ed_journals/modules/io/models/
async_iter.rs1use 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#[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}