From 542f306fedb32d69b4f53948c68383a045f30e87 Mon Sep 17 00:00:00 2001 From: Volodymyr Sazheniuk Date: Sat, 21 Feb 2026 15:58:20 +0100 Subject: [PATCH 1/3] Arbitrary S3 object versions support --- Cargo.lock | 10 ++ Cargo.toml | 1 + .../examples/client_benchmark.rs | 2 +- mountpoint-s3-client/examples/download.rs | 2 +- mountpoint-s3-client/src/failure_client.rs | 8 +- mountpoint-s3-client/src/lib.rs | 2 +- mountpoint-s3-client/src/mock_client.rs | 80 ++++++--- .../src/mock_client/throughput_client.rs | 8 +- mountpoint-s3-client/src/object_client.rs | 5 + mountpoint-s3-client/src/s3_crt_client.rs | 6 +- .../src/s3_crt_client/get_object.rs | 15 +- .../src/s3_crt_client/head_object.rs | 14 +- mountpoint-s3-crt-sys/crt/aws-c-auth | 2 +- mountpoint-s3-crt-sys/crt/aws-c-cal | 2 +- mountpoint-s3-crt-sys/crt/aws-c-common | 2 +- mountpoint-s3-crt-sys/crt/aws-c-compression | 2 +- mountpoint-s3-crt-sys/crt/aws-c-http | 2 +- mountpoint-s3-crt-sys/crt/aws-c-io | 2 +- mountpoint-s3-crt-sys/crt/aws-c-s3 | 2 +- mountpoint-s3-crt-sys/crt/aws-checksums | 2 +- mountpoint-s3-crt-sys/crt/aws-lc | 2 +- mountpoint-s3-crt-sys/crt/s2n-tls | 2 +- mountpoint-s3-fs/Cargo.toml | 3 +- .../examples/prefetch_benchmark.rs | 2 +- .../src/data_cache/express_data_cache.rs | 1 + mountpoint-s3-fs/src/fs.rs | 40 ++++- mountpoint-s3-fs/src/fs/handles.rs | 7 +- mountpoint-s3-fs/src/fuse.rs | 28 +++- mountpoint-s3-fs/src/manifest/metablock.rs | 4 + mountpoint-s3-fs/src/metablock.rs | 4 + mountpoint-s3-fs/src/metablock/path.rs | 15 +- mountpoint-s3-fs/src/metablock/stat.rs | 5 + mountpoint-s3-fs/src/object.rs | 17 +- mountpoint-s3-fs/src/prefetch/part_stream.rs | 2 +- mountpoint-s3-fs/src/superblock.rs | 101 ++++++++++- mountpoint-s3-fs/src/superblock/inode.rs | 47 +++++- mountpoint-s3-fs/src/superblock/readdir.rs | 157 ++++++++++++++---- mountpoint-s3-fs/src/upload/incremental.rs | 16 +- mountpoint-s3-fs/tests/fs.rs | 2 +- mountpoint-s3-ioctl/Cargo.toml | 15 ++ mountpoint-s3-ioctl/src/error.rs | 5 + mountpoint-s3-ioctl/src/lib.rs | 41 +++++ mountpoint-s3-ioctl/src/main.rs | 42 +++++ 43 files changed, 620 insertions(+), 107 deletions(-) create mode 100644 mountpoint-s3-ioctl/Cargo.toml create mode 100644 mountpoint-s3-ioctl/src/error.rs create mode 100644 mountpoint-s3-ioctl/src/lib.rs create mode 100644 mountpoint-s3-ioctl/src/main.rs diff --git a/Cargo.lock b/Cargo.lock index c0753faf71..316fec8bbf 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2911,6 +2911,7 @@ dependencies = [ "mime_guess", "mountpoint-s3-client", "mountpoint-s3-fuser", + "mountpoint-s3-ioctl", "nix 0.31.3", "opentelemetry", "opentelemetry-otlp", @@ -2967,6 +2968,15 @@ dependencies = [ "zerocopy", ] +[[package]] +name = "mountpoint-s3-ioctl" +version = "0.1.0" +dependencies = [ + "clap", + "nix 0.30.1", + "thiserror", +] + [[package]] name = "nix" version = "0.29.0" diff --git a/Cargo.toml b/Cargo.toml index 88dd6b7a17..cd57a53e54 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -6,6 +6,7 @@ members = [ "mountpoint-s3", "mountpoint-s3-fuser", "mountpoint-s3-fs", + "mountpoint-s3-ioctl", ] resolver = "2" diff --git a/mountpoint-s3-client/examples/client_benchmark.rs b/mountpoint-s3-client/examples/client_benchmark.rs index 5eca50e84b..e3c28a5368 100644 --- a/mountpoint-s3-client/examples/client_benchmark.rs +++ b/mountpoint-s3-client/examples/client_benchmark.rs @@ -62,7 +62,7 @@ fn run_benchmark( futures::executor::block_on(async move { let mut received_obj_len = 0u64; let mut request = client - .get_object(bucket, key, &GetObjectParams::new()) + .get_object(bucket, key, None, &GetObjectParams::new()) .await .expect("couldn't create get request"); let mut backpressure_handle = request.backpressure_handle().cloned(); diff --git a/mountpoint-s3-client/examples/download.rs b/mountpoint-s3-client/examples/download.rs index dd45820cbe..0640e8389a 100644 --- a/mountpoint-s3-client/examples/download.rs +++ b/mountpoint-s3-client/examples/download.rs @@ -88,7 +88,7 @@ fn main() { let last_offset_clone = Arc::clone(&last_offset); futures::executor::block_on(async move { let mut request = client - .get_object(bucket, key, &GetObjectParams::new().range(range)) + .get_object(bucket, key, None, &GetObjectParams::new().range(range)) .await .expect("couldn't create get request"); loop { diff --git a/mountpoint-s3-client/src/failure_client.rs b/mountpoint-s3-client/src/failure_client.rs index 45720511b7..5cf9769e14 100644 --- a/mountpoint-s3-client/src/failure_client.rs +++ b/mountpoint-s3-client/src/failure_client.rs @@ -115,6 +115,7 @@ where &self, bucket: &str, key: &str, + version: Option<&str>, params: &GetObjectParams, ) -> ObjectClientResult { let failure_mode = (self.get_object_cb)(&mut *self.state.lock().unwrap(), bucket, key, params); @@ -122,7 +123,7 @@ where return Err(err); } - let request = self.client.get_object(bucket, key, params).await?; + let request = self.client.get_object(bucket, key, version, params).await?; Ok(FailureGetResponse { request, poll_count: 0, @@ -156,10 +157,11 @@ where &self, bucket: &str, key: &str, + version: Option<&str>, params: &HeadObjectParams, ) -> ObjectClientResult { (self.head_object_cb)(&mut *self.state.lock().unwrap(), bucket, key)?; - self.client.head_object(bucket, key, params).await + self.client.head_object(bucket, key, version, params).await } async fn put_object( @@ -491,7 +493,7 @@ mod tests { let fail_set = HashSet::from([2, 4, 5]); for i in 1..=6 { - let r = fail_client.get_object(bucket, key, &GetObjectParams::new()).await; + let r = fail_client.get_object(bucket, key, None, &GetObjectParams::new()).await; if fail_set.contains(&i) { assert!(r.is_err()); } else { diff --git a/mountpoint-s3-client/src/lib.rs b/mountpoint-s3-client/src/lib.rs index 181d161902..19d8fbb534 100644 --- a/mountpoint-s3-client/src/lib.rs +++ b/mountpoint-s3-client/src/lib.rs @@ -22,7 +22,7 @@ //! //! let client = S3CrtClient::new(Default::default()).expect("client construction failed"); //! -//! let response = client.get_object("my-bucket", "my-key", &GetObjectParams::new()).await.expect("get_object failed"); +//! let response = client.get_object("my-bucket", "my-key", Some("object-version"), &GetObjectParams::new()).await.expect("get_object failed"); //! let body = response.map_ok(|part| part.data.to_vec()).try_concat().await.expect("body streaming failed"); //! # } //! ``` diff --git a/mountpoint-s3-client/src/mock_client.rs b/mountpoint-s3-client/src/mock_client.rs index 28bff22666..b5ce215eb5 100644 --- a/mountpoint-s3-client/src/mock_client.rs +++ b/mountpoint-s3-client/src/mock_client.rs @@ -936,6 +936,7 @@ impl ObjectClient for MockClient { &self, bucket: &str, key: &str, + _version: Option<&str>, params: &GetObjectParams, ) -> ObjectClientResult { trace!(bucket, key, ?params.range, ?params.if_match, "GetObject"); @@ -999,6 +1000,7 @@ impl ObjectClient for MockClient { &self, bucket: &str, key: &str, + _version: Option<&str>, params: &HeadObjectParams, ) -> ObjectClientResult { trace!(bucket, key, "HeadObject"); @@ -1025,6 +1027,7 @@ impl ObjectClient for MockClient { checksum, sse_type: None, sse_kms_key_id: None, + version_id: None, }) } else { Err(ObjectClientError::ServiceError(HeadObjectError::NotFound)) @@ -1455,7 +1458,7 @@ mod tests { client.add_object(key, object); let mut get_request = client - .get_object("test_bucket", key, &GetObjectParams::new().range(range.clone())) + .get_object("test_bucket", key, None, &GetObjectParams::new().range(range.clone())) .await .expect("should not fail"); @@ -1509,7 +1512,7 @@ mod tests { client.add_object(key, MockObject::from_bytes(&body, ETag::for_tests())); let mut get_request = client - .get_object("test_bucket", key, &GetObjectParams::new().range(range.clone())) + .get_object("test_bucket", key, None, &GetObjectParams::new().range(range.clone())) .await .expect("should not fail"); let mut backpressure_handle = get_request @@ -1553,44 +1556,71 @@ mod tests { client.add_object("key1", body[..].into()); assert!(matches!( - client.get_object("wrong_bucket", "key1", &GetObjectParams::new()).await, + client + .get_object("wrong_bucket", "key1", None, &GetObjectParams::new()) + .await, Err(ObjectClientError::ServiceError(GetObjectError::NoSuchBucket(_))) )); assert!(matches!( client - .get_object("test_bucket", "wrong_key", &GetObjectParams::new()) + .get_object("test_bucket", "wrong_key", None, &GetObjectParams::new()) .await, Err(ObjectClientError::ServiceError(GetObjectError::NoSuchKey(_))) )); assert_client_error!( client - .get_object("test_bucket", "key1", &GetObjectParams::new().range(Some(0..2001))) + .get_object( + "test_bucket", + "key1", + None, + &GetObjectParams::new().range(Some(0..2001)) + ) .await, "invalid range, length=2000" ); assert_client_error!( client - .get_object("test_bucket", "key1", &GetObjectParams::new().range(Some(2000..2000))) + .get_object( + "test_bucket", + "key1", + None, + &GetObjectParams::new().range(Some(2000..2000)) + ) .await, "invalid range, length=2000" ); assert_client_error!( client - .get_object("test_bucket", "key1", &GetObjectParams::new().range(Some(500..2001))) + .get_object( + "test_bucket", + "key1", + None, + &GetObjectParams::new().range(Some(500..2001)) + ) .await, "invalid range, length=2000" ); assert_client_error!( client - .get_object("test_bucket", "key1", &GetObjectParams::new().range(Some(5000..2001))) + .get_object( + "test_bucket", + "key1", + None, + &GetObjectParams::new().range(Some(5000..2001)) + ) .await, "invalid range, length=2000" ); assert_client_error!( client - .get_object("test_bucket", "key1", &GetObjectParams::new().range(Some(5000..1))) + .get_object( + "test_bucket", + "key1", + None, + &GetObjectParams::new().range(Some(5000..1)) + ) .await, "invalid range, length=2000" ); @@ -1619,7 +1649,12 @@ mod tests { client.add_object(key, MockObject::from_bytes(&expected_body, ETag::for_tests())); let mut get_request = client - .get_object("test_bucket", key, &GetObjectParams::new().range(Some(range.clone()))) + .get_object( + "test_bucket", + key, + None, + &GetObjectParams::new().range(Some(range.clone())), + ) .await .expect("should not fail"); @@ -1652,7 +1687,7 @@ mod tests { .expect("Should not fail"); client - .get_object(bucket, dst_key, &GetObjectParams::new()) + .get_object(bucket, dst_key, None, &GetObjectParams::new()) .await .expect("get_object should succeed"); } @@ -2059,7 +2094,7 @@ mod tests { put_request.complete().await.expect("put_object failed"); let mut get_request = client - .get_object("test_bucket", "key1", &GetObjectParams::new()) + .get_object("test_bucket", "key1", None, &GetObjectParams::new()) .await .expect("get_object failed"); @@ -2092,7 +2127,7 @@ mod tests { .expect("put_object failed"); let get_request = client - .get_object("test_bucket", "key1", &GetObjectParams::new()) + .get_object("test_bucket", "key1", None, &GetObjectParams::new()) .await .expect("get_object failed"); @@ -2284,12 +2319,12 @@ mod tests { .expect("rename should succeed"); assert!(matches!( client - .head_object("test_bucket", "key1", &HeadObjectParams::new()) + .head_object("test_bucket", "key1", None, &HeadObjectParams::new()) .await, Err(ObjectClientError::ServiceError(HeadObjectError::NotFound)) )); client - .head_object("test_bucket", "new_key1", &HeadObjectParams::new()) + .head_object("test_bucket", "new_key1", None, &HeadObjectParams::new()) .await .expect("object should now exist with new key"); } @@ -2347,12 +2382,12 @@ mod tests { // Assert that key2 is accessible, while key1 is not assert!(matches!( client - .head_object("test_bucket", "key1", &HeadObjectParams::new()) + .head_object("test_bucket", "key1", None, &HeadObjectParams::new()) .await, Err(ObjectClientError::ServiceError(HeadObjectError::NotFound)) )); client - .head_object("test_bucket", "key2", &HeadObjectParams::new()) + .head_object("test_bucket", "key2", None, &HeadObjectParams::new()) .await .expect("object should now exist with new key"); } @@ -2391,7 +2426,10 @@ mod tests { put_request.complete().await.unwrap(); // head_object returns storage class - let head_result = client.head_object(bucket, key, &HeadObjectParams::new()).await.unwrap(); + let head_result = client + .head_object(bucket, key, None, &HeadObjectParams::new()) + .await + .unwrap(); assert_eq!(head_result.storage_class.as_deref(), storage_class); // list_objects returns storage class @@ -2409,14 +2447,14 @@ mod tests { let head_counter_1 = client.new_counter(Operation::HeadObject); let delete_counter_1 = client.new_counter(Operation::DeleteObject); - let _result = client.head_object(bucket, "key", &HeadObjectParams::new()).await; + let _result = client.head_object(bucket, "key", None, &HeadObjectParams::new()).await; assert_eq!(1, head_counter_1.count()); assert_eq!(0, delete_counter_1.count()); let head_counter_2 = client.new_counter(Operation::HeadObject); assert_eq!(0, head_counter_2.count()); - let _result = client.head_object(bucket, "key", &HeadObjectParams::new()).await; + let _result = client.head_object(bucket, "key", None, &HeadObjectParams::new()).await; let _result = client.delete_object(bucket, "key").await; let _result = client.delete_object(bucket, "key").await; let _result = client.delete_object(bucket, "key").await; @@ -2547,7 +2585,7 @@ mod tests { .expect("append failed"); let get_request = client - .get_object(bucket, key, &GetObjectParams::default()) + .get_object(bucket, key, None, &GetObjectParams::default()) .await .expect("get_object failed"); diff --git a/mountpoint-s3-client/src/mock_client/throughput_client.rs b/mountpoint-s3-client/src/mock_client/throughput_client.rs index 34f0ac8ebb..974a755de8 100644 --- a/mountpoint-s3-client/src/mock_client/throughput_client.rs +++ b/mountpoint-s3-client/src/mock_client/throughput_client.rs @@ -164,9 +164,10 @@ impl ObjectClient for ThroughputMockClient { &self, bucket: &str, key: &str, + version: Option<&str>, params: &GetObjectParams, ) -> ObjectClientResult { - let request = self.inner.get_object(bucket, key, params).await?; + let request = self.inner.get_object(bucket, key, version, params).await?; let rate_limiter = self.rate_limiter.clone(); Ok(ThroughputGetObjectResponse { request, rate_limiter }) } @@ -188,9 +189,10 @@ impl ObjectClient for ThroughputMockClient { &self, bucket: &str, key: &str, + version: Option<&str>, params: &HeadObjectParams, ) -> ObjectClientResult { - self.inner.head_object(bucket, key, params).await + self.inner.head_object(bucket, key, version, params).await } async fn put_object( @@ -267,7 +269,7 @@ mod tests { let num_bytes = block_on(async move { let mut num_bytes = 0; let mut get = client - .get_object("test_bucket", "testfile", &GetObjectParams::new()) + .get_object("test_bucket", "testfile", None, &GetObjectParams::new()) .await .unwrap(); while let Some(part) = get.next().await { diff --git a/mountpoint-s3-client/src/object_client.rs b/mountpoint-s3-client/src/object_client.rs index 1aee960033..8fcb8e25a4 100644 --- a/mountpoint-s3-client/src/object_client.rs +++ b/mountpoint-s3-client/src/object_client.rs @@ -78,6 +78,7 @@ pub trait ObjectClient { &self, bucket: &str, key: &str, + version: Option<&str>, params: &GetObjectParams, ) -> ObjectClientResult; @@ -96,6 +97,7 @@ pub trait ObjectClient { &self, bucket: &str, key: &str, + version: Option<&str>, params: &HeadObjectParams, ) -> ObjectClientResult; @@ -366,6 +368,9 @@ pub struct HeadObjectResult { /// Server-side encryption KMS key ID that was used to store the object. pub sse_kms_key_id: Option, + + /// Version ID of the object if versioning is enabled on the bucket. + pub version_id: Option, } /// Errors returned by a [`head_object`](ObjectClient::head_object) request diff --git a/mountpoint-s3-client/src/s3_crt_client.rs b/mountpoint-s3-client/src/s3_crt_client.rs index 6e61fa2fe3..fd200dafcc 100644 --- a/mountpoint-s3-client/src/s3_crt_client.rs +++ b/mountpoint-s3-client/src/s3_crt_client.rs @@ -1628,9 +1628,10 @@ impl ObjectClient for S3CrtClient { &self, bucket: &str, key: &str, + version: Option<&str>, params: &GetObjectParams, ) -> ObjectClientResult { - self.get_object(bucket, key, params).await + self.get_object(bucket, key, version, params).await } async fn list_objects( @@ -1649,9 +1650,10 @@ impl ObjectClient for S3CrtClient { &self, bucket: &str, key: &str, + version: Option<&str>, params: &HeadObjectParams, ) -> ObjectClientResult { - self.head_object(bucket, key, params).await + self.head_object(bucket, key, version, params).await } async fn put_object( diff --git a/mountpoint-s3-client/src/s3_crt_client/get_object.rs b/mountpoint-s3-client/src/s3_crt_client/get_object.rs index 3feabef42e..634b8c2779 100644 --- a/mountpoint-s3-client/src/s3_crt_client/get_object.rs +++ b/mountpoint-s3-client/src/s3_crt_client/get_object.rs @@ -21,7 +21,10 @@ use crate::object_client::{ ObjectChecksumError, ObjectClientError, ObjectClientResult, ObjectMetadata, }; -use super::{CancellingMetaRequest, ResponseHeadersError, S3CrtClient, S3Operation, S3RequestError, parse_checksum}; +use super::{ + CancellingMetaRequest, QueryFragment, ResponseHeadersError, S3CrtClient, S3Operation, S3RequestError, + parse_checksum, +}; impl S3CrtClient { /// Create and begin a new GetObject request. The returned [S3GetObjectResponse] is a [Stream] of @@ -30,6 +33,7 @@ impl S3CrtClient { &self, bucket: &str, key: &str, + version: Option<&str>, params: &GetObjectParams, ) -> Result> { let requested_checksums = params.checksum_mode.as_ref() == Some(&ChecksumMode::Enabled); @@ -71,9 +75,14 @@ impl S3CrtClient { .map_err(S3RequestError::construction_failure)?; } - let key = format!("/{key}"); + let query = if let Some(version) = version { + QueryFragment::Query(&[("versionId", version)]) + } else { + Default::default() + }; + message - .set_request_path(key) + .set_request_path_and_query(format!("/{key}"), query) .map_err(S3RequestError::construction_failure)?; let mut options = message.into_options(S3Operation::GetObject); diff --git a/mountpoint-s3-client/src/s3_crt_client/head_object.rs b/mountpoint-s3-client/src/s3_crt_client/head_object.rs index d94517af71..5cbb94a09c 100644 --- a/mountpoint-s3-client/src/s3_crt_client/head_object.rs +++ b/mountpoint-s3-client/src/s3_crt_client/head_object.rs @@ -10,7 +10,7 @@ use time::format_description::well_known::Rfc2822; use crate::object_client::{HeadObjectError, HeadObjectParams, HeadObjectResult, ObjectClientResult, RestoreStatus}; -use super::{ChecksumMode, S3CrtClient, S3Operation, S3RequestError, parse_checksum}; +use super::{ChecksumMode, QueryFragment, S3CrtClient, S3Operation, S3RequestError, parse_checksum}; #[derive(Error, Debug)] #[non_exhaustive] @@ -74,6 +74,7 @@ impl HeadObjectResult { let restore_status = Self::parse_restore_status(headers)?; let sse_type = headers.get_as_optional_string("x-amz-server-side-encryption")?; let sse_kms_key_id = headers.get_as_optional_string("x-amz-server-side-encryption-aws-kms-key-id")?; + let version_id = headers.get_as_optional_string("x-amz-version-id")?; let checksum = parse_checksum(headers)?; let result = HeadObjectResult { size, @@ -84,6 +85,7 @@ impl HeadObjectResult { checksum, sse_type, sse_kms_key_id, + version_id, }; Ok(result) } @@ -94,6 +96,7 @@ impl S3CrtClient { &self, bucket: &str, key: &str, + version: Option<&str>, params: &HeadObjectParams, ) -> ObjectClientResult { let request = { @@ -102,9 +105,14 @@ impl S3CrtClient { .new_request_template("HEAD", bucket) .map_err(S3RequestError::construction_failure)?; - let key = key.to_string(); + let query = if let Some(version) = version { + QueryFragment::Query(&[("versionId", version)]) + } else { + Default::default() + }; + message - .set_request_path(format!("/{key}")) + .set_request_path_and_query(format!("/{key}"), query) .map_err(S3RequestError::construction_failure)?; let bucket = bucket.to_owned(); diff --git a/mountpoint-s3-crt-sys/crt/aws-c-auth b/mountpoint-s3-crt-sys/crt/aws-c-auth index 4b5d524bf1..16c7289432 160000 --- a/mountpoint-s3-crt-sys/crt/aws-c-auth +++ b/mountpoint-s3-crt-sys/crt/aws-c-auth @@ -1 +1 @@ -Subproject commit 4b5d524bf1a511b05e0fffe5bdc51800770b9427 +Subproject commit 16c7289432839ca7e49132b68ad5a776e7279394 diff --git a/mountpoint-s3-crt-sys/crt/aws-c-cal b/mountpoint-s3-crt-sys/crt/aws-c-cal index 9edd8eac2b..9abb8e31c4 160000 --- a/mountpoint-s3-crt-sys/crt/aws-c-cal +++ b/mountpoint-s3-crt-sys/crt/aws-c-cal @@ -1 +1 @@ -Subproject commit 9edd8eac2b21ca6a04535b91d60d361c2f1bb60f +Subproject commit 9abb8e31c43b705377639095d616f56611d428a6 diff --git a/mountpoint-s3-crt-sys/crt/aws-c-common b/mountpoint-s3-crt-sys/crt/aws-c-common index a9d57d2d37..dc44b5a53f 160000 --- a/mountpoint-s3-crt-sys/crt/aws-c-common +++ b/mountpoint-s3-crt-sys/crt/aws-c-common @@ -1 +1 @@ -Subproject commit a9d57d2d372582fa18bf1d81741e107c59a23226 +Subproject commit dc44b5a53f0503eecd7a043459d06879d5f21acd diff --git a/mountpoint-s3-crt-sys/crt/aws-c-compression b/mountpoint-s3-crt-sys/crt/aws-c-compression index d8264e64f6..f951ab2b81 160000 --- a/mountpoint-s3-crt-sys/crt/aws-c-compression +++ b/mountpoint-s3-crt-sys/crt/aws-c-compression @@ -1 +1 @@ -Subproject commit d8264e64f698341eb03039b96b4f44702a9b3f83 +Subproject commit f951ab2b819fc6993b6e5e6cfef64b1a1554bfc8 diff --git a/mountpoint-s3-crt-sys/crt/aws-c-http b/mountpoint-s3-crt-sys/crt/aws-c-http index 8aefd899fc..f32d2de349 160000 --- a/mountpoint-s3-crt-sys/crt/aws-c-http +++ b/mountpoint-s3-crt-sys/crt/aws-c-http @@ -1 +1 @@ -Subproject commit 8aefd899fc3210bfd0e3fd414011a3cb708bf6e4 +Subproject commit f32d2de3497ad03aa316d64d7db3c6773399cc8e diff --git a/mountpoint-s3-crt-sys/crt/aws-c-io b/mountpoint-s3-crt-sys/crt/aws-c-io index 8bda5cf0fe..fbac3c30fd 160000 --- a/mountpoint-s3-crt-sys/crt/aws-c-io +++ b/mountpoint-s3-crt-sys/crt/aws-c-io @@ -1 +1 @@ -Subproject commit 8bda5cf0fe7f075ad879f27232d5ee83e1b67431 +Subproject commit fbac3c30fd8c50c05168f41486403a69d91f7600 diff --git a/mountpoint-s3-crt-sys/crt/aws-c-s3 b/mountpoint-s3-crt-sys/crt/aws-c-s3 index 448ec5e49e..1e9c52ca20 160000 --- a/mountpoint-s3-crt-sys/crt/aws-c-s3 +++ b/mountpoint-s3-crt-sys/crt/aws-c-s3 @@ -1 +1 @@ -Subproject commit 448ec5e49eb181687f9af39b0d3b83a563b15039 +Subproject commit 1e9c52ca2059ab72f8c066d12960a5e33a7ac043 diff --git a/mountpoint-s3-crt-sys/crt/aws-checksums b/mountpoint-s3-crt-sys/crt/aws-checksums index 1d5f2f1f3e..5c0d51dbf2 160000 --- a/mountpoint-s3-crt-sys/crt/aws-checksums +++ b/mountpoint-s3-crt-sys/crt/aws-checksums @@ -1 +1 @@ -Subproject commit 1d5f2f1f3e5d013aae8810878ceb5b3f6f258c4e +Subproject commit 5c0d51dbf239fe845353cd5c08ee077f8f2fe9fb diff --git a/mountpoint-s3-crt-sys/crt/aws-lc b/mountpoint-s3-crt-sys/crt/aws-lc index 6283365b1d..ea70f681ce 160000 --- a/mountpoint-s3-crt-sys/crt/aws-lc +++ b/mountpoint-s3-crt-sys/crt/aws-lc @@ -1 +1 @@ -Subproject commit 6283365b1d43abadfa9b997812cef06b095f7f04 +Subproject commit ea70f681ce48c3996b7584be573355c2ebdc56e3 diff --git a/mountpoint-s3-crt-sys/crt/s2n-tls b/mountpoint-s3-crt-sys/crt/s2n-tls index f5f6c6c2ce..627cfd188d 160000 --- a/mountpoint-s3-crt-sys/crt/s2n-tls +++ b/mountpoint-s3-crt-sys/crt/s2n-tls @@ -1 +1 @@ -Subproject commit f5f6c6c2ce2370de1aa3ade6899a7321d1127bb8 +Subproject commit 627cfd188df62a0e5c299b74ca9228d2977bba4b diff --git a/mountpoint-s3-fs/Cargo.toml b/mountpoint-s3-fs/Cargo.toml index dca0ba3b6e..cf4b3043de 100644 --- a/mountpoint-s3-fs/Cargo.toml +++ b/mountpoint-s3-fs/Cargo.toml @@ -10,6 +10,7 @@ description = "Mountpoint S3 main library" [dependencies] mountpoint-s3-fuser = { workspace = true } mountpoint-s3-client = { workspace = true } +mountpoint-s3-ioctl = { path = "../mountpoint-s3-ioctl", version = "0.1.0" } anyhow = { version = "1.0.103", features = ["backtrace"] } async-channel = "2.5.0" @@ -34,7 +35,7 @@ libc = "0.2.186" linked-hash-map = "0.5.6" metrics = "0.24.6" mime_guess = "2.0.5" -nix = { version = "0.31.3", default-features = false, features = ["fs", "process", "signal", "user"] } +nix = { version = "0.31.3", default-features = false, features = ["fs", "process", "signal", "user", "ioctl"] } rand = "0.10.1" regex = "1.12.4" rusqlite = { version = "0.40.1", features = ["bundled", "fallible_uint"], optional = true } diff --git a/mountpoint-s3-fs/examples/prefetch_benchmark.rs b/mountpoint-s3-fs/examples/prefetch_benchmark.rs index 0cd0904c39..57172a4e4d 100644 --- a/mountpoint-s3-fs/examples/prefetch_benchmark.rs +++ b/mountpoint-s3-fs/examples/prefetch_benchmark.rs @@ -178,7 +178,7 @@ fn main() -> anyhow::Result<()> { .s3_keys .iter() .map(|key| { - let head_result = block_on(client.head_object(bucket, key, &HeadObjectParams::new())) + let head_result = block_on(client.head_object(bucket, key, None, &HeadObjectParams::new())) .with_context(|| format!("HeadObject failed for {key}"))?; Ok((ObjectId::new(key.to_string(), head_result.etag), head_result.size)) }) diff --git a/mountpoint-s3-fs/src/data_cache/express_data_cache.rs b/mountpoint-s3-fs/src/data_cache/express_data_cache.rs index d304df2d6d..2fcf6f7181 100644 --- a/mountpoint-s3-fs/src/data_cache/express_data_cache.rs +++ b/mountpoint-s3-fs/src/data_cache/express_data_cache.rs @@ -185,6 +185,7 @@ where .get_object( &self.config.bucket_name, &object_key, + None, &GetObjectParams::new().checksum_mode(Some(ChecksumMode::Enabled)), ) .await diff --git a/mountpoint-s3-fs/src/fs.rs b/mountpoint-s3-fs/src/fs.rs index 96c52d54f3..4317f174ea 100644 --- a/mountpoint-s3-fs/src/fs.rs +++ b/mountpoint-s3-fs/src/fs.rs @@ -16,7 +16,9 @@ use crate::async_util::Runtime; use crate::logging; use crate::mem_limiter::MemoryLimiter; use crate::memory::PagedPool; -use crate::metablock::{AddDirEntry, AddDirEntryResult, InodeInformation, Metablock, PendingUploadHook, ReadWriteMode}; +use crate::metablock::{ + AddDirEntry, AddDirEntryResult, InodeInformation, Metablock, NewHandle, PendingUploadHook, ReadWriteMode, +}; pub use crate::metablock::{InodeError, InodeKind, InodeNo}; use crate::prefetch::{Prefetcher, PrefetcherBuilder}; use crate::sync::atomic::{AtomicU64, Ordering}; @@ -780,6 +782,42 @@ where tracing::trace!("rename complete"); Ok(()) } + + pub async fn set_inode_version(&self, ino: InodeNo, version: Option<&str>) -> Result<(), Error> { + trace!(inode = ino, ?version, "fs:set_inode_version"); + + let mut handles_write = self.file_handles.write().await; + + // First, update the inode itself. + self.metablock.set_inode_version(ino, version).await?; + + if let std::collections::hash_map::Entry::Occupied(mut entry) = handles_write.entry(ino) { + let handle = entry.get(); + if matches!(*entry.get().state.lock().await, FileHandleState::Write { .. }) { + // We cannot change the inode version while there is a write handle open. + return Err(err!( + libc::EBUSY, + "Cannot set inode version while file is open for writing." + )); + } + trace!(pid = handle.open_pid, "Recreating read handle with new inode version."); + // Recreate the read handle with the new version. + let lookup = self.metablock.getattr(ino, false).await?; + let new_handle = NewHandle::read(lookup.clone()); + let handle_state = FileHandleState::new(&new_handle, OpenFlags::empty(), self).await?; + let handle = Arc::new(FileHandle { + ino, + open_pid: handle.open_pid, + location: lookup.try_into_s3_location()?, + state: AsyncMutex::new(handle_state), + }); + *entry.get_mut() = handle; + } + + trace!("fs:set_inode_version complete"); + + Ok(()) + } } #[cfg(test)] diff --git a/mountpoint-s3-fs/src/fs/handles.rs b/mountpoint-s3-fs/src/fs/handles.rs index b7b8596f2e..e685719a1e 100644 --- a/mountpoint-s3-fs/src/fs/handles.rs +++ b/mountpoint-s3-fs/src/fs/handles.rs @@ -67,6 +67,7 @@ where let stat = handle.lookup.stat(); let location = handle.lookup.s3_location()?; let full_key = location.full_key(); + let object_version = location.version(); let bucket = location.bucket_name(); match handle.mode { @@ -76,7 +77,11 @@ where None => return Err(err!(libc::EBADF, "no E-Tag for inode {}", ino)), Some(etag) => ETag::from_str(etag).expect("E-Tag should be set"), }; - let object_id = ObjectId::new(full_key.into(), etag); + let object_id = if let Some(object_version) = object_version { + ObjectId::new_with_version(full_key.into(), Some(object_version.to_owned()), etag) + } else { + ObjectId::new(full_key.into(), etag) + }; let request = fs .prefetcher .prefetch(bucket.to_string(), object_id, HandleId::new(fh), object_size); diff --git a/mountpoint-s3-fs/src/fuse.rs b/mountpoint-s3-fs/src/fuse.rs index 2b6b9df9f5..e80b4cb00d 100644 --- a/mountpoint-s3-fs/src/fuse.rs +++ b/mountpoint-s3-fs/src/fuse.rs @@ -566,19 +566,35 @@ where fuse_unsupported!("bmap", reply); } - #[instrument(level="warn", skip_all, fields(req=_req.unique(), ino=_ino, fh=_fh, cmd=_cmd, pid=_req.pid()))] + #[instrument(level="warn", skip_all, fields(req=req.unique(), ino=ino, fh=_fh, cmd=cmd, pid=req.pid()))] fn ioctl( &self, - _req: &Request<'_>, - _ino: u64, + req: &Request<'_>, + ino: u64, _fh: u64, _flags: u32, - _cmd: u32, - _in_data: &[u8], + cmd: u32, + in_data: &[u8], _out_size: u32, reply: ReplyIoctl, ) { - fuse_unsupported!("ioctl", reply, libc::ENOSYS, tracing::Level::DEBUG); + match cmd as u64 { + mountpoint_s3_ioctl::MOUNT_S3_IOC_TYPE_SET_VERSION => { + let version = { + let len = in_data.iter().position(|b| *b == 0).map_or(in_data.len(), |pos| pos); + (len > 0).then(|| String::from_utf8_lossy(&in_data[..len])) + }; + match block_on(self.fs.set_inode_version(ino, version.as_deref()).in_current_span()) { + Ok(()) => { + reply.ioctl(0, &[]); + } + Err(e) => fuse_error!("ioctl", reply, e, self, req), + } + } + _ => { + fuse_unsupported!("ioctl", reply, libc::ENOSYS, tracing::Level::DEBUG); + } + } } #[instrument(level="warn", skip_all, fields(req=_req.unique(), ino=_ino, fh=_fh, offset=_offset, length=_length, pid=_req.pid()))] diff --git a/mountpoint-s3-fs/src/manifest/metablock.rs b/mountpoint-s3-fs/src/manifest/metablock.rs index 548ddf8f80..1ea188fa55 100644 --- a/mountpoint-s3-fs/src/manifest/metablock.rs +++ b/mountpoint-s3-fs/src/manifest/metablock.rs @@ -318,4 +318,8 @@ impl Metablock for ManifestMetablock { bucket: None, })) } + + async fn set_inode_version(&self, _ino: InodeNo, _version: Option<&str>) -> Result<(), InodeError> { + Ok(()) + } } diff --git a/mountpoint-s3-fs/src/metablock.rs b/mountpoint-s3-fs/src/metablock.rs index a3795a7713..d36a9821b4 100644 --- a/mountpoint-s3-fs/src/metablock.rs +++ b/mountpoint-s3-fs/src/metablock.rs @@ -152,6 +152,10 @@ pub trait Metablock: Send + Sync { /// Unlink the entry described by `parent_ino` and `name`. async fn unlink(&self, parent_ino: InodeNo, name: &OsStr) -> Result<(), InodeError>; + + /// Explicitly set S3 object version for the given inode. + /// If `version` is `None`, any existing version information shall be removed. + async fn set_inode_version(&self, ino: InodeNo, version: Option<&str>) -> Result<(), InodeError>; } /// Callback to the file system which adds directory entries to the reply buffer. diff --git a/mountpoint-s3-fs/src/metablock/path.rs b/mountpoint-s3-fs/src/metablock/path.rs index 295da2b21e..4afac95a8b 100644 --- a/mountpoint-s3-fs/src/metablock/path.rs +++ b/mountpoint-s3-fs/src/metablock/path.rs @@ -168,12 +168,21 @@ impl TryFrom for ValidKey { #[derive(Debug, Clone)] pub struct S3Location { pub path: Arc, + pub version: Option, pub partial_key: ValidKey, } impl S3Location { pub fn new(path: Arc, partial_key: ValidKey) -> Self { - Self { path, partial_key } + Self::new_with_object_version(path, None, partial_key) + } + + pub fn new_with_object_version(path: Arc, version: Option, partial_key: ValidKey) -> Self { + Self { + path, + version, + partial_key, + } } /// Get the bucket name @@ -189,6 +198,10 @@ impl S3Location { pub fn name(&self) -> &str { self.partial_key.name() } + + pub fn version(&self) -> Option<&str> { + self.version.as_deref() + } } impl Display for S3Location { diff --git a/mountpoint-s3-fs/src/metablock/stat.rs b/mountpoint-s3-fs/src/metablock/stat.rs index 811a42f778..4ea4b76447 100644 --- a/mountpoint-s3-fs/src/metablock/stat.rs +++ b/mountpoint-s3-fs/src/metablock/stat.rs @@ -51,6 +51,8 @@ pub struct InodeStat { /// are only readable after restoration. For objects with other storage classes /// this field should be always `true`. pub is_readable: bool, + /// Version ID of the object, if versioning is enabled + pub version_id: Option>, } impl InodeStat { @@ -96,6 +98,7 @@ impl InodeStat { storage_class: Option<&str>, restore_status: Option, validity: Duration, + version_id: Option>, ) -> InodeStat { let is_readable = Self::is_readable(storage_class, restore_status); InodeStat { @@ -106,6 +109,7 @@ impl InodeStat { mtime: datetime, etag, is_readable, + version_id, } } @@ -119,6 +123,7 @@ impl InodeStat { mtime: datetime, etag: None, is_readable: true, + version_id: None, } } } diff --git a/mountpoint-s3-fs/src/object.rs b/mountpoint-s3-fs/src/object.rs index cd733b215d..54d37e5fbd 100644 --- a/mountpoint-s3-fs/src/object.rs +++ b/mountpoint-s3-fs/src/object.rs @@ -23,13 +23,24 @@ impl Debug for ObjectId { #[derive(Debug, Hash, PartialEq, Eq)] struct InnerObjectId { key: String, + version: Option, etag: ETag, } impl ObjectId { pub fn new(key: String, etag: ETag) -> Self { Self { - inner: Arc::new(InnerObjectId { key, etag }), + inner: Arc::new(InnerObjectId { + key, + version: None, + etag, + }), + } + } + + pub fn new_with_version(key: String, version: Option, etag: ETag) -> Self { + Self { + inner: Arc::new(InnerObjectId { key, version, etag }), } } @@ -37,6 +48,10 @@ impl ObjectId { &self.inner.key } + pub fn version(&self) -> Option<&str> { + self.inner.version.as_deref() + } + pub fn etag(&self) -> &ETag { &self.inner.etag } diff --git a/mountpoint-s3-fs/src/prefetch/part_stream.rs b/mountpoint-s3-fs/src/prefetch/part_stream.rs index 3736ae4ac2..85fc69f062 100644 --- a/mountpoint-s3-fs/src/prefetch/part_stream.rs +++ b/mountpoint-s3-fs/src/prefetch/part_stream.rs @@ -437,7 +437,7 @@ fn read_from_request<'a, Client: ObjectClient + 'a>( ) -> impl Stream> + 'a { try_stream! { let mut request = client - .get_object(&bucket, id.key(), &GetObjectParams::new().range(Some(request_range.clone())).if_match(Some(id.etag().clone())).custom_id(Some(handle_id.as_raw()))) + .get_object(&bucket, id.key(), id.version(), &GetObjectParams::new().range(Some(request_range.clone())).if_match(Some(id.etag().clone())).custom_id(Some(handle_id.as_raw()))) .await .inspect_err(|e| error!(key=id.key(), error=?e, "GetObject request failed")) .map_err(|err| PrefetchReadError::get_request_failed(err, &bucket, id.key()))?; diff --git a/mountpoint-s3-fs/src/superblock.rs b/mountpoint-s3-fs/src/superblock.rs index ff10caacc9..b33f04f99e 100644 --- a/mountpoint-s3-fs/src/superblock.rs +++ b/mountpoint-s3-fs/src/superblock.rs @@ -118,6 +118,9 @@ struct SuperblockInner { client: OC, dir_handles: RwLock>>, next_dir_handle_id: AtomicU64, + /// Map of S3 objects (expressed by inodes in the superblock) to their version IDs. + /// If an object is not present in the map, the most recent version is assumed. + versioned_objects: Arc, Box>>>, } /// Configuration for superblock operations @@ -246,6 +249,7 @@ impl Superblock { client, next_dir_handle_id: AtomicU64::new(1), dir_handles: Default::default(), + versioned_objects: Arc::new(RwLock::new(HashMap::new())), }; Self { inner: Arc::new(inner) } } @@ -860,11 +864,19 @@ impl Metablock for Superblock { } }; + let version = matches!(mode, ReadWriteMode::Read) + .then(|| self.inner.get_inode_version(&inode)) + .transpose()? + .flatten(); let inode_lookup = Lookup::new( inode.ino(), locked_inode.stat.clone(), inode.kind(), - Some(S3Location::new(self.inner.s3_path.clone(), inode.valid_key().clone())), + Some(S3Location::new_with_object_version( + self.inner.s3_path.clone(), + version, + inode.valid_key().clone(), + )), ); (pending_upload_hook, inode_lookup) @@ -1223,6 +1235,7 @@ impl Metablock for Superblock { return Ok(()); } if is_readdirplus { + trace!(inode = ?next.inode, "remembering inode from readdirplus"); self.inner.remember(&next.inode) } dir_handle.next_offset(); @@ -1277,6 +1290,7 @@ impl Metablock for Superblock { None, None, self.inner.config.cache_config.file_ttl, + None, ), InodeKind::Directory => { InodeStat::for_directory(self.inner.mount_time, self.inner.config.cache_config.dir_ttl) @@ -1385,6 +1399,10 @@ impl Metablock for Superblock { } Ok(false) } + + async fn set_inode_version(&self, ino: InodeNo, version: Option<&str>) -> Result<(), InodeError> { + self.inner.set_inode_version(ino, version).await + } } impl SuperblockInner { @@ -1439,7 +1457,11 @@ impl SuperblockInner { let lookup = match lookup { Some(lookup) => lookup?, None => { - let remote = self.remote_lookup(parent_ino, name).await?; + let full_key = self.inode_full_key(parent_ino, &name)?; + let version = self.versioned_objects.read().unwrap().get(full_key.as_str()).cloned(); + let remote = self + .remote_lookup(parent_ino, name, version.as_ref().map(|v| v.as_ref())) + .await?; self.update_from_remote(parent_ino, name, remote)? } }; @@ -1505,6 +1527,7 @@ impl SuperblockInner { &self, parent_ino: InodeNo, name: ValidName<'_>, + version: Option<&str>, ) -> Result, InodeError> { let parent = self.get(parent_ino)?; let full_path: String = self @@ -1543,7 +1566,7 @@ impl SuperblockInner { let head_object_params = HeadObjectParams::new(); let mut file_lookup = self .client - .head_object(&self.s3_path.bucket, object_key, &head_object_params) + .head_object(&self.s3_path.bucket, object_key, version, &head_object_params) .fuse(); let mut dir_lookup = self .client @@ -1556,8 +1579,8 @@ impl SuperblockInner { select_biased! { result = file_lookup => { match result { - Ok(HeadObjectResult { size, last_modified, restore_status, etag, storage_class, .. }) => { - let stat = InodeStat::for_file(size as usize, last_modified, Some(etag.into_inner().into_boxed_str()), storage_class.as_deref(), restore_status, self.config.cache_config.file_ttl); + Ok(HeadObjectResult { size, last_modified, restore_status, etag, storage_class, version_id, .. }) => { + let stat = InodeStat::for_file(size as usize, last_modified, Some(etag.into_inner().into_boxed_str()), storage_class.as_deref(), restore_status, self.config.cache_config.file_ttl, version_id.map(|v| v.into_boxed_str())); file_state = Some(stat); } // If the object is not found, might be a directory, so keep going @@ -1882,6 +1905,63 @@ impl SuperblockInner { Ok(inode) } + + async fn set_inode_version(&self, ino: InodeNo, version: Option<&str>) -> Result<(), InodeError> { + let inode = self + .inodes + .read() + .unwrap() + .get_inode(&ino) + .ok_or(InodeError::InodeDoesNotExist(ino))? + .clone(); + // Forcibly lookup the inode remotely to get its new stat that corresponds to the given version. + let lookup = self + .remote_lookup(inode.parent(), inode.name().try_into()?, version) + .await? + .ok_or(InodeError::InodeDoesNotExist(ino))?; + + // Now update the inode in place. + inode.set_version(version, lookup.stat)?; + + // Update the map of versioned objects. + let full_key = self.inode_full_key(inode.parent(), inode.name())?; + let mut versioned_objects = self.versioned_objects.write().unwrap(); + if let Some(version) = version { + trace!(%full_key, %version, "set versioned object in superblock"); + versioned_objects.insert(full_key.into(), version.into()); + } else { + trace!(%full_key, "removed versioned object from superblock"); + versioned_objects.remove(full_key.as_str()); + } + + Ok(()) + } + + /// Return a full key that does not contain the trailing '/'. + fn inode_full_key(&self, parent_ino: InodeNo, key: &str) -> Result { + let parent_ino = self + .inodes + .read() + .unwrap() + .get_inode(&parent_ino) + .ok_or(InodeError::InodeDoesNotExist(parent_ino))? + .clone(); + let full_path: String = self + .full_key_for_inode(&parent_ino) + .new_child(key.try_into()?, InodeKind::Directory) + .map_err(|_| InodeError::NotADirectory(parent_ino.err())) + .map(|p| p.into())?; + Ok(full_path[..(full_path.len() - 1)].to_string()) + } + + fn get_inode_version(&self, ino: &Inode) -> Result, InodeError> { + Ok(self + .versioned_objects + .read() + .unwrap() + .get(ino.valid_key().as_ref()) + .map(|v| v.to_string())) + } } /// Data from a remote object. @@ -1917,7 +1997,11 @@ impl From for InodeInformation { impl From for Lookup { fn from(val: LookedUpInode) -> Self { - let location = Some(S3Location::new(val.path.clone(), val.inode.valid_key().clone())); + let location = Some(S3Location::new_with_object_version( + val.path.clone(), + val.inode.version(), + val.inode.valid_key().clone(), + )); Lookup::new_from_info_and_loc(val.into(), location) } @@ -2172,7 +2256,7 @@ mod tests { // Grab last modified time according to mock S3 let full_key = file.s3_location().expect("should have location").full_key(); let modified_time = client - .head_object(&bucket, full_key.as_ref(), &HeadObjectParams::new()) + .head_object(&bucket, full_key.as_ref(), None, &HeadObjectParams::new()) .await .expect("object should exist") .last_modified; @@ -2409,6 +2493,7 @@ mod tests { None, None, std::time::Duration::from_secs(24 * 60 * 60), + None, ); if invalidate_stat { // For testing the expired stat cases @@ -3248,7 +3333,7 @@ mod tests { #[test] fn test_inodestat_constructors() { let ts = OffsetDateTime::UNIX_EPOCH + Duration::days(90); - let file_inodestat = InodeStat::for_file(128, ts, None, None, None, Default::default()); + let file_inodestat = InodeStat::for_file(128, ts, None, None, None, Default::default(), None); assert_eq!(file_inodestat.size, 128); assert_eq!(file_inodestat.atime, ts); assert_eq!(file_inodestat.ctime, ts); diff --git a/mountpoint-s3-fs/src/superblock/inode.rs b/mountpoint-s3-fs/src/superblock/inode.rs index d4607326cd..cdc9f81364 100644 --- a/mountpoint-s3-fs/src/superblock/inode.rs +++ b/mountpoint-s3-fs/src/superblock/inode.rs @@ -54,6 +54,14 @@ impl Deref for InodeLockedForReading<'_> { } } +/// Custom (not the latest one) S3 object version identifier. +/// If specified in [`InodeInner`], all operations (metadata, reads) performed with this inode, +/// will actually use this version ID when accessing S3 objects. +#[derive(Debug, Clone)] +pub struct S3ObjectVersion { + pub version: Box, +} + #[derive(Debug)] struct InodeInner { // Immutable inode state -- any changes to these requires a new inode @@ -103,6 +111,12 @@ impl Inode { &self.inner.valid_key } + pub fn version(&self) -> Option { + self.get_inode_state() + .ok() + .and_then(|inode| inode.version.as_ref().map(|v| v.version.to_string())) + } + /// return Inode State with read lock after checking whether the directory inode is deleted or not. pub fn get_inode_state(&self) -> Result, InodeError> { let inode_state = self.inner.sync.read().unwrap(); @@ -160,6 +174,7 @@ impl Inode { write_status: WriteStatus::Remote, kind_data: InodeKindData::default_for(InodeKind::Directory), pending_upload_hook: None, + version: None, }, ) } @@ -186,10 +201,12 @@ impl Inode { None, None, new_validity, + old_inode_state.stat.version_id.clone(), ), write_status: WriteStatus::Remote, kind_data: InodeKindData::default_for(InodeKind::File), pending_upload_hook: None, + version: None, }; Ok(Self::new(self.ino(), new_parent, new_key, prefix, new_inode_state)) @@ -221,6 +238,18 @@ impl Inode { } } + /// Set S3 object version and inode stat. + /// [None] means the most recent object version should be used. + pub fn set_version(&self, version: Option<&str>, stat: InodeStat) -> Result<(), InodeError> { + let mut mut_inode = self.get_mut_inode_state()?; + mut_inode.state.stat = stat; + mut_inode.state.version = version.map(|version| S3ObjectVersion { + version: version.into(), + }); + + Ok(()) + } + fn compute_checksum(ino: InodeNo, prefix: &Prefix, key: &str) -> Crc32c { let mut hasher = crc32c::Hasher::new(); hasher.update(ino.to_be_bytes().as_ref()); @@ -245,6 +274,7 @@ pub struct InodeState { pub write_status: WriteStatus, pub kind_data: InodeKindData, pub pending_upload_hook: Option, + pub version: Option, } impl InodeState { @@ -254,6 +284,7 @@ impl InodeState { kind_data: InodeKindData::default_for(kind), write_status, pending_upload_hook: None, + version: stat.version_id.as_ref().map(|v| S3ObjectVersion { version: v.clone() }), } } } @@ -341,9 +372,10 @@ mod tests { &superblock.inner.s3_path.prefix, InodeState { write_status: WriteStatus::Remote, - stat: InodeStat::for_file(0, OffsetDateTime::now_utc(), None, None, None, Default::default()), + stat: InodeStat::for_file(0, OffsetDateTime::now_utc(), None, None, None, Default::default(), None), kind_data: InodeKindData::File {}, pending_upload_hook: None, + version: None, }, ); superblock.inner.inodes.write().unwrap().insert(ino, inode.clone(), 5); @@ -483,10 +515,12 @@ mod tests { None, None, NEVER_EXPIRE_TTL, + None, ), write_status: WriteStatus::Remote, kind_data: InodeKindData::File {}, pending_upload_hook: None, + version: None, }), }), }; @@ -542,9 +576,18 @@ mod tests { checksum, sync: RwLock::new(InodeState { write_status: WriteStatus::LocalOpenForWriting, - stat: InodeStat::for_file(0, OffsetDateTime::UNIX_EPOCH, None, None, None, Default::default()), + stat: InodeStat::for_file( + 0, + OffsetDateTime::UNIX_EPOCH, + None, + None, + None, + Default::default(), + None, + ), kind_data: InodeKindData::File {}, pending_upload_hook: None, + version: None, }), }), }; diff --git a/mountpoint-s3-fs/src/superblock/readdir.rs b/mountpoint-s3-fs/src/superblock/readdir.rs index 7f568bb4fe..7dad2ab5fb 100644 --- a/mountpoint-s3-fs/src/superblock/readdir.rs +++ b/mountpoint-s3-fs/src/superblock/readdir.rs @@ -41,16 +41,16 @@ //! These children are listed only once, at the start of the readdir operation, and so are a //! snapshot in time of the directory. -use std::collections::VecDeque; -use std::ffi::OsString; - use super::{InodeKindData, LookedUpInode, RemoteLookup, SuperblockInner}; use crate::metablock::{InodeError, InodeKind, InodeNo, InodeStat}; use crate::superblock::ValidName; use crate::sync::atomic::{AtomicI64, Ordering}; -use crate::sync::{AsyncMutex, Mutex}; +use crate::sync::{AsyncMutex, Mutex, RwLock}; use mountpoint_s3_client::ObjectClient; -use mountpoint_s3_client::types::RestoreStatus; +use mountpoint_s3_client::types::{HeadObjectParams, HeadObjectResult, RestoreStatus}; +use std::collections::{HashMap, VecDeque}; +use std::ffi::OsString; +use std::sync::Arc; use time::OffsetDateTime; use tracing::{error, trace, warn}; @@ -106,9 +106,21 @@ impl ReaddirHandle { }; let iter = if inner.config.s3_personality.is_list_ordered() { - ReaddirIter::ordered(&inner.s3_path.bucket, &full_path, page_size, local_entries.into()) + ReaddirIter::ordered( + &inner.s3_path.bucket, + &full_path, + page_size, + local_entries.into(), + inner.versioned_objects.clone(), + ) } else { - ReaddirIter::unordered(&inner.s3_path.bucket, &full_path, page_size, local_entries.into()) + ReaddirIter::unordered( + &inner.s3_path.bucket, + &full_path, + page_size, + local_entries.into(), + inner.versioned_objects.clone(), + ) }; Ok(Self { @@ -188,6 +200,7 @@ impl ReaddirHandle { etag, storage_class, restore_status, + version_id, .. } => { let stat = InodeStat::for_file( @@ -197,6 +210,7 @@ impl ReaddirHandle { storage_class.as_deref(), *restore_status, inner.config.cache_config.file_ttl, + version_id.as_ref().map(|v| v.clone().into_boxed_str()), ); RemoteLookup { stat, @@ -234,6 +248,9 @@ enum ReaddirEntry { restore_status: Option, /// Entity tag of this object. etag: String, + /// Optional version ID of this S3 object. + /// Will be [None] if versioning is not enabled for the bucket or the most recent version is used. + version_id: Option, }, LocalInode { lookup: LookedUpInode, @@ -271,8 +288,13 @@ impl ReaddirEntry { Self::RemotePrefix { name } => { format!("directory '{name}'") } - Self::RemoteObject { name, full_key, .. } => { - format!("file '{name}' (full key {full_key:?})") + Self::RemoteObject { + name, + full_key, + version_id, + .. + } => { + format!("file '{name}' (full key {full_key:?}, version id {version_id:?})") } Self::LocalInode { lookup } => { let kind = match lookup.inode.kind() { @@ -321,12 +343,36 @@ enum ReaddirIter { } impl ReaddirIter { - fn ordered(bucket: &str, full_path: &str, page_size: usize, local_entries: VecDeque) -> Self { - Self::Ordered(ordered::ReaddirIter::new(bucket, full_path, page_size, local_entries)) + fn ordered( + bucket: &str, + full_path: &str, + page_size: usize, + local_entries: VecDeque, + versioned_objects: Arc, Box>>>, + ) -> Self { + Self::Ordered(ordered::ReaddirIter::new( + bucket, + full_path, + page_size, + local_entries, + versioned_objects, + )) } - fn unordered(bucket: &str, full_path: &str, page_size: usize, local_entries: VecDeque) -> Self { - Self::Unordered(unordered::ReaddirIter::new(bucket, full_path, page_size, local_entries)) + fn unordered( + bucket: &str, + full_path: &str, + page_size: usize, + local_entries: VecDeque, + versioned_objects: Arc, Box>>>, + ) -> Self { + Self::Unordered(unordered::ReaddirIter::new( + bucket, + full_path, + page_size, + local_entries, + versioned_objects, + )) } async fn next(&mut self, client: &impl ObjectClient) -> Result, InodeError> { @@ -362,10 +408,17 @@ struct RemoteIter { state: RemoteIterState, /// Does the S3 implementation return ordered results? ordered: bool, + versioned_objects: Arc, Box>>>, } impl RemoteIter { - fn new(bucket: &str, full_path: &str, page_size: usize, ordered: bool) -> Self { + fn new( + bucket: &str, + full_path: &str, + page_size: usize, + ordered: bool, + versioned_objects: Arc, Box>>>, + ) -> Self { Self { entries: VecDeque::new(), bucket: bucket.to_owned(), @@ -373,6 +426,7 @@ impl RemoteIter { page_size, state: RemoteIterState::InProgress(None), ordered, + versioned_objects, } } @@ -411,18 +465,63 @@ impl RemoteIter { name: prefix[self.full_path.len()..prefix.len() - 1].to_owned(), }); - let objects = result - .objects - .into_iter() - .map(|object_info| ReaddirEntry::RemoteObject { - name: object_info.key[self.full_path.len()..].to_owned(), - full_key: object_info.key, - size: object_info.size, - last_modified: object_info.last_modified, - storage_class: object_info.storage_class, - restore_status: object_info.restore_status, - etag: object_info.etag, - }); + let mut objects = vec![]; + for object_info in result.objects { + let version = self + .versioned_objects + .read() + .unwrap() + .get(object_info.key.as_str()) + .cloned(); + + // `ListObjectsV2` returns metadata only for latest versions of objects, so if a version + // id is specified we need to make a separate `HeadObject` call to get the correct metadata. + let entry = if let Some(version) = version { + let HeadObjectResult { + size, + last_modified, + etag, + storage_class, + restore_status, + version_id, + .. + } = client + .head_object(&self.bucket, &object_info.key, Some(&version), &HeadObjectParams::new()) + .await + .map_err(|e| { + InodeError::client_error( + e, + "HeadObject failed for versioned object", + &self.bucket, + &object_info.key, + ) + })?; + + ReaddirEntry::RemoteObject { + name: object_info.key[self.full_path.len()..].to_owned(), + full_key: object_info.key, + size, + last_modified, + storage_class, + restore_status, + etag: etag.into_inner(), + version_id, + } + } else { + ReaddirEntry::RemoteObject { + name: object_info.key[self.full_path.len()..].to_owned(), + full_key: object_info.key, + size: object_info.size, + last_modified: object_info.last_modified, + storage_class: object_info.storage_class, + restore_status: object_info.restore_status, + etag: object_info.etag, + version_id: None, + } + }; + + objects.push(entry); + } if self.ordered { // ListObjectsV2 results are sorted, so ideally we'd just merge-sort the two streams. @@ -466,9 +565,10 @@ mod ordered { full_path: &str, page_size: usize, local_entries: VecDeque, + versioned_objects: Arc, Box>>>, ) -> Self { Self { - remote: RemoteIter::new(bucket, full_path, page_size, true), + remote: RemoteIter::new(bucket, full_path, page_size, true, versioned_objects), local: LocalIter::new(local_entries), next_remote: None, next_local: None, @@ -573,6 +673,7 @@ mod unordered { full_path: &str, page_size: usize, local_entries: VecDeque, + versioned_objects: Arc, Box>>>, ) -> Self { let local_map = local_entries .into_iter() @@ -585,7 +686,7 @@ mod unordered { .collect::>(); Self { - remote: RemoteIter::new(bucket, full_path, page_size, false), + remote: RemoteIter::new(bucket, full_path, page_size, false, versioned_objects), local: local_map, local_iter: VecDeque::new(), } diff --git a/mountpoint-s3-fs/src/upload/incremental.rs b/mountpoint-s3-fs/src/upload/incremental.rs index 7dcb9501a5..0f3c4d15ab 100644 --- a/mountpoint-s3-fs/src/upload/incremental.rs +++ b/mountpoint-s3-fs/src/upload/incremental.rs @@ -335,6 +335,7 @@ where .head_object( ¶ms.bucket, ¶ms.key, + None, &HeadObjectParams::new().checksum_mode(Some(ChecksumMode::Enabled)), ) .await?; @@ -632,7 +633,7 @@ mod tests { // Verify content of the object let get_request = client - .get_object(bucket, key, &GetObjectParams::default()) + .get_object(bucket, key, None, &GetObjectParams::default()) .await .expect("get_object failed"); let actual = get_request.collect().await.expect("failed to collect body"); @@ -698,7 +699,7 @@ mod tests { // Verify content of the object let get_request = client - .get_object(bucket, key, &GetObjectParams::default()) + .get_object(bucket, key, None, &GetObjectParams::default()) .await .expect("get_object failed"); let actual = get_request.collect().await.expect("failed to collect body"); @@ -767,6 +768,7 @@ mod tests { .get_object( bucket, key, + None, &GetObjectParams::default().checksum_mode(Some(ChecksumMode::Enabled)), ) .await @@ -812,7 +814,7 @@ mod tests { // Verify content of the object let get_request = client - .get_object(bucket, key, &GetObjectParams::default()) + .get_object(bucket, key, None, &GetObjectParams::default()) .await .expect("get_object failed"); let actual = get_request.collect().await.expect("failed to collect body"); @@ -894,7 +896,7 @@ mod tests { // Verify that object is partially appended from the first request let get_request = client - .get_object(bucket, key, &GetObjectParams::default()) + .get_object(bucket, key, None, &GetObjectParams::default()) .await .expect("get_object failed"); let actual = get_request.collect().await.expect("failed to collect body"); @@ -1018,7 +1020,7 @@ mod tests { // Verify that object is partially appended from the first request let get_request = client - .get_object(bucket, key, &GetObjectParams::default()) + .get_object(bucket, key, None, &GetObjectParams::default()) .await .expect("get_object failed"); let actual = get_request.collect().await.expect("failed to collect body"); @@ -1205,7 +1207,7 @@ mod tests { // Verify content of the object let get_request = client - .get_object(bucket, key, &GetObjectParams::default()) + .get_object(bucket, key, None, &GetObjectParams::default()) .await .expect("get_object failed"); let actual = get_request.collect().await.expect("failed to collect body"); @@ -1254,7 +1256,7 @@ mod tests { // Verify content of the object let get_request = client - .get_object(bucket, key, &GetObjectParams::default()) + .get_object(bucket, key, None, &GetObjectParams::default()) .await .expect("get_object failed"); let actual = get_request.collect().await.expect("failed to collect body"); diff --git a/mountpoint-s3-fs/tests/fs.rs b/mountpoint-s3-fs/tests/fs.rs index 0cde0587ba..1cd6f71de8 100644 --- a/mountpoint-s3-fs/tests/fs.rs +++ b/mountpoint-s3-fs/tests/fs.rs @@ -570,7 +570,7 @@ async fn test_sequential_write(write_size: usize) { // Check that the object made it to S3 as we expected let get = client - .get_object(BUCKET_NAME, "dir1/file2.bin", &GetObjectParams::new()) + .get_object(BUCKET_NAME, "dir1/file2.bin", None, &GetObjectParams::new()) .await .unwrap(); let actual = get.collect().await.unwrap(); diff --git a/mountpoint-s3-ioctl/Cargo.toml b/mountpoint-s3-ioctl/Cargo.toml new file mode 100644 index 0000000000..3cb1637bf8 --- /dev/null +++ b/mountpoint-s3-ioctl/Cargo.toml @@ -0,0 +1,15 @@ +[package] +name = "mountpoint-s3-ioctl" +version = "0.1.0" +edition = "2024" +license = "Apache-2.0" +publish = false + +[[bin]] +name = "mount-s3-ioctl" +path = "src/main.rs" + +[dependencies] +clap = { version = "4.5.53", features = ["derive"] } +nix = { version = "0.30.1", features = ["ioctl"] } +thiserror = "2.0.17" \ No newline at end of file diff --git a/mountpoint-s3-ioctl/src/error.rs b/mountpoint-s3-ioctl/src/error.rs new file mode 100644 index 0000000000..730272b216 --- /dev/null +++ b/mountpoint-s3-ioctl/src/error.rs @@ -0,0 +1,5 @@ +#[derive(Debug, thiserror::Error)] +pub enum Error { + #[error("Version ID is too long ({provided}) to fit in the buffer ({max})")] + VersionIdTooLong { provided: usize, max: usize }, +} diff --git a/mountpoint-s3-ioctl/src/lib.rs b/mountpoint-s3-ioctl/src/lib.rs new file mode 100644 index 0000000000..d9d3f33296 --- /dev/null +++ b/mountpoint-s3-ioctl/src/lib.rs @@ -0,0 +1,41 @@ +mod error; + +use nix::sys::ioctl::ioctl_num_type; +use nix::{ioctl_write_buf, request_code_write}; + +pub struct S3ObjectVersionBuffer { + pub data: [u8; 64], +} + +impl TryFrom<&str> for S3ObjectVersionBuffer { + type Error = error::Error; + + fn try_from(value: &str) -> Result { + let mut buffer = S3ObjectVersionBuffer { data: [0; 64] }; + + if value.len() > buffer.data.len() { + return Err(error::Error::VersionIdTooLong { + max: buffer.data.len(), + provided: value.len(), + }); + } + + let bytes = value.as_bytes(); + let len = bytes.len().min(buffer.data.len()); + buffer.data[..len].copy_from_slice(&bytes[..len]); + Ok(buffer) + } +} + +const MOUNT_S3_IOC_MAGIC: u8 = b'm'; + +/// IOCTL to set the S3 object version ID for an inode. Once set, +/// this version id will be used for all S3 operations on the inode. +pub const MOUNT_S3_IOC_TYPE_SET_VERSION: ioctl_num_type = + request_code_write!(MOUNT_S3_IOC_MAGIC, 0, size_of::()); +ioctl_write_buf!( + ioctl_mount_s3_set_inode_version, + MOUNT_S3_IOC_MAGIC, + 0, + S3ObjectVersionBuffer +); diff --git a/mountpoint-s3-ioctl/src/main.rs b/mountpoint-s3-ioctl/src/main.rs new file mode 100644 index 0000000000..5fae9428e3 --- /dev/null +++ b/mountpoint-s3-ioctl/src/main.rs @@ -0,0 +1,42 @@ +use clap::Parser; +use mountpoint_s3_ioctl::{S3ObjectVersionBuffer, ioctl_mount_s3_set_inode_version}; +use std::os::fd::AsRawFd; + +#[derive(Parser)] +enum Args { + /// Set the S3 object version ID for the given file. + /// This version ID will be used for all S3 operations on the file. + /// If the version ID is an empty string, any existing version information will be removed and + /// the latest version will be used. + SetVersion { + /// Path to the file on the mountpoint handled by `mountpoint-s3`. + file: String, + /// S3 object version ID to set for the file. + version: String, + }, +} + +fn main() { + let args = Args::parse(); + + match args { + Args::SetVersion { file, version } => { + set_inode_version(&file, &version); + } + } +} + +fn set_inode_version(file: &str, version: &str) { + let fd = std::fs::File::open(file).unwrap_or_else(|e| { + eprintln!("Failed to open file '{}': {}", file, e); + std::process::exit(1); + }); + let version_buffer = S3ObjectVersionBuffer::try_from(version).unwrap(); + + let result = unsafe { ioctl_mount_s3_set_inode_version(fd.as_raw_fd(), std::slice::from_ref(&version_buffer)) }; + + if let Err(e) = result { + eprintln!("Failed to set S3 object version for file '{}': {}", file, e); + std::process::exit(1); + } +} From 7baa0e263cc07edb531927df9fb33e6a65ceb1d2 Mon Sep 17 00:00:00 2001 From: zoryamba-elastio Date: Wed, 29 Jul 2026 20:45:30 +0300 Subject: [PATCH 2/3] fix rebase --- Cargo.lock | 12 ++++++++++++ mountpoint-s3-crt-sys/crt/aws-c-compression | 2 +- mountpoint-s3-crt-sys/crt/aws-c-sdkutils | 2 +- mountpoint-s3-fs/src/fs.rs | 3 ++- 4 files changed, 16 insertions(+), 3 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 316fec8bbf..57eba5115a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2989,6 +2989,18 @@ dependencies = [ "libc", ] +[[package]] +name = "nix" +version = "0.30.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "74523f3a35e05aba87a1d978330aef40f67b0304ac79c1c00b294c9830543db6" +dependencies = [ + "bitflags 2.13.1", + "cfg-if", + "cfg_aliases", + "libc", +] + [[package]] name = "nix" version = "0.31.3" diff --git a/mountpoint-s3-crt-sys/crt/aws-c-compression b/mountpoint-s3-crt-sys/crt/aws-c-compression index f951ab2b81..48d8a47d20 160000 --- a/mountpoint-s3-crt-sys/crt/aws-c-compression +++ b/mountpoint-s3-crt-sys/crt/aws-c-compression @@ -1 +1 @@ -Subproject commit f951ab2b819fc6993b6e5e6cfef64b1a1554bfc8 +Subproject commit 48d8a47d20f6a1ced4f4f2325803c95b494d02c7 diff --git a/mountpoint-s3-crt-sys/crt/aws-c-sdkutils b/mountpoint-s3-crt-sys/crt/aws-c-sdkutils index cb14fea362..6b8f56012e 160000 --- a/mountpoint-s3-crt-sys/crt/aws-c-sdkutils +++ b/mountpoint-s3-crt-sys/crt/aws-c-sdkutils @@ -1 +1 @@ -Subproject commit cb14fea362c82c995eebd34e2e96590ab4e0ed58 +Subproject commit 6b8f56012e6420b6a942c9c19c91db802e24507c diff --git a/mountpoint-s3-fs/src/fs.rs b/mountpoint-s3-fs/src/fs.rs index 4317f174ea..6b44f73450 100644 --- a/mountpoint-s3-fs/src/fs.rs +++ b/mountpoint-s3-fs/src/fs.rs @@ -802,9 +802,10 @@ where } trace!(pid = handle.open_pid, "Recreating read handle with new inode version."); // Recreate the read handle with the new version. + let fh = self.next_handle(); let lookup = self.metablock.getattr(ino, false).await?; let new_handle = NewHandle::read(lookup.clone()); - let handle_state = FileHandleState::new(&new_handle, OpenFlags::empty(), self).await?; + let handle_state = FileHandleState::new(fh, &new_handle, OpenFlags::empty(), self).await?; let handle = Arc::new(FileHandle { ino, open_pid: handle.open_pid, From 95c5038ff70978eaacc97b21ac028e3447d0f539 Mon Sep 17 00:00:00 2001 From: zoryamba-elastio Date: Thu, 30 Jul 2026 21:35:44 +0300 Subject: [PATCH 3/3] update submodule --- .gitmodules | 2 +- mountpoint-s3-crt-sys/crt/aws-c-s3 | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/.gitmodules b/.gitmodules index 687383c75a..5aa8cf24a7 100644 --- a/.gitmodules +++ b/.gitmodules @@ -30,4 +30,4 @@ url = https://github.com/awslabs/aws-c-auth.git [submodule "aws-c-s3"] path = mountpoint-s3-crt-sys/crt/aws-c-s3 - url = https://github.com/awslabs/aws-c-s3.git + url = https://github.com/elastio/aws-c-s3.git diff --git a/mountpoint-s3-crt-sys/crt/aws-c-s3 b/mountpoint-s3-crt-sys/crt/aws-c-s3 index 1e9c52ca20..8dbbb2d466 160000 --- a/mountpoint-s3-crt-sys/crt/aws-c-s3 +++ b/mountpoint-s3-crt-sys/crt/aws-c-s3 @@ -1 +1 @@ -Subproject commit 1e9c52ca2059ab72f8c066d12960a5e33a7ac043 +Subproject commit 8dbbb2d46670da65e44e1a2296ad1f4ff3c4417d