1use std::fs::{File, OpenOptions};
28use std::io::{self, BufRead, BufReader, BufWriter, Write};
29use std::marker::PhantomData;
30use std::path::Path;
31
32use serde::de::DeserializeOwned;
33use serde::Serialize;
34
35use crate::error::IoError;
36
37pub type JsonlResult<T> = Result<T, IoError>;
39
40pub struct JsonlReader<T> {
63 inner: BufReader<File>,
64 line_buf: String,
65 _marker: PhantomData<T>,
66}
67
68impl<T: DeserializeOwned> JsonlReader<T> {
69 pub fn open(path: &Path) -> JsonlResult<Self> {
71 let file = File::open(path).map_err(|e| IoError::FileError(e.to_string()))?;
72 Ok(Self {
73 inner: BufReader::new(file),
74 line_buf: String::new(),
75 _marker: PhantomData,
76 })
77 }
78
79 pub fn next_record(&mut self) -> JsonlResult<Option<T>> {
84 loop {
85 self.line_buf.clear();
86 let n = self
87 .inner
88 .read_line(&mut self.line_buf)
89 .map_err(|e| IoError::FileError(e.to_string()))?;
90
91 if n == 0 {
92 return Ok(None);
93 }
94
95 let trimmed = self.line_buf.trim();
96 if trimmed.is_empty() || trimmed.starts_with('#') {
97 continue;
98 }
99
100 let record = serde_json::from_str::<T>(trimmed)
101 .map_err(|e| IoError::ParseError(format!("JSON parse error: {e}")))?;
102
103 return Ok(Some(record));
104 }
105 }
106
107 pub fn collect_all(&mut self) -> JsonlResult<Vec<T>> {
109 let mut out = Vec::new();
110 while let Some(r) = self.next_record()? {
111 out.push(r);
112 }
113 Ok(out)
114 }
115}
116
117pub struct JsonlWriter {
138 inner: BufWriter<File>,
139}
140
141impl JsonlWriter {
142 pub fn create(path: &Path) -> JsonlResult<Self> {
144 let file = File::create(path).map_err(|e| IoError::FileError(e.to_string()))?;
145 Ok(Self {
146 inner: BufWriter::new(file),
147 })
148 }
149
150 pub fn append(path: &Path) -> JsonlResult<Self> {
152 let file = OpenOptions::new()
153 .create(true)
154 .append(true)
155 .open(path)
156 .map_err(|e| IoError::FileError(e.to_string()))?;
157 Ok(Self {
158 inner: BufWriter::new(file),
159 })
160 }
161
162 pub fn write_record<T: Serialize>(&mut self, record: &T) -> JsonlResult<()> {
164 let json = serde_json::to_string(record)
165 .map_err(|e| IoError::SerializationError(format!("JSON serialization: {e}")))?;
166 self.inner
167 .write_all(json.as_bytes())
168 .map_err(|e| IoError::FileError(e.to_string()))?;
169 self.inner
170 .write_all(b"\n")
171 .map_err(|e| IoError::FileError(e.to_string()))?;
172 Ok(())
173 }
174
175 pub fn flush(&mut self) -> JsonlResult<()> {
177 self.inner
178 .flush()
179 .map_err(|e| IoError::FileError(e.to_string()))
180 }
181}
182
183pub fn read_jsonl<T: DeserializeOwned>(path: &Path) -> JsonlResult<Vec<T>> {
190 JsonlReader::open(path)?.collect_all()
191}
192
193pub fn write_jsonl<T: Serialize>(records: &[T], path: &Path) -> JsonlResult<()> {
197 let mut writer = JsonlWriter::create(path)?;
198 for record in records {
199 writer.write_record(record)?;
200 }
201 writer.flush()
202}
203
204pub fn stream_jsonl<T: DeserializeOwned>(path: &Path) -> JsonlStreamIter<T> {
224 JsonlStreamIter::new(path)
225}
226
227pub struct JsonlStreamIter<T> {
231 reader: Option<BufReader<File>>,
232 line_buf: String,
233 _marker: PhantomData<T>,
234}
235
236impl<T: DeserializeOwned> JsonlStreamIter<T> {
237 fn new(path: &Path) -> Self {
238 match File::open(path) {
239 Ok(f) => Self {
240 reader: Some(BufReader::new(f)),
241 line_buf: String::new(),
242 _marker: PhantomData,
243 },
244 Err(e) => {
245 let _ = e; Self {
252 reader: None,
253 line_buf: String::new(),
254 _marker: PhantomData,
255 }
256 }
257 }
258 }
259}
260
261impl<T: DeserializeOwned> Iterator for JsonlStreamIter<T> {
262 type Item = JsonlResult<T>;
263
264 fn next(&mut self) -> Option<Self::Item> {
265 let reader = self.reader.as_mut()?;
266
267 loop {
268 self.line_buf.clear();
269 let n = match reader.read_line(&mut self.line_buf) {
270 Ok(n) => n,
271 Err(e) => return Some(Err(IoError::FileError(e.to_string()))),
272 };
273
274 if n == 0 {
275 return None;
276 }
277
278 let trimmed = self.line_buf.trim();
279 if trimmed.is_empty() || trimmed.starts_with('#') {
280 continue;
281 }
282
283 return Some(
284 serde_json::from_str::<T>(trimmed)
285 .map_err(|e| IoError::ParseError(format!("JSON parse: {e}"))),
286 );
287 }
288 }
289}
290
291#[cfg(test)]
294mod tests {
295 use super::*;
296 use serde::{Deserialize, Serialize};
297 use std::env::temp_dir;
298
299 #[derive(Debug, Serialize, Deserialize, PartialEq)]
300 struct Point {
301 x: f64,
302 y: f64,
303 label: String,
304 }
305
306 fn sample_points() -> Vec<Point> {
307 vec![
308 Point {
309 x: 1.0,
310 y: 2.0,
311 label: "A".to_string(),
312 },
313 Point {
314 x: -3.5,
315 y: 0.0,
316 label: "B".to_string(),
317 },
318 Point {
319 x: 100.0,
320 y: -100.0,
321 label: "C".to_string(),
322 },
323 ]
324 }
325
326 fn tmp_path(name: &str) -> std::path::PathBuf {
327 temp_dir().join(name)
328 }
329
330 #[test]
331 fn test_write_and_read_jsonl() {
332 let path = tmp_path("test_points.jsonl");
333 let pts = sample_points();
334 write_jsonl(&pts, &path).expect("write failed");
335 let loaded: Vec<Point> = read_jsonl(&path).expect("read failed");
336 assert_eq!(loaded.len(), pts.len());
337 assert_eq!(loaded[0], pts[0]);
338 assert_eq!(loaded[2], pts[2]);
339 }
340
341 #[test]
342 fn test_jsonl_reader_next_record() {
343 let path = tmp_path("test_next.jsonl");
344 let pts = sample_points();
345 write_jsonl(&pts, &path).expect("write failed");
346
347 let mut reader = JsonlReader::<Point>::open(&path).expect("open failed");
348 let first = reader
349 .next_record()
350 .expect("read error")
351 .expect("should have record");
352 assert_eq!(first, pts[0]);
353 let second = reader
354 .next_record()
355 .expect("read error")
356 .expect("should have record");
357 assert_eq!(second, pts[1]);
358 let third = reader
359 .next_record()
360 .expect("read error")
361 .expect("should have record");
362 assert_eq!(third, pts[2]);
363 let eof = reader.next_record().expect("read error");
364 assert!(eof.is_none());
365 }
366
367 #[test]
368 fn test_jsonl_writer_append() {
369 let path = tmp_path("test_append.jsonl");
370 let batch1 = vec![Point {
372 x: 0.0,
373 y: 0.0,
374 label: "Origin".to_string(),
375 }];
376 write_jsonl(&batch1, &path).expect("write batch1 failed");
377
378 let batch2 = [Point {
380 x: 1.0,
381 y: 1.0,
382 label: "Unit".to_string(),
383 }];
384 let mut writer = JsonlWriter::append(&path).expect("append open failed");
385 writer
386 .write_record(&batch2[0])
387 .expect("write record failed");
388 writer.flush().expect("flush failed");
389
390 let all: Vec<Point> = read_jsonl(&path).expect("read failed");
391 assert_eq!(all.len(), 2);
392 assert_eq!(all[0].label, "Origin");
393 assert_eq!(all[1].label, "Unit");
394 }
395
396 #[test]
397 fn test_stream_jsonl_iterator() {
398 let path = tmp_path("test_stream.jsonl");
399 let pts = sample_points();
400 write_jsonl(&pts, &path).expect("write failed");
401
402 let collected: Vec<Point> = stream_jsonl::<Point>(&path)
403 .map(|r| r.expect("stream error"))
404 .collect();
405 assert_eq!(collected.len(), pts.len());
406 for (a, b) in collected.iter().zip(pts.iter()) {
407 assert_eq!(a, b);
408 }
409 }
410
411 #[test]
412 fn test_empty_file() {
413 let path = tmp_path("test_empty.jsonl");
414 write_jsonl::<Point>(&[], &path).expect("write failed");
415 let loaded: Vec<Point> = read_jsonl(&path).expect("read failed");
416 assert!(loaded.is_empty());
417 }
418
419 #[test]
420 fn test_jsonl_skips_blank_lines_and_comments() {
421 let path = tmp_path("test_comments.jsonl");
422 {
424 let mut f = File::create(&path).expect("create failed");
425 writeln!(f, "# This is a comment").expect("write failed");
426 writeln!(f, r#"{{"x":1.0,"y":2.0,"label":"A"}}"#).expect("write failed");
427 writeln!(f).expect("write failed"); writeln!(f, r#"{{"x":3.0,"y":4.0,"label":"B"}}"#).expect("write failed");
429 }
430 let loaded: Vec<Point> = read_jsonl(&path).expect("read failed");
431 assert_eq!(loaded.len(), 2);
432 assert_eq!(loaded[0].label, "A");
433 assert_eq!(loaded[1].label, "B");
434 }
435
436 #[test]
437 fn test_large_dataset() {
438 let path = tmp_path("test_large.jsonl");
439 let n = 10_000usize;
440 let records: Vec<Point> = (0..n)
441 .map(|i| Point {
442 x: i as f64,
443 y: -(i as f64),
444 label: format!("item_{i}"),
445 })
446 .collect();
447 write_jsonl(&records, &path).expect("write failed");
448 let loaded: Vec<Point> = read_jsonl(&path).expect("read failed");
449 assert_eq!(loaded.len(), n);
450 assert_eq!(loaded[9999].label, "item_9999");
451 }
452
453 #[test]
454 fn test_collect_all_via_reader() {
455 let path = tmp_path("test_collect.jsonl");
456 let pts = sample_points();
457 write_jsonl(&pts, &path).expect("write failed");
458 let mut reader = JsonlReader::<Point>::open(&path).expect("open failed");
459 let all = reader.collect_all().expect("collect_all failed");
460 assert_eq!(all.len(), pts.len());
461 }
462
463 #[test]
464 fn test_parse_error_propagated() {
465 let path = tmp_path("test_parse_err.jsonl");
466 {
467 let mut f = File::create(&path).expect("create");
468 writeln!(f, "not valid json {{{{").expect("write");
469 }
470 let result: Result<Vec<Point>, _> = read_jsonl(&path);
471 assert!(result.is_err());
472 }
473}