diff --git a/bench/bench_ocaml.ml b/bench/bench_ocaml.ml index a4e2278..dd6ceda 100644 --- a/bench/bench_ocaml.ml +++ b/bench/bench_ocaml.ml @@ -382,7 +382,7 @@ let logseq_get_page_data db name = let build_storage_db size = let storage = memory_storage () in let db = db_with (people size) (empty_db ~schema ~storage ()) in - store db; + ignore (store db); match restore storage with | Some db -> db | None -> failwith "storage-backed benchmark db should restore" diff --git a/bench/codec_parity.ml b/bench/codec_parity.ml new file mode 100644 index 0000000..0ae721c --- /dev/null +++ b/bench/codec_parity.ml @@ -0,0 +1,38 @@ +(* Byte-parity check: decode cljs-written kvs rows, re-encode with our + codec, and compare bytes. Usage: codec_parity [max-rows] *) + +let () = + let db_path = Sys.argv.(1) in + let max_rows = if Array.length Sys.argv > 2 then int_of_string Sys.argv.(2) else 50 in + let db = Sqlite3.db_open db_path in + let rows = ref [] in + (match + Sqlite3.exec db ~cb:(fun row _headers -> + rows := (Option.value ~default:"" row.(0), Option.value ~default:"" row.(1)) :: !rows) + "select addr, content from kvs order by rowid" + with + | Sqlite3.Rc.OK -> () + | rc -> Printf.eprintf "sqlite error: %s\n" (Sqlite3.Rc.to_string rc); exit 1); + let rows = List.rev !rows in + let total = ref 0 and identical = ref 0 and diffs = ref [] in + List.iteri + (fun i (addr, content) -> + if i < max_rows then ( + incr total; + match Datascript_sqlite_codec.decode content with + | payload -> + let re = Datascript_sqlite_codec.encode payload in + if String.equal re content then incr identical + else + diffs := (addr, String.length content, String.length re, content, re) :: !diffs + | exception exn -> + diffs := (addr, -1, -1, content, Printexc.to_string exn) :: !diffs)) + rows; + Printf.printf "%d/%d rows byte-identical\n" !identical !total; + List.iter + (fun (addr, clen, rlen, c, r) -> + Printf.printf "\n--- addr %s (cljs %dB, ours %dB)\ncljs: %s\nours: %s\n" addr clen + rlen + (if String.length c > 300 then String.sub c 0 300 ^ "..." else c) + (if String.length r > 300 then String.sub r 0 300 ^ "..." else r)) + (List.rev (if List.length !diffs > 5 then List.filteri (fun i _ -> i < 5) !diffs else !diffs)) diff --git a/bench/dune b/bench/dune index 67c4a83..6815c7a 100644 --- a/bench/dune +++ b/bench/dune @@ -70,3 +70,15 @@ ../script/benchmark_gate_vs_cljs_check_test.sh) (action (run bash %{dep:../script/benchmark_gate_vs_cljs_check_test.sh}))) + +(executable + (name write_probe) + (modules write_probe) + (modes exe) + (libraries datascript-ocaml-native datascript_sqlite unix sqlite3)) + +(executable + (name codec_parity) + (modules codec_parity) + (modes exe) + (libraries datascript-ocaml-native datascript_sqlite unix sqlite3)) diff --git a/bench/persistent_sqlite.ml b/bench/persistent_sqlite.ml index 15717eb..803f4c5 100644 --- a/bench/persistent_sqlite.ml +++ b/bench/persistent_sqlite.ml @@ -126,7 +126,7 @@ let run_size size = let persistent_build, persistent_db = time "snapshot-build-and-store" (fun () -> let db = db_with tx (empty_db ~schema ~storage ()) in - store db; + ignore (store db); db) in print_timing persistent_build; @@ -152,7 +152,7 @@ let run_size size = (add_block_tx "persistent-new" (Float.of_int (size + 1))) restored_db in - store db; + ignore (store db); db) in print_timing persistent_add; @@ -163,7 +163,7 @@ let run_size size = let persistent_update, restored_db = time "snapshot-update-one-and-store-after-add" (fun () -> let db = db_with (update_content_tx "block-00001" "Edited") restored_db in - store db; + ignore (store db); db) in print_timing persistent_update; diff --git a/bench/write_probe.ml b/bench/write_probe.ml new file mode 100644 index 0000000..0daa49f --- /dev/null +++ b/bench/write_probe.ml @@ -0,0 +1,67 @@ +open Datascript + +let indexed = + { + cardinality = One; + unique = None; + indexed = true; + is_component = false; + no_history = false; + doc = None; + value_type = None; + tuple_attrs = None; + tuple_types = None; + } + +let unique_identity = { indexed with unique = Some Identity } + +let schema = + [ + ("block/id", unique_identity); + ("block/content", indexed); + ("block/order", indexed); + ] + +let block_tx base count = + List.init count (fun index -> + let i = base + index in + Entity + { + db_id = Some (Temp_id (Printf.sprintf "b-%d" i)); + attrs = + [ + ("block/id", One_value (String (Printf.sprintf "b-%d" i))); + ("block/content", One_value (String (Printf.sprintf "Block %d" i))); + ("block/order", One_value (Float (Float.of_int i))); + ]; + }) + +let kvs_rows db_path = + let db = Sqlite3.db_open db_path in + let stmt = Sqlite3.prepare db "select count(*) from kvs" in + let n = + match Sqlite3.step stmt with + | Sqlite3.Rc.ROW -> Sqlite3.column_int stmt 0 + | rc -> failwith (Sqlite3.Rc.to_string rc) + in + ignore (Sqlite3.finalize stmt); + ignore (Sqlite3.db_close db); + n + +let () = + let batches = int_of_string Sys.argv.(1) in + let per_batch = int_of_string Sys.argv.(2) in + let db_path = Filename.temp_file "write-probe" ".sqlite3" in + let session = Datascript_sqlite.open_session db_path in + let storage = Datascript_sqlite.storage session in + let conn = create_conn ~schema ~storage () in + let mutable_rows = ref (kvs_rows db_path) in + Printf.printf "start rows=%d\n%!" !mutable_rows; + for b = 1 to batches do + ignore (transact_conn conn (block_tx (b * 10000) per_batch)); + let rows = kvs_rows db_path in + Printf.printf "batch %d: rows=%d (+%d)\n%!" b rows (rows - !mutable_rows); + mutable_rows := rows + done; + Datascript_sqlite.close session; + Sys.remove db_path diff --git a/examples/logseq_sqlite_storage.ml b/examples/logseq_sqlite_storage.ml index 1769136..0c0fe9a 100644 --- a/examples/logseq_sqlite_storage.ml +++ b/examples/logseq_sqlite_storage.ml @@ -218,7 +218,7 @@ let transit_of_tuple_attrs attrs = let transit_of_tuple_types types = Transit.Array (List.map transit_of_value_type types) -let schema_attr_to_transit attr = +let schema_attr_to_transit ~ident_backed attr ident = let entries = ref [] in let add key value = entries := (Transit.Keyword key, value) :: !entries in if attr.cardinality <> One then add "db/cardinality" (transit_of_cardinality attr.cardinality); @@ -230,13 +230,20 @@ let schema_attr_to_transit attr = Option.iter (fun value_type -> add "db/valueType" (transit_of_value_type value_type)) attr.value_type; Option.iter (fun attrs -> add "db/tupleAttrs" (transit_of_tuple_attrs attrs)) attr.tuple_attrs; Option.iter (fun types -> add "db/tupleTypes" (transit_of_tuple_types types)) attr.tuple_types; + (* cljs update-schema merges {:db/ident ident} into the attr's spec map *) + if ident_backed then add "db/ident" (Transit.Keyword ident); Transit.Map (List.rev !entries) -let schema_to_transit schema = +let schema_to_transit ?(eids = []) schema = Transit.Map - (schema - |> List.map (fun (attr, schema_attr) -> - Transit.Keyword attr, schema_attr_to_transit schema_attr)) + ((schema + |> List.map (fun (attr, schema_attr) -> + ( Transit.Keyword attr + , schema_attr_to_transit + ~ident_backed:(List.exists (fun (_, ident) -> String.equal ident attr) eids) + schema_attr + attr ))) + @ List.map (fun (eid, ident) -> Transit.Int eid, Transit.Keyword ident) eids) let rec value_to_transit = function | Nil -> Transit.Null @@ -283,7 +290,7 @@ let optional_metadata_entry key = function let storage_root_to_transit root = Transit.Map - ([ Transit.Keyword "schema", schema_to_transit root.storage_schema + ([ Transit.Keyword "schema", schema_to_transit ~eids:root.storage_schema_idents root.storage_schema ; Transit.Keyword "max-eid", Transit.Int root.storage_max_eid ; Transit.Keyword "max-tx", Transit.Int root.storage_max_tx ; Transit.Keyword "eavt", Transit.Int (sqlite_addr_of_storage_address root.storage_eavt) @@ -389,6 +396,15 @@ let schema_of_transit = function | None -> None) | _ -> [] +let schema_eids_of_transit = function + | Transit.Map entries -> + entries + |> List.filter_map (fun (key, value) -> + match int_of_transit_value key, keyword_of_transit value with + | Some eid, Some ident -> Some (eid, ident) + | _ -> None) + | _ -> [] + let ref_type_of_transit = function | Transit.Keyword "weak" -> PSet.Weak | Transit.Keyword "strong" | _ -> PSet.Strong @@ -472,6 +488,7 @@ let storage_root_of_transit entries = | _ -> None in { storage_schema = schema_of_transit (find "schema") + ; storage_schema_idents = schema_eids_of_transit (find "schema") ; storage_max_eid = int_of_transit "storage root :max-eid" (find "max-eid") ; storage_max_tx = int_of_transit "storage root :max-tx" (find "max-tx") ; storage_eavt = diff --git a/examples/sqlite_storage_example.ml b/examples/sqlite_storage_example.ml index 401be70..234bfb9 100644 --- a/examples/sqlite_storage_example.ml +++ b/examples/sqlite_storage_example.ml @@ -43,7 +43,7 @@ let run_roundtrip db_path = ~schema:[ "name", indexed ] [ datom ~e:1 ~a:"name" ~v:(String "SQLite example") () ] in - store ~storage db; + ignore (store ~storage db); match restore (Storage.storage db_path) with | None -> failwith "failed to restore SQLite-backed db" | Some restored -> diff --git a/impl/conn.ml b/impl/conn.ml index b3ec9c6..491c246 100644 --- a/impl/conn.ml +++ b/impl/conn.ml @@ -11,11 +11,11 @@ type t = type creation_context = { empty_db : ?schema:schema -> ?storage:storage -> unit -> db ; init_db : ?schema:schema -> ?storage:storage -> datom list -> db - ; store : ?storage:storage -> db -> unit + ; store : ?storage:storage -> db -> db } type schema_context = - { store : ?storage:storage -> db -> unit + { store : ?storage:storage -> db -> db ; with_schema : db -> schema -> db } @@ -25,7 +25,7 @@ type restore_context = } type transact_context = - { store : ?storage:storage -> db -> unit + { store : ?storage:storage -> db -> db ; store_tail : storage -> datom list list -> unit ; storage_tail_datom_count : datom list list -> int ; storage_tail_compaction_threshold : db -> int @@ -33,14 +33,14 @@ type transact_context = } type reset_context = - { store : ?storage:storage -> db -> unit + { store : ?storage:storage -> db -> db ; datoms : db -> datom list } type context = { empty_db : ?schema:schema -> ?storage:storage -> unit -> db ; init_db : ?schema:schema -> ?storage:storage -> datom list -> db - ; store : ?storage:storage -> db -> unit + ; store : ?storage:storage -> db -> db ; store_tail : storage -> datom list list -> unit ; restore : storage -> db option ; restore_tail_groups : storage -> datom list list @@ -78,8 +78,7 @@ let create (context : creation_context) ?schema ?storage () = match storage with | None -> db | Some storage -> - context.store ~storage db; - { db with storage_ref = Some storage } + { (context.store ~storage db) with storage_ref = Some storage } in make ?storage db @@ -87,8 +86,7 @@ let from_db (context : creation_context) db = match db.storage_ref with | None -> make db | Some storage -> - context.store ~storage db; - make ~storage db + make ~storage (context.store ~storage db) let from_datoms (context : creation_context) ?schema ?storage datoms = from_db context (context.init_db ?schema ?storage datoms) @@ -129,7 +127,7 @@ let reset_schema (context : schema_context) conn schema = (match conn.storage with | None -> () | Some storage -> - context.store ~storage db; + conn.db <- context.store ~storage db; conn.storage_tail <- []); db @@ -149,7 +147,7 @@ let transact (context : transact_context) ?(tx_meta = []) conn tx_data = if report.tx_data <> [] then begin let tail = conn.storage_tail @ [ report.tx_data ] in if context.storage_tail_datom_count tail > context.storage_tail_compaction_threshold report.db_after then begin - context.store ~storage report.db_after; + conn.db <- context.store ~storage report.db_after; conn.storage_tail <- [] end else begin conn.storage_tail <- tail; @@ -177,7 +175,7 @@ let apply_report (context : transact_context) conn (report : tx_report) = if report.tx_data <> [] then begin let tail = conn.storage_tail @ [ report.tx_data ] in if context.storage_tail_datom_count tail > context.storage_tail_compaction_threshold db_after then begin - context.store ~storage db_after; + conn.db <- context.store ~storage db_after; conn.storage_tail <- [] end else begin conn.storage_tail <- tail; @@ -203,7 +201,7 @@ let reset (context : reset_context) ?(tx_meta = []) conn db = (match conn.storage with | None -> () | Some storage -> - context.store ~storage db; + conn.db <- context.store ~storage db; conn.storage_tail <- []); notify_listeners conn report; db diff --git a/impl/conn.mli b/impl/conn.mli index 94b6772..bc45688 100644 --- a/impl/conn.mli +++ b/impl/conn.mli @@ -5,11 +5,11 @@ type t type creation_context = { empty_db : ?schema:schema -> ?storage:storage -> unit -> db ; init_db : ?schema:schema -> ?storage:storage -> datom list -> db - ; store : ?storage:storage -> db -> unit + ; store : ?storage:storage -> db -> db } type schema_context = - { store : ?storage:storage -> db -> unit + { store : ?storage:storage -> db -> db ; with_schema : db -> schema -> db } @@ -19,7 +19,7 @@ type restore_context = } type transact_context = - { store : ?storage:storage -> db -> unit + { store : ?storage:storage -> db -> db ; store_tail : storage -> datom list list -> unit ; storage_tail_datom_count : datom list list -> int ; storage_tail_compaction_threshold : db -> int @@ -27,14 +27,14 @@ type transact_context = } type reset_context = - { store : ?storage:storage -> db -> unit + { store : ?storage:storage -> db -> db ; datoms : db -> datom list } type context = { empty_db : ?schema:schema -> ?storage:storage -> unit -> db ; init_db : ?schema:schema -> ?storage:storage -> datom list -> db - ; store : ?storage:storage -> db -> unit + ; store : ?storage:storage -> db -> db ; store_tail : storage -> datom list list -> unit ; restore : storage -> db option ; restore_tail_groups : storage -> datom list list diff --git a/impl/datascript.ml b/impl/datascript.ml index 7b72d4c..9527f9b 100644 --- a/impl/datascript.ml +++ b/impl/datascript.ml @@ -900,15 +900,19 @@ let tx_meta_skips_store tx_meta = tx_meta let persist_transact_tail ~tx_meta db tx_data = - if tx_data <> [] && not (tx_meta_skips_store tx_meta) then + if tx_data = [] || tx_meta_skips_store tx_meta then + db + else match db.storage_ref with - | None -> () + | None -> db | Some storage -> let tail = restore_tail_groups storage @ [ tx_data ] in if storage_tail_datom_count tail > storage_tail_compaction_threshold db then store ~storage db - else - store_tail storage tail + else begin + store_tail storage tail; + db + end let transact_report ?(tx_meta = []) db tx_ops = let db_after, tempids, tx_data = apply_tx tx_ops db in @@ -916,8 +920,8 @@ let transact_report ?(tx_meta = []) db tx_ops = let transact ?(tx_meta = []) db tx_ops = let report = transact_report ~tx_meta db tx_ops in - persist_transact_tail ~tx_meta report.db_after report.tx_data; - report + let db_after = persist_transact_tail ~tx_meta report.db_after report.tx_data in + { report with db_after } let with_tx ?tx_meta db tx_ops = transact ?tx_meta db tx_ops diff --git a/impl/datascript.mli b/impl/datascript.mli index 68f43cc..76b39d8 100644 --- a/impl/datascript.mli +++ b/impl/datascript.mli @@ -81,11 +81,11 @@ module Conn : sig type creation_context = { empty_db : ?schema:schema -> ?storage:storage -> unit -> db ; init_db : ?schema:schema -> ?storage:storage -> datom list -> db - ; store : ?storage:storage -> db -> unit + ; store : ?storage:storage -> db -> db } type schema_context = - { store : ?storage:storage -> db -> unit + { store : ?storage:storage -> db -> db ; with_schema : db -> schema -> db } @@ -95,7 +95,7 @@ module Conn : sig } type transact_context = - { store : ?storage:storage -> db -> unit + { store : ?storage:storage -> db -> db ; store_tail : storage -> datom list list -> unit ; storage_tail_datom_count : datom list list -> int ; storage_tail_compaction_threshold : db -> int @@ -103,7 +103,7 @@ module Conn : sig } type reset_context = - { store : ?storage:storage -> db -> unit + { store : ?storage:storage -> db -> db ; datoms : db -> datom list } @@ -248,7 +248,7 @@ module Storage : sig val tail_address : storage_address val memory_storage : unit -> storage val file_storage : string -> storage - val store : ?storage:storage -> db -> unit + val store : ?storage:storage -> db -> db val store_tail : storage -> datom list list -> unit val tail_compaction_threshold : db -> int val tail_datom_count : datom list list -> int @@ -414,7 +414,7 @@ val from_serializable : serializable_db -> db val db_from_reader_string : string -> db val memory_storage : unit -> storage val file_storage : string -> storage -val store : ?storage:storage -> db -> unit +val store : ?storage:storage -> db -> db val store_tail : storage -> datom list list -> unit val restore : storage -> db option val db_with_tail : db -> datom list list -> db diff --git a/impl/storage.ml b/impl/storage.ml index 14d919e..bd63217 100644 --- a/impl/storage.ml +++ b/impl/storage.ml @@ -51,8 +51,12 @@ let file_storage = Platform.file_storage let buffered_node_storage pending_entries = { PSet.store_node = - (fun node -> - let address = next_storage_address () in + (fun ?address node -> + let address = + match address with + | Some address -> address + | None -> next_storage_address () + in pending_entries := (address, Storage_node node) :: !pending_entries; address) ; restore_node = (fun _address -> None) @@ -105,8 +109,12 @@ let normalize_stored_tail schema = let restoring_node_storage ?schema storage = { PSet.store_node = - (fun node -> - let address = next_storage_address () in + (fun ?address node -> + let address = + match address with + | Some address -> address + | None -> next_storage_address () + in storage.storage_store [ address, Storage_node node ]; address) ; restore_node = @@ -157,7 +165,18 @@ let index_metadata pending_entries storage index_set root_address = let root_of_stored_indexes db ~eavt_metadata ~aevt_metadata ~avet_metadata eavt_address aevt_address avet_address = let settings = PSet.settings db.eavt_index in + let schema_idents = + (* eid -> :db/ident pairs, matching cljs's schema map entries *) + let first_ident = { e = 0; a = "db/ident"; v = Nil; tx = 0; added = true } in + db.aevt_index + |> PSet.slice_seq ~from_:first_ident + |> PSet.to_seq + |> Seq.take_while (fun d -> String.equal d.a "db/ident") + |> Seq.filter_map (fun d -> match d.v with Keyword ident -> Some (d.e, ident) | _ -> None) + |> List.of_seq + in { storage_schema = db.schema + ; storage_schema_idents = schema_idents ; storage_max_eid = db.max_eid ; storage_max_tx = db.max_tx ; storage_eavt = eavt_address @@ -202,8 +221,7 @@ let store_index node_storage index index_set = match PSet.store index_set with | address, _ -> address | exception Invalid_argument message when String.equal message "store requires a storage-backed set" -> - let storage_backed = storage_backed_index node_storage index index_set in - fst (PSet.store storage_backed) + fst (PSet.store (storage_backed_index node_storage index index_set)) let store_to_storage db storage = let pending_entries = ref [] in @@ -222,7 +240,29 @@ let store_to_storage db storage = (List.rev !pending_entries @ [ root_address, Storage_root root ; tail_address, Storage_tail [] - ]) + ]); + (* store_node_tree can't write addresses back into immutable nodes, so + reinstall each index as a lazy Deferred set over the freshly stored + root — cljs achieves the same by flagging nodes stored in place. + Rebuild via PSet.restore with the real (readable) node storage: + the set PSet.store returns inherits whatever storage the input set + carried, which for sets built by storage_backed_index is this call's + throwaway buffered storage — adopting it would route later node + writes into a dead buffer and lose them. *) + let adopt_node_storage = restoring_node_storage ~schema:db.schema storage in + let adopt index address metadata = + match + PSet.restore ~cmp:(Util.compare_datom index) ~settings:(settings_of_root root) + ~count:metadata.storage_index_count adopt_node_storage address + with + | Some index -> index + | None -> invalid_arg "stored index failed to restore" + in + { db with + eavt_index = adopt Eavt eavt_address eavt_metadata + ; aevt_index = adopt Aevt aevt_address aevt_metadata + ; avet_index = adopt Avet avet_address avet_metadata + } let store ?storage db = match storage, db.storage_ref with diff --git a/impl/storage.mli b/impl/storage.mli index 9cdc992..afc094e 100644 --- a/impl/storage.mli +++ b/impl/storage.mli @@ -13,7 +13,7 @@ val root_address : storage_address val tail_address : storage_address val memory_storage : unit -> storage val file_storage : string -> storage -val store : ?storage:storage -> db -> unit +val store : ?storage:storage -> db -> db val store_tail : storage -> datom list list -> unit val normalize_stored_datom : schema -> datom -> datom val tail_compaction_threshold : db -> int diff --git a/melange/datascript_melange_storage.ml b/melange/datascript_melange_storage.ml index bc3a1d6..8b7caea 100644 --- a/melange/datascript_melange_storage.ml +++ b/melange/datascript_melange_storage.ml @@ -109,7 +109,7 @@ let transit_of_tuple_attrs attrs = let transit_of_tuple_types types = Transit.Array (List.map transit_of_value_type types) -let schema_attr_to_transit attr = +let schema_attr_to_transit ~ident_backed attr ident = let entries = ref [] in let add key value = entries := (Transit.Keyword key, value) :: !entries in (match attr.cardinality with @@ -125,13 +125,21 @@ let schema_attr_to_transit attr = attr.value_type; Option.iter (fun attrs -> add "db/tupleAttrs" (transit_of_tuple_attrs attrs)) attr.tuple_attrs; Option.iter (fun types -> add "db/tupleTypes" (transit_of_tuple_types types)) attr.tuple_types; + (* cljs update-schema merges {:db/ident ident} into the attr's spec map *) + if ident_backed then add "db/ident" (Transit.Keyword ident); Transit.Map (List.rev !entries) -let schema_to_transit schema = +let schema_to_transit ?(eids = []) schema = Transit.Map (List.map - (fun (attr, schema_attr) -> (Transit.Keyword attr, schema_attr_to_transit schema_attr)) - schema) + (fun (attr, schema_attr) -> + ( Transit.Keyword attr + , schema_attr_to_transit + ~ident_backed:(List.exists (fun (_, ident) -> String.equal ident attr) eids) + schema_attr + attr )) + schema + @ List.map (fun (eid, ident) -> (Transit.Int eid, Transit.Keyword ident)) eids) let tuple_attrs_of_transit = function | Transit.Array values | Transit.List values -> Some (List.filter_map keyword_of_transit values) @@ -174,6 +182,16 @@ let schema_of_transit = function entries | _ -> [] +let schema_eids_of_transit = function + | Transit.Map entries -> + List.filter_map + (fun (key, value) -> + match int_of_transit_value key, keyword_of_transit value with + | Some eid, Some ident -> Some (eid, ident) + | _ -> None) + entries + | _ -> [] + let rec value_to_transit = function | Ds.Nil -> Transit.Null | Int64 value -> Transit.Int64 value @@ -267,7 +285,7 @@ let optional_metadata_entry key = function let storage_root_to_transit root = Transit.Map ([ - (Transit.Keyword "schema", schema_to_transit root.storage_schema); + (Transit.Keyword "schema", schema_to_transit ~eids:root.storage_schema_idents root.storage_schema); (Transit.Keyword "max-eid", Transit.Int root.storage_max_eid); (Transit.Keyword "max-tx", Transit.Int root.storage_max_tx); (Transit.Keyword "eavt", address_to_transit root.storage_eavt); @@ -328,6 +346,7 @@ let optional_metadata key entries = let storage_root_of_transit entries = { Ds.storage_schema = schema_of_transit (require_key "schema" entries); + storage_schema_idents = schema_eids_of_transit (require_key "schema" entries); storage_max_eid = int_of_transit "storage root :max-eid" (require_key "max-eid" entries); storage_max_tx = int_of_transit "storage root :max-tx" (require_key "max-tx" entries); storage_eavt = address_of_transit "storage root :eavt" (require_key "eavt" entries); diff --git a/sqlite/datascript_sqlite_codec.ml b/sqlite/datascript_sqlite_codec.ml index 12bc229..4bd89d4 100644 --- a/sqlite/datascript_sqlite_codec.ml +++ b/sqlite/datascript_sqlite_codec.ml @@ -109,7 +109,7 @@ let transit_of_tuple_attrs attrs = let transit_of_tuple_types types = Transit.Array (List.map transit_of_value_type types) -let schema_attr_to_transit attr = +let schema_attr_to_transit ~ident_backed attr ident = let entries = ref [] in let add key value = entries := (Transit.Keyword key, value) :: !entries in (match attr.cardinality with @@ -125,13 +125,21 @@ let schema_attr_to_transit attr = attr.value_type; Option.iter (fun attrs -> add "db/tupleAttrs" (transit_of_tuple_attrs attrs)) attr.tuple_attrs; Option.iter (fun types -> add "db/tupleTypes" (transit_of_tuple_types types)) attr.tuple_types; + (* cljs update-schema merges {:db/ident ident} into the attr's spec map *) + if ident_backed then add "db/ident" (Transit.Keyword ident); Transit.Map (List.rev !entries) -let schema_to_transit schema = +let schema_to_transit ?(eids = []) schema = Transit.Map (List.map - (fun (attr, schema_attr) -> (Transit.Keyword attr, schema_attr_to_transit schema_attr)) - schema) + (fun (attr, schema_attr) -> + ( Transit.Keyword attr + , schema_attr_to_transit + ~ident_backed:(List.exists (fun (_, ident) -> String.equal ident attr) eids) + schema_attr + attr )) + schema + @ List.map (fun (eid, ident) -> (Transit.Int eid, Transit.Keyword ident)) eids) let tuple_attrs_of_transit = function | Transit.Array values | Transit.List values -> Some (List.filter_map keyword_of_transit values) @@ -174,6 +182,16 @@ let schema_of_transit = function entries | _ -> [] +let schema_eids_of_transit = function + | Transit.Map entries -> + List.filter_map + (fun (key, value) -> + match int_of_transit_value key, keyword_of_transit value with + | Some eid, Some ident -> Some (eid, ident) + | _ -> None) + entries + | _ -> [] + let rec value_to_transit = function | Ds.Nil -> Transit.Null | Int64 value -> Transit.Int64 value @@ -267,7 +285,7 @@ let optional_metadata_entry key = function let storage_root_to_transit root = Transit.Map ([ - (Transit.Keyword "schema", schema_to_transit root.storage_schema); + (Transit.Keyword "schema", schema_to_transit ~eids:root.storage_schema_idents root.storage_schema); (Transit.Keyword "max-eid", Transit.Int root.storage_max_eid); (Transit.Keyword "max-tx", Transit.Int root.storage_max_tx); (Transit.Keyword "eavt", address_to_transit root.storage_eavt); @@ -328,6 +346,7 @@ let optional_metadata key entries = let storage_root_of_transit entries = { Ds.storage_schema = schema_of_transit (require_key "schema" entries); + storage_schema_idents = schema_eids_of_transit (require_key "schema" entries); storage_max_eid = int_of_transit "storage root :max-eid" (require_key "max-eid" entries); storage_max_tx = int_of_transit "storage root :max-tx" (require_key "max-tx" entries); storage_eavt = address_of_transit "storage root :eavt" (require_key "eavt" entries); @@ -367,5 +386,7 @@ let payload_of_transit = function | (Transit.Array _ | Transit.List _) as tail -> Storage_tail (storage_tail_of_transit tail) | _ -> invalid_arg "unknown storage payload" -let encode payload = payload |> payload_to_transit |> Transit.to_string ~mode:Transit.Verbose +(* cljs datascript writes storage payloads in transit's non-verbose JSON + mode: maps as ["^ ",k,v] arrays with the write cache's ^N refs *) +let encode payload = payload |> payload_to_transit |> Transit.to_string ~mode:Transit.Normal let decode content = content |> Transit.of_string |> payload_of_transit diff --git a/test/js_smoke.ml b/test/js_smoke.ml index d675a8e..cf92a67 100644 --- a/test/js_smoke.ml +++ b/test/js_smoke.ml @@ -40,7 +40,7 @@ let () = incr reads; base.storage_restore address) } in let schema = schema_of_edn_string "{:name {:db/index true}}" in let db = init_db ~schema [datom ~e:1 ~a:"name" ~v:(String "v") ()] in - store ~storage db; + ignore (store ~storage db); let restored = Option.get (restore storage) in reads := 0; List.iter (fun set -> diff --git a/test/melange_smoke.ml b/test/melange_smoke.ml index 29024ab..0a233d8 100644 --- a/test/melange_smoke.ml +++ b/test/melange_smoke.ml @@ -33,7 +33,7 @@ let () = incr reads; base.storage_restore address) } in let schema = schema_of_edn_string "{:name {:db/index true}}" in let db = init_db ~schema [datom ~e:1 ~a:"name" ~v:(String "v") ()] in - store ~storage db; + ignore (store ~storage db); let restored = Option.get (restore storage) in reads := 0; List.iter (fun set -> diff --git a/test/test_int64.ml b/test/test_int64.ml index 60c87de..bd635ee 100644 --- a/test/test_int64.ml +++ b/test/test_int64.ml @@ -150,7 +150,7 @@ let test_int64_storage_migration () = ; Db.datom ~e:2 ~a:"count" ~v:(Int64 9_223_372_036_854_775_000L) () ] in - store ~storage db; + ignore (store ~storage db); (match restore storage with | None -> failwith "restore failed" | Some restored -> diff --git a/test/test_sqlite_package.ml b/test/test_sqlite_package.ml index 9fde404..41f892c 100644 --- a/test/test_sqlite_package.ml +++ b/test/test_sqlite_package.ml @@ -32,7 +32,7 @@ let test_storage_roundtrip () = ; Add (Temp_id "todo-1", "todo/title", String "Move storage into datascript") ] in - store ~storage report.db_after; + ignore (store ~storage report.db_after); let restored = match restore storage with | Some db -> db diff --git a/test/test_sqlite_storage.ml b/test/test_sqlite_storage.ml index c356b3f..bf2dd7a 100644 --- a/test/test_sqlite_storage.ml +++ b/test/test_sqlite_storage.ml @@ -625,7 +625,7 @@ let test_sqlite_storage_round_trips_ocaml_payloads () = ~schema:[ "name", indexed ] [ datom ~e:1 ~a:"name" ~v:(String "Ada") () ] in - store ~storage db; + ignore (store ~storage db); assert_equal "kvs schema" "CREATE TABLE kvs (addr INTEGER primary key, content TEXT, addresses JSON)" diff --git a/test/test_storage.ml b/test/test_storage.ml index fdd96a0..f4eefc7 100644 --- a/test/test_storage.ml +++ b/test/test_storage.ml @@ -86,7 +86,7 @@ let reset_writes writes = writes := [] let test_storage__test_basics () = let storage = memory_storage () in let db = small_db () in - store ~storage db; + ignore (store ~storage db); assert_upstream_storage_addresses "store writes upstream storage addresses" (storage_addresses storage); (match restore storage with | None -> failwith "restore should read stored db" @@ -99,7 +99,7 @@ let test_storage__test_basics () = failwith "settings should expose storage attachment"); let attached_storage = memory_storage () in let attached = empty_db ~schema:[ "name", indexed ] ~storage:attached_storage () in - store attached; + ignore (store attached); (match restore attached_storage with | None -> failwith "store should use db-attached storage" | Some restored -> @@ -108,7 +108,7 @@ let test_storage__test_basics () = let test_storage__test_upstream_wire_addresses () = let storage = memory_storage () in let db = small_db () in - store ~storage db; + ignore (store ~storage db); let addresses = storage_addresses storage in if List.mem "datascript/root" addresses || List.mem "datascript/tail" addresses then failwith "storage should not use OCaml snapshot address names"; @@ -134,7 +134,7 @@ let test_storage__test_file_storage () = (fun () -> let storage = file_storage dir in let db = small_db () in - store ~storage db; + ignore (store ~storage db); store_tail storage [ [ datom ~tx:(tx0 + 2) ~e:1 ~a:"name" ~v:(String "Alex") () ] ]; let restored_storage = file_storage dir in assert_upstream_storage_addresses "file_storage lists persisted addresses" (storage_addresses restored_storage); @@ -149,7 +149,7 @@ let test_storage__test_file_storage () = let test_storage__test_gc () = let storage = memory_storage () in let db = small_db () in - store ~storage db; + ignore (store ~storage db); store_tail storage [ [ datom ~tx:(tx0 + 2) ~e:1 ~a:"name" ~v:(String "Alex") () ] ]; storage.storage_store [ "stale/node", Storage_tail [] ]; collect_garbage storage; @@ -165,7 +165,7 @@ let test_storage__test_gc () = let test_storage__test_restored_db_addresses () = let storage = memory_storage () in let db = small_db () in - store ~storage db; + ignore (store ~storage db); let restored = match restore storage with | Some db -> db @@ -176,14 +176,14 @@ let test_storage__test_restored_db_addresses () = let test_storage__test_restored_incremental_store_reuses_index_nodes () = let storage, writes = counting_storage () in let db = large_db () in - store ~storage db; + ignore (store ~storage db); let restored = match restore storage with | Some db -> db | None -> failwith "restore should read stored large db" in reset_writes writes; - store ~storage restored; + ignore (store ~storage restored); assert_int_at_most "storing an unchanged restored db should not rewrite index nodes" 2 @@ -192,7 +192,7 @@ let test_storage__test_restored_incremental_store_reuses_index_nodes () = let db_after = db_with [ Add (Entity_id 1001, "str", String "1001") ] restored in - store ~storage db_after; + ignore (store ~storage db_after); assert_int_at_most "storing an incrementally changed restored db should write only changed index paths" 8 @@ -205,7 +205,7 @@ let test_storage__test_restored_incremental_store_reuses_index_nodes () = let db_after_replacement = db_with [ Add (Entity_id 1, "str", String "changed") ] restored in - store ~storage db_after_replacement; + ignore (store ~storage db_after_replacement); assert_int_at_most "storing a cardinality-one replacement should write only changed index paths" 16 @@ -215,9 +215,36 @@ let test_storage__test_restored_incremental_store_reuses_index_nodes () = [ 1, "str", String "changed" ] (datoms db_after_replacement Eavt ~e:1 ()) +let test_storage__test_conn_repeated_transacts_store_incrementally () = + let storage, writes = counting_storage () in + let conn = create_conn ~schema:[ "name", indexed ] ~storage () in + for i = 1 to 300 do + ignore + (transact_conn conn + [ Add (Entity_id i, "name", String (string_of_int i)) ]) + done; + if List.length !writes = 0 then + failwith "storage-backed conn should compact its tail into index nodes"; + (* each compaction should only write the index nodes created since the + previous store; rewriting the whole tree on every store made total + writes grow quadratically with the number of transactions. At 300 + transactions the buggy code writes well over 3000 entries while the + incremental store stays under 1000. *) + assert_int_at_most + "repeated transacts should not rewrite the whole index tree on every store" + 1500 + (List.length !writes); + (match restore storage with + | None -> failwith "storage-backed conn should restore" + | Some restored -> + assert_equal_int + "all committed datoms remain restorable" + 300 + (List.length (datoms restored Eavt ()))) + let test_storage__test_restore_is_lazy () = let storage = memory_storage () in - large_db () |> store ~storage; + large_db () |> store ~storage |> ignore; let address_count = List.length (storage_addresses storage) in if address_count < 20 then failf "large stored db should have many index nodes, got %d" address_count; @@ -240,7 +267,7 @@ let test_storage__test_restore_is_lazy () = let test_storage__test_restore_with_tail_is_lazy () = let storage = memory_storage () in - large_db () |> store ~storage; + large_db () |> store ~storage |> ignore; let address_count = List.length (storage_addresses storage) in store_tail storage [ @@ -266,7 +293,7 @@ let test_storage__test_restore_with_tail_is_lazy () = let test_storage__test_transact_after_restore_uses_index_slices () = let storage = memory_storage () in - large_db () |> store ~storage; + large_db () |> store ~storage |> ignore; let baseline_storage, baseline_reads = restore_counting_storage storage in let baseline = match restore baseline_storage with @@ -390,7 +417,7 @@ let count_fixture () = datom ~e:(i / 2 + 1) ~a:(if i mod 2 = 0 then "probe/indexed" else "probe/plain") ~v:(String (Printf.sprintf "v%04d" i)) ()) in - store ~storage (init_db ~schema datoms); + ignore (store ~storage (init_db ~schema datoms)); storage let index_counts db = @@ -455,13 +482,13 @@ let test_storage__test_restore_counts_after_tail () = check_cached_counts "restore_conn counts after tail" [4096;4096;2048] (conn_db conn) reads; let tx = transact_conn conn [ Add (Entity_id 2050, "probe/indexed", String "later") ] in check_cached_counts "incremental conn transaction count" [4097;4097;2049] tx.db_after reads; - store tx.db_after; + ignore (store tx.db_after); let restored_again = Option.get (restore measured) in check_cached_counts "new snapshot count excludes cleared tail" [4097;4097;2049] restored_again reads let test_storage__test_empty_and_missing_count_metadata () = let storage = memory_storage () in - store ~storage (empty_db ()); + ignore (store ~storage (empty_db ())); let measured, reads = restore_counting_storage storage in let empty = Option.get (restore measured) in check_cached_counts "empty snapshot count" [0;0;0] empty reads; @@ -496,5 +523,6 @@ let () = test_storage__test_restore_is_lazy (); test_storage__test_restore_with_tail_is_lazy (); test_storage__test_transact_after_restore_uses_index_slices (); + test_storage__test_conn_repeated_transacts_store_incrementally (); test_storage__test_conn (); test_storage__test_db_with_tail (); diff --git a/type/datascript_types.ml b/type/datascript_types.ml index 2202e9e..b990671 100644 --- a/type/datascript_types.ml +++ b/type/datascript_types.ml @@ -85,6 +85,10 @@ type storage_index_metadata = type storage_root = { storage_schema : schema + ; (* cljs datascript's schema map also carries eid -> :db/ident entries + for every attribute installed through a :db/ident datom; keeping + them lets the stored root match cljs's payload shape *) + storage_schema_idents : (entity_id * attr) list ; storage_max_eid : entity_id ; storage_max_tx : tx ; storage_eavt : storage_address