Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion bench/bench_ocaml.ml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
38 changes: 38 additions & 0 deletions bench/codec_parity.ml
Original file line number Diff line number Diff line change
@@ -0,0 +1,38 @@
(* Byte-parity check: decode cljs-written kvs rows, re-encode with our
codec, and compare bytes. Usage: codec_parity <db.sqlite> [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))
12 changes: 12 additions & 0 deletions bench/dune
Original file line number Diff line number Diff line change
Expand Up @@ -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))
6 changes: 3 additions & 3 deletions bench/persistent_sqlite.ml
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand All @@ -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;
Expand Down
67 changes: 67 additions & 0 deletions bench/write_probe.ml
Original file line number Diff line number Diff line change
@@ -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
29 changes: 23 additions & 6 deletions examples/logseq_sqlite_storage.ml
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand All @@ -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
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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 =
Expand Down
2 changes: 1 addition & 1 deletion examples/sqlite_storage_example.ml
Original file line number Diff line number Diff line change
Expand Up @@ -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 ->
Expand Down
24 changes: 11 additions & 13 deletions impl/conn.ml
Original file line number Diff line number Diff line change
Expand Up @@ -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
}

Expand All @@ -25,22 +25,22 @@ 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
; transact : tx_meta:tx_meta -> db -> tx_op list -> tx_report
}

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
Expand Down Expand Up @@ -78,17 +78,15 @@ 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

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)
Expand Down Expand Up @@ -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

Expand All @@ -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;
Expand Down Expand Up @@ -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;
Expand All @@ -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
10 changes: 5 additions & 5 deletions impl/conn.mli
Original file line number Diff line number Diff line change
Expand Up @@ -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
}

Expand All @@ -19,22 +19,22 @@ 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
; transact : tx_meta:tx_meta -> db -> tx_op list -> tx_report
}

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
Expand Down
Loading
Loading