Skip to content

Commit 2fe8092

Browse files
committed
Implement JSTableIterator for streaming document reading
1 parent 5ba50ea commit 2fe8092

1 file changed

Lines changed: 101 additions & 32 deletions

File tree

src/jstable.rs

Lines changed: 101 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -50,51 +50,85 @@ impl JSTable {
5050
}
5151
}
5252

53-
pub fn read_jstable(path: &str) -> io::Result<JSTable> {
54-
let file = File::open(path)?;
55-
let mut reader = BufReader::new(file);
53+
pub struct JSTableIterator {
54+
reader: BufReader<File>,
55+
pub timestamp: u64,
56+
pub schema: Schema,
57+
}
5658

57-
// Read Header Length
58-
let mut len_buf = [0u8; 4];
59-
reader.read_exact(&mut len_buf)?;
60-
let header_len = u32::from_le_bytes(len_buf) as usize;
59+
impl JSTableIterator {
60+
pub fn new(path: &str) -> io::Result<Self> {
61+
let file = File::open(path)?;
62+
let mut reader = BufReader::new(file);
6163

62-
// Read Header Blob
63-
let mut header_blob = vec![0u8; header_len];
64-
reader.read_exact(&mut header_blob)?;
65-
66-
let header_val = jsonb::from_slice(&header_blob).map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
67-
// Convert jsonb::Value -> String -> T
68-
let header_str = header_val.to_string();
69-
let header: JSTableHeader = serde_json::from_str(&header_str).map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
64+
// Read Header Length
65+
let mut len_buf = [0u8; 4];
66+
reader.read_exact(&mut len_buf)?;
67+
let header_len = u32::from_le_bytes(len_buf) as usize;
7068

71-
let mut documents = BTreeMap::new();
72-
73-
// Read Records
74-
loop {
75-
match reader.read_exact(&mut len_buf) {
69+
// Read Header Blob
70+
let mut header_blob = vec![0u8; header_len];
71+
reader.read_exact(&mut header_blob)?;
72+
73+
let header_val = jsonb::from_slice(&header_blob).map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
74+
// Convert jsonb::Value -> String -> T
75+
let header_str = header_val.to_string();
76+
let header: JSTableHeader = serde_json::from_str(&header_str).map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
77+
78+
Ok(Self {
79+
reader,
80+
timestamp: header.timestamp,
81+
schema: header.schema,
82+
})
83+
}
84+
}
85+
86+
impl Iterator for JSTableIterator {
87+
type Item = io::Result<(String, Value)>;
88+
89+
fn next(&mut self) -> Option<Self::Item> {
90+
let mut len_buf = [0u8; 4];
91+
match self.reader.read_exact(&mut len_buf) {
7692
Ok(_) => {
7793
let record_len = u32::from_le_bytes(len_buf) as usize;
7894
let mut record_blob = vec![0u8; record_len];
79-
reader.read_exact(&mut record_blob)?;
95+
if let Err(e) = self.reader.read_exact(&mut record_blob) {
96+
return Some(Err(e));
97+
}
8098

81-
let record_val = jsonb::from_slice(&record_blob).map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
99+
let record_val = match jsonb::from_slice(&record_blob).map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e)) {
100+
Ok(v) => v,
101+
Err(e) => return Some(Err(e)),
102+
};
82103
let record_str = record_val.to_string();
83-
let record: (String, Value) = serde_json::from_str(&record_str).map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
104+
let record: (String, Value) = match serde_json::from_str(&record_str).map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e)) {
105+
Ok(v) => v,
106+
Err(e) => return Some(Err(e)),
107+
};
84108

85-
documents.insert(record.0, record.1);
86-
}
87-
Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => {
88-
break;
109+
Some(Ok(record))
89110
}
90-
Err(e) => return Err(e),
111+
Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => None,
112+
Err(e) => Some(Err(e)),
91113
}
92114
}
115+
}
93116

94-
Ok(JSTable {
95-
timestamp: header.timestamp,
96-
schema: header.schema,
97-
documents
117+
pub fn read_jstable(path: &str) -> io::Result<JSTable> {
118+
let mut iterator = JSTableIterator::new(path)?;
119+
let timestamp = iterator.timestamp;
120+
let schema = iterator.schema.clone();
121+
122+
let mut documents = BTreeMap::new();
123+
for result in iterator {
124+
let (id, doc) = result?;
125+
documents.insert(id, doc);
126+
}
127+
128+
Ok(JSTable {
129+
timestamp,
130+
schema,
131+
documents,
98132
})
99133
}
100134

@@ -155,6 +189,41 @@ mod tests {
155189
Ok(())
156190
}
157191

192+
#[test]
193+
fn test_jstable_iterator() -> Result<(), Box<dyn std::error::Error>> {
194+
let schema = Schema {
195+
types: vec![SchemaType::Object],
196+
properties: Some(BTreeMap::from([
197+
("a".to_string(), Schema::new(SchemaType::Integer)),
198+
])),
199+
items: None,
200+
};
201+
let mut documents = BTreeMap::new();
202+
documents.insert("id1".to_string(), json!({"a": 1}));
203+
documents.insert("id2".to_string(), json!({"a": 2}));
204+
let jstable = JSTable::new(12345, schema.clone(), documents.clone());
205+
206+
let file = NamedTempFile::new().unwrap();
207+
jstable.write(file.path().to_str().unwrap()).unwrap();
208+
209+
let mut iterator = JSTableIterator::new(file.path().to_str().unwrap())?;
210+
assert_eq!(iterator.timestamp, 12345);
211+
212+
let mut count = 0;
213+
let mut ids = Vec::new();
214+
for result in iterator {
215+
let (id, doc) = result?;
216+
count += 1;
217+
ids.push(id);
218+
assert!(doc == json!({"a": 1}) || doc == json!({"a": 2}));
219+
}
220+
assert_eq!(count, 2);
221+
assert!(ids.contains(&"id1".to_string()));
222+
assert!(ids.contains(&"id2".to_string()));
223+
224+
Ok(())
225+
}
226+
158227
#[test]
159228
fn test_merge_jstables_conflict_resolution() {
160229
let schema = Schema::new(SchemaType::Object);

0 commit comments

Comments
 (0)