@@ -3,7 +3,7 @@ use serde::{Deserialize, Serialize};
33use serde_json:: Value ;
44use std:: collections:: BTreeMap ;
55use std:: fs:: File ;
6- use std:: io:: { self , BufReader , Read , Write } ;
6+ use std:: io:: { self , BufReader , Read , Seek , SeekFrom , Write } ;
77use xorf:: BinaryFuse8 ;
88
99pub struct JSTable {
@@ -77,16 +77,41 @@ impl JSTable {
7777 summary_file. write_all ( & filter_len. to_le_bytes ( ) ) ?;
7878 summary_file. write_all ( & filter_bytes) ?;
7979
80- // Write Documents to data
80+ // Write Documents to data and build index
81+ let mut index: Vec < ( String , u64 ) > = Vec :: new ( ) ;
82+ let mut current_offset: u64 = 0 ;
83+ let mut bytes_since_last_index: u64 = 0 ;
84+ let mut first = true ;
85+
8186 for ( id, doc) in & self . documents {
87+ // Add index entry if needed
88+ if first || bytes_since_last_index >= 1024 {
89+ index. push ( ( id. clone ( ) , current_offset) ) ;
90+ bytes_since_last_index = 0 ;
91+ first = false ;
92+ }
93+
8294 let record: ( String , & Value ) = ( id. clone ( ) , doc) ;
8395 let record_blob = jsonb:: to_owned_jsonb ( & record)
8496 . map_err ( |e| io:: Error :: new ( io:: ErrorKind :: InvalidData , e) ) ?;
8597 let record_bytes = record_blob. to_vec ( ) ;
8698 let record_len = record_bytes. len ( ) as u32 ;
99+
87100 data_file. write_all ( & record_len. to_le_bytes ( ) ) ?;
88101 data_file. write_all ( & record_bytes) ?;
102+
103+ let written = 4 + record_bytes. len ( ) as u64 ;
104+ current_offset += written;
105+ bytes_since_last_index += written;
89106 }
107+
108+ // Write Index to summary
109+ let index_bytes = serde_json:: to_vec ( & index)
110+ . map_err ( |e| io:: Error :: new ( io:: ErrorKind :: InvalidData , e) ) ?;
111+ let index_len = index_bytes. len ( ) as u32 ;
112+ summary_file. write_all ( & index_len. to_le_bytes ( ) ) ?;
113+ summary_file. write_all ( & index_bytes) ?;
114+
90115 Ok ( ( ) )
91116 }
92117}
@@ -122,7 +147,7 @@ impl JSTableIterator {
122147 let header: JSTableHeader = serde_json:: from_str ( & header_str)
123148 . map_err ( |e| io:: Error :: new ( io:: ErrorKind :: InvalidData , e) ) ?;
124149
125- // We don't need to read the filter here, so we are done with summary file
150+ // We don't need to read the filter or index here
126151
127152 let data_file = File :: open ( data_path) ?;
128153 let data_reader = BufReader :: new ( data_file) ;
@@ -134,6 +159,11 @@ impl JSTableIterator {
134159 schema : header. schema ,
135160 } )
136161 }
162+
163+ pub fn seek ( & mut self , offset : u64 ) -> io:: Result < ( ) > {
164+ self . reader . seek ( SeekFrom :: Start ( offset) ) ?;
165+ Ok ( ( ) )
166+ }
137167}
138168
139169impl Iterator for JSTableIterator {
@@ -223,6 +253,49 @@ pub fn read_filter(path: &str) -> io::Result<BinaryFuse8> {
223253 Ok ( filter)
224254}
225255
256+ pub fn read_index ( path : & str ) -> io:: Result < Vec < ( String , u64 ) > > {
257+ let summary_path = format ! ( "{}.summary" , path) ;
258+ let file = File :: open ( summary_path) ?;
259+ let mut reader = BufReader :: new ( file) ;
260+
261+ // Read Header Length
262+ let mut len_buf = [ 0u8 ; 4 ] ;
263+ reader. read_exact ( & mut len_buf) ?;
264+ let header_len = u32:: from_le_bytes ( len_buf) as usize ;
265+
266+ // Skip Header Blob
267+ io:: copy (
268+ & mut reader. by_ref ( ) . take ( header_len as u64 ) ,
269+ & mut io:: sink ( ) ,
270+ ) ?;
271+
272+ // Read Filter Length
273+ let mut len_buf = [ 0u8 ; 4 ] ;
274+ reader. read_exact ( & mut len_buf) ?;
275+ let filter_len = u32:: from_le_bytes ( len_buf) as usize ;
276+
277+ // Skip Filter Blob
278+ io:: copy (
279+ & mut reader. by_ref ( ) . take ( filter_len as u64 ) ,
280+ & mut io:: sink ( ) ,
281+ ) ?;
282+
283+ // Read Index Length
284+ let mut len_buf = [ 0u8 ; 4 ] ;
285+ reader. read_exact ( & mut len_buf) ?;
286+ let index_len = u32:: from_le_bytes ( len_buf) as usize ;
287+
288+ // Read Index Blob
289+ let mut index_blob = vec ! [ 0u8 ; index_len] ;
290+ reader. read_exact ( & mut index_blob) ?;
291+
292+ // Deserialize
293+ let index: Vec < ( String , u64 ) > = serde_json:: from_slice ( & index_blob)
294+ . map_err ( |e| io:: Error :: new ( io:: ErrorKind :: InvalidData , e) ) ?;
295+
296+ Ok ( index)
297+ }
298+
226299pub fn merge_jstables ( tables : & [ JSTable ] ) -> JSTable {
227300 let mut sorted_tables: Vec < & JSTable > = tables. iter ( ) . collect ( ) ;
228301 sorted_tables. sort_by_key ( |t| t. timestamp ) ;
@@ -435,4 +508,63 @@ mod tests {
435508 assert_eq ! ( keys, vec![ "a" , "b" , "c" ] ) ;
436509 Ok ( ( ) )
437510 }
511+
512+ #[ test]
513+ fn test_read_index ( ) -> Result < ( ) , Box < dyn std:: error:: Error > > {
514+ let schema = Schema :: new ( SchemaType :: Object ) ;
515+ let mut documents = BTreeMap :: new ( ) ;
516+ // Insert enough data to trigger indexing (threshold 1024 bytes)
517+ // Each entry: 4 bytes length + record bytes
518+ // Record: ["id", "val..."]
519+ // We want at least one entry after the first one.
520+
521+ let large_val = "x" . repeat ( 500 ) ; // ~500 bytes
522+ documents. insert ( "a" . to_string ( ) , json ! ( large_val) ) ;
523+ documents. insert ( "b" . to_string ( ) , json ! ( large_val) ) ;
524+ documents. insert ( "c" . to_string ( ) , json ! ( large_val) ) ;
525+ // a: offset 0. write ~500+ -> offset ~500+.
526+ // b: offset ~500+. bytes_written since a ~ 500+. < 1024.
527+ // c: offset ~1000+. bytes_written since a ~ 1000+. >= 1024?
528+ // Let's make it larger.
529+ let larger_val = "x" . repeat ( 1100 ) ;
530+ documents. insert ( "d" . to_string ( ) , json ! ( larger_val) ) ;
531+ documents. insert ( "e" . to_string ( ) , json ! ( 1 ) ) ;
532+
533+ let jstable = JSTable :: new ( 123 , "idx_test" . to_string ( ) , schema, documents) ;
534+ let dir = tempdir ( ) ?;
535+ let path = dir. path ( ) . join ( "idx_table" ) ;
536+ jstable. write ( path. to_str ( ) . unwrap ( ) ) ?;
537+
538+ let index = read_index ( path. to_str ( ) . unwrap ( ) ) ?;
539+
540+ // Should contain at least "a" (first) and "e" (after "d" which is large)
541+ // actually "d" is ~1100.
542+ // a (0), b (large), c (large), d (larger), e (1)
543+ // sorted: a, b, c, d, e
544+
545+ // "a": offset 0.
546+ // write "a" (large). bytes=1100.
547+ // next is "b". bytes_since >= 1024. so "b" is indexed?
548+ // Logic:
549+ // if first || bytes_since >= 1024 { push; bytes=0 }
550+ // "a": first. push ("a", 0). bytes=0.
551+ // write "a" (1100). bytes=1100.
552+ // "b": bytes >= 1024. push ("b", off_b). bytes=0.
553+ // write "b" (1100). bytes=1100.
554+ // "c": bytes >= 1024. push ("c", off_c).
555+
556+ assert ! ( !index. is_empty( ) ) ;
557+ assert_eq ! ( index[ 0 ] . 0 , "a" ) ;
558+ assert_eq ! ( index[ 0 ] . 1 , 0 ) ;
559+
560+ // Check seeking
561+ let mut iter = JSTableIterator :: new ( path. to_str ( ) . unwrap ( ) ) ?;
562+ // Seek to last index entry
563+ let last = index. last ( ) . unwrap ( ) ;
564+ iter. seek ( last. 1 ) ?;
565+ let ( key, _) = iter. next ( ) . unwrap ( ) ?;
566+ assert_eq ! ( key, last. 0 ) ;
567+
568+ Ok ( ( ) )
569+ }
438570}
0 commit comments