Skip to main content

rc_core/
traits.rs

1//! ObjectStore trait definition
2//!
3//! This trait defines the interface for S3-compatible storage operations.
4//! It allows the CLI to be decoupled from the specific S3 SDK implementation.
5
6use std::collections::HashMap;
7
8use async_trait::async_trait;
9use jiff::Timestamp;
10use serde::{Deserialize, Serialize};
11use tokio::io::{AsyncWrite, AsyncWriteExt};
12
13use crate::cors::CorsRule;
14use crate::encryption::{BucketEncryption, ObjectEncryptionRequest};
15use crate::error::{Error, Result};
16use crate::lifecycle::LifecycleRule;
17use crate::multipart_copy::{
18    MultipartCopyCancellation, MultipartCopyOptions, MultipartCopyProgress, MultipartCopyResult,
19};
20use crate::object_lock::{
21    BucketObjectLockConfiguration, LegalHoldStatus, ObjectLockOptions, ObjectRetention,
22};
23use crate::path::RemotePath;
24use crate::replication::{
25    ReplicationCheckResult, ReplicationConfiguration, ReplicationResyncStartOptions,
26    ReplicationResyncStartResult, ReplicationResyncStatus,
27};
28use crate::select::SelectOptions;
29use crate::transfer_options::{
30    ObjectTransferMetadata, ObjectWriteOptions, TransferCopyOptions, TransferReadOptions,
31};
32
33/// Requested behavior for bucket creation.
34#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
35pub struct CreateBucketOptions {
36    /// Explicit S3 location constraint. `None` omits the request body.
37    pub region: Option<String>,
38    /// Whether the resulting bucket must have versioning enabled.
39    pub versioning_enabled: bool,
40    /// Whether Object Lock must be enabled in the create request.
41    pub object_lock_enabled: bool,
42}
43
44impl CreateBucketOptions {
45    /// Build CLI options while applying Object Lock's required versioning invariant.
46    pub fn for_cli(
47        region: Option<String>,
48        versioning_enabled: bool,
49        object_lock_enabled: bool,
50    ) -> Result<Self> {
51        let options = Self {
52            region,
53            versioning_enabled: versioning_enabled || object_lock_enabled,
54            object_lock_enabled,
55        };
56        options.validate()?;
57        Ok(options)
58    }
59
60    /// Reject request states that cannot produce the promised bucket state.
61    pub fn validate(&self) -> Result<()> {
62        if self.object_lock_enabled && !self.versioning_enabled {
63            return Err(Error::InvalidPath(
64                "Bucket Object Lock requires versioning to be enabled".to_string(),
65            ));
66        }
67        if self
68            .region
69            .as_deref()
70            .is_some_and(|region| region.trim().is_empty())
71        {
72            return Err(Error::InvalidPath(
73                "Bucket region cannot be empty".to_string(),
74            ));
75        }
76        Ok(())
77    }
78}
79
80/// Metadata for an object version
81#[derive(Debug, Clone, Serialize, Deserialize)]
82pub struct ObjectVersion {
83    /// Object key
84    pub key: String,
85
86    /// Version ID
87    pub version_id: String,
88
89    /// Whether this is the latest version
90    pub is_latest: bool,
91
92    /// Whether this is a delete marker
93    pub is_delete_marker: bool,
94
95    /// Last modified timestamp
96    #[serde(skip_serializing_if = "Option::is_none")]
97    pub last_modified: Option<Timestamp>,
98
99    /// Size in bytes
100    #[serde(skip_serializing_if = "Option::is_none")]
101    pub size_bytes: Option<i64>,
102
103    /// ETag
104    #[serde(skip_serializing_if = "Option::is_none")]
105    pub etag: Option<String>,
106}
107
108/// Result of an object version list operation
109#[derive(Debug, Clone, Serialize, Deserialize)]
110pub struct ObjectVersionListResult {
111    /// Listed object versions and delete markers
112    pub items: Vec<ObjectVersion>,
113
114    /// Whether the result is truncated (more items available)
115    pub truncated: bool,
116
117    /// Continuation key marker for pagination
118    #[serde(skip_serializing_if = "Option::is_none")]
119    pub continuation_token: Option<String>,
120
121    /// Continuation version marker for pagination
122    #[serde(skip_serializing_if = "Option::is_none")]
123    pub version_id_marker: Option<String>,
124}
125
126/// Options for selecting an object version during read and metadata operations.
127#[derive(Debug, Clone, Default, PartialEq, Eq)]
128pub struct ObjectReadOptions {
129    /// Exact object version to select. `None` selects the current object.
130    pub version_id: Option<String>,
131}
132
133/// Request options for a server-side object copy.
134#[derive(Debug, Clone, Default, PartialEq, Eq)]
135pub struct CopyObjectOptions {
136    /// Exact historical source version to copy. `None` selects the current source.
137    pub source_version_id: Option<String>,
138    /// Copy only when the selected source still has this ETag.
139    pub source_etag: Option<String>,
140}
141
142impl CopyObjectOptions {
143    /// Build copy options while rejecting an ambiguous empty source version identifier.
144    pub fn for_source_version(source_version_id: Option<String>) -> Result<Self> {
145        Self::for_source_identity(source_version_id, None)
146    }
147
148    /// Build copy options pinned to an exact source version or ETag.
149    pub fn for_source_identity(
150        source_version_id: Option<String>,
151        source_etag: Option<String>,
152    ) -> Result<Self> {
153        if source_version_id.as_deref().is_some_and(str::is_empty) {
154            return Err(Error::InvalidPath(
155                "Source version ID cannot be empty".to_string(),
156            ));
157        }
158        if source_etag.as_deref().is_some_and(str::is_empty) {
159            return Err(Error::InvalidPath(
160                "Source ETag cannot be empty".to_string(),
161            ));
162        }
163        Ok(Self {
164            source_version_id,
165            source_etag,
166        })
167    }
168}
169
170impl ObjectReadOptions {
171    /// Build read options while rejecting ambiguous empty version identifiers.
172    pub fn for_version(version_id: Option<String>) -> Result<Self> {
173        if version_id.as_deref().is_some_and(str::is_empty) {
174            return Err(Error::InvalidPath("Version ID cannot be empty".to_string()));
175        }
176        Ok(Self { version_id })
177    }
178}
179
180/// Pagination options for listing object versions and delete markers.
181#[derive(Debug, Clone, Default, PartialEq, Eq)]
182pub struct ListObjectVersionsOptions {
183    /// Maximum number of entries to return.
184    pub max_keys: Option<i32>,
185    /// Key marker returned by the previous page.
186    pub key_marker: Option<String>,
187    /// Version marker returned by the previous page.
188    pub version_id_marker: Option<String>,
189}
190
191/// Request-level options for object deletion.
192#[derive(Debug, Clone, Default, PartialEq, Eq)]
193pub struct DeleteRequestOptions {
194    /// Exact object version to delete. `None` targets the current object state.
195    pub version_id: Option<String>,
196    /// Explicitly bypass Object Lock governance retention.
197    pub bypass_governance: bool,
198    /// Ask RustFS to permanently delete data instead of creating delete markers.
199    pub force_delete: bool,
200}
201
202/// An object key and optional historical version selected for deletion.
203#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
204pub struct ObjectVersionIdentifier {
205    /// Object key.
206    pub key: String,
207    /// Exact version to delete, when version-aware deletion is requested.
208    #[serde(skip_serializing_if = "Option::is_none")]
209    pub version_id: Option<String>,
210    /// Whether the selected version was listed as a delete marker.
211    pub is_delete_marker: bool,
212}
213
214/// A version-aware delete result returned by the object store.
215#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
216pub struct DeletedObject {
217    /// Deleted object key.
218    pub key: String,
219    /// Deleted object version, when reported by the backend.
220    #[serde(skip_serializing_if = "Option::is_none")]
221    pub version_id: Option<String>,
222    /// Whether the deleted entry is a delete marker.
223    pub is_delete_marker: bool,
224}
225
226/// A per-object failure returned by a multi-object delete request.
227#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
228pub struct DeleteObjectFailure {
229    /// Object key that could not be deleted.
230    pub key: String,
231    /// Requested version, when present.
232    #[serde(skip_serializing_if = "Option::is_none")]
233    pub version_id: Option<String>,
234    /// S3 error code, when provided.
235    #[serde(skip_serializing_if = "Option::is_none")]
236    pub code: Option<String>,
237    /// Backend error message, when provided.
238    #[serde(skip_serializing_if = "Option::is_none")]
239    pub message: Option<String>,
240}
241
242/// Result of deleting multiple version-aware object identifiers.
243#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
244pub struct DeleteObjectsResult {
245    /// Successfully deleted entries.
246    pub deleted: Vec<DeletedObject>,
247    /// Entries rejected by the backend.
248    pub failures: Vec<DeleteObjectFailure>,
249}
250
251/// Metadata for an object or bucket
252#[derive(Debug, Clone, Serialize, Deserialize)]
253pub struct ObjectInfo {
254    /// Object key or bucket name
255    pub key: String,
256
257    /// Size in bytes (None for buckets)
258    #[serde(skip_serializing_if = "Option::is_none")]
259    pub size_bytes: Option<i64>,
260
261    /// Human-readable size
262    #[serde(skip_serializing_if = "Option::is_none")]
263    pub size_human: Option<String>,
264
265    /// Last modified timestamp
266    #[serde(skip_serializing_if = "Option::is_none")]
267    pub last_modified: Option<Timestamp>,
268
269    /// ETag (usually MD5 for single-part uploads)
270    #[serde(skip_serializing_if = "Option::is_none")]
271    pub etag: Option<String>,
272
273    /// Storage class
274    #[serde(skip_serializing_if = "Option::is_none")]
275    pub storage_class: Option<String>,
276
277    /// Content type
278    #[serde(skip_serializing_if = "Option::is_none")]
279    pub content_type: Option<String>,
280
281    /// User-defined metadata
282    #[serde(skip_serializing_if = "Option::is_none")]
283    pub metadata: Option<HashMap<String, String>>,
284
285    /// Object version selected or created by the operation.
286    #[serde(skip_serializing_if = "Option::is_none")]
287    pub version_id: Option<String>,
288
289    /// Source object version used by a copy operation.
290    #[serde(skip_serializing_if = "Option::is_none")]
291    pub source_version_id: Option<String>,
292
293    /// Whether the selected version is a delete marker.
294    #[serde(skip_serializing_if = "Option::is_none")]
295    pub is_delete_marker: Option<bool>,
296
297    /// Whether this is a directory/prefix
298    pub is_dir: bool,
299}
300
301impl ObjectInfo {
302    /// Create a new ObjectInfo for a file
303    pub fn file(key: impl Into<String>, size: i64) -> Self {
304        Self {
305            key: key.into(),
306            size_bytes: Some(size),
307            size_human: Some(humansize::format_size(size as u64, humansize::BINARY)),
308            last_modified: None,
309            etag: None,
310            storage_class: None,
311            content_type: None,
312            metadata: None,
313            version_id: None,
314            source_version_id: None,
315            is_delete_marker: None,
316            is_dir: false,
317        }
318    }
319
320    /// Create a new ObjectInfo for a directory/prefix
321    pub fn dir(key: impl Into<String>) -> Self {
322        Self {
323            key: key.into(),
324            size_bytes: None,
325            size_human: None,
326            last_modified: None,
327            etag: None,
328            storage_class: None,
329            content_type: None,
330            metadata: None,
331            version_id: None,
332            source_version_id: None,
333            is_delete_marker: None,
334            is_dir: true,
335        }
336    }
337
338    /// Create a new ObjectInfo for a bucket
339    pub fn bucket(name: impl Into<String>) -> Self {
340        Self {
341            key: name.into(),
342            size_bytes: None,
343            size_human: None,
344            last_modified: None,
345            etag: None,
346            storage_class: None,
347            content_type: None,
348            metadata: None,
349            version_id: None,
350            source_version_id: None,
351            is_delete_marker: None,
352            is_dir: true,
353        }
354    }
355}
356
357/// Result of a list operation
358#[derive(Debug, Clone, Serialize, Deserialize)]
359pub struct ListResult {
360    /// Listed objects
361    pub items: Vec<ObjectInfo>,
362
363    /// Whether the result is truncated (more items available)
364    pub truncated: bool,
365
366    /// Continuation token for pagination
367    #[serde(skip_serializing_if = "Option::is_none")]
368    pub continuation_token: Option<String>,
369}
370
371/// Options for list operations
372#[derive(Debug, Clone, Default)]
373pub struct ListOptions {
374    /// Maximum number of keys to return per request
375    pub max_keys: Option<i32>,
376
377    /// Delimiter for grouping (usually "/")
378    pub delimiter: Option<String>,
379
380    /// Prefix to filter by
381    pub prefix: Option<String>,
382
383    /// Continuation token for pagination
384    pub continuation_token: Option<String>,
385
386    /// Whether to list recursively (ignore delimiter)
387    pub recursive: bool,
388}
389
390/// S3 identity attached to an incomplete multipart upload.
391#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
392pub struct MultipartIdentity {
393    /// Stable account or principal identifier, when returned by the server.
394    #[serde(skip_serializing_if = "Option::is_none")]
395    pub id: Option<String>,
396
397    /// Human-readable principal name, when returned by the server.
398    #[serde(skip_serializing_if = "Option::is_none")]
399    pub display_name: Option<String>,
400}
401
402/// Metadata for an incomplete multipart upload.
403///
404/// `size_bytes` is nullable because the S3 list-multipart-uploads response does not expose the
405/// uploaded part total. Keeping the field in the typed record makes the representation compatible
406/// with richer backends and the output v3 mapping without issuing one request per upload.
407#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
408pub struct MultipartUpload {
409    /// Bucket containing the pending upload.
410    pub bucket: String,
411
412    /// Object key being uploaded.
413    pub key: String,
414
415    /// Server-assigned multipart upload identifier.
416    pub upload_id: String,
417
418    /// Time at which the upload was initiated, when returned by the server.
419    pub initiated: Option<Timestamp>,
420
421    /// Bytes uploaded so far, when the backend provides this value.
422    pub size_bytes: Option<i64>,
423
424    /// Requested storage class, when returned by the server.
425    pub storage_class: Option<String>,
426
427    /// Principal that initiated the upload.
428    #[serde(skip_serializing_if = "Option::is_none")]
429    pub initiator: Option<MultipartIdentity>,
430
431    /// Owner of the pending object.
432    #[serde(skip_serializing_if = "Option::is_none")]
433    pub owner: Option<MultipartIdentity>,
434
435    /// Checksum algorithm selected when the upload was created.
436    #[serde(skip_serializing_if = "Option::is_none")]
437    pub checksum_algorithm: Option<String>,
438
439    /// Checksum aggregation type selected when the upload was created.
440    #[serde(skip_serializing_if = "Option::is_none")]
441    pub checksum_type: Option<String>,
442}
443
444/// Options for one page of incomplete multipart uploads.
445#[derive(Debug, Clone, Default, PartialEq, Eq)]
446pub struct MultipartUploadListOptions {
447    /// Only return object keys beginning with this prefix.
448    pub prefix: Option<String>,
449
450    /// Group keys containing this delimiter into common prefixes.
451    pub delimiter: Option<String>,
452
453    /// Key marker returned by the previous page.
454    pub key_marker: Option<String>,
455
456    /// Upload ID marker returned by the previous page.
457    pub upload_id_marker: Option<String>,
458
459    /// Maximum number of uploads and common prefixes to return.
460    pub max_uploads: Option<i32>,
461}
462
463/// Result of one incomplete multipart upload listing page.
464#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
465pub struct MultipartUploadListResult {
466    /// Uploads returned on this page.
467    pub uploads: Vec<MultipartUpload>,
468
469    /// Grouped prefixes returned when a delimiter was supplied.
470    #[serde(default, skip_serializing_if = "Vec::is_empty")]
471    pub common_prefixes: Vec<String>,
472
473    /// Whether another page is available.
474    pub truncated: bool,
475
476    /// Key marker for the next page.
477    #[serde(skip_serializing_if = "Option::is_none")]
478    pub next_key_marker: Option<String>,
479
480    /// Upload ID marker for the next page.
481    #[serde(skip_serializing_if = "Option::is_none")]
482    pub next_upload_id_marker: Option<String>,
483}
484
485/// Identifies one incomplete multipart upload to abort.
486#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
487pub struct AbortMultipartUploadRequest {
488    /// Bucket containing the upload.
489    pub bucket: String,
490
491    /// Object key being uploaded.
492    pub key: String,
493
494    /// Server-assigned multipart upload identifier.
495    pub upload_id: String,
496}
497
498/// Backend capability information
499#[derive(Debug, Clone, Default)]
500pub struct Capabilities {
501    /// Supports bucket versioning
502    pub versioning: bool,
503
504    /// Supports object lock/retention
505    pub object_lock: bool,
506
507    /// Supports object tagging
508    pub tagging: bool,
509
510    /// Supports anonymous bucket access policies
511    pub anonymous: bool,
512
513    /// S3 Select (`SelectObjectContent`).
514    ///
515    /// This remains `false` in generic capability hints because support is determined by issuing
516    /// a real request against the target object.
517    pub select: bool,
518
519    /// Supports event notifications
520    pub notifications: bool,
521
522    /// Supports lifecycle configuration
523    pub lifecycle: bool,
524
525    /// Supports bucket replication
526    pub replication: bool,
527
528    /// Supports bucket CORS configuration
529    pub cors: bool,
530}
531
532/// Bucket notification target type
533#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
534#[serde(rename_all = "lowercase")]
535pub enum NotificationTarget {
536    /// SQS queue target
537    Queue,
538    /// SNS topic target
539    Topic,
540    /// Lambda function target
541    Lambda,
542}
543
544/// Bucket notification rule
545#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
546pub struct BucketNotification {
547    /// Optional rule id
548    #[serde(skip_serializing_if = "Option::is_none")]
549    pub id: Option<String>,
550    /// Notification target type
551    pub target: NotificationTarget,
552    /// Target ARN
553    pub arn: String,
554    /// Event patterns
555    pub events: Vec<String>,
556    /// Optional key prefix filter
557    #[serde(skip_serializing_if = "Option::is_none")]
558    pub prefix: Option<String>,
559    /// Optional key suffix filter
560    #[serde(skip_serializing_if = "Option::is_none")]
561    pub suffix: Option<String>,
562}
563
564/// Trait for S3-compatible storage operations
565///
566/// This trait is implemented by the S3 adapter and can be mocked for testing.
567#[async_trait]
568pub trait ObjectStore: Send + Sync {
569    /// List buckets
570    async fn list_buckets(&self) -> Result<Vec<ObjectInfo>>;
571
572    /// List objects in a bucket or prefix
573    async fn list_objects(&self, path: &RemotePath, options: ListOptions) -> Result<ListResult>;
574
575    /// List one page of incomplete multipart uploads.
576    async fn list_multipart_uploads(
577        &self,
578        bucket: &str,
579        options: MultipartUploadListOptions,
580    ) -> Result<MultipartUploadListResult> {
581        let _ = (bucket, options);
582        Err(Error::UnsupportedFeature(
583            "Multipart upload listing is not supported by this object store".to_string(),
584        ))
585    }
586
587    /// Abort one incomplete multipart upload.
588    async fn abort_multipart_upload(&self, request: &AbortMultipartUploadRequest) -> Result<()> {
589        let _ = request;
590        Err(Error::UnsupportedFeature(
591            "Multipart upload cleanup is not supported by this object store".to_string(),
592        ))
593    }
594
595    /// Get object metadata
596    async fn head_object(&self, path: &RemotePath) -> Result<ObjectInfo>;
597
598    /// Get metadata for the current object or an exact historical version.
599    async fn head_object_with_options(
600        &self,
601        path: &RemotePath,
602        options: &ObjectReadOptions,
603    ) -> Result<ObjectInfo> {
604        if options.version_id.is_some() {
605            return Err(Error::UnsupportedFeature(
606                "Exact-version metadata reads are not implemented by this object store".to_string(),
607            ));
608        }
609        self.head_object(path).await
610    }
611
612    /// Get metadata with checksum-mode and SSE-C support when implemented by the backend.
613    ///
614    /// The default preserves legacy version selection and rejects advanced options explicitly.
615    async fn head_object_with_transfer_options(
616        &self,
617        path: &RemotePath,
618        options: &TransferReadOptions,
619    ) -> Result<ObjectInfo> {
620        let legacy = options.legacy_read_options()?;
621        self.head_object_with_options(path, &legacy).await
622    }
623
624    /// Read complete transfer metadata without changing ObjectInfo's stable output contract.
625    async fn head_object_transfer_metadata(
626        &self,
627        _path: &RemotePath,
628        options: &TransferReadOptions,
629    ) -> Result<ObjectTransferMetadata> {
630        options.validate()?;
631        Err(Error::UnsupportedFeature(
632            "Complete transfer metadata is not implemented by this object store".to_string(),
633        ))
634    }
635
636    /// Check if a bucket exists
637    async fn bucket_exists(&self, bucket: &str) -> Result<bool>;
638
639    /// Create a bucket
640    async fn create_bucket(&self, bucket: &str) -> Result<()>;
641
642    /// Create a bucket with explicit region, versioning, and Object Lock intent.
643    ///
644    /// The default preserves existing implementations for the option-free request and rejects
645    /// advanced behavior instead of silently ignoring it.
646    async fn create_bucket_with_options(
647        &self,
648        bucket: &str,
649        options: &CreateBucketOptions,
650    ) -> Result<()> {
651        options.validate()?;
652        if options != &CreateBucketOptions::default() {
653            return Err(Error::UnsupportedFeature(
654                "Bucket creation options are not implemented by this object store".to_string(),
655            ));
656        }
657        self.create_bucket(bucket).await
658    }
659
660    /// Return the effective location reported by the service.
661    ///
662    /// `None` is the S3 representation for the default `us-east-1` location.
663    async fn get_bucket_location(&self, _bucket: &str) -> Result<Option<String>> {
664        Err(Error::UnsupportedFeature(
665            "Bucket location inspection is not implemented by this object store".to_string(),
666        ))
667    }
668
669    /// Delete a bucket
670    async fn delete_bucket(&self, bucket: &str) -> Result<()>;
671
672    /// Get backend capabilities
673    async fn capabilities(&self) -> Result<Capabilities>;
674
675    /// Get object content as bytes
676    async fn get_object(&self, path: &RemotePath) -> Result<Vec<u8>>;
677
678    /// Get current object content or an exact historical version as bytes.
679    async fn get_object_with_options(
680        &self,
681        path: &RemotePath,
682        options: &ObjectReadOptions,
683    ) -> Result<Vec<u8>> {
684        if options.version_id.is_some() {
685            return Err(Error::UnsupportedFeature(
686                "Exact-version object reads are not implemented by this object store".to_string(),
687            ));
688        }
689        self.get_object(path).await
690    }
691
692    /// Read an object with checksum-mode and SSE-C support when implemented by the backend.
693    ///
694    /// The default preserves legacy version selection and rejects advanced options explicitly.
695    async fn get_object_with_transfer_options(
696        &self,
697        path: &RemotePath,
698        options: &TransferReadOptions,
699    ) -> Result<Vec<u8>> {
700        let legacy = options.legacy_read_options()?;
701        self.get_object_with_options(path, &legacy).await
702    }
703
704    /// Stream current object content or an exact historical version to a writer.
705    async fn write_object_to_with_options(
706        &self,
707        path: &RemotePath,
708        options: &ObjectReadOptions,
709        writer: &mut (dyn AsyncWrite + Send + Unpin),
710        max_bytes: Option<u64>,
711    ) -> Result<u64> {
712        let data = self.get_object_with_options(path, options).await?;
713        let write_len = max_bytes
714            .and_then(|limit| usize::try_from(limit).ok())
715            .map(|limit| limit.min(data.len()))
716            .unwrap_or(data.len());
717        writer.write_all(&data[..write_len]).await?;
718        writer.flush().await?;
719        Ok(write_len as u64)
720    }
721
722    /// Stream an object with checksum-mode and SSE-C support when implemented by the backend.
723    ///
724    /// The default preserves legacy version selection and rejects advanced options explicitly.
725    async fn write_object_to_with_transfer_options(
726        &self,
727        path: &RemotePath,
728        options: &TransferReadOptions,
729        writer: &mut (dyn AsyncWrite + Send + Unpin),
730        max_bytes: Option<u64>,
731    ) -> Result<u64> {
732        let legacy = options.legacy_read_options()?;
733        self.write_object_to_with_options(path, &legacy, writer, max_bytes)
734            .await
735    }
736
737    /// Upload object from bytes
738    async fn put_object(
739        &self,
740        path: &RemotePath,
741        data: Vec<u8>,
742        content_type: Option<&str>,
743        encryption: Option<&ObjectEncryptionRequest>,
744    ) -> Result<ObjectInfo>;
745
746    /// Upload an object with complete transfer-fidelity options.
747    ///
748    /// The default delegates Content-Type and managed encryption to the legacy API and rejects
749    /// every option that cannot be represented there.
750    async fn put_object_with_options(
751        &self,
752        path: &RemotePath,
753        data: Vec<u8>,
754        options: &ObjectWriteOptions,
755    ) -> Result<ObjectInfo> {
756        let (content_type, encryption) = options.legacy_put_arguments()?;
757        self.put_object(path, data, content_type, encryption).await
758    }
759
760    /// Delete an object
761    async fn delete_object(&self, path: &RemotePath) -> Result<()>;
762
763    /// Delete the current object state or one exact version with explicit request options.
764    async fn delete_object_with_options(
765        &self,
766        path: &RemotePath,
767        options: DeleteRequestOptions,
768    ) -> Result<DeletedObject> {
769        if options.version_id.is_some() || options.bypass_governance || options.force_delete {
770            return Err(Error::UnsupportedFeature(
771                "Version-aware or policy-bypassing deletion is not implemented by this object store"
772                    .to_string(),
773            ));
774        }
775        self.delete_object(path).await?;
776        Ok(DeletedObject {
777            key: path.key.clone(),
778            version_id: None,
779            is_delete_marker: false,
780        })
781    }
782
783    /// Delete multiple objects (batch delete)
784    async fn delete_objects(&self, bucket: &str, keys: Vec<String>) -> Result<Vec<String>>;
785
786    /// Delete exact object versions and delete markers in one request.
787    async fn delete_object_versions(
788        &self,
789        _bucket: &str,
790        _objects: Vec<ObjectVersionIdentifier>,
791        _options: DeleteRequestOptions,
792    ) -> Result<DeleteObjectsResult> {
793        Err(Error::UnsupportedFeature(
794            "Multi-object version deletion is not implemented by this object store".to_string(),
795        ))
796    }
797
798    /// Copy object within S3 (server-side copy)
799    async fn copy_object(
800        &self,
801        src: &RemotePath,
802        dst: &RemotePath,
803        encryption: Option<&ObjectEncryptionRequest>,
804    ) -> Result<ObjectInfo>;
805
806    /// Copy the current object or one exact historical source version.
807    async fn copy_object_with_options(
808        &self,
809        src: &RemotePath,
810        dst: &RemotePath,
811        options: &CopyObjectOptions,
812        encryption: Option<&ObjectEncryptionRequest>,
813    ) -> Result<ObjectInfo> {
814        if options.source_version_id.is_some() || options.source_etag.is_some() {
815            return Err(Error::UnsupportedFeature(
816                "Exact-source server-side copy is not implemented by this object store".to_string(),
817            ));
818        }
819        self.copy_object(src, dst, encryption).await
820    }
821
822    /// Copy an object with explicit metadata, tags, checksum, SSE-C, and lock intent.
823    ///
824    /// The default preserves legacy source-version and managed-encryption behavior while
825    /// rejecting every advanced option that the original API cannot represent.
826    async fn copy_object_with_transfer_options(
827        &self,
828        src: &RemotePath,
829        dst: &RemotePath,
830        options: &TransferCopyOptions,
831    ) -> Result<ObjectInfo> {
832        let (legacy, encryption) = options.legacy_copy_arguments()?;
833        self.copy_object_with_options(src, dst, &legacy, encryption)
834            .await
835    }
836
837    /// Copy an object with S3's multipart server-side copy lifecycle.
838    ///
839    /// Backends that do not implement multipart copy fail explicitly. The
840    /// callback receives cumulative logical bytes after each successful part.
841    async fn multipart_copy(
842        &self,
843        _src: &RemotePath,
844        _dst: &RemotePath,
845        _options: &MultipartCopyOptions,
846        _cancellation: &MultipartCopyCancellation,
847        _encryption: Option<&ObjectEncryptionRequest>,
848        _on_progress: &MultipartCopyProgress<'_>,
849    ) -> Result<MultipartCopyResult> {
850        Err(Error::UnsupportedFeature(
851            "Multipart server-side copy is not implemented by this object store".to_string(),
852        ))
853    }
854
855    /// Copy an object through the multipart lifecycle with transfer-fidelity options.
856    ///
857    /// The default delegates only when all requested fidelity can be represented by the legacy
858    /// multipart API. Implementations must override this method before accepting advanced fields.
859    async fn multipart_copy_with_transfer_options(
860        &self,
861        src: &RemotePath,
862        dst: &RemotePath,
863        multipart: &MultipartCopyOptions,
864        transfer: &TransferCopyOptions,
865        cancellation: &MultipartCopyCancellation,
866        on_progress: &MultipartCopyProgress<'_>,
867    ) -> Result<MultipartCopyResult> {
868        transfer.validate_multipart_source_version(multipart.source_version_id.as_deref())?;
869        let (_, encryption) = transfer.legacy_copy_arguments()?;
870        self.multipart_copy(src, dst, multipart, cancellation, encryption, on_progress)
871            .await
872    }
873
874    /// Generate a presigned URL for an object
875    async fn presign_get(&self, path: &RemotePath, expires_secs: u64) -> Result<String>;
876
877    /// Generate a presigned URL for uploading an object
878    async fn presign_put(
879        &self,
880        path: &RemotePath,
881        expires_secs: u64,
882        content_type: Option<&str>,
883    ) -> Result<String>;
884
885    // Phase 5: Optional operations (capability-dependent)
886
887    /// Get bucket versioning status
888    async fn get_versioning(&self, bucket: &str) -> Result<Option<bool>>;
889
890    /// Set bucket versioning status
891    async fn set_versioning(&self, bucket: &str, enabled: bool) -> Result<()>;
892
893    /// Get a bucket's Object Lock configuration.
894    ///
895    /// `None` means the bucket exists but has no Object Lock configuration.
896    async fn get_bucket_object_lock_configuration(
897        &self,
898        _bucket: &str,
899    ) -> Result<Option<BucketObjectLockConfiguration>> {
900        Err(Error::UnsupportedFeature(
901            "Bucket Object Lock configuration is not implemented by this object store".to_string(),
902        ))
903    }
904
905    /// Update an Object Lock enabled bucket's configuration.
906    async fn put_bucket_object_lock_configuration(
907        &self,
908        _bucket: &str,
909        _configuration: BucketObjectLockConfiguration,
910    ) -> Result<()> {
911        Err(Error::UnsupportedFeature(
912            "Bucket Object Lock configuration is not implemented by this object store".to_string(),
913        ))
914    }
915
916    /// Get retention applied to the selected object version.
917    async fn get_object_retention(
918        &self,
919        _path: &RemotePath,
920        _options: &ObjectLockOptions,
921    ) -> Result<Option<ObjectRetention>> {
922        Err(Error::UnsupportedFeature(
923            "Object retention is not implemented by this object store".to_string(),
924        ))
925    }
926
927    /// Set or clear retention on the selected object version.
928    async fn put_object_retention(
929        &self,
930        _path: &RemotePath,
931        _retention: Option<ObjectRetention>,
932        _options: &ObjectLockOptions,
933    ) -> Result<()> {
934        Err(Error::UnsupportedFeature(
935            "Object retention is not implemented by this object store".to_string(),
936        ))
937    }
938
939    /// Get legal-hold status for the selected object version.
940    async fn get_object_legal_hold(
941        &self,
942        _path: &RemotePath,
943        _options: &ObjectLockOptions,
944    ) -> Result<LegalHoldStatus> {
945        Err(Error::UnsupportedFeature(
946            "Object legal hold is not implemented by this object store".to_string(),
947        ))
948    }
949
950    /// Set legal-hold status for the selected object version.
951    async fn put_object_legal_hold(
952        &self,
953        _path: &RemotePath,
954        _status: LegalHoldStatus,
955        _options: &ObjectLockOptions,
956    ) -> Result<()> {
957        Err(Error::UnsupportedFeature(
958            "Object legal hold is not implemented by this object store".to_string(),
959        ))
960    }
961
962    /// Get bucket default encryption. Returns None when encryption is not configured.
963    async fn get_bucket_encryption(&self, bucket: &str) -> Result<Option<BucketEncryption>>;
964
965    /// Set bucket default encryption.
966    async fn set_bucket_encryption(&self, bucket: &str, encryption: BucketEncryption)
967    -> Result<()>;
968
969    /// Delete bucket default encryption.
970    async fn delete_bucket_encryption(&self, bucket: &str) -> Result<()>;
971
972    /// List object versions
973    async fn list_object_versions(
974        &self,
975        path: &RemotePath,
976        max_keys: Option<i32>,
977    ) -> Result<Vec<ObjectVersion>>;
978
979    /// List one page of object versions and delete markers with both S3 pagination markers.
980    async fn list_object_versions_page_with_options(
981        &self,
982        _path: &RemotePath,
983        _options: &ListObjectVersionsOptions,
984    ) -> Result<ObjectVersionListResult> {
985        Err(Error::UnsupportedFeature(
986            "Paginated object version listing is not implemented by this object store".to_string(),
987        ))
988    }
989
990    /// Get object tags
991    async fn get_object_tags(
992        &self,
993        path: &RemotePath,
994    ) -> Result<std::collections::HashMap<String, String>>;
995
996    /// Get bucket tags
997    async fn get_bucket_tags(
998        &self,
999        bucket: &str,
1000    ) -> Result<std::collections::HashMap<String, String>>;
1001
1002    /// Set object tags
1003    async fn set_object_tags(
1004        &self,
1005        path: &RemotePath,
1006        tags: std::collections::HashMap<String, String>,
1007    ) -> Result<()>;
1008
1009    /// Set bucket tags
1010    async fn set_bucket_tags(
1011        &self,
1012        bucket: &str,
1013        tags: std::collections::HashMap<String, String>,
1014    ) -> Result<()>;
1015
1016    /// Delete object tags
1017    async fn delete_object_tags(&self, path: &RemotePath) -> Result<()>;
1018
1019    /// Delete bucket tags
1020    async fn delete_bucket_tags(&self, bucket: &str) -> Result<()>;
1021
1022    /// Get bucket policy as raw JSON string. Returns `None` when no policy exists.
1023    async fn get_bucket_policy(&self, bucket: &str) -> Result<Option<String>>;
1024
1025    /// Replace bucket policy using raw JSON string.
1026    async fn set_bucket_policy(&self, bucket: &str, policy: &str) -> Result<()>;
1027
1028    /// Remove bucket policy (set anonymous access to private).
1029    async fn delete_bucket_policy(&self, bucket: &str) -> Result<()>;
1030
1031    /// Get bucket notification configuration as flat rules.
1032    async fn get_bucket_notifications(&self, bucket: &str) -> Result<Vec<BucketNotification>>;
1033
1034    /// Replace bucket notification configuration with flat rules.
1035    async fn set_bucket_notifications(
1036        &self,
1037        bucket: &str,
1038        notifications: Vec<BucketNotification>,
1039    ) -> Result<()>;
1040
1041    // Lifecycle operations (capability-dependent)
1042
1043    /// Get bucket lifecycle rules. Returns empty vec if no lifecycle config exists.
1044    async fn get_bucket_lifecycle(&self, bucket: &str) -> Result<Vec<LifecycleRule>>;
1045
1046    /// Set bucket lifecycle configuration (replaces all rules).
1047    async fn set_bucket_lifecycle(&self, bucket: &str, rules: Vec<LifecycleRule>) -> Result<()>;
1048
1049    /// Delete bucket lifecycle configuration.
1050    async fn delete_bucket_lifecycle(&self, bucket: &str) -> Result<()>;
1051
1052    /// Restore a transitioned (archived) object.
1053    async fn restore_object(&self, path: &RemotePath, days: i32) -> Result<()>;
1054
1055    // Replication operations (capability-dependent)
1056
1057    /// Get bucket replication configuration. Returns None if not configured.
1058    async fn get_bucket_replication(
1059        &self,
1060        bucket: &str,
1061    ) -> Result<Option<ReplicationConfiguration>>;
1062
1063    /// Set bucket replication configuration.
1064    async fn set_bucket_replication(
1065        &self,
1066        bucket: &str,
1067        config: ReplicationConfiguration,
1068    ) -> Result<()>;
1069
1070    /// Delete bucket replication configuration.
1071    async fn delete_bucket_replication(&self, bucket: &str) -> Result<()>;
1072
1073    /// Actively validate configured replication targets using the legacy
1074    /// success/failure contract.
1075    async fn check_bucket_replication(&self, bucket: &str) -> Result<()>;
1076
1077    /// Actively validate configured replication targets and retain structured
1078    /// per-target and cleanup outcomes when the server provides them.
1079    async fn check_bucket_replication_detailed(
1080        &self,
1081        bucket: &str,
1082    ) -> Result<ReplicationCheckResult> {
1083        self.check_bucket_replication(bucket).await?;
1084        Ok(ReplicationCheckResult::legacy_success())
1085    }
1086
1087    /// Start a server-side bucket replication resync.
1088    async fn start_bucket_replication_resync(
1089        &self,
1090        bucket: &str,
1091        options: ReplicationResyncStartOptions,
1092    ) -> Result<ReplicationResyncStartResult>;
1093
1094    /// Read persisted server-side bucket replication resync status.
1095    async fn bucket_replication_resync_status(
1096        &self,
1097        bucket: &str,
1098        target_arn: Option<&str>,
1099    ) -> Result<ReplicationResyncStatus>;
1100
1101    /// Get bucket CORS rules. Returns empty vec if no CORS config exists.
1102    async fn get_bucket_cors(&self, bucket: &str) -> Result<Vec<CorsRule>>;
1103
1104    /// Set bucket CORS configuration (replaces all rules).
1105    async fn set_bucket_cors(&self, bucket: &str, rules: Vec<CorsRule>) -> Result<()>;
1106
1107    /// Delete bucket CORS configuration.
1108    async fn delete_bucket_cors(&self, bucket: &str) -> Result<()>;
1109
1110    /// Run S3 Select on an object and stream result payloads to `writer`.
1111    async fn select_object_content(
1112        &self,
1113        path: &RemotePath,
1114        options: &SelectOptions,
1115        writer: &mut (dyn AsyncWrite + Send + Unpin),
1116    ) -> Result<()>;
1117    // async fn get_versioning(&self, bucket: &str) -> Result<bool>;
1118    // async fn set_versioning(&self, bucket: &str, enabled: bool) -> Result<()>;
1119    // async fn get_tags(&self, path: &RemotePath) -> Result<HashMap<String, String>>;
1120    // async fn set_tags(&self, path: &RemotePath, tags: HashMap<String, String>) -> Result<()>;
1121}
1122
1123#[cfg(test)]
1124mod tests {
1125    use super::*;
1126
1127    #[test]
1128    fn create_bucket_options_reject_lock_without_versioning() {
1129        let options = CreateBucketOptions {
1130            region: Some("us-east-1".to_string()),
1131            versioning_enabled: false,
1132            object_lock_enabled: true,
1133        };
1134
1135        let error = options
1136            .validate()
1137            .expect_err("Object Lock without versioning must be rejected");
1138
1139        assert!(matches!(error, crate::Error::InvalidPath(_)));
1140    }
1141
1142    #[test]
1143    fn create_bucket_options_normalize_cli_lock_to_versioning() {
1144        let options = CreateBucketOptions::for_cli(Some("us-east-1".to_string()), false, true)
1145            .expect("CLI Object Lock options should be valid");
1146
1147        assert!(options.object_lock_enabled);
1148        assert!(options.versioning_enabled);
1149        assert_eq!(options.region.as_deref(), Some("us-east-1"));
1150    }
1151
1152    #[test]
1153    fn test_object_info_file() {
1154        let info = ObjectInfo::file("test.txt", 1024);
1155        assert_eq!(info.key, "test.txt");
1156        assert_eq!(info.size_bytes, Some(1024));
1157        assert!(!info.is_dir);
1158        assert_eq!(info.version_id, None);
1159        assert_eq!(info.source_version_id, None);
1160        assert_eq!(info.is_delete_marker, None);
1161    }
1162
1163    #[test]
1164    fn object_read_options_reject_empty_version_ids() {
1165        let error = ObjectReadOptions::for_version(Some(String::new()))
1166            .expect_err("empty version IDs must be rejected");
1167
1168        assert!(matches!(error, crate::Error::InvalidPath(_)));
1169    }
1170
1171    #[test]
1172    fn versioned_delete_targets_are_serializable_for_structured_output() {
1173        let target = ObjectVersionIdentifier {
1174            key: "reports/a.csv".to_string(),
1175            version_id: Some("v1".to_string()),
1176            is_delete_marker: true,
1177        };
1178
1179        let json = serde_json::to_value(target).expect("serialize versioned delete target");
1180        assert_eq!(json["key"], "reports/a.csv");
1181        assert_eq!(json["version_id"], "v1");
1182        assert_eq!(json["is_delete_marker"], true);
1183    }
1184
1185    #[test]
1186    fn object_info_version_fields_are_optional_and_serializable() {
1187        let current = serde_json::to_value(ObjectInfo::file("current.txt", 1))
1188            .expect("serialize current object info");
1189        assert!(current.get("version_id").is_none());
1190        assert!(current.get("source_version_id").is_none());
1191        assert!(current.get("is_delete_marker").is_none());
1192
1193        let mut copied = ObjectInfo::file("copy.txt", 1);
1194        copied.version_id = Some("destination-v2".to_string());
1195        copied.source_version_id = Some("source-v1".to_string());
1196        let copied = serde_json::to_value(copied).expect("serialize copy object info");
1197        assert_eq!(copied["version_id"], "destination-v2");
1198        assert_eq!(copied["source_version_id"], "source-v1");
1199    }
1200
1201    #[test]
1202    fn test_object_info_dir() {
1203        let info = ObjectInfo::dir("path/to/dir/");
1204        assert_eq!(info.key, "path/to/dir/");
1205        assert!(info.is_dir);
1206        assert!(info.size_bytes.is_none());
1207    }
1208
1209    #[test]
1210    fn test_object_info_bucket() {
1211        let info = ObjectInfo::bucket("my-bucket");
1212        assert_eq!(info.key, "my-bucket");
1213        assert!(info.is_dir);
1214    }
1215
1216    #[test]
1217    fn test_object_info_metadata_default_none() {
1218        let info = ObjectInfo::file("test.txt", 1024);
1219        assert!(info.metadata.is_none());
1220    }
1221
1222    #[test]
1223    fn test_object_info_metadata_set() {
1224        let mut info = ObjectInfo::file("test.txt", 1024);
1225        let mut meta = HashMap::new();
1226        meta.insert("content-disposition".to_string(), "attachment".to_string());
1227        meta.insert("custom-key".to_string(), "custom-value".to_string());
1228        info.metadata = Some(meta);
1229
1230        let metadata = info.metadata.as_ref().expect("metadata should be Some");
1231        assert_eq!(metadata.len(), 2);
1232        assert_eq!(metadata.get("content-disposition").unwrap(), "attachment");
1233        assert_eq!(metadata.get("custom-key").unwrap(), "custom-value");
1234    }
1235
1236    #[test]
1237    fn multipart_upload_serializes_stable_identity_and_server_metadata() {
1238        let upload = MultipartUpload {
1239            bucket: "archive".to_string(),
1240            key: "backups/data.tar".to_string(),
1241            upload_id: "upload-123".to_string(),
1242            initiated: Some(
1243                "2026-07-21T04:00:00Z"
1244                    .parse()
1245                    .expect("test timestamp should be valid"),
1246            ),
1247            size_bytes: None,
1248            storage_class: Some("STANDARD".to_string()),
1249            initiator: Some(MultipartIdentity {
1250                id: Some("user-1".to_string()),
1251                display_name: Some("backup-agent".to_string()),
1252            }),
1253            owner: None,
1254            checksum_algorithm: Some("CRC64NVME".to_string()),
1255            checksum_type: Some("FULL_OBJECT".to_string()),
1256        };
1257
1258        let value = serde_json::to_value(upload).expect("multipart upload should serialize");
1259        assert_eq!(value["bucket"], "archive");
1260        assert_eq!(value["key"], "backups/data.tar");
1261        assert_eq!(value["upload_id"], "upload-123");
1262        assert_eq!(value["initiated"], "2026-07-21T04:00:00Z");
1263        assert!(value["size_bytes"].is_null());
1264        assert_eq!(value["storage_class"], "STANDARD");
1265        assert_eq!(value["initiator"]["id"], "user-1");
1266    }
1267
1268    #[test]
1269    fn multipart_list_options_default_to_the_first_unfiltered_page() {
1270        let options = MultipartUploadListOptions::default();
1271
1272        assert!(options.prefix.is_none());
1273        assert!(options.delimiter.is_none());
1274        assert!(options.key_marker.is_none());
1275        assert!(options.upload_id_marker.is_none());
1276        assert!(options.max_uploads.is_none());
1277    }
1278}