11use crate :: schema:: Schema ;
2+ use serde:: { Deserialize , Serialize } ;
23use serde_json:: Value ;
34use std:: collections:: BTreeMap ;
45use std:: io:: { self , BufReader , BufRead , Write } ;
56use std:: fs:: File ;
67
78pub struct JSTable {
9+ pub timestamp : u64 ,
810 pub schema : Schema ,
911 pub documents : BTreeMap < String , Value > ,
1012}
1113
14+ #[ derive( Serialize , Deserialize ) ]
15+ struct JSTableHeader {
16+ timestamp : u64 ,
17+ schema : Schema ,
18+ }
19+
1220impl JSTable {
13- pub fn new ( schema : Schema , documents : BTreeMap < String , Value > ) -> Self {
14- JSTable { schema, documents }
21+ pub fn new ( timestamp : u64 , schema : Schema , documents : BTreeMap < String , Value > ) -> Self {
22+ JSTable { timestamp , schema, documents }
1523 }
1624
1725 pub fn write ( & self , path : & str ) -> io:: Result < ( ) > {
1826 let mut file = File :: create ( path) ?;
19- let schema_json = serde_json:: to_string ( & self . schema ) ?;
20- writeln ! ( file, "{}" , schema_json) ?;
27+ let header = JSTableHeader {
28+ timestamp : self . timestamp ,
29+ schema : self . schema . clone ( ) ,
30+ } ;
31+ let header_json = serde_json:: to_string ( & header) ?;
32+ writeln ! ( file, "{}" , header_json) ?;
2133 for ( id, doc) in & self . documents {
2234 let record: ( String , & Value ) = ( id. clone ( ) , doc) ;
2335 let record_json = serde_json:: to_string ( & record) ?;
@@ -32,8 +44,8 @@ pub fn read_jstable(path: &str) -> io::Result<JSTable> {
3244 let reader = BufReader :: new ( file) ;
3345 let mut lines = reader. lines ( ) ;
3446
35- let schema_line = lines. next ( ) . ok_or_else ( || io:: Error :: new ( io:: ErrorKind :: InvalidData , "Missing schema line" ) ) ??;
36- let schema : Schema = serde_json:: from_str ( & schema_line ) ?;
47+ let header_line = lines. next ( ) . ok_or_else ( || io:: Error :: new ( io:: ErrorKind :: InvalidData , "Missing header line" ) ) ??;
48+ let header : JSTableHeader = serde_json:: from_str ( & header_line ) ?;
3749
3850 let mut documents = BTreeMap :: new ( ) ;
3951 for line_result in lines {
@@ -45,14 +57,25 @@ pub fn read_jstable(path: &str) -> io::Result<JSTable> {
4557 documents. insert ( record. 0 , record. 1 ) ;
4658 }
4759
48- Ok ( JSTable { schema, documents } )
60+ Ok ( JSTable {
61+ timestamp : header. timestamp ,
62+ schema : header. schema ,
63+ documents
64+ } )
4965}
5066
5167pub fn merge_jstables ( tables : & [ JSTable ] ) -> JSTable {
68+ let mut sorted_tables: Vec < & JSTable > = tables. iter ( ) . collect ( ) ;
69+ sorted_tables. sort_by_key ( |t| t. timestamp ) ;
70+
5271 let mut merged_schema = Schema :: new ( crate :: schema:: SchemaType :: Object ) ;
5372 let mut merged_documents = BTreeMap :: new ( ) ;
73+ let mut max_timestamp = 0 ;
5474
55- for table in tables {
75+ for table in sorted_tables {
76+ if table. timestamp > max_timestamp {
77+ max_timestamp = table. timestamp ;
78+ }
5679 merged_schema. merge ( table. schema . clone ( ) ) ;
5780 for ( id, doc) in & table. documents {
5881 merged_documents. insert ( id. clone ( ) , doc. clone ( ) ) ;
@@ -61,7 +84,7 @@ pub fn merge_jstables(tables: &[JSTable]) -> JSTable {
6184
6285 merged_documents. retain ( |_, v| !v. is_null ( ) ) ;
6386
64- JSTable :: new ( merged_schema, merged_documents)
87+ JSTable :: new ( max_timestamp , merged_schema, merged_documents)
6588}
6689
6790#[ cfg( test) ]
@@ -75,19 +98,25 @@ mod tests {
7598 #[ test]
7699 fn test_read_jstable ( ) -> Result < ( ) , Box < dyn std:: error:: Error > > {
77100 let mut file = NamedTempFile :: new ( ) . unwrap ( ) ;
78- let schema_json = serde_json :: to_string ( & Schema {
101+ let schema = Schema {
79102 types : vec ! [ SchemaType :: Object ] ,
80103 properties : Some ( BTreeMap :: from ( [
81104 ( "a" . to_string ( ) , Schema :: new ( SchemaType :: Integer ) ) ,
82105 ] ) ) ,
83106 items : None ,
84- } ) . unwrap ( ) ;
85- writeln ! ( file, "{}" , schema_json) . unwrap ( ) ;
107+ } ;
108+ let header = JSTableHeader {
109+ timestamp : 12345 ,
110+ schema : schema,
111+ } ;
112+ let header_json = serde_json:: to_string ( & header) . unwrap ( ) ;
113+ writeln ! ( file, "{}" , header_json) . unwrap ( ) ;
86114 writeln ! ( file, "{}" , serde_json:: to_string( & json!( [ "id1" , { "a" : 1 } ] ) ) . unwrap( ) ) . unwrap ( ) ;
87115 writeln ! ( file, "{}" , serde_json:: to_string( & json!( [ "id2" , { "a" : 2 } ] ) ) . unwrap( ) ) . unwrap ( ) ;
88116
89117 let jstable = read_jstable ( file. path ( ) . to_str ( ) . unwrap ( ) ) . unwrap ( ) ;
90118
119+ assert_eq ! ( jstable. timestamp, 12345 ) ;
91120 assert_eq ! ( jstable. schema. types, vec![ SchemaType :: Object ] ) ;
92121 assert_eq ! ( jstable. documents. len( ) , 2 ) ;
93122 assert_eq ! ( * jstable. documents. get( "id1" ) . unwrap( ) , json!( { "a" : 1 } ) ) ;
@@ -107,16 +136,53 @@ mod tests {
107136 let mut documents = BTreeMap :: new ( ) ;
108137 documents. insert ( "id1" . to_string ( ) , json ! ( { "a" : 1 } ) ) ;
109138 documents. insert ( "id2" . to_string ( ) , json ! ( { "a" : 2 } ) ) ;
110- let jstable = JSTable :: new ( schema, documents) ;
139+ let jstable = JSTable :: new ( 67890 , schema, documents) ;
111140
112141 let file = NamedTempFile :: new ( ) . unwrap ( ) ;
113142 jstable. write ( file. path ( ) . to_str ( ) . unwrap ( ) ) . unwrap ( ) ;
114143
115144 let content = std:: fs:: read_to_string ( file. path ( ) ) . unwrap ( ) ;
116- assert ! ( content. contains( "{\" type\" :\" object\" ,\" properties\" :{\" a\" :{\" type\" :[\" integer\" ]}}}" ) ) ;
145+ let lines: Vec < & str > = content. lines ( ) . collect ( ) ;
146+ let header: JSTableHeader = serde_json:: from_str ( lines[ 0 ] ) ?;
147+
148+ assert_eq ! ( header. timestamp, 67890 ) ;
149+ assert_eq ! ( header. schema. types, vec![ SchemaType :: Object ] ) ;
150+
117151 assert ! ( content. contains( "[\" id1\" ,{\" a\" :1}]" ) ) ;
118152 assert ! ( content. contains( "[\" id2\" ,{\" a\" :2}]" ) ) ;
119153
120154 Ok ( ( ) )
121155 }
156+
157+ #[ test]
158+ fn test_merge_jstables_conflict_resolution ( ) {
159+ let schema = Schema :: new ( SchemaType :: Object ) ;
160+
161+ let mut docs1 = BTreeMap :: new ( ) ;
162+ docs1. insert ( "id1" . to_string ( ) , json ! ( { "v" : 1 } ) ) ;
163+ let t1 = JSTable :: new ( 100 , schema. clone ( ) , docs1) ;
164+
165+ let mut docs2 = BTreeMap :: new ( ) ;
166+ docs2. insert ( "id1" . to_string ( ) , json ! ( { "v" : 2 } ) ) ;
167+ let t2 = JSTable :: new ( 200 , schema. clone ( ) , docs2) ;
168+
169+ // Case 1: t1 (older) then t2 (newer) in the slice
170+ // Note: Creating array [t1, t2] moves them.
171+ let merged = merge_jstables ( & [ t1, t2] ) ;
172+ assert_eq ! ( * merged. documents. get( "id1" ) . unwrap( ) , json!( { "v" : 2 } ) ) ;
173+ assert_eq ! ( merged. timestamp, 200 ) ;
174+
175+ // Case 2: Reverse order
176+ let mut docs1 = BTreeMap :: new ( ) ;
177+ docs1. insert ( "id1" . to_string ( ) , json ! ( { "v" : 1 } ) ) ;
178+ let t1b = JSTable :: new ( 100 , schema. clone ( ) , docs1) ;
179+
180+ let mut docs2 = BTreeMap :: new ( ) ;
181+ docs2. insert ( "id1" . to_string ( ) , json ! ( { "v" : 2 } ) ) ;
182+ let t2b = JSTable :: new ( 200 , schema. clone ( ) , docs2) ;
183+
184+ let merged_reverse = merge_jstables ( & [ t2b, t1b] ) ;
185+ assert_eq ! ( * merged_reverse. documents. get( "id1" ) . unwrap( ) , json!( { "v" : 2 } ) ) ;
186+ assert_eq ! ( merged_reverse. timestamp, 200 ) ;
187+ }
122188}
0 commit comments