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());
    }
}