Skip to main content

ed_journals/modules/io/models/
async_iter.rs

1use crate::logs::LogEvent;
2use crate::modules::io::error::LogIOError;
3use futures::{AsyncRead, AsyncReadExt, FutureExt, Stream};
4use serde::de::DeserializeOwned;
5use std::marker::PhantomData;
6use std::pin::{pin, Pin};
7use std::task::{Context, Poll};
8
9/// Asynchronous iterator for iterating over some [AsyncRead] and returning [serde_json::Value]s
10/// which can then manually be parsed using [serde_json::from_value]. To automatically parse
11/// entries, use [AsyncIter](crate::io::AsyncIter) instead.
12#[derive(Debug)]
13pub struct AsyncIter<T, R = LogEvent>
14where
15    T: AsyncRead + Unpin,
16    R: DeserializeOwned + Unpin,
17{
18    inner: T,
19    _p: PhantomData<R>,
20}
21
22impl<T> AsyncIter<T>
23where
24    T: AsyncRead + Unpin,
25{
26    pub fn new(inner: T) -> AsyncIter<T, LogEvent> {
27        AsyncIter {
28            inner,
29            _p: PhantomData,
30        }
31    }
32
33    pub fn new_raw(inner: T) -> AsyncIter<T, serde_json::Value> {
34        AsyncIter {
35            inner,
36            _p: PhantomData,
37        }
38    }
39
40    pub fn new_typed<R>(inner: T) -> AsyncIter<T, R>
41    where
42        R: DeserializeOwned + Unpin,
43    {
44        AsyncIter {
45            inner,
46            _p: PhantomData,
47        }
48    }
49}
50
51impl<T, R> AsyncIter<T, R>
52where
53    T: AsyncRead + Unpin,
54    R: DeserializeOwned + Unpin,
55{
56    async fn inner_next(&mut self) -> Option<Result<R, LogIOError>> {
57        let mut line = Vec::with_capacity(64);
58
59        loop {
60            let mut buf: [u8; 1] = [0; 1];
61            let result = self.inner.read(&mut buf).await;
62
63            match result {
64                Ok(0) => break,
65                Ok(_) => {}
66                Err(e) => return Some(Err(e.into())),
67            }
68
69            let byte = buf[0];
70
71            if byte == b'\n' && !line.is_empty() {
72                break;
73            }
74
75            if byte == 0x00 || (line.is_empty() && byte == b' ') {
76                continue;
77            }
78
79            line.push(byte);
80        }
81
82        if line.is_empty() {
83            return None;
84        }
85
86        Some(Ok(match serde_json::from_slice(&line) {
87            Ok(event) => event,
88            Err(e) => return Some(Err(e.into())),
89        }))
90    }
91}
92
93impl<T, R> From<T> for AsyncIter<T, R>
94where
95    T: AsyncRead + Unpin,
96    R: DeserializeOwned + Unpin,
97{
98    fn from(inner: T) -> Self {
99        AsyncIter {
100            inner,
101            _p: PhantomData,
102        }
103    }
104}
105
106#[cfg(feature = "tokio")]
107#[cfg_attr(docsrs, doc(cfg(feature = "tokio")))]
108impl<A> From<A> for AsyncIter<tokio_util::compat::Compat<A>>
109where
110    A: tokio::io::AsyncRead + Unpin,
111{
112    fn from(value: A) -> Self {
113        let compat = tokio_util::compat::TokioAsyncReadCompatExt::compat(value);
114        AsyncIter::from(compat)
115    }
116}
117
118impl<T, R> Stream for AsyncIter<T, R>
119where
120    T: AsyncRead + Unpin,
121    R: DeserializeOwned + Unpin,
122{
123    type Item = Result<R, LogIOError>;
124
125    fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
126        pin!(self.inner_next()).poll_unpin(cx)
127    }
128}
129
130#[cfg(test)]
131mod tests {
132    use crate::logs::LogEventContentKind;
133    use crate::modules::io::AsyncIter;
134    use async_fs::File;
135    use futures::io::Cursor;
136    use futures::StreamExt;
137    use smol::fs;
138
139    fn async_reader_reads_complete_file_correctly() {
140        smol::block_on(async {
141            let data = r#"{ "timestamp":"2020-09-21T19:04:44Z", "event":"Repair", "Item":"Paint", "Cost":1 }
142{ "timestamp":"2020-09-21T19:04:51Z", "event":"Repair", "Item":"Wear", "Cost":10 }"#;
143
144            let cursor = Cursor::new(data);
145            let buf_reader = futures::io::BufReader::new(cursor);
146
147            let mut reader = AsyncIter::new(buf_reader);
148
149            assert!(reader.next().await.is_some());
150            assert!(reader.next().await.is_some());
151
152            assert!(dbg!(reader.next().await).is_none());
153        });
154    }
155
156    fn last_lines_are_read_correctly() {
157        smol::block_on(async {
158            fs::write("c.tmp", "").await.unwrap();
159
160            let file = File::open("c.tmp").await.unwrap();
161
162            let buf_reader = futures::io::BufReader::new(file);
163
164            let mut reader = AsyncIter::new(buf_reader);
165
166            assert!(reader.next().await.is_none());
167
168            fs::write(
169                "c.tmp",
170                r#"{"timestamp":"2022-10-22T15:10:41Z","event":"Fileheader","part":1,"language":"English/UK","Odyssey":true,"gameversion":"4.0.0.1450","build":"r286858/r0 "}"#,
171            )
172                .await
173                .unwrap();
174
175            assert_eq!(
176                reader.next().await.unwrap().unwrap().content.kind(),
177                LogEventContentKind::FileHeader
178            );
179
180            fs::write("c.tmp", r#"{"timestamp":"2022-10-22T15:10:41Z","event":"Fileheader","part":1,"language":"English/UK","Odyssey":true,"gameversion":"4.0.0.1450","build":"r286858/r0 "}
181{"timestamp":"2022-10-22T15:12:05Z","event":"Commander","FID":"F123456789","Name":"TEST"}"#)
182                .await
183                .unwrap();
184
185            assert_eq!(
186                reader.next().await.unwrap().unwrap().content.kind(),
187                LogEventContentKind::Commander
188            );
189
190            fs::remove_file("c.tmp").await.unwrap();
191        });
192    }
193}