koprogo_api/infrastructure/storage/
s3_storage.rs1use 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
16fn 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#[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 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 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
308pub type SharedS3Storage = Arc<S3Storage>;