Skip to main content

coven_storage/cloud/s3/
exact.rs

1use super::*;
2
3#[async_trait]
4impl ExactSlotStorage for S3CloudHome {
5    async fn provider_binding(
6        &self,
7    ) -> Result<coven_protocol::objects::ResolvedProviderBinding, CloudHomeError> {
8        use coven_protocol::objects::{
9            ProviderDeviceBinding, ProviderPrincipalId, ResolvedProviderBinding, S3EndpointBinding,
10            StoreProviderBinding,
11        };
12
13        if self.bucket.is_empty() || self.region.is_empty() || self.access_key.is_empty() {
14            return Err(CloudHomeError::Configuration(
15                "S3 provider binding requires a bucket, region, and access-key id".to_string(),
16            ));
17        }
18        let (endpoint, principal) = match self.endpoint.as_deref() {
19            None => {
20                let client = self.sts_client.clone().ok_or_else(|| {
21                    CloudHomeError::Configuration("AWS S3 adapter has no STS client".to_string())
22                })?;
23                let identity = self
24                    .runtime
25                    .run_cloud(move || async move {
26                        client
27                            .get_caller_identity()
28                            .send()
29                            .await
30                            .map_err(sts_request_error)
31                    })
32                    .await?;
33                let account = identity.account().ok_or_else(|| {
34                    CloudHomeError::Configuration(
35                        "STS GetCallerIdentity returned no account id".to_string(),
36                    )
37                })?;
38                let arn = identity.arn().ok_or_else(|| {
39                    CloudHomeError::Configuration(
40                        "STS GetCallerIdentity returned no caller ARN".to_string(),
41                    )
42                })?;
43                let user_id = identity.user_id().ok_or_else(|| {
44                    CloudHomeError::Configuration(
45                        "STS GetCallerIdentity returned no user id".to_string(),
46                    )
47                })?;
48                let (partition, principal) = aws_caller_identity(account, arn, user_id)?;
49                (
50                    S3EndpointBinding::Aws { partition },
51                    ProviderPrincipalId::Aws {
52                        account_id: account.to_string(),
53                        principal,
54                    },
55                )
56            }
57            Some(endpoint) => (
58                S3EndpointBinding::Custom {
59                    origin: coven_protocol::provider::canonical_custom_s3_origin(endpoint)
60                        .map_err(|error| {
61                            CloudHomeError::configuration("validate custom S3 endpoint", error)
62                        })?,
63                },
64                ProviderPrincipalId::CustomS3Credential {
65                    access_key_id_hash: s3_access_key_id_hash(&self.access_key),
66                },
67            ),
68        };
69        let binding = ResolvedProviderBinding {
70            store: StoreProviderBinding::S3 {
71                endpoint,
72                region: self.region.to_ascii_lowercase(),
73                bucket: self.bucket.clone(),
74                key_prefix: self.key_prefix.clone(),
75            },
76            device: ProviderDeviceBinding { principal },
77        };
78        binding.validate().map_err(|error| {
79            CloudHomeError::configuration("validate S3 provider binding", error)
80        })?;
81        Ok(binding)
82    }
83
84    async fn create_at(
85        &self,
86        upload: &crate::cloud::ExactUpload<'_>,
87        control: &UploadControl,
88    ) -> Result<crate::cloud::ExactCreateOutcome, CloudHomeError> {
89        let checksum = matches!(
90            self.exact_upload_verification,
91            coven_foundation::config::ExactUploadVerification::UploadChecksum
92                | coven_foundation::config::ExactUploadVerification::MetadataHash
93        )
94        .then(|| sha256_base64(upload.object().stored_hash()));
95        let operation = match self.google_xml.is_none() {
96            true => {
97                S3CloudHome::create_at_slot(
98                    self,
99                    upload.object().slot(),
100                    upload.body().await?,
101                    checksum,
102                    control,
103                )
104                .await
105            }
106            false => {
107                upload.object().slot().require_logical_key_for("S3")?;
108                let source = match upload.source() {
109                    crate::cloud::ExactUploadSource::Bytes(bytes) => {
110                        GoogleUploadSource::Bytes(bytes.to_vec())
111                    }
112                    crate::cloud::ExactUploadSource::File(path) => {
113                        GoogleUploadSource::File(path.to_path_buf())
114                    }
115                };
116                let result = self
117                    .put_google_exact_create_only(
118                        upload.object().slot().logical_key(),
119                        source,
120                        upload.object().stored_size(),
121                        hex::encode(upload.object().stored_hash().as_bytes()),
122                        control.clone(),
123                    )
124                    .await;
125                if result.is_ok() {
126                    // The body already reported every chunk it handed over; this
127                    // settles the count on the exact stored size once the
128                    // provider has acknowledged the whole object.
129                    control.report(upload.object().stored_size());
130                }
131                result
132            }
133        };
134        crate::cloud::exact_upload::settle_exact_create(operation, |observed| {
135            self.verify_exact_upload(upload, observed)
136        })
137        .await
138    }
139    async fn list_slots(&self, prefix: &str) -> Result<Vec<ObjectSlot>, CloudHomeError> {
140        crate::cloud::logical_slots(CloudHome::list(self, prefix).await?)
141    }
142    async fn read_at(&self, slot: &ObjectSlot) -> Result<Vec<u8>, CloudHomeError> {
143        slot.require_logical_key_for("S3")?;
144        S3CloudHome::read(self, slot.logical_key()).await
145    }
146    async fn read_versioned_at(
147        &self,
148        slot: &ObjectSlot,
149    ) -> Result<crate::cloud::CloudVersionedObject, CloudHomeError> {
150        slot.require_logical_key_for("S3")?;
151        let key = slot.logical_key().to_string();
152        if let Some(google_xml) = self.google_xml.clone() {
153            let endpoint = self.endpoint.clone().ok_or_else(|| {
154                CloudHomeError::Configuration("Google Cloud Storage endpoint is absent".to_string())
155            })?;
156            let bucket = self.bucket.clone();
157            let region = self.region.clone();
158            let access_key = self.access_key.clone();
159            let secret_key = self.secret_key.clone();
160            let full = self.full_key(&key);
161            let now = self.clock.now();
162            return self
163                .runtime
164                .run(move || async move {
165                    google_xml
166                        .read_versioned(
167                            &endpoint,
168                            &bucket,
169                            &region,
170                            &access_key,
171                            &secret_key,
172                            &full,
173                            now,
174                        )
175                        .await
176                })
177                .await
178                .map_err(|error| {
179                    CloudHomeError::transport("run Google Cloud Storage versioned read", error)
180                })?;
181        }
182        let full = self.full_key(&key);
183        let client = self.client.clone();
184        let bucket = self.bucket.clone();
185        self.runtime
186            .run_cloud(move || async move {
187                let response = client
188                    .get_object()
189                    .bucket(&bucket)
190                    .key(full)
191                    .send()
192                    .await
193                    .map_err(|error| get_object_error(&key, error))?;
194                let version = response.e_tag().ok_or_else(|| {
195                    CloudHomeError::Transport(format!("read versioned {key}: S3 returned no ETag"))
196                })?;
197                let version = crate::cloud::CloudObjectVersion::from_provider(version.to_string())?;
198                let bytes = response
199                    .body
200                    .collect()
201                    .await
202                    .map_err(|error| body_read_error("read versioned body", &key, error))?
203                    .into_bytes()
204                    .to_vec();
205                Ok(crate::cloud::CloudVersionedObject { bytes, version })
206            })
207            .await
208    }
209    async fn replace_at_if_version(
210        &self,
211        slot: &ObjectSlot,
212        expected: &crate::cloud::CloudObjectVersion,
213        bytes: Vec<u8>,
214    ) -> Result<crate::cloud::ConditionalWriteOutcome, CloudHomeError> {
215        slot.require_logical_key_for("S3")?;
216        let key = slot.logical_key().to_string();
217        if let Some(google_xml) = self.google_xml.clone() {
218            let endpoint = self.endpoint.clone().ok_or_else(|| {
219                CloudHomeError::Configuration("Google Cloud Storage endpoint is absent".to_string())
220            })?;
221            let bucket = self.bucket.clone();
222            let region = self.region.clone();
223            let access_key = self.access_key.clone();
224            let secret_key = self.secret_key.clone();
225            let full = self.full_key(&key);
226            let expected = expected.clone();
227            let now = self.clock.now();
228            return self
229                .runtime
230                .run(move || async move {
231                    google_xml
232                        .replace_if_generation(
233                            &endpoint,
234                            &bucket,
235                            &region,
236                            &access_key,
237                            &secret_key,
238                            &full,
239                            &expected,
240                            bytes,
241                            now,
242                        )
243                        .await
244                })
245                .await
246                .map_err(|error| {
247                    CloudHomeError::transport(
248                        "run Google Cloud Storage conditional replacement",
249                        error,
250                    )
251                })?;
252        }
253        let full = self.full_key(&key);
254        let client = self.client.clone();
255        let bucket = self.bucket.clone();
256        let expected = expected.as_provider().to_string();
257        self.runtime
258            .run_cloud(move || async move {
259                let result = client
260                    .put_object()
261                    .bucket(&bucket)
262                    .key(full)
263                    .if_match(expected)
264                    .body(bytes.into())
265                    .send()
266                    .await;
267                let response = match result {
268                    Ok(response) => response,
269                    Err(error) if create_only_put_failed(&error) => {
270                        return Ok(crate::cloud::ConditionalWriteOutcome::VersionChanged)
271                    }
272                    Err(error) => return Err(put_object_error(&key, error)),
273                };
274                let version = response.e_tag().ok_or_else(|| {
275                    CloudHomeError::Transport(format!(
276                        "replace versioned {key}: S3 returned no ETag"
277                    ))
278                })?;
279                Ok(crate::cloud::ConditionalWriteOutcome::Replaced(
280                    crate::cloud::CloudObjectVersion::from_provider(version.to_string())?,
281                ))
282            })
283            .await
284    }
285    async fn read_range_at(
286        &self,
287        slot: &ObjectSlot,
288        start: u64,
289        end: u64,
290    ) -> Result<Vec<u8>, CloudHomeError> {
291        slot.require_logical_key_for("S3")?;
292        S3CloudHome::read_range(self, slot.logical_key(), start, end).await
293    }
294    async fn read_at_to_file(
295        &self,
296        slot: &ObjectSlot,
297        destination: &std::path::Path,
298        progress: crate::cloud::DownloadProgress,
299    ) -> Result<(), crate::cloud::CloudFileReadError> {
300        S3CloudHome::read_exact_to_file(self, slot, destination, progress).await
301    }
302    async fn delete_at(&self, slot: &ObjectSlot) -> Result<(), CloudHomeError> {
303        slot.require_logical_key_for("S3")?;
304        S3CloudHome::delete(self, slot.logical_key()).await
305    }
306}