1use std::fs;
25use std::io::{Read, Seek, SeekFrom, Write};
26use std::ops::Range;
27use std::path::PathBuf;
28use std::sync::Arc;
29
30use async_trait::async_trait;
31use bytes::Bytes;
32use futures::StreamExt;
33use futures::stream::BoxStream;
34use serde::{Deserialize, Serialize};
35
36use crate::error::invalid_data;
37use crate::io::{
38 FileMetadata, FileRead, FileWrite, InputFile, OutputFile, Storage, StorageConfig,
39 StorageFactory,
40};
41use crate::{Error, ErrorKind, Result};
42
43#[derive(Debug, Clone, Default, Serialize, Deserialize)]
55pub struct LocalFsStorage;
56
57impl LocalFsStorage {
58 pub fn new() -> Self {
60 Self
61 }
62
63 pub(crate) fn normalize_path(path: &str) -> PathBuf {
71 let path = if let Some(stripped) = path.strip_prefix("file://") {
72 if stripped.starts_with('/') {
74 stripped.to_string()
75 } else {
76 format!("/{stripped}")
77 }
78 } else if let Some(stripped) = path.strip_prefix("file:") {
79 if stripped.starts_with('/') {
81 stripped.to_string()
82 } else {
83 format!("/{stripped}")
84 }
85 } else {
86 path.to_string()
87 };
88 PathBuf::from(path)
89 }
90}
91
92#[async_trait]
93#[typetag::serde]
94impl Storage for LocalFsStorage {
95 async fn exists(&self, path: &str) -> Result<bool> {
96 let path = Self::normalize_path(path);
97 Ok(path.exists())
98 }
99
100 async fn metadata(&self, path: &str) -> Result<FileMetadata> {
101 let path = Self::normalize_path(path);
102 let metadata = fs::metadata(&path)
103 .map_err(|e| invalid_data!("Failed to get metadata for {}: {}", path.display(), e))?;
104 Ok(FileMetadata {
105 size: metadata.len(),
106 })
107 }
108
109 async fn read(&self, path: &str) -> Result<Bytes> {
110 let path = Self::normalize_path(path);
111 let content = fs::read(&path)
112 .map_err(|e| invalid_data!("Failed to read file {}: {}", path.display(), e))?;
113 Ok(Bytes::from(content))
114 }
115
116 async fn reader(&self, path: &str) -> Result<Box<dyn FileRead>> {
117 let path = Self::normalize_path(path);
118 let file = fs::File::open(&path)
119 .map_err(|e| invalid_data!("Failed to open file {}: {}", path.display(), e))?;
120 Ok(Box::new(LocalFsFileRead::new(file)))
121 }
122
123 async fn write(&self, path: &str, bs: Bytes) -> Result<()> {
124 let path = Self::normalize_path(path);
125
126 if let Some(parent) = path.parent() {
128 fs::create_dir_all(parent).map_err(|e| {
129 Error::new(
130 ErrorKind::Unexpected,
131 format!("Failed to create directory {}: {}", parent.display(), e),
132 )
133 })?;
134 }
135
136 fs::write(&path, &bs).map_err(|e| {
137 Error::new(
138 ErrorKind::Unexpected,
139 format!("Failed to write file {}: {}", path.display(), e),
140 )
141 })?;
142 Ok(())
143 }
144
145 async fn writer(&self, path: &str) -> Result<Box<dyn FileWrite>> {
146 let path = Self::normalize_path(path);
147
148 if let Some(parent) = path.parent() {
150 fs::create_dir_all(parent).map_err(|e| {
151 Error::new(
152 ErrorKind::Unexpected,
153 format!("Failed to create directory {}: {}", parent.display(), e),
154 )
155 })?;
156 }
157
158 let file = fs::File::create(&path).map_err(|e| {
159 Error::new(
160 ErrorKind::Unexpected,
161 format!("Failed to create file {}: {}", path.display(), e),
162 )
163 })?;
164 Ok(Box::new(LocalFsFileWrite::new(file)))
165 }
166
167 async fn delete(&self, path: &str) -> Result<()> {
168 let path = Self::normalize_path(path);
169 if path.exists() {
170 fs::remove_file(&path).map_err(|e| {
171 Error::new(
172 ErrorKind::Unexpected,
173 format!("Failed to delete file {}: {}", path.display(), e),
174 )
175 })?;
176 }
177 Ok(())
178 }
179
180 async fn delete_prefix(&self, path: &str) -> Result<()> {
181 let path = Self::normalize_path(path);
182 if path.is_dir() {
183 fs::remove_dir_all(&path).map_err(|e| {
184 Error::new(
185 ErrorKind::Unexpected,
186 format!("Failed to delete directory {}: {}", path.display(), e),
187 )
188 })?;
189 }
190 Ok(())
191 }
192
193 async fn delete_stream(&self, mut paths: BoxStream<'static, String>) -> Result<()> {
194 while let Some(path) = paths.next().await {
195 self.delete(&path).await?;
196 }
197 Ok(())
198 }
199
200 fn new_input(&self, path: &str) -> Result<InputFile> {
201 Ok(InputFile::new(Arc::new(self.clone()), path.to_string()))
202 }
203
204 fn new_output(&self, path: &str) -> Result<OutputFile> {
205 Ok(OutputFile::new(Arc::new(self.clone()), path.to_string()))
206 }
207}
208
209#[derive(Debug)]
211pub struct LocalFsFileRead {
212 file: std::sync::Mutex<fs::File>,
213}
214
215impl LocalFsFileRead {
216 pub fn new(file: fs::File) -> Self {
218 Self {
219 file: std::sync::Mutex::new(file),
220 }
221 }
222}
223
224#[async_trait]
225impl FileRead for LocalFsFileRead {
226 async fn read(&self, range: Range<u64>) -> Result<Bytes> {
227 let mut file = self.file.lock().map_err(|e| {
228 Error::new(
229 ErrorKind::Unexpected,
230 format!("Failed to acquire file lock: {e}"),
231 )
232 })?;
233
234 file.seek(SeekFrom::Start(range.start))
235 .map_err(|e| invalid_data!("Failed to seek to position {}: {}", range.start, e))?;
236
237 let len = (range.end - range.start) as usize;
238 let mut buffer = vec![0u8; len];
239 file.read_exact(&mut buffer)
240 .map_err(|e| invalid_data!("Failed to read {len} bytes: {e}"))?;
241
242 Ok(Bytes::from(buffer))
243 }
244}
245
246#[derive(Debug)]
250pub struct LocalFsFileWrite {
251 file: Option<fs::File>,
252 bytes_written: u64,
253}
254
255impl LocalFsFileWrite {
256 pub fn new(file: fs::File) -> Self {
258 Self {
259 file: Some(file),
260 bytes_written: 0,
261 }
262 }
263}
264
265#[async_trait]
266impl FileWrite for LocalFsFileWrite {
267 async fn write(&mut self, bs: Bytes) -> Result<()> {
268 let file = self
269 .file
270 .as_mut()
271 .ok_or_else(|| invalid_data!("Cannot write to closed file"))?;
272
273 file.write_all(&bs).map_err(|e| {
274 Error::new(
275 ErrorKind::Unexpected,
276 format!("Failed to write to file: {e}"),
277 )
278 })?;
279 self.bytes_written += bs.len() as u64;
280
281 Ok(())
282 }
283
284 async fn close(&mut self) -> Result<FileMetadata> {
285 let file = self
286 .file
287 .take()
288 .ok_or_else(|| invalid_data!("File already closed"))?;
289
290 file.sync_all()
291 .map_err(|e| Error::new(ErrorKind::Unexpected, format!("Failed to sync file: {e}")))?;
292
293 Ok(FileMetadata {
294 size: self.bytes_written,
295 })
296 }
297}
298
299#[derive(Clone, Debug, Default, Serialize, Deserialize)]
314pub struct LocalFsStorageFactory;
315
316#[typetag::serde]
317impl StorageFactory for LocalFsStorageFactory {
318 fn build(&self, _config: &StorageConfig) -> Result<Arc<dyn Storage>> {
319 Ok(Arc::new(LocalFsStorage::new()))
320 }
321}
322
323#[cfg(test)]
324mod tests {
325 use tempfile::TempDir;
326
327 use super::*;
328
329 #[test]
330 fn test_normalize_path() {
331 assert_eq!(
333 LocalFsStorage::normalize_path("file:///path/to/file"),
334 PathBuf::from("/path/to/file")
335 );
336
337 assert_eq!(
339 LocalFsStorage::normalize_path("file://path/to/file"),
340 PathBuf::from("/path/to/file")
341 );
342
343 assert_eq!(
345 LocalFsStorage::normalize_path("file:/path/to/file"),
346 PathBuf::from("/path/to/file")
347 );
348
349 assert_eq!(
351 LocalFsStorage::normalize_path("/path/to/file"),
352 PathBuf::from("/path/to/file")
353 );
354 }
355
356 #[tokio::test]
357 async fn test_local_fs_storage_write_read() {
358 let tmp_dir = TempDir::new().unwrap();
359 let storage = LocalFsStorage::new();
360 let path = tmp_dir.path().join("test.txt");
361 let path_str = path.to_str().unwrap();
362 let content = Bytes::from("Hello, World!");
363
364 storage.write(path_str, content.clone()).await.unwrap();
366
367 let read_content = storage.read(path_str).await.unwrap();
369 assert_eq!(read_content, content);
370 }
371
372 #[tokio::test]
373 async fn test_local_fs_storage_exists() {
374 let tmp_dir = TempDir::new().unwrap();
375 let storage = LocalFsStorage::new();
376 let path = tmp_dir.path().join("test.txt");
377 let path_str = path.to_str().unwrap();
378
379 assert!(!storage.exists(path_str).await.unwrap());
381
382 storage.write(path_str, Bytes::from("test")).await.unwrap();
384
385 assert!(storage.exists(path_str).await.unwrap());
387 }
388
389 #[tokio::test]
390 async fn test_local_fs_storage_metadata() {
391 let tmp_dir = TempDir::new().unwrap();
392 let storage = LocalFsStorage::new();
393 let path = tmp_dir.path().join("test.txt");
394 let path_str = path.to_str().unwrap();
395 let content = Bytes::from("Hello, World!");
396
397 storage.write(path_str, content.clone()).await.unwrap();
398
399 let metadata = storage.metadata(path_str).await.unwrap();
400 assert_eq!(metadata.size, content.len() as u64);
401 }
402
403 #[tokio::test]
404 async fn test_local_fs_storage_delete() {
405 let tmp_dir = TempDir::new().unwrap();
406 let storage = LocalFsStorage::new();
407 let path = tmp_dir.path().join("test.txt");
408 let path_str = path.to_str().unwrap();
409
410 storage.write(path_str, Bytes::from("test")).await.unwrap();
411 assert!(storage.exists(path_str).await.unwrap());
412
413 storage.delete(path_str).await.unwrap();
414 assert!(!storage.exists(path_str).await.unwrap());
415 }
416
417 #[tokio::test]
418 async fn test_local_fs_storage_delete_prefix() {
419 let tmp_dir = TempDir::new().unwrap();
420 let storage = LocalFsStorage::new();
421 let dir_path = tmp_dir.path().join("subdir");
422 let file1 = dir_path.join("file1.txt");
423 let file2 = dir_path.join("file2.txt");
424
425 storage
427 .write(file1.to_str().unwrap(), Bytes::from("1"))
428 .await
429 .unwrap();
430 storage
431 .write(file2.to_str().unwrap(), Bytes::from("2"))
432 .await
433 .unwrap();
434
435 storage
437 .delete_prefix(dir_path.to_str().unwrap())
438 .await
439 .unwrap();
440
441 assert!(!dir_path.exists());
443 }
444
445 #[tokio::test]
446 async fn test_local_fs_storage_reader() {
447 let tmp_dir = TempDir::new().unwrap();
448 let storage = LocalFsStorage::new();
449 let path = tmp_dir.path().join("test.txt");
450 let path_str = path.to_str().unwrap();
451 let content = Bytes::from("Hello, World!");
452
453 storage.write(path_str, content.clone()).await.unwrap();
454
455 let reader = storage.reader(path_str).await.unwrap();
456 let read_content = reader.read(0..content.len() as u64).await.unwrap();
457 assert_eq!(read_content, content);
458
459 let partial = reader.read(0..5).await.unwrap();
461 assert_eq!(partial, Bytes::from("Hello"));
462 }
463
464 #[tokio::test]
465 async fn test_local_fs_storage_writer() {
466 let tmp_dir = TempDir::new().unwrap();
467 let storage = LocalFsStorage::new();
468 let path = tmp_dir.path().join("test.txt");
469 let path_str = path.to_str().unwrap();
470
471 let mut writer = storage.writer(path_str).await.unwrap();
472 writer.write(Bytes::from("Hello, ")).await.unwrap();
473 writer.write(Bytes::from("World!")).await.unwrap();
474 let metadata = writer.close().await.unwrap();
475
476 let content = storage.read(path_str).await.unwrap();
477 assert_eq!(content, Bytes::from("Hello, World!"));
478 assert_eq!(metadata.size, content.len() as u64);
479 }
480
481 #[tokio::test]
482 async fn test_local_fs_file_write_double_close() {
483 let tmp_dir = TempDir::new().unwrap();
484 let storage = LocalFsStorage::new();
485 let path = tmp_dir.path().join("test.txt");
486 let path_str = path.to_str().unwrap();
487
488 let mut writer = storage.writer(path_str).await.unwrap();
489 writer.write(Bytes::from("test")).await.unwrap();
490 writer.close().await.unwrap();
491
492 let result = writer.close().await;
494 assert!(result.is_err());
495 }
496
497 #[tokio::test]
498 async fn test_local_fs_file_write_after_close() {
499 let tmp_dir = TempDir::new().unwrap();
500 let storage = LocalFsStorage::new();
501 let path = tmp_dir.path().join("test.txt");
502 let path_str = path.to_str().unwrap();
503
504 let mut writer = storage.writer(path_str).await.unwrap();
505 assert_eq!(writer.close().await.unwrap().size, 0);
506
507 let result = writer.write(Bytes::from("test")).await;
509 assert!(result.is_err());
510 }
511
512 #[test]
513 fn test_local_fs_storage_factory() {
514 let factory = LocalFsStorageFactory;
515 let config = StorageConfig::new();
516 let storage = factory.build(&config).unwrap();
517
518 assert!(format!("{storage:?}").contains("LocalFsStorage"));
520 }
521
522 #[tokio::test]
523 async fn test_local_fs_creates_parent_directories() {
524 let tmp_dir = TempDir::new().unwrap();
525 let storage = LocalFsStorage::new();
526 let path = tmp_dir.path().join("a/b/c/test.txt");
527 let path_str = path.to_str().unwrap();
528
529 storage.write(path_str, Bytes::from("test")).await.unwrap();
531
532 assert!(path.exists());
533 }
534
535 #[tokio::test]
536 async fn test_local_fs_storage_delete_stream() {
537 use futures::stream;
538
539 let tmp_dir = TempDir::new().unwrap();
540 let storage = LocalFsStorage::new();
541
542 let file1 = tmp_dir.path().join("file1.txt");
544 let file2 = tmp_dir.path().join("file2.txt");
545 let file3 = tmp_dir.path().join("file3.txt");
546
547 storage
548 .write(file1.to_str().unwrap(), Bytes::from("1"))
549 .await
550 .unwrap();
551 storage
552 .write(file2.to_str().unwrap(), Bytes::from("2"))
553 .await
554 .unwrap();
555 storage
556 .write(file3.to_str().unwrap(), Bytes::from("3"))
557 .await
558 .unwrap();
559
560 assert!(storage.exists(file1.to_str().unwrap()).await.unwrap());
562 assert!(storage.exists(file2.to_str().unwrap()).await.unwrap());
563 assert!(storage.exists(file3.to_str().unwrap()).await.unwrap());
564
565 let paths = vec![
567 file1.to_str().unwrap().to_string(),
568 file2.to_str().unwrap().to_string(),
569 ];
570 let path_stream = stream::iter(paths).boxed();
571 storage.delete_stream(path_stream).await.unwrap();
572
573 assert!(!storage.exists(file1.to_str().unwrap()).await.unwrap());
575 assert!(!storage.exists(file2.to_str().unwrap()).await.unwrap());
576
577 assert!(storage.exists(file3.to_str().unwrap()).await.unwrap());
579 }
580
581 #[tokio::test]
582 async fn test_local_fs_storage_delete_stream_empty() {
583 use futures::stream;
584
585 let storage = LocalFsStorage::new();
586
587 let path_stream = stream::iter(Vec::<String>::new()).boxed();
589 storage.delete_stream(path_stream).await.unwrap();
590 }
591}