Skip to main content

koprogo_api/infrastructure/storage/
s3_storage.rs

1use super::{metrics::record_storage_operation, StorageProvider};
2use async_trait::async_trait;
3use aws_config::meta::region::RegionProviderChain;
4use aws_config::BehaviorVersion;
5use aws_credential_types::provider::SharedCredentialsProvider;
6use aws_credential_types::Credentials;
7use aws_sdk_s3::config::{Builder as S3ConfigBuilder, Region};
8use aws_sdk_s3::error::SdkError;
9use aws_sdk_s3::primitives::ByteStream;
10use aws_sdk_s3::Client;
11use std::env;
12use std::sync::Arc;
13use std::time::Instant;
14use uuid::Uuid;
15
16/// Deplie la chaine de causes d'une erreur du SDK AWS.
17///
18/// `SdkError` implemente `Display` de facon quasi muette : un echec
19/// d'authentification, un bucket absent et un refus de politique rendent tous
20/// les trois la meme chaine, « service error ». Le detail exploitable
21/// (`NoSuchBucket`, `AccessDenied`, `InvalidAccessKeyId`, code HTTP) vit dans
22/// `source()`.
23///
24/// Consequence concrete, constatee le 2026-08-26 : l'upload de documents
25/// echouait en production et l'API ne savait dire que
26/// « Failed to upload object: service error ». Impossible de distinguer une
27/// erreur de configuration d'une panne, ni cote exploitation ni cote test.
28/// Un message qui ne discrimine rien equivaut a une absence de message.
29fn describe_error(err: &dyn std::error::Error) -> String {
30    let mut parts = vec![err.to_string()];
31    let mut cur = err.source();
32    while let Some(e) = cur {
33        let s = e.to_string();
34        if !s.is_empty() && !parts.contains(&s) {
35            parts.push(s);
36        }
37        cur = e.source();
38    }
39    parts.join(" <- ")
40}
41
42/// Configuration holder for the S3/MinIO storage backend.
43#[derive(Clone, Debug)]
44pub struct S3StorageConfig {
45    pub bucket: String,
46    pub region: Option<String>,
47    pub endpoint: Option<String>,
48    pub access_key: String,
49    pub secret_key: String,
50    pub force_path_style: bool,
51    pub key_prefix: Option<String>,
52}
53
54impl S3StorageConfig {
55    /// Load configuration from environment variables.
56    pub fn from_env() -> Result<Self, String> {
57        let bucket =
58            env::var("S3_BUCKET").map_err(|_| "S3_BUCKET is required when using s3 storage")?;
59        let access_key = env::var("S3_ACCESS_KEY")
60            .map_err(|_| "S3_ACCESS_KEY is required when using s3 storage")?;
61        let secret_key = env::var("S3_SECRET_KEY")
62            .map_err(|_| "S3_SECRET_KEY is required when using s3 storage")?;
63
64        let region = env::var("S3_REGION").ok();
65        let endpoint = env::var("S3_ENDPOINT").ok();
66        let key_prefix = env::var("S3_KEY_PREFIX").ok();
67        let force_path_style = env::var("S3_FORCE_PATH_STYLE")
68            .unwrap_or_else(|_| "true".to_string())
69            .parse::<bool>()
70            .unwrap_or(true);
71
72        Ok(Self {
73            bucket,
74            region,
75            endpoint,
76            access_key,
77            secret_key,
78            force_path_style,
79            key_prefix,
80        })
81    }
82}
83
84pub struct S3Storage {
85    client: Client,
86    bucket: String,
87    key_prefix: Option<String>,
88}
89
90impl S3Storage {
91    /// Build a storage instance from a configuration object.
92    pub async fn from_config(config: S3StorageConfig) -> Result<Self, String> {
93        let S3StorageConfig {
94            bucket,
95            region,
96            endpoint,
97            access_key,
98            secret_key,
99            force_path_style,
100            key_prefix,
101        } = config;
102
103        let region_provider = if let Some(ref region) = region {
104            RegionProviderChain::first_try(Region::new(region.clone()))
105        } else {
106            RegionProviderChain::default_provider()
107        };
108
109        let credentials = SharedCredentialsProvider::new(Credentials::new(
110            access_key,
111            secret_key,
112            None,
113            None,
114            "koprogo-storage",
115        ));
116
117        let shared_config = aws_config::defaults(BehaviorVersion::latest())
118            .region(region_provider)
119            .credentials_provider(credentials)
120            .load()
121            .await;
122
123        let mut builder = S3ConfigBuilder::from(&shared_config);
124
125        if let Some(region) = region {
126            builder = builder.region(Region::new(region));
127        }
128
129        if let Some(endpoint) = endpoint {
130            builder = builder.endpoint_url(endpoint);
131        }
132
133        if force_path_style {
134            builder = builder.force_path_style(true);
135        }
136
137        let client = Client::from_conf(builder.build());
138
139        Self::ensure_bucket(&client, &bucket).await?;
140
141        Ok(Self {
142            client,
143            bucket,
144            key_prefix,
145        })
146    }
147
148    fn build_key(&self, building_id: Uuid, original_name: &str) -> String {
149        let sanitized = Self::sanitize_filename(original_name);
150        let unique = format!("{}_{}", Uuid::new_v4(), sanitized);
151        let key = format!("{}/{}", building_id, unique);
152        if let Some(prefix) = &self.key_prefix {
153            format!("{}/{}", prefix.trim_end_matches('/'), key)
154        } else {
155            key
156        }
157    }
158
159    fn sanitize_filename(filename: &str) -> String {
160        filename.replace("..", "_").replace(['/', '\\'], "_")
161    }
162}
163
164impl S3Storage {
165    async fn ensure_bucket(client: &Client, bucket: &str) -> Result<(), String> {
166        match client.head_bucket().bucket(bucket).send().await {
167            Ok(_) => Ok(()),
168            Err(SdkError::ServiceError(err)) if err.err().is_not_found() => {
169                match client.create_bucket().bucket(bucket).send().await {
170                    Ok(_) => Ok(()),
171                    Err(SdkError::ServiceError(err))
172                        if err.err().is_bucket_already_exists()
173                            || err.err().is_bucket_already_owned_by_you() =>
174                    {
175                        Ok(())
176                    }
177                    Err(e) => Err(format!(
178                        "Failed to create bucket `{}`: {}",
179                        bucket,
180                        describe_error(&e)
181                    )),
182                }
183            }
184            Err(e) => Err(format!(
185                "Failed to verify bucket `{}`: {}",
186                bucket,
187                describe_error(&e)
188            )),
189        }
190    }
191}
192
193#[async_trait]
194impl StorageProvider for S3Storage {
195    async fn save_file(
196        &self,
197        building_id: Uuid,
198        filename: &str,
199        content: &[u8],
200    ) -> Result<String, String> {
201        let start = Instant::now();
202        let key = self.build_key(building_id, filename);
203
204        let result = self
205            .client
206            .put_object()
207            .bucket(&self.bucket)
208            .key(&key)
209            .body(ByteStream::from(content.to_vec()))
210            .send()
211            .await
212            .map_err(|e| format!("Failed to upload object: {}", describe_error(&e)))
213            .map(|_| key.clone());
214
215        record_storage_operation(
216            "s3",
217            "save_file",
218            start.elapsed(),
219            result.as_ref().map(|_| ()).map_err(|e| e.as_str()),
220        );
221
222        result
223    }
224
225    async fn read_file(&self, relative_path: &str) -> Result<Vec<u8>, String> {
226        let start = Instant::now();
227        let result = match self
228            .client
229            .get_object()
230            .bucket(&self.bucket)
231            .key(relative_path)
232            .send()
233            .await
234        {
235            Ok(output) => match output.body.collect().await {
236                Ok(data) => Ok(data.into_bytes().to_vec()),
237                Err(e) => Err(format!(
238                    "Failed to read object body: {}",
239                    describe_error(&e)
240                )),
241            },
242            Err(e) => Err(format!("Failed to fetch object: {}", describe_error(&e))),
243        };
244
245        record_storage_operation(
246            "s3",
247            "read_file",
248            start.elapsed(),
249            result.as_ref().map(|_| ()).map_err(|e| e.as_str()),
250        );
251
252        result
253    }
254
255    async fn delete_file(&self, relative_path: &str) -> Result<(), String> {
256        let start = Instant::now();
257        let result = self
258            .client
259            .delete_object()
260            .bucket(&self.bucket)
261            .key(relative_path)
262            .send()
263            .await
264            .map_err(|e| format!("Failed to delete object: {}", describe_error(&e)))
265            .map(|_| ());
266
267        record_storage_operation(
268            "s3",
269            "delete_file",
270            start.elapsed(),
271            result.as_ref().map(|_| ()).map_err(|e| e.as_str()),
272        );
273
274        result
275    }
276
277    async fn file_exists(&self, relative_path: &str) -> bool {
278        let start = Instant::now();
279        let result = self
280            .client
281            .head_object()
282            .bucket(&self.bucket)
283            .key(relative_path)
284            .send()
285            .await;
286
287        let exists = match result {
288            Ok(_) => true,
289            Err(SdkError::ServiceError(err)) if err.err().is_not_found() => false,
290            Err(_) => false,
291        };
292
293        record_storage_operation("s3", "file_exists", start.elapsed(), Ok(()));
294
295        exists
296    }
297}
298
299impl std::fmt::Debug for S3Storage {
300    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
301        f.debug_struct("S3Storage")
302            .field("bucket", &self.bucket)
303            .field("key_prefix", &self.key_prefix)
304            .finish()
305    }
306}
307
308/// Convenient alias for sharing the S3 storage provider.
309pub type SharedS3Storage = Arc<S3Storage>;