File size: 3,852 Bytes
1851bae | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 | use std::path::{Path, PathBuf};
use anyhow::Result;
use forge_app::FileReaderInfra;
use futures::{StreamExt, stream};
#[derive(Clone)]
pub struct ForgeFileReadService;
impl Default for ForgeFileReadService {
fn default() -> Self {
Self
}
}
impl ForgeFileReadService {
pub fn new() -> Self {
Self
}
}
#[async_trait::async_trait]
impl FileReaderInfra for ForgeFileReadService {
async fn read_utf8(&self, path: &Path) -> Result<String> {
forge_fs::ForgeFS::read_utf8(path).await
}
fn read_batch_utf8(
&self,
batch_size: usize,
paths: Vec<PathBuf>,
) -> impl futures::Stream<Item = (PathBuf, anyhow::Result<String>)> + Send {
let batches: Vec<Vec<PathBuf>> = paths
.chunks(batch_size)
.map(|chunk| chunk.to_vec())
.collect();
stream::iter(batches)
.then(move |batch| async move {
let futures = batch.into_iter().map(|path| async move {
let result = self.read_utf8(&path).await;
(path, result)
});
futures::future::join_all(futures).await
})
.flat_map(stream::iter)
}
async fn read(&self, path: &Path) -> Result<Vec<u8>> {
forge_fs::ForgeFS::read(path).await
}
async fn range_read_utf8(
&self,
path: &Path,
start_line: u64,
end_line: u64,
) -> Result<(String, forge_domain::FileInfo)> {
forge_fs::ForgeFS::read_range_utf8(path, start_line, end_line).await
}
}
#[cfg(test)]
mod tests {
use std::io::Write;
use futures::StreamExt;
use tempfile::NamedTempFile;
use super::*;
#[tokio::test]
async fn test_read_batch_utf8() {
let fixture = ForgeFileReadService::new();
// Create temporary test files
let mut file1 = NamedTempFile::new().unwrap();
let mut file2 = NamedTempFile::new().unwrap();
let mut file3 = NamedTempFile::new().unwrap();
writeln!(file1, "content1").unwrap();
writeln!(file2, "content2").unwrap();
writeln!(file3, "content3").unwrap();
let paths = vec![
file1.path().to_path_buf(),
file2.path().to_path_buf(),
file3.path().to_path_buf(),
];
// Read with batch size of 2
let stream = fixture.read_batch_utf8(2, paths.clone());
futures::pin_mut!(stream);
let item1 = stream.next().await.unwrap();
assert_eq!(item1.0, paths[0]);
assert_eq!(item1.1.as_deref().unwrap().trim(), "content1");
let item2 = stream.next().await.unwrap();
assert_eq!(item2.0, paths[1]);
assert_eq!(item2.1.as_deref().unwrap().trim(), "content2");
let item3 = stream.next().await.unwrap();
assert_eq!(item3.0, paths[2]);
assert_eq!(item3.1.as_deref().unwrap().trim(), "content3");
// No more items
assert!(stream.next().await.is_none());
}
#[tokio::test]
async fn test_read_batch_utf8_single_batch() {
let fixture = ForgeFileReadService::new();
let mut file1 = NamedTempFile::new().unwrap();
let mut file2 = NamedTempFile::new().unwrap();
writeln!(file1, "test1").unwrap();
writeln!(file2, "test2").unwrap();
let paths = vec![file1.path().to_path_buf(), file2.path().to_path_buf()];
// Read with batch size larger than number of files
let stream = fixture.read_batch_utf8(10, paths.clone());
futures::pin_mut!(stream);
let item1 = stream.next().await.unwrap();
assert_eq!(item1.0, paths[0]);
let item2 = stream.next().await.unwrap();
assert_eq!(item2.0, paths[1]);
// No more items
assert!(stream.next().await.is_none());
}
}
|