From 0c8253e1b966d475a3fb5376a9cdb769ea0d0b40 Mon Sep 17 00:00:00 2001 From: Zixuan Chen Date: Thu, 1 Oct 2026 05:00:52 +0800 Subject: [PATCH 1/2] fix: reduce overlapping import costs See context/import-peer-id-reuse.md for comparison and allocation rationale. Co-Authored-By: GPT-6 (Codex) --- .changeset/fix-import-overlap-cost.md | 7 + context/arena-parent-links.md | 22 +- .../benchmarks/import-overlap-cost-fat1.csv | 82 ++ .../benchmarks/import-overlap-cost-fat64.csv | 244 ++++ context/import-batch-atomicity.md | 11 +- context/import-peer-id-reuse.md | 162 ++- context/internal-encoding.md | 7 +- .../examples/import_scaling_stress.rs | 1044 +++++++++++++++++ crates/examples/import_overlap_measure.py | 90 ++ crates/loro-internal/src/arena.rs | 7 + .../src/encoding/fast_snapshot.rs | 13 +- crates/loro-internal/src/loro.rs | 2 +- crates/loro-internal/src/oplog.rs | 31 +- .../loro-internal/src/oplog/change_store.rs | 256 +++- .../src/oplog/change_store/block_encode.rs | 110 +- .../loro-internal/src/oplog/known_history.rs | 263 ++++- .../src/oplog/pending_changes.rs | 4 +- .../src/tests/import_atomicity.rs | 44 + crates/loro/tests/import_reused_peer_id.rs | 42 + 19 files changed, 2376 insertions(+), 65 deletions(-) create mode 100644 .changeset/fix-import-overlap-cost.md create mode 100644 context/benchmarks/import-overlap-cost-fat1.csv create mode 100644 context/benchmarks/import-overlap-cost-fat64.csv create mode 100644 crates/examples/examples/import_scaling_stress.rs create mode 100644 crates/examples/import_overlap_measure.py diff --git a/.changeset/fix-import-overlap-cost.md b/.changeset/fix-import-overlap-cost.md new file mode 100644 index 000000000..27ee3a0c4 --- /dev/null +++ b/.changeset/fix-import-overlap-cost.md @@ -0,0 +1,7 @@ +--- +"loro-crdt": patch +--- + +Reduce known-history comparison costs for snapshot-loaded documents while preserving peer-id reuse detection. Snapshot updates allocate only new history in the document arena and avoid cloning discarded changes. + +Avoid unnecessary rollback scopes for detached Text imports while retaining validation of pending movable-list operations. diff --git a/context/arena-parent-links.md b/context/arena-parent-links.md index 0d0d4665c..d3eee0769 100644 --- a/context/arena-parent-links.md +++ b/context/arena-parent-links.md @@ -1,6 +1,6 @@ # Arena Parent Links -Verified against code 2026-09-30. +Verified against code 2026-10-01. `SharedArena` (`crates/loro-internal/src/arena.rs`) stores each container's parent. Liveness (`DocState::is_deleted`), paths (`DocState::get_path`, @@ -142,8 +142,10 @@ Limits: from the change store and panics with "unparsed vv don't match with change store" (`loro_dag.rs`, `ensure_lazy_load_node`). - A read that returns before any failure was recorded is not undone: an - `import` or `checkout` that is the first to read the block finishes on the - partial history. Only `export` checks again at the end. + `import` or `checkout` using a general history reader can finish on partial + history when it is the first to read the block. Only `export` checks again at + the end. The direct cold Text comparison in `known_history.rs` returns its + decode error immediately and also records it in `parse_failures`. - Local edits still succeed after a failure was recorded, but they cannot be exported from this document any more. What can be salvaged is the current state (`get_deep_value`). @@ -182,6 +184,20 @@ A vv shortcut for IDs beyond the history would have to be exact for every path that writes the KV store, or a live container would read as deleted. Measured 2026-09-28 (loro-dev/loro#1159). +## Snapshot overlap decoding + +`ChangeStore::decode_snapshot_for_updates` uses a temporary arena for unmatched +incoming blocks, rather than allocating a second known prefix in the document +arena. Known-history comparison resolves container IDs across the two arenas; +only checked, trimmed new ops are converted into document indices/value slices. +`register_container_and_parent_link` runs on those converted changes before they +leave the decoder. Temporary indices and parent links never escape into the +document. A rejected comparison follows the existing arena rollback path. + +Cold Text-insert-only local blocks can also be compared directly from bytes +without registering anything in the document arena. Other blocks still use the +normal lazy reader. See [import-peer-id-reuse.md](import-peer-id-reuse.md). + ## Import rollback A failed import can parse old change blocks while it computes its diff diff --git a/context/benchmarks/import-overlap-cost-fat1.csv b/context/benchmarks/import-overlap-cost-fat1.csv new file mode 100644 index 000000000..b3cb37bf7 --- /dev/null +++ b/context/benchmarks/import-overlap-cost-fat1.csv @@ -0,0 +1,82 @@ +scenario,n,rep,version,ms,heap_kb,delta_kb +overlap_snap,8000,0,old,0.268,0,0 +overlap_snap,8000,0,main,1.071,0,0 +overlap_snap,8000,0,fixed,0.594,0,0 +overlap_snap,8000,1,old,0.247,0,0 +overlap_snap,8000,1,main,1.095,0,0 +overlap_snap,8000,1,fixed,0.584,0,0 +overlap_snap,8000,2,old,0.266,0,0 +overlap_snap,8000,2,main,1.058,0,0 +overlap_snap,8000,2,fixed,0.595,0,0 +overlap_snap,32000,0,old,1.056,0,0 +overlap_snap,32000,0,main,4.379,0,0 +overlap_snap,32000,0,fixed,2.36,0,0 +overlap_snap,32000,1,old,1.049,0,0 +overlap_snap,32000,1,main,4.438,0,0 +overlap_snap,32000,1,fixed,2.308,0,0 +overlap_snap,32000,2,old,1.066,0,0 +overlap_snap,32000,2,main,4.461,0,0 +overlap_snap,32000,2,fixed,2.422,0,0 +overlap_snap,64000,0,old,2.209,0,0 +overlap_snap,64000,0,main,9.024,0,0 +overlap_snap,64000,0,fixed,4.867,0,0 +overlap_snap,64000,1,old,2.125,0,0 +overlap_snap,64000,1,main,9.219,0,0 +overlap_snap,64000,1,fixed,4.75,0,0 +overlap_snap,64000,2,old,2.17,0,0 +overlap_snap,64000,2,main,9.08,0,0 +overlap_snap,64000,2,fixed,4.828,0,0 +overlap_partial,8000,0,old,0.249,0,0 +overlap_partial,8000,0,main,0.678,0,0 +overlap_partial,8000,0,fixed,0.445,0,0 +overlap_partial,8000,1,old,0.247,0,0 +overlap_partial,8000,1,main,0.673,0,0 +overlap_partial,8000,1,fixed,0.444,0,0 +overlap_partial,8000,2,old,0.247,0,0 +overlap_partial,8000,2,main,0.672,0,0 +overlap_partial,8000,2,fixed,0.45,0,0 +overlap_partial,32000,0,old,1.028,0,0 +overlap_partial,32000,0,main,2.735,0,0 +overlap_partial,32000,0,fixed,1.735,0,0 +overlap_partial,32000,1,old,1.043,0,0 +overlap_partial,32000,1,main,2.694,0,0 +overlap_partial,32000,1,fixed,1.666,0,0 +overlap_partial,32000,2,old,1.13,0,0 +overlap_partial,32000,2,main,2.725,0,0 +overlap_partial,32000,2,fixed,1.667,0,0 +overlap_partial,64000,0,old,2.082,0,0 +overlap_partial,64000,0,main,5.479,0,0 +overlap_partial,64000,0,fixed,3.372,0,0 +overlap_partial,64000,1,old,2.197,0,0 +overlap_partial,64000,1,main,5.506,0,0 +overlap_partial,64000,1,fixed,3.339,0,0 +overlap_partial,64000,2,old,2.181,0,0 +overlap_partial,64000,2,main,5.421,0,0 +overlap_partial,64000,2,fixed,3.512,0,0 +snap_plus,8000,0,old,0.793,0,354 +snap_plus,8000,0,main,1.731,0,1508 +snap_plus,8000,0,fixed,0.364,0,334 +snap_plus,8000,1,old,0.801,0,353 +snap_plus,8000,1,main,1.745,0,1507 +snap_plus,8000,1,fixed,0.356,0,334 +snap_plus,8000,2,old,0.807,0,353 +snap_plus,8000,2,main,1.746,0,1508 +snap_plus,8000,2,fixed,0.361,0,334 +snap_plus,32000,0,old,3.274,0,1604 +snap_plus,32000,0,main,7.571,0,6673 +snap_plus,32000,0,fixed,1.137,0,1571 +snap_plus,32000,1,old,3.274,0,1605 +snap_plus,32000,1,main,7.692,0,6673 +snap_plus,32000,1,fixed,1.171,0,1571 +snap_plus,32000,2,old,3.296,0,1604 +snap_plus,32000,2,main,7.555,0,6673 +snap_plus,32000,2,fixed,1.12,0,1570 +snap_plus,64000,0,old,6.568,0,3211 +snap_plus,64000,0,main,15.173,0,13346 +snap_plus,64000,0,fixed,2.283,0,3142 +snap_plus,64000,1,old,6.454,0,3211 +snap_plus,64000,1,main,15.927,0,13346 +snap_plus,64000,1,fixed,2.213,0,3143 +snap_plus,64000,2,old,6.492,0,3211 +snap_plus,64000,2,main,15.34,0,13347 +snap_plus,64000,2,fixed,2.236,0,3142 diff --git a/context/benchmarks/import-overlap-cost-fat64.csv b/context/benchmarks/import-overlap-cost-fat64.csv new file mode 100644 index 000000000..e1d6e2015 --- /dev/null +++ b/context/benchmarks/import-overlap-cost-fat64.csv @@ -0,0 +1,244 @@ +scenario,n,rep,version,ms,heap_kb,delta_kb +overlap_snap,8000,0,old,1.757,0,0 +overlap_snap,8000,0,main,3.095,0,0 +overlap_snap,8000,0,fixed,2.295,0,0 +overlap_snap,8000,1,old,1.742,0,0 +overlap_snap,8000,1,main,3.136,0,0 +overlap_snap,8000,1,fixed,2.328,0,0 +overlap_snap,8000,2,old,1.812,0,0 +overlap_snap,8000,2,main,3.066,0,0 +overlap_snap,8000,2,fixed,2.338,0,0 +overlap_snap,32000,0,old,7.377,0,0 +overlap_snap,32000,0,main,13.0,0,0 +overlap_snap,32000,0,fixed,9.504,0,0 +overlap_snap,32000,1,old,7.318,0,0 +overlap_snap,32000,1,main,12.734,0,0 +overlap_snap,32000,1,fixed,9.317,0,0 +overlap_snap,32000,2,old,7.389,0,0 +overlap_snap,32000,2,main,12.803,0,0 +overlap_snap,32000,2,fixed,9.68,0,0 +overlap_snap,64000,0,old,15.539,0,0 +overlap_snap,64000,0,main,26.87,0,0 +overlap_snap,64000,0,fixed,19.33,0,0 +overlap_snap,64000,1,old,15.462,0,0 +overlap_snap,64000,1,main,26.566,0,0 +overlap_snap,64000,1,fixed,19.493,0,0 +overlap_snap,64000,2,old,15.123,0,0 +overlap_snap,64000,2,main,26.516,0,0 +overlap_snap,64000,2,fixed,19.44,0,0 +overlap_partial,8000,0,old,1.548,0,0 +overlap_partial,8000,0,main,2.24,0,0 +overlap_partial,8000,0,fixed,1.79,0,0 +overlap_partial,8000,1,old,1.612,0,0 +overlap_partial,8000,1,main,2.249,0,0 +overlap_partial,8000,1,fixed,1.852,0,0 +overlap_partial,8000,2,old,1.568,0,0 +overlap_partial,8000,2,main,2.26,0,0 +overlap_partial,8000,2,fixed,1.869,0,0 +overlap_partial,32000,0,old,6.546,0,0 +overlap_partial,32000,0,main,9.233,0,0 +overlap_partial,32000,0,fixed,7.506,0,0 +overlap_partial,32000,1,old,6.509,0,0 +overlap_partial,32000,1,main,9.33,0,0 +overlap_partial,32000,1,fixed,7.676,0,0 +overlap_partial,32000,2,old,6.605,0,0 +overlap_partial,32000,2,main,9.203,0,0 +overlap_partial,32000,2,fixed,7.682,0,0 +overlap_partial,64000,0,old,13.401,0,0 +overlap_partial,64000,0,main,19.404,0,0 +overlap_partial,64000,0,fixed,15.429,0,0 +overlap_partial,64000,1,old,13.599,0,0 +overlap_partial,64000,1,main,19.139,0,0 +overlap_partial,64000,1,fixed,15.51,0,0 +overlap_partial,64000,2,old,13.391,0,0 +overlap_partial,64000,2,main,18.976,0,0 +overlap_partial,64000,2,fixed,15.316,0,0 +overlap_mem,8000,0,old,0.394,0,0 +overlap_mem,8000,0,main,0.398,0,0 +overlap_mem,8000,0,fixed,0.4,0,0 +overlap_mem,8000,1,old,0.387,0,0 +overlap_mem,8000,1,main,0.387,0,0 +overlap_mem,8000,1,fixed,0.393,0,0 +overlap_mem,8000,2,old,0.385,0,0 +overlap_mem,8000,2,main,0.397,0,0 +overlap_mem,8000,2,fixed,0.385,0,0 +overlap_mem,32000,0,old,1.568,0,0 +overlap_mem,32000,0,main,1.567,0,0 +overlap_mem,32000,0,fixed,1.578,0,0 +overlap_mem,32000,1,old,1.543,0,0 +overlap_mem,32000,1,main,1.565,0,0 +overlap_mem,32000,1,fixed,1.567,0,0 +overlap_mem,32000,2,old,1.549,0,0 +overlap_mem,32000,2,main,1.606,0,0 +overlap_mem,32000,2,fixed,1.553,0,0 +overlap_mem,64000,0,old,3.112,0,0 +overlap_mem,64000,0,main,3.125,0,0 +overlap_mem,64000,0,fixed,3.108,0,0 +overlap_mem,64000,1,old,3.028,0,0 +overlap_mem,64000,1,main,3.118,0,0 +overlap_mem,64000,1,fixed,3.082,0,0 +overlap_mem,64000,2,old,3.107,0,0 +overlap_mem,64000,2,main,3.152,0,0 +overlap_mem,64000,2,fixed,3.082,0,0 +snap_plus,8000,0,old,2.473,0,2164 +snap_plus,8000,0,main,3.867,0,4598 +snap_plus,8000,0,fixed,1.597,0,1604 +snap_plus,8000,1,old,2.501,0,2165 +snap_plus,8000,1,main,3.947,0,4598 +snap_plus,8000,1,fixed,1.603,0,1605 +snap_plus,8000,2,old,2.474,0,2164 +snap_plus,8000,2,main,3.925,0,4598 +snap_plus,8000,2,fixed,1.592,0,1604 +snap_plus,32000,0,old,10.024,0,8682 +snap_plus,32000,0,main,16.703,0,19192 +snap_plus,32000,0,fixed,6.526,0,6452 +snap_plus,32000,1,old,10.229,0,8682 +snap_plus,32000,1,main,16.324,0,19192 +snap_plus,32000,1,fixed,6.411,0,6453 +snap_plus,32000,2,old,10.213,0,8681 +snap_plus,32000,2,main,16.615,0,19192 +snap_plus,32000,2,fixed,6.475,0,6452 +snap_plus,64000,0,old,20.766,0,21282 +snap_plus,64000,0,main,34.489,0,41202 +snap_plus,64000,0,fixed,13.397,0,16827 +snap_plus,64000,1,old,20.615,0,21283 +snap_plus,64000,1,main,34.882,0,41229 +snap_plus,64000,1,fixed,13.568,0,16827 +snap_plus,64000,2,old,20.481,0,21273 +snap_plus,64000,2,main,34.627,0,41193 +snap_plus,64000,2,fixed,13.282,0,16832 +snap_mem,8000,0,old,1.2105000000000001,7727,5382 +snap_mem,8000,0,main,1.911,13190,10839 +snap_mem,8000,0,fixed,0.203,3955,1604 +snap_mem,8000,1,old,1.1484999999999999,7727,5382 +snap_mem,8000,1,main,1.958,13189,10838 +snap_mem,8000,1,fixed,0.20400000000000001,3956,1605 +snap_mem,8000,2,old,1.138,7727,5382 +snap_mem,8000,2,main,1.919,13190,10839 +snap_mem,8000,2,fixed,0.2025,3956,1606 +snap_mem,32000,0,old,4.887,30590,21548 +snap_mem,32000,0,main,7.8675,52443,43383 +snap_mem,32000,0,fixed,0.71,15516,6456 +snap_mem,32000,1,old,4.8255,30589,21547 +snap_mem,32000,1,main,7.9475,52445,43385 +snap_mem,32000,1,fixed,0.722,15514,6454 +snap_mem,32000,2,old,4.7844999999999995,30587,21545 +snap_mem,32000,2,main,7.8875,52444,43384 +snap_mem,32000,2,fixed,0.7264999999999999,15515,6455 +snap_mem,64000,0,old,9.486,64813,47011 +snap_mem,64000,0,main,15.8765,108979,91143 +snap_mem,64000,0,fixed,1.589,34664,16828 +snap_mem,64000,1,old,9.5715,64808,47006 +snap_mem,64000,1,main,16.019,108932,91096 +snap_mem,64000,1,fixed,1.512,34665,16829 +snap_mem,64000,2,old,9.811,64809,47007 +snap_mem,64000,2,main,15.917000000000002,108973,91137 +snap_mem,64000,2,fixed,1.5805,34663,16827 +detached,8000,0,old,7.117,0,0 +detached,8000,0,main,8.931,0,0 +detached,8000,0,fixed,7.216,0,0 +detached,8000,1,old,6.865,0,0 +detached,8000,1,main,8.897,0,0 +detached,8000,1,fixed,6.954,0,0 +detached,8000,2,old,6.831,0,0 +detached,8000,2,main,8.863,0,0 +detached,8000,2,fixed,6.909,0,0 +detached,32000,0,old,27.968,0,0 +detached,32000,0,main,35.847,0,0 +detached,32000,0,fixed,27.75,0,0 +detached,32000,1,old,28.11,0,0 +detached,32000,1,main,35.118,0,0 +detached,32000,1,fixed,28.886,0,0 +detached,32000,2,old,28.873,0,0 +detached,32000,2,main,35.053,0,0 +detached,32000,2,fixed,28.423,0,0 +detached,64000,0,old,57.046,0,0 +detached,64000,0,main,73.145,0,0 +detached,64000,0,fixed,59.277,0,0 +detached,64000,1,old,57.146,0,0 +detached,64000,1,main,71.227,0,0 +detached,64000,1,fixed,57.281,0,0 +detached,64000,2,old,56.329,0,0 +detached,64000,2,main,72.767,0,0 +detached,64000,2,fixed,56.579,0,0 +mlist_batch,8000,0,old,7.456,0,0 +mlist_batch,8000,0,main,8.909,0,0 +mlist_batch,8000,0,fixed,9.053,0,0 +mlist_batch,8000,1,old,7.46,0,0 +mlist_batch,8000,1,main,9.06,0,0 +mlist_batch,8000,1,fixed,8.993,0,0 +mlist_batch,8000,2,old,7.563,0,0 +mlist_batch,8000,2,main,8.9,0,0 +mlist_batch,8000,2,fixed,8.982,0,0 +mlist_batch,32000,0,old,29.524,0,0 +mlist_batch,32000,0,main,38.558,0,0 +mlist_batch,32000,0,fixed,39.529,0,0 +mlist_batch,32000,1,old,29.442,0,0 +mlist_batch,32000,1,main,36.899,0,0 +mlist_batch,32000,1,fixed,37.298,0,0 +mlist_batch,32000,2,old,29.867,0,0 +mlist_batch,32000,2,main,37.508,0,0 +mlist_batch,32000,2,fixed,37.358,0,0 +mlist_batch,64000,0,old,79.806,0,0 +mlist_batch,64000,0,main,96.3,0,0 +mlist_batch,64000,0,fixed,95.375,0,0 +mlist_batch,64000,1,old,79.049,0,0 +mlist_batch,64000,1,main,100.955,0,0 +mlist_batch,64000,1,fixed,96.268,0,0 +mlist_batch,64000,2,old,83.69,0,0 +mlist_batch,64000,2,main,96.798,0,0 +mlist_batch,64000,2,fixed,96.028,0,0 +text_stream,8000,0,old,18.094,0,0 +text_stream,8000,0,main,20.952,0,0 +text_stream,8000,0,fixed,20.439,0,0 +text_stream,8000,1,old,18.364,0,0 +text_stream,8000,1,main,20.088,0,0 +text_stream,8000,1,fixed,20.701,0,0 +text_stream,8000,2,old,18.2,0,0 +text_stream,8000,2,main,21.394,0,0 +text_stream,8000,2,fixed,20.529,0,0 +text_stream,32000,0,old,76.572,0,0 +text_stream,32000,0,main,84.157,0,0 +text_stream,32000,0,fixed,84.151,0,0 +text_stream,32000,1,old,71.38,0,0 +text_stream,32000,1,main,85.664,0,0 +text_stream,32000,1,fixed,82.378,0,0 +text_stream,32000,2,old,73.324,0,0 +text_stream,32000,2,main,86.208,0,0 +text_stream,32000,2,fixed,81.682,0,0 +text_stream,64000,0,old,149.703,0,0 +text_stream,64000,0,main,172.67,0,0 +text_stream,64000,0,fixed,168.454,0,0 +text_stream,64000,1,old,151.261,0,0 +text_stream,64000,1,main,165.534,0,0 +text_stream,64000,1,fixed,166.158,0,0 +text_stream,64000,2,old,149.038,0,0 +text_stream,64000,2,main,162.693,0,0 +text_stream,64000,2,fixed,164.711,0,0 +batch,8000,0,old,29.818,0,0 +batch,8000,0,main,29.709,0,0 +batch,8000,0,fixed,29.698,0,0 +batch,8000,1,old,29.492,0,0 +batch,8000,1,main,30.685,0,0 +batch,8000,1,fixed,29.245,0,0 +batch,8000,2,old,30.727,0,0 +batch,8000,2,main,29.681,0,0 +batch,8000,2,fixed,29.886,0,0 +batch,32000,0,old,124.071,0,0 +batch,32000,0,main,127.269,0,0 +batch,32000,0,fixed,124.147,0,0 +batch,32000,1,old,122.213,0,0 +batch,32000,1,main,123.015,0,0 +batch,32000,1,fixed,126.228,0,0 +batch,32000,2,old,125.755,0,0 +batch,32000,2,main,129.023,0,0 +batch,32000,2,fixed,125.134,0,0 +batch,64000,0,old,260.088,0,0 +batch,64000,0,main,255.777,0,0 +batch,64000,0,fixed,255.857,0,0 +batch,64000,1,old,257.137,0,0 +batch,64000,1,main,256.648,0,0 +batch,64000,1,fixed,255.895,0,0 +batch,64000,2,old,254.521,0,0 +batch,64000,2,main,255.13,0,0 +batch,64000,2,fixed,255.577,0,0 diff --git a/context/import-batch-atomicity.md b/context/import-batch-atomicity.md index 41e9d8c0d..83017d7d1 100644 --- a/context/import-batch-atomicity.md +++ b/context/import-batch-atomicity.md @@ -1,6 +1,6 @@ # `import_batch` Atomicity and the Detached-Mode Invariant -Verified against code 2026-08-09. +Verified against code 2026-10-01. `LoroDoc::import_batch` (`crates/loro-internal/src/loro.rs`) does not import blobs the way `import` does. It stops the auto-commit txn, keeps the txn mutex for the whole @@ -95,6 +95,15 @@ regression on out-of-order batches, where every blob parks and is later unlocked does not re-scan the pending set the earlier blobs grew, once per blob. Doing it eagerly made such a batch quadratic (12k blobs: 3.4s vs 0.45s). +Standalone detached imports do not apply state, so preflight uses +`import_op_can_reject(op, true)`: only MovableList `Move`/`Set` element validation +requires a rollback scope. Text/List/Tree state validation still opens a scope +when attached. An applied Text or scalar update can unlock a pre-existing bad +MovableList pending change, so preflight includes pending `Move`/`Set` ops when +`applies_to_dag` is true. The batch-wide scope remains unconditional and is not +replaced or nested by this optimization. Regression: +`detached_text_import_that_unlocks_invalid_movable_list_ops_rolls_back`. + `crates/examples/examples/import_batch_perf.rs` is the ad-hoc probe for these shapes (not part of CI; run it on two revisions and compare). diff --git a/context/import-peer-id-reuse.md b/context/import-peer-id-reuse.md index 638f8f612..ebe8cd479 100644 --- a/context/import-peer-id-reuse.md +++ b/context/import-peer-id-reuse.md @@ -1,6 +1,6 @@ # Imports That Reuse Local Op Ids -Verified against code 2026-09-30. +Verified against code 2026-10-01. An import skips the part of each change the doc already has by version vector and applies the rest on top of the local history. When two clients shared a peer id @@ -33,9 +33,19 @@ path, which rolls back only the arena (decode may have registered containers). that pass on counter ends only. Parked pending changes already wait for their own predecessor (`remote_change_apply_state`). -The decoders (`ChangeStore::decode_block_bytes`, `decode_snapshot_for_updates`) -therefore no longer trim. A snapshot with nothing new still returns no changes, -without cloning them. +Binary decoding can skip a fully known block whose bytes match the current +local block exactly (`ChangeStore::contains_encoded_block`). A dirty parsed +block shadows the KV copy, so stale bytes never authorize this shortcut. +Different bytes still use the semantic check: encoding, change boundaries, +timestamps, and messages are not required to match. + +For snapshots imported into a loaded document, +`ChangeStore::decode_snapshot_for_updates` decodes unmatched blocks in a temporary +arena, checks and trims there, and moves only the retained suffix into the document +arena. Frontier blocks parsed by `import_all` are taken, not cloned. No dropped +prefix strings or list values are allocated in the document arena, including when +snapshot and local history use different change boundaries. The regular import +entry then runs its counter/preflight checks as usual. ## What "equal" means @@ -74,11 +84,26 @@ is per atom range, not per change: ## Cost -Imports that don't overlap local history, or overlap only a little, cost the same -as before. An import that re-sends all known history *and* brings new changes -pays for comparing the overlap. On a 10k-change doc that was about +15% (55 โ†’ 64 ms) -for both "all updates + 1 change" and "snapshot + 1 change". Probe: -`crates/examples/examples/import_known_history_perf.rs`. +The scope of conflict detection is unchanged: all imported overlapping peers +are compared when the import brings new ops, including a fully known peer whose +history another peer depends on. No new detection limit is introduced. + +Cold blocks containing only Text inserts can be compared directly from their +encoded container IDs, positions, UTF-8 strings and dependency boundaries +(`OpLog::check_known_text_in_cold_block`, `block_encode::visit_text_insert_block`). +This does not parse changes or allocate strings in the local arena. Both change +and op boundaries remain significant to the comparison. Unsupported op kinds +fall back to the general check; already parsed blocks keep the existing path. +A short imported change that ends inside a cold block uses the general path so +many short changes cannot repeatedly scan the same large block. + +The comparison still costs O(overlap); it is not a content-hashed version vector. +The release probe is `crates/examples/examples/import_scaling_stress.rs`. +Use its `overlap_snap`, `overlap_partial`, `overlap_mem`, `snap_plus` and `snap_mem` +scenarios with separate old/main/fixed binaries and interleaved runs. Heap is a +process-wide allocator statistic; inspect both the one-import delta and the +repeated-import series rather than interpreting one allocator growth step as a +universal ratio. ## Tests @@ -88,3 +113,122 @@ for both "all updates + 1 change" and "snapshot + 1 change". Probe: and re-imports that must still succeed (piecewise vs merged change stores, every op kind including `NaN` and one-atom reversed deletes, shallow docs). - `crates/loro-wasm/tests/import_reused_peer_id.test.ts`. + +- `oplog::known_history::tests`: cross-arena container identity and Unicode atom + comparisons, and a dependency mismatch inside a merged import's cold prefix. +- `oplog::change_store::test`: identical blocks stay lazy, dirty caches shadow + old KV bytes, and repeated snapshots allocate only new text/list values. +- `cold_history_conflict_inside_a_large_known_prefix_is_rejected` in the public + reused-peer tests: exact mismatch ID across multiple cold blocks, through + updates, snapshots, and batches, with subsequent usability checks. + +## Measurements + +Measured on macOS arm64, `rustc 1.96.0 (ac68faa20 2026-05-25)`, release builds, +`CARGO_BUILD_JOBS=4`. Baselines: 1.16.3 `ad5b2a6d4546473d9f4a96412d2ae6808c8403ec`, +main `c00c9fa501f8d32f68d6255eacb7035a67fb6ab6`, fixed = this change. +The verifier directory was read only; its probe and NOSEED/SEEDSYNC patch were +copied here. Each cell has 3 interleaved old/main/fixed process runs, each reporting +3 timed rounds after a warmup. Tables show medians of process medians. +The machine is shared; recorded one-minute load ranged 3.48-4.37. +Use ratios, not an individual wall time. + +The first five scenarios use `FAT=64 REPEATS=4`. Other scenarios use `FAT=1`; +`text_stream` uses `NOSEED=1` (no concurrent seed). `snap_mem` has no warmup: +its time column is the median of each process's four growing-snapshot imports, +then the median over 3 processes. Its first-import latency is separately covered +by `snap_plus`. + +| Scenario | Changes | 1.16.3 ms | main ms | fixed ms | fixed/main | fixed/1.16.3 | +|---|---:|---:|---:|---:|---:|---:| +| overlap_snap | 8,000 | 1.757 | 3.095 | 2.328 | 0.75x | 1.32x | +| overlap_snap | 32,000 | 7.377 | 12.803 | 9.504 | 0.74x | 1.29x | +| overlap_snap | 64,000 | 15.462 | 26.566 | 19.440 | 0.73x | 1.26x | +| overlap_partial | 8,000 | 1.568 | 2.249 | 1.852 | 0.82x | 1.18x | +| overlap_partial | 32,000 | 6.546 | 9.233 | 7.676 | 0.83x | 1.17x | +| overlap_partial | 64,000 | 13.401 | 19.139 | 15.429 | 0.81x | 1.15x | +| overlap_mem | 8,000 | 0.387 | 0.397 | 0.393 | 0.99x | 1.02x | +| overlap_mem | 32,000 | 1.549 | 1.567 | 1.567 | 1.00x | 1.01x | +| overlap_mem | 64,000 | 3.107 | 3.125 | 3.082 | 0.99x | 0.99x | +| snap_plus | 8,000 | 2.474 | 3.925 | 1.597 | 0.41x | 0.65x | +| snap_plus | 32,000 | 10.213 | 16.615 | 6.475 | 0.39x | 0.63x | +| snap_plus | 64,000 | 20.615 | 34.627 | 13.397 | 0.39x | 0.65x | +| snap_mem | 8,000 | 1.148 | 1.919 | 0.203 | 0.11x | 0.18x | +| snap_mem | 32,000 | 4.825 | 7.888 | 0.722 | 0.09x | 0.15x | +| snap_mem | 64,000 | 9.572 | 15.917 | 1.581 | 0.10x | 0.17x | +| detached | 8,000 | 6.865 | 8.897 | 6.954 | 0.78x | 1.01x | +| detached | 32,000 | 28.110 | 35.118 | 28.423 | 0.81x | 1.01x | +| detached | 64,000 | 57.046 | 72.767 | 57.281 | 0.79x | 1.00x | +| mlist_batch | 8,000 | 7.460 | 8.909 | 8.993 | 1.01x | 1.21x | +| mlist_batch | 32,000 | 29.524 | 37.508 | 37.358 | 1.00x | 1.27x | +| mlist_batch | 64,000 | 79.806 | 96.798 | 96.028 | 0.99x | 1.20x | +| text_stream | 8,000 | 18.200 | 20.952 | 20.529 | 0.98x | 1.13x | +| text_stream | 32,000 | 73.324 | 85.664 | 82.378 | 0.96x | 1.12x | +| text_stream | 64,000 | 149.703 | 165.534 | 166.158 | 1.00x | 1.11x | +| batch | 8,000 | 29.818 | 29.709 | 29.698 | 1.00x | 1.00x | +| batch | 32,000 | 124.071 | 127.269 | 125.134 | 0.98x | 1.01x | +| batch | 64,000 | 257.137 | 255.777 | 255.857 | 1.00x | 1.00x | + +### Retained allocator memory + +Process-wide `malloc_zone_statistics().size_in_use`, in KiB. `snap_plus` is the +one-import delta; `snap_mem` is the delta from the loaded base to import 4. +These statistics have allocator growth steps and include state/history caches; +the exact arena regression test verifies that no incoming known prefix survives. + +| Scenario | Changes | 1.16.3 KiB | main KiB | fixed KiB | fixed/main | +|---|---:|---:|---:|---:|---:| +| snap_plus | 8,000 | 2,164 | 4,598 | 1,604 | 0.35x | +| snap_plus | 32,000 | 8,682 | 19,192 | 6,452 | 0.34x | +| snap_plus | 64,000 | 21,282 | 41,202 | 16,827 | 0.41x | +| snap_mem | 8,000 | 5,382 | 10,839 | 1,605 | 0.15x | +| snap_mem | 32,000 | 21,547 | 43,384 | 6,455 | 0.15x | +| snap_mem | 64,000 | 47,007 | 91,137 | 16,828 | 0.18x | + +Fixed `snap_mem` retained-heap deltas from its loaded base after imports 1/2/3/4 +(median of 3 processes, KiB): + +| Changes | Import 1 | Import 2 | Import 3 | Import 4 | +|---|---:|---:|---:|---:| +| 8,000 | 1,605 | 1,605 | 1,605 | 1,605 | +| 32,000 | 6,455 | 6,455 | 6,455 | 6,455 | +| 64,000 | 16,826 | 16,827 | 16,828 | 16,828 | + +### FAT=1 controls for the original regression shape + +Same release/interleaving/round/repeat settings; these reproduce the originally +reported 9 ms and 15 ms regressions without large inserted strings. + +| Scenario | Changes | 1.16.3 ms | main ms | fixed ms | fixed/main | fixed/1.16.3 | +|---|---:|---:|---:|---:|---:|---:| +| overlap_snap | 8,000 | 0.266 | 1.071 | 0.594 | 0.55x | 2.23x | +| overlap_snap | 32,000 | 1.056 | 4.438 | 2.360 | 0.53x | 2.23x | +| overlap_snap | 64,000 | 2.170 | 9.080 | 4.828 | 0.53x | 2.22x | +| overlap_partial | 8,000 | 0.247 | 0.673 | 0.445 | 0.66x | 1.80x | +| overlap_partial | 32,000 | 1.043 | 2.725 | 1.667 | 0.61x | 1.60x | +| overlap_partial | 64,000 | 2.181 | 5.479 | 3.372 | 0.62x | 1.55x | +| snap_plus | 8,000 | 0.801 | 1.745 | 0.361 | 0.21x | 0.45x | +| snap_plus | 32,000 | 3.274 | 7.571 | 1.137 | 0.15x | 0.35x | +| snap_plus | 64,000 | 6.492 | 15.340 | 2.236 | 0.15x | 0.34x | + +R1 retains some comparison cost above 1.16.3, which did not check known history. +Movable-list element validation is unchanged: `mlist_batch` retains its roughly +1.2x cost relative to 1.16.3. Text streaming and batch controls are comparable to +main. Snapshot decoding is faster than both baselines and repeated imports plateau. + +Reproduce with separate binaries named `old`, `main`, `fixed`: + +```sh +export CARGO_BUILD_JOBS=4 +cargo build -p examples --release --example import_scaling_stress +# Copy each revision's binary to $BENCH_BIN_DIR/{old,main,fixed}. +python3 crates/examples/import_overlap_measure.py \ + --bin-dir "$BENCH_BIN_DIR" --output-dir /tmp/import-overlap-fat64 +python3 crates/examples/import_overlap_measure.py \ + --bin-dir "$BENCH_BIN_DIR" --output-dir /tmp/import-overlap-fat1 \ + --fat 1 --scenarios overlap_snap overlap_partial snap_plus +``` + +Samples: [FAT=64 matrix](benchmarks/import-overlap-cost-fat64.csv), +[FAT=1 controls](benchmarks/import-overlap-cost-fat1.csv). +The runner also saves full probe output and load readings to `raw.jsonl`. diff --git a/context/internal-encoding.md b/context/internal-encoding.md index bf2772c40..a68b3a2ec 100644 --- a/context/internal-encoding.md +++ b/context/internal-encoding.md @@ -1,6 +1,6 @@ # Internal Encoding Context -Verified against code 2026-09-28. +Verified against code 2026-10-01. Loro has one binary blob envelope, two current binary body formats, two recognized-but-unsupported legacy top-level modes, and a separate JSON updates @@ -83,6 +83,11 @@ still contains current helpers including `import_changes_to_oplog`, `encode_op`, document. If a snapshot is imported into a non-empty document, `LoroDoc::_import_with` routes through decoded oplog changes instead. Failed direct snapshot import must reset both state and oplog. +The non-empty path compares byte-identical known blocks without parsing them; +unmatched incoming blocks use a temporary arena, and only checked new suffixes +are moved into the document arena. This changes no on-wire format. See +[import-peer-id-reuse.md](import-peer-id-reuse.md) for the semantic fallback and +[arena-parent-links.md](arena-parent-links.md) for converted-op parent links. Default snapshot export does not walk the alive-container graph. It relies on a write-time invariant: every container brought alive by an applied op or diff diff --git a/crates/examples/examples/import_scaling_stress.rs b/crates/examples/examples/import_scaling_stress.rs new file mode 100644 index 000000000..709ecba5a --- /dev/null +++ b/crates/examples/examples/import_scaling_stress.rs @@ -0,0 +1,1044 @@ +//! Import-path scaling probe for loro-crdt 1.16.4 vs 1.16.3 (`ad5b2a6d`). +//! +//! Each process prints one `RESULT` line. `rounds` (default 3) are timed after one +//! warmup; the line is median / min / mean. `FAT` is the inserted string length +//! (default 1). History uses merge interval -1 so each commit is its own change. +//! +//! ```text +//! CARGO_BUILD_JOBS=4 cargo run -p examples --release --example import_scaling_stress -- [rounds] +//! ``` +//! +//! Scenarios: `text_peers`, `map_peers`, `text_stream`, `overlap_mem`, `overlap_snap`, +//! `overlap_partial`, `snap_plus`, `snap_mem`, `batch`, `batch_rev`, `ooo`, `detached`, +//! `mlist_batch`, `mlist_stream`, `check`. +use loro::{ExportMode, LoroDoc, ID}; +use std::time::{Duration, Instant}; + +fn main() { + let mut args = std::env::args().skip(1); + let scenario = args.next().unwrap_or_else(|| "help".into()); + if scenario == "help" || scenario == "-h" { + eprintln!( + "usage: import_scaling_stress [rounds]\n\ + env: FAT= IMPORTS=" + ); + std::process::exit(2); + } + let n: usize = args + .next() + .unwrap_or_else(|| "1000".into()) + .parse() + .expect("n"); + let rounds: usize = args.next().and_then(|s| s.parse().ok()).unwrap_or(3); + let fat: usize = std::env::var("FAT") + .ok() + .and_then(|s| s.parse().ok()) + .unwrap_or(1) + .max(1); + println!( + "META scenario={scenario} n={n} rounds={rounds} fat={fat} load={}", + loadavg() + ); + match scenario.as_str() { + "text_peers" => peers(n, rounds, true), + "map_peers" => peers(n, rounds, false), + "text_stream" => text_stream(n, rounds, fat), + "overlap_mem" => overlap(n, rounds, fat, OverlapMode::Mem), + "overlap_snap" => overlap(n, rounds, fat, OverlapMode::Snap), + "overlap_partial" => overlap(n, rounds, fat, OverlapMode::Partial), + "snap_plus" => snap_plus(n, rounds, fat), + "snap_mem" => snap_mem(n, fat), + "batch" => batch(n, rounds, false), + "batch_rev" => batch(n, rounds, true), + "ooo" => ooo(n, rounds), + "detached" => detached(n, rounds), + "mlist_batch" => mlist(n, rounds, false), + "mlist_stream" => mlist(n, rounds, true), + "mlist_hole" => mlist_hole(n, true), + "check" => check(), + other => { + eprintln!("unknown scenario {other}"); + std::process::exit(2); + } + } +} + +fn loadavg() -> String { + std::fs::read_to_string("/proc/loadavg") + .ok() + .or_else(|| { + std::process::Command::new("sysctl") + .args(["-n", "vm.loadavg"]) + .output() + .ok() + .map(|o| String::from_utf8_lossy(&o.stdout).trim().to_string()) + }) + .unwrap_or_default() + .split_whitespace() + .take(3) + .collect::>() + .join(",") +} + +fn rss_kb() -> u64 { + std::process::Command::new("ps") + .args(["-o", "rss=", "-p", &std::process::id().to_string()]) + .output() + .ok() + .and_then(|o| String::from_utf8(o.stdout).ok()) + .and_then(|s| s.trim().parse().ok()) + .unwrap_or(0) +} + +#[repr(C)] +struct MallocStats { + blocks_in_use: u32, + size_in_use: usize, + max_size_in_use: usize, + size_allocated: usize, +} + +extern "C" { + fn malloc_default_zone() -> *mut std::ffi::c_void; + fn malloc_zone_statistics(zone: *mut std::ffi::c_void, stats: *mut MallocStats); +} + +fn heap_kb() -> u64 { + unsafe { + let mut stats = MallocStats { + blocks_in_use: 0, + size_in_use: 0, + max_size_in_use: 0, + size_allocated: 0, + }; + malloc_zone_statistics(malloc_default_zone(), &mut stats); + (stats.size_in_use / 1024) as u64 + } +} + +struct Samples { + median: Duration, + min: Duration, + mean: Duration, +} + +fn samples_of(mut times: Vec) -> Samples { + times.sort(); + let total: Duration = times.iter().copied().sum(); + Samples { + median: times[times.len() / 2], + min: times[0], + mean: total / times.len() as u32, + } +} + +fn report(scenario: &str, n: usize, samples: &Samples, extra: &str) { + let ms = |d: Duration| d.as_secs_f64() * 1e3; + println!( + "RESULT scenario={scenario} n={n} median_ms={:.3} min_ms={:.3} mean_ms={:.3} rss_kb={} heap_kb={} {extra}", + ms(samples.median), + ms(samples.min), + ms(samples.mean), + rss_kb(), + heap_kb(), + ); +} + +fn chunk(fat: usize) -> String { + if fat == 1 { + "a".to_string() + } else { + "x".repeat(fat) + } +} + +fn append_text(doc: &LoroDoc, s: &str) { + let text = doc.get_text("t"); + let len = text.len_unicode(); + text.insert(len, s).unwrap(); + doc.commit(); +} + +/// `n` commits from `peer`. Merge interval -1 keeps each commit its own change +/// (0 still merges commits that share a timestamp second). +fn text_history(n: usize, peer: u64, fat: usize) -> LoroDoc { + let doc = LoroDoc::new(); + doc.set_peer_id(peer).unwrap(); + doc.set_change_merge_interval(-1); + let piece = chunk(fat); + for _ in 0..n { + append_text(&doc, &piece); + } + doc +} + +fn peers(peer_count: usize, rounds: usize, text_imports: bool) { + let imports: usize = std::env::var("IMPORTS") + .ok() + .and_then(|s| s.parse().ok()) + .unwrap_or(200); + let setup = Instant::now(); + let base = LoroDoc::new(); + for p in 1..=peer_count { + let d = LoroDoc::new(); + d.set_peer_id(p as u64).unwrap(); + d.get_map("m").insert("k", p as i64).unwrap(); + d.commit(); + base.import(&d.export(ExportMode::all_updates()).unwrap()) + .unwrap(); + } + let frontiers = base.oplog_frontiers().len(); + let src = LoroDoc::new(); + src.set_peer_id(1).unwrap(); + src.set_change_merge_interval(-1); + src.get_map("m").insert("k", 1i64).unwrap(); + src.commit(); + let mut vv = src.oplog_vv(); + let mut blobs = Vec::with_capacity(imports); + for i in 0..imports { + if text_imports { + src.get_text("t").insert(i, "a").unwrap(); + } else { + src.get_map("m").insert(&format!("n{i}"), i as i64).unwrap(); + } + src.commit(); + blobs.push(src.export(ExportMode::updates(&vv)).unwrap()); + vv = src.oplog_vv(); + } + let expect_text = { + let d = base.fork(); + for b in &blobs { + d.import(b).unwrap(); + } + if text_imports { + d.get_text("t").to_string() + } else { + String::new() + } + }; + println!( + "SETUP peers={peer_count} frontiers={frontiers} imports={imports} setup_ms={:.1}", + setup.elapsed().as_secs_f64() * 1e3 + ); + let mut times = Vec::with_capacity(rounds + 1); + for round in 0..=rounds { + let doc = base.fork(); + let start = Instant::now(); + for b in &blobs { + doc.import(b).unwrap(); + } + let elapsed = start.elapsed(); + if text_imports { + assert_eq!(doc.get_text("t").to_string(), expect_text); + } + std::hint::black_box(doc.state_frontiers()); + if round > 0 { + times.push(elapsed); + } + } + report( + if text_imports { + "text_peers" + } else { + "map_peers" + }, + peer_count, + &samples_of(times), + &format!("imports={imports} frontiers={frontiers}"), + ); +} + +fn text_stream(n: usize, rounds: usize, fat: usize) { + let piece = chunk(fat); + // Remote imports merge when the timestamp delta is <= 0. Spreading the + // commit timestamps keeps each update as its own op instead of one growing + // insert, which is what makes a per-import clone of the last op quadratic. + let spread = std::env::var("SPREAD").ok().as_deref() == Some("1"); + let src = LoroDoc::new(); + src.set_peer_id(2).unwrap(); + src.set_change_merge_interval(-1); + let seedsync = std::env::var("SEEDSYNC").ok().as_deref() == Some("1"); + let mut seed_blob = Vec::new(); + if seedsync { + let sd = LoroDoc::new(); + sd.set_peer_id(1).unwrap(); + sd.get_text("t").insert(0, "seed").unwrap(); + sd.commit(); + seed_blob = sd.export(ExportMode::all_updates()).unwrap(); + src.import(&seed_blob).unwrap(); + } + let mut vv = src.oplog_vv(); + let mut blobs = Vec::with_capacity(n); + for i in 0..n { + let text = src.get_text("t"); + let len = text.len_unicode(); + text.insert(len, &piece).unwrap(); + if spread { + src.set_next_commit_timestamp((i as i64) * 10); + } + src.commit(); + blobs.push(src.export(ExportMode::updates(&vv)).unwrap()); + vv = src.oplog_vv(); + } + let expected = format!("seed{}", src.get_text("t").to_string()); + let expected = if seedsync { + src.get_text("t").to_string() + } else { + expected + }; + let mut times = Vec::with_capacity(rounds + 1); + let mut recv_changes = 0; + for round in 0..=rounds { + let doc = LoroDoc::new(); + doc.set_peer_id(1).unwrap(); + let noseed = std::env::var("NOSEED").ok().as_deref() == Some("1"); + if seedsync { + doc.import(&seed_blob).unwrap(); + } else if !noseed { + doc.get_text("t").insert(0, "seed").unwrap(); + doc.commit(); + } + let start = Instant::now(); + for b in &blobs { + doc.import(b).unwrap(); + } + let elapsed = start.elapsed(); + if noseed || seedsync { + assert_eq!(doc.get_text("t").to_string(), src.get_text("t").to_string()); + } else { + assert_eq!(doc.get_text("t").to_string(), expected); + } + recv_changes = doc.len_changes(); + if round > 0 { + times.push(elapsed); + } + } + report( + "text_stream", + n, + &samples_of(times), + &format!( + "fat={fat} spread={spread} src_changes={} recv_changes={recv_changes}", + src.len_changes() + ), + ); +} + +#[derive(Clone, Copy)] +enum OverlapMode { + Mem, + Snap, + Partial, +} + +fn overlap(n: usize, rounds: usize, fat: usize, mode: OverlapMode) { + let src = text_history(n, 1, fat); + let base_vv = src.oplog_vv(); + let base_end = *base_vv.get(&1).expect("peer 1"); + let changes = src.len_changes(); + let base_updates = src.export(ExportMode::all_updates()).unwrap(); + let base_snap = src.export(ExportMode::Snapshot).unwrap(); + append_text(&src, "Z"); + let all_plus = src.export(ExportMode::all_updates()).unwrap(); + let expected = src.get_text("t").to_string(); + let mut mid = base_vv; + mid.set_end(ID::new(1, base_end / 2)); + let partial = src.export(ExportMode::updates(&mid)).unwrap(); + let (name, blob_len) = match mode { + OverlapMode::Mem => ("overlap_mem", all_plus.len()), + OverlapMode::Snap => ("overlap_snap", all_plus.len()), + OverlapMode::Partial => ("overlap_partial", partial.len()), + }; + let mut times = Vec::with_capacity(rounds + 1); + for round in 0..=rounds { + let doc = LoroDoc::new(); + match mode { + OverlapMode::Mem => { + doc.import(&base_updates).unwrap(); + } + OverlapMode::Snap | OverlapMode::Partial => { + doc.import(&base_snap).unwrap(); + } + } + // Leave the loaded doc cold: do not touch len_changes / get_change here. + let blob = match mode { + OverlapMode::Partial => &partial, + _ => &all_plus, + }; + let start = Instant::now(); + doc.import(blob).unwrap(); + let elapsed = start.elapsed(); + assert_eq!(doc.get_text("t").to_string(), expected); + if round > 0 { + times.push(elapsed); + } + } + report( + name, + n, + &samples_of(times), + &format!( + "fat={fat} changes={changes} blob={blob_len} snap={}", + base_snap.len() + ), + ); +} + +fn snap_plus(n: usize, rounds: usize, fat: usize) { + let src = text_history(n, 1, fat); + let changes = src.len_changes(); + let base_snap = src.export(ExportMode::Snapshot).unwrap(); + append_text(&src, "Z"); + let snap_plus = src.export(ExportMode::Snapshot).unwrap(); + let expected = src.get_text("t").to_string(); + let mut times = Vec::with_capacity(rounds + 1); + let mut last_heap = 0; + for round in 0..=rounds { + let doc = LoroDoc::new(); + doc.import(&base_snap).unwrap(); + let before = heap_kb(); + let start = Instant::now(); + doc.import(&snap_plus).unwrap(); + let elapsed = start.elapsed(); + last_heap = heap_kb().saturating_sub(before); + assert_eq!(doc.get_text("t").to_string(), expected); + if round > 0 { + times.push(elapsed); + } + drop(doc); + } + report( + "snap_plus", + n, + &samples_of(times), + &format!( + "fat={fat} changes={changes} snap={} snap_plus={} heap_delta_kb={last_heap}", + base_snap.len(), + snap_plus.len() + ), + ); +} + +/// Import a growing snapshot (full history + one new op) into the same doc. +/// Heap after each import is the retained-memory signal. +fn snap_mem(n: usize, fat: usize) { + let repeats: usize = std::env::var("REPEATS") + .ok() + .and_then(|s| s.parse().ok()) + .unwrap_or(6); + let src = text_history(n, 1, fat); + let base_snap = src.export(ExportMode::Snapshot).unwrap(); + let mut snaps = Vec::with_capacity(repeats); + let mut only_new = Vec::with_capacity(repeats); + for _ in 0..repeats { + let vv = src.oplog_vv(); + append_text(&src, "Z"); + snaps.push(src.export(ExportMode::Snapshot).unwrap()); + only_new.push(src.export(ExportMode::updates(&vv)).unwrap()); + } + let doc = LoroDoc::new(); + doc.import(&base_snap).unwrap(); + println!( + "MEM kind=snap_base heap_kb={} rss_kb={} snap={}", + heap_kb(), + rss_kb(), + base_snap.len() + ); + for (i, snap) in snaps.iter().enumerate() { + let before = heap_kb(); + let start = Instant::now(); + doc.import(snap).unwrap(); + let ms = start.elapsed().as_secs_f64() * 1e3; + println!( + "MEM kind=snap i={i} import_ms={ms:.3} heap_kb={} delta_kb={} rss_kb={} bytes={}", + heap_kb(), + heap_kb() as i64 - before as i64, + rss_kb(), + snap.len() + ); + } + assert_eq!(doc.get_text("t").to_string(), src.get_text("t").to_string()); + + let control = LoroDoc::new(); + control.import(&base_snap).unwrap(); + let before = heap_kb(); + for (i, upd) in only_new.iter().enumerate() { + let start = Instant::now(); + control.import(upd).unwrap(); + println!( + "MEM kind=update i={i} import_ms={:.3} heap_kb={} rss_kb={} bytes={}", + start.elapsed().as_secs_f64() * 1e3, + heap_kb(), + rss_kb(), + upd.len() + ); + } + println!( + "MEM kind=update_total delta_kb={}", + heap_kb() as i64 - before as i64 + ); + assert_eq!( + control.get_text("t").to_string(), + src.get_text("t").to_string() + ); +} + +fn one_peer_blobs(n: usize) -> Vec> { + let src = LoroDoc::new(); + src.set_peer_id(2).unwrap(); + src.set_change_merge_interval(-1); + let mut vv = src.oplog_vv(); + let mut blobs = Vec::with_capacity(n); + for _ in 0..n { + append_text(&src, "a"); + blobs.push(src.export(ExportMode::updates(&vv)).unwrap()); + vv = src.oplog_vv(); + } + blobs +} + +fn seeded() -> LoroDoc { + let doc = LoroDoc::new(); + doc.set_peer_id(1).unwrap(); + doc.get_text("t").insert(0, "seed").unwrap(); + doc.commit(); + doc +} + +fn batch(n: usize, rounds: usize, reverse: bool) { + let mut blobs = one_peer_blobs(n); + if reverse { + blobs.reverse(); + } + let expected = { + let doc = seeded(); + doc.import_batch(&blobs).unwrap(); + doc.get_text("t").to_string() + }; + let mut times = Vec::with_capacity(rounds + 1); + for round in 0..=rounds { + let doc = seeded(); + let start = Instant::now(); + doc.import_batch(&blobs).unwrap(); + let elapsed = start.elapsed(); + assert_eq!(doc.get_text("t").to_string(), expected); + if round > 0 { + times.push(elapsed); + } + } + report( + if reverse { "batch_rev" } else { "batch" }, + n, + &samples_of(times), + &format!("blobs={n}"), + ); +} + +fn ooo(n: usize, rounds: usize) { + let mut blobs = one_peer_blobs(n); + blobs.reverse(); + let expected = { + let doc = seeded(); + for b in &blobs { + doc.import(b).unwrap(); + } + doc.get_text("t").to_string() + }; + let mut times = Vec::with_capacity(rounds + 1); + for round in 0..=rounds { + let doc = seeded(); + let start = Instant::now(); + for b in &blobs { + doc.import(b).unwrap(); + } + let elapsed = start.elapsed(); + assert_eq!(doc.get_text("t").to_string(), expected); + if round > 0 { + times.push(elapsed); + } + } + report("ooo", n, &samples_of(times), &format!("blobs={n}")); +} + +fn detached(n: usize, rounds: usize) { + let blobs = one_peer_blobs(n); + let expected = { + let doc = seeded(); + doc.detach(); + for b in &blobs { + doc.import(b).unwrap(); + } + doc.attach(); + doc.get_text("t").to_string() + }; + let mut import_times = Vec::with_capacity(rounds + 1); + let mut attach_times = Vec::with_capacity(rounds + 1); + for round in 0..=rounds { + let doc = seeded(); + doc.detach(); + let start = Instant::now(); + for b in &blobs { + doc.import(b).unwrap(); + } + let import_elapsed = start.elapsed(); + let start = Instant::now(); + doc.attach(); + let attach_elapsed = start.elapsed(); + assert_eq!(doc.get_text("t").to_string(), expected); + assert!(!doc.is_detached()); + if round > 0 { + import_times.push(import_elapsed); + attach_times.push(attach_elapsed); + } + } + let imports = samples_of(import_times); + let attaches = samples_of(attach_times); + report( + "detached", + n, + &imports, + &format!( + "attach_median_ms={:.3} attach_min_ms={:.3}", + attaches.median.as_secs_f64() * 1e3, + attaches.min.as_secs_f64() * 1e3 + ), + ); +} + +fn mlist_elements(n: usize) -> (LoroDoc, Vec) { + let src = LoroDoc::new(); + src.set_peer_id(1).unwrap(); + src.set_change_merge_interval(-1); + let list = src.get_movable_list("m"); + for i in 0..n { + list.insert(i, i as i64).unwrap(); + src.commit(); + } + let snap = src.export(ExportMode::Snapshot).unwrap(); + (src, snap) +} + +fn mlist(n: usize, rounds: usize, stream: bool) { + let (src, snap) = mlist_elements(n); + let changes = src.len_changes(); + let base_vv = src.oplog_vv(); + let editor = LoroDoc::new(); + editor.import(&snap).unwrap(); + editor.set_peer_id(2).unwrap(); + editor.set_change_merge_interval(-1); + let list = editor.get_movable_list("m"); + let mut per_op = Vec::with_capacity(n); + let mut vv = editor.oplog_vv(); + for i in 0..n { + list.set(i, (i as i64) + 1).unwrap(); + editor.commit(); + per_op.push(editor.export(ExportMode::updates(&vv)).unwrap()); + vv = editor.oplog_vv(); + } + let batch_blob = editor.export(ExportMode::updates(&base_vv)).unwrap(); + let expected = editor.get_movable_list("m").get_deep_value(); + let mut times = Vec::with_capacity(rounds + 1); + for round in 0..=rounds { + let doc = LoroDoc::new(); + doc.import(&snap).unwrap(); + let start = Instant::now(); + if stream { + for b in &per_op { + doc.import(b).unwrap(); + } + } else { + doc.import(&batch_blob).unwrap(); + } + let elapsed = start.elapsed(); + assert_eq!(doc.get_movable_list("m").get_deep_value(), expected); + if round > 0 { + times.push(elapsed); + } + } + report( + if stream { + "mlist_stream" + } else { + "mlist_batch" + }, + n, + &samples_of(times), + &format!( + "changes={changes} snap={} blob={}", + snap.len(), + if stream { + per_op.iter().map(Vec::len).sum::() + } else { + batch_blob.len() + } + ), + ); +} + +fn check() { + let mut failed = 0; + let mut run = |name: &str, f: fn()| { + let start = Instant::now(); + match std::panic::catch_unwind(f) { + Ok(()) => println!("PASS {name} {:.1}ms", start.elapsed().as_secs_f64() * 1e3), + Err(payload) => { + failed += 1; + let msg = payload + .downcast_ref::() + .map(String::as_str) + .or_else(|| payload.downcast_ref::<&str>().copied()) + .unwrap_or("panic"); + println!("FAIL {name}: {msg}"); + } + } + }; + run("overlap_snap_equals", overlap_snap_equals); + run("overlap_partial_equals", overlap_partial_equals); + run("merged_change_prefix", merged_change_prefix); + run("deletes_roundtrip_overlap", deletes_roundtrip_overlap); + run("random_text_overlap", random_text_overlap); + run("snap_repeat_equals", snap_repeat_equals); + run("mlist_sets_after_snapshot", mlist_sets_after_snapshot); + run("mlist_moves_after_snapshot", mlist_moves_after_snapshot); + run("mlist_after_get_change", mlist_after_get_change); + run("mlist_set_before_insert", mlist_set_before_insert); + run("batch_rev_equals_forward", batch_rev_equals_forward); + run("ooo_equals_forward", ooo_equals_forward); + run("detached_equals_attached", detached_equals_attached); + run("two_peers_snapshot_overlap", two_peers_snapshot_overlap); + run("conflict_same_peer", conflict_same_peer); + if failed > 0 { + eprintln!("CHECK_FAILS {failed}"); + std::process::exit(1); + } + println!("CHECK_OK"); +} + +fn overlap_snap_equals() { + let src = text_history(80, 1, 8); + let snap = src.export(ExportMode::Snapshot).unwrap(); + append_text(&src, "Q"); + let doc = LoroDoc::new(); + doc.import(&snap).unwrap(); + doc.import(&src.export(ExportMode::all_updates()).unwrap()) + .unwrap(); + assert_eq!(doc.get_text("t").to_string(), src.get_text("t").to_string()); +} + +fn overlap_partial_equals() { + let src = text_history(80, 1, 4); + let vv = src.oplog_vv(); + let end = *vv.get(&1).unwrap(); + let snap = src.export(ExportMode::Snapshot).unwrap(); + append_text(&src, "Q"); + let mut mid = vv; + mid.set_end(ID::new(1, end / 2)); + let doc = LoroDoc::new(); + doc.import(&snap).unwrap(); + doc.import(&src.export(ExportMode::updates(&mid)).unwrap()) + .unwrap(); + assert_eq!(doc.get_text("t").to_string(), src.get_text("t").to_string()); +} + +fn merged_change_prefix() { + let src = LoroDoc::new(); + src.set_peer_id(1).unwrap(); + src.set_change_merge_interval(10_000); + for _ in 0..40 { + append_text(&src, "m"); + } + assert_eq!(src.len_changes(), 1, "expected one merged change"); + let end = *src.oplog_vv().get(&1).unwrap(); + let receiver = LoroDoc::new(); + receiver.set_peer_id(1).unwrap(); + receiver.set_change_merge_interval(10_000); + for _ in 0..(end as usize / 2) { + append_text(&receiver, "m"); + } + append_text(&src, "Z"); + receiver + .import(&src.export(ExportMode::all_updates()).unwrap()) + .unwrap(); + assert_eq!( + receiver.get_text("t").to_string(), + src.get_text("t").to_string() + ); +} + +fn deletes_roundtrip_overlap() { + let src = LoroDoc::new(); + src.set_peer_id(1).unwrap(); + src.set_change_merge_interval(-1); + for _ in 0..30 { + append_text(&src, "abcd"); + } + let text = src.get_text("t"); + for i in 0..10 { + let len = text.len_unicode(); + text.delete(len / 2, (i % 3) + 1).unwrap(); + src.commit(); + } + let snap = src.export(ExportMode::Snapshot).unwrap(); + append_text(&src, "Z"); + let doc = LoroDoc::new(); + doc.import(&snap).unwrap(); + doc.import(&src.export(ExportMode::all_updates()).unwrap()) + .unwrap(); + assert_eq!(doc.get_text("t").to_string(), src.get_text("t").to_string()); +} + +fn random_text_overlap() { + for seed in [1u64, 2, 3, 7, 11, 19, 23, 42] { + let mut state = seed; + let mut next = || { + state = state.wrapping_mul(6364136223846793005).wrapping_add(1); + state + }; + let src = LoroDoc::new(); + src.set_peer_id(1).unwrap(); + src.set_change_merge_interval(-1); + for _ in 0..60 { + let text = src.get_text("t"); + let len = text.len_unicode(); + if len > 4 && next() % 5 == 0 { + let at = (next() as usize) % len; + let n = (1 + (next() as usize) % 3).min(len - at); + text.delete(at, n).unwrap(); + } else { + let at = if len == 0 { + 0 + } else { + (next() as usize) % (len + 1) + }; + let s = if next() % 2 == 0 { "xy" } else { "z" }; + text.insert(at, s).unwrap(); + } + src.commit(); + } + let snap = src.export(ExportMode::Snapshot).unwrap(); + append_text(&src, "Q"); + let doc = LoroDoc::new(); + doc.import(&snap).unwrap(); + doc.import(&src.export(ExportMode::Snapshot).unwrap()) + .unwrap(); + assert_eq!( + doc.get_text("t").to_string(), + src.get_text("t").to_string(), + "seed {seed}" + ); + } +} + +fn snap_repeat_equals() { + let src = text_history(40, 1, 8); + let doc = LoroDoc::new(); + doc.import(&src.export(ExportMode::Snapshot).unwrap()) + .unwrap(); + for _ in 0..5 { + append_text(&src, "Z"); + doc.import(&src.export(ExportMode::Snapshot).unwrap()) + .unwrap(); + assert_eq!(doc.get_text("t").to_string(), src.get_text("t").to_string()); + } +} + +fn mlist_sets_after_snapshot() { + let (src, snap) = mlist_elements(120); + let vv = src.oplog_vv(); + let editor = LoroDoc::new(); + editor.import(&snap).unwrap(); + editor.set_peer_id(2).unwrap(); + let list = editor.get_movable_list("m"); + for i in 0..120 { + list.set(i, (i as i64) * 10).unwrap(); + } + editor.commit(); + let doc = LoroDoc::new(); + doc.import(&snap).unwrap(); + doc.import(&editor.export(ExportMode::updates(&vv)).unwrap()) + .unwrap(); + assert_eq!( + doc.get_movable_list("m").get_deep_value(), + editor.get_movable_list("m").get_deep_value() + ); +} + +fn mlist_moves_after_snapshot() { + let (src, snap) = mlist_elements(60); + let vv = src.oplog_vv(); + let editor = LoroDoc::new(); + editor.import(&snap).unwrap(); + editor.set_peer_id(2).unwrap(); + let list = editor.get_movable_list("m"); + for i in 0..20 { + list.mov(i, 59 - i).unwrap(); + } + editor.commit(); + let doc = LoroDoc::new(); + doc.import(&snap).unwrap(); + doc.import(&editor.export(ExportMode::updates(&vv)).unwrap()) + .unwrap(); + assert_eq!( + doc.get_movable_list("m").get_deep_value(), + editor.get_movable_list("m").get_deep_value() + ); +} + +fn mlist_after_get_change() { + // Enough inserts that the history is more than one 4KB block. + let (src, snap) = mlist_elements(800); + assert!(src.len_changes() > 1); + let vv = src.oplog_vv(); + let editor = LoroDoc::new(); + editor.import(&snap).unwrap(); + editor.set_peer_id(2).unwrap(); + let list = editor.get_movable_list("m"); + for i in (0..800).step_by(7) { + list.set(i, -1i64).unwrap(); + } + editor.commit(); + let blob = editor.export(ExportMode::updates(&vv)).unwrap(); + let doc = LoroDoc::new(); + doc.import(&snap).unwrap(); + // Parse an early block, then import sets of later elements. + assert!(doc.get_change(ID::new(1, 0)).is_some()); + doc.import(&blob).unwrap(); + assert_eq!( + doc.get_movable_list("m").get_deep_value(), + editor.get_movable_list("m").get_deep_value() + ); +} + +/// Snapshot-load `n` elements, optionally parse an early block, then import a +/// set of every element. A lamport lookup that stops on the cached early block +/// rejects the later sets. +fn mlist_hole(n: usize, prime: bool) { + let (src, snap) = mlist_elements(n); + let vv = src.oplog_vv(); + let editor = LoroDoc::new(); + editor.import(&snap).unwrap(); + editor.set_peer_id(2).unwrap(); + let list = editor.get_movable_list("m"); + for i in 0..n { + list.set(i, -1i64).unwrap(); + } + editor.commit(); + let blob = editor.export(ExportMode::updates(&vv)).unwrap(); + let doc = LoroDoc::new(); + doc.import(&snap).unwrap(); + if prime { + assert!(doc.get_change(ID::new(1, 0)).is_some()); + } + doc.import(&blob).unwrap(); + assert_eq!( + doc.get_movable_list("m").get_deep_value(), + editor.get_movable_list("m").get_deep_value() + ); + println!( + "PASS mlist_hole n={n} prime={prime} changes={}", + src.len_changes() + ); +} + +fn mlist_set_before_insert() { + let src = LoroDoc::new(); + src.set_peer_id(1).unwrap(); + src.get_movable_list("m").insert(0, "a").unwrap(); + src.commit(); + let insert_blob = src.export(ExportMode::all_updates()).unwrap(); + let vv = src.oplog_vv(); + src.get_movable_list("m").set(0, "b").unwrap(); + src.commit(); + let set_blob = src.export(ExportMode::updates(&vv)).unwrap(); + let doc = LoroDoc::new(); + doc.import(&set_blob).unwrap(); + doc.import(&insert_blob).unwrap(); + assert_eq!( + doc.get_movable_list("m").get_deep_value(), + src.get_movable_list("m").get_deep_value() + ); +} + +fn batch_rev_equals_forward() { + let blobs = one_peer_blobs(40); + let mut rev = blobs.clone(); + rev.reverse(); + let a = seeded(); + a.import_batch(&blobs).unwrap(); + let b = seeded(); + b.import_batch(&rev).unwrap(); + assert_eq!(a.get_text("t").to_string(), b.get_text("t").to_string()); +} + +fn ooo_equals_forward() { + let blobs = one_peer_blobs(40); + let mut rev = blobs.clone(); + rev.reverse(); + let a = seeded(); + for b in &blobs { + a.import(b).unwrap(); + } + let b = seeded(); + for blob in &rev { + b.import(blob).unwrap(); + } + assert_eq!(a.get_text("t").to_string(), b.get_text("t").to_string()); +} + +fn detached_equals_attached() { + let blobs = one_peer_blobs(30); + let attached = seeded(); + for b in &blobs { + attached.import(b).unwrap(); + } + let detached = seeded(); + detached.detach(); + for b in &blobs { + detached.import(b).unwrap(); + } + detached.attach(); + assert_eq!( + attached.get_text("t").to_string(), + detached.get_text("t").to_string() + ); +} + +fn two_peers_snapshot_overlap() { + let a = text_history(40, 1, 4); + let b = text_history(40, 2, 4); + let both = LoroDoc::new(); + both.import(&a.export(ExportMode::all_updates()).unwrap()) + .unwrap(); + both.import(&b.export(ExportMode::all_updates()).unwrap()) + .unwrap(); + append_text(&both, "Z"); + let receiver = LoroDoc::new(); + receiver + .import(&a.export(ExportMode::Snapshot).unwrap()) + .unwrap(); + receiver + .import(&both.export(ExportMode::Snapshot).unwrap()) + .unwrap(); + assert_eq!( + receiver.get_text("t").to_string(), + both.get_text("t").to_string() + ); +} + +fn conflict_same_peer() { + let a = LoroDoc::new(); + a.set_peer_id(1).unwrap(); + a.get_text("t").insert(0, "aaa").unwrap(); + a.commit(); + let b = LoroDoc::new(); + b.set_peer_id(1).unwrap(); + b.get_text("t").insert(0, "bbb").unwrap(); + b.commit(); + b.get_text("t").insert(3, "X").unwrap(); + b.commit(); + let doc = LoroDoc::new(); + doc.import(&a.export(ExportMode::all_updates()).unwrap()) + .unwrap(); + let result = doc.import(&b.export(ExportMode::all_updates()).unwrap()); + println!("INFO conflict_same_peer {result:?}"); +} diff --git a/crates/examples/import_overlap_measure.py b/crates/examples/import_overlap_measure.py new file mode 100644 index 000000000..ef8c0bd3a --- /dev/null +++ b/crates/examples/import_overlap_measure.py @@ -0,0 +1,90 @@ +#!/usr/bin/env python3 +"""Interleave release binaries of the import scaling probe on a shared host. + +Build/copy each revision's import_scaling_stress executable as old/main/fixed, +then run with --bin-dir and --output-dir. All dependencies are Python stdlib. +""" +import argparse +import csv +import json +import os +import re +import statistics +import subprocess +from pathlib import Path + +SCENARIOS = [ + "overlap_snap", "overlap_partial", "overlap_mem", "snap_plus", "snap_mem", + "detached", "mlist_batch", "text_stream", "batch", +] + + +def main(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--bin-dir", type=Path, required=True) + parser.add_argument("--output-dir", type=Path, required=True) + parser.add_argument("--fat", type=int, default=64) + parser.add_argument("--scenarios", nargs="+", choices=SCENARIOS, default=SCENARIOS) + args = parser.parse_args() + args.output_dir.mkdir(parents=True, exist_ok=True) + rows = [] + with (args.output_dir / "raw.jsonl").open("w") as raw: + for scenario in args.scenarios: + for n in [8000, 32000, 64000]: + for rep in range(3): + for version in ["old", "main", "fixed"]: + env = os.environ.copy() + env.update(CARGO_BUILD_JOBS="4", FAT="1", REPEATS="4") + env.pop("SPREAD", None) + env.pop("SEEDSYNC", None) + env.pop("NOSEED", None) + if scenario in SCENARIOS[:5]: + env["FAT"] = str(args.fat) + if scenario == "text_stream": + env["NOSEED"] = "1" + run = subprocess.run( + [str(args.bin_dir.resolve() / version), scenario, str(n), "3"], + env=env, capture_output=True, text=True, + ) + load = subprocess.run( + ["uptime"], capture_output=True, text=True, + ).stdout.strip() + record = dict( + scenario=scenario, n=n, rep=rep, version=version, + load=load, stdout=run.stdout, stderr=run.stderr, + exit=run.returncode, + ) + raw.write(json.dumps(record) + "\n") + raw.flush() + if run.returncode: + raise RuntimeError(record) + if scenario == "snap_mem": + samples = re.findall( + r"MEM kind=snap i=\d+ import_ms=([\d.]+) " + r"heap_kb=(\d+) delta_kb=(-?\d+)", run.stdout, + ) + base = int(re.search( + r"MEM kind=snap_base heap_kb=(\d+)", run.stdout, + )[1]) + ms = statistics.median(float(x[0]) for x in samples) + heap = int(samples[-1][1]) + delta = heap - base + else: + ms = float(re.search(r"median_ms=([\d.]+)", run.stdout)[1]) + heap = 0 + match = re.search(r"heap_delta_kb=(\d+)", run.stdout) + delta = int(match[1]) if match else 0 + rows.append(dict( + scenario=scenario, n=n, rep=rep, version=version, + ms=ms, heap_kb=heap, delta_kb=delta, + )) + print(f"{scenario} {n} rep {rep + 1}/3", flush=True) + with (args.output_dir / "samples.csv").open("w") as output: + writer = csv.DictWriter(output, fieldnames=rows[0].keys()) + writer.writeheader() + writer.writerows(rows) + print("MEASUREMENT_COMPLETE", flush=True) + + +if __name__ == "__main__": + main() diff --git a/crates/loro-internal/src/arena.rs b/crates/loro-internal/src/arena.rs index ab7588ede..871faf6f2 100644 --- a/crates/loro-internal/src/arena.rs +++ b/crates/loro-internal/src/arena.rs @@ -122,6 +122,13 @@ pub(crate) struct ArenaExtent { str_bytes: usize, } +#[cfg(test)] +impl ArenaExtent { + pub(crate) fn values_for_test(&self) -> usize { + self.values + } +} + impl SharedArenaRollback { /// Whether rolling back to this checkpoint keeps everything within `extent`. pub(crate) fn keeps(&self, extent: ArenaExtent) -> bool { diff --git a/crates/loro-internal/src/encoding/fast_snapshot.rs b/crates/loro-internal/src/encoding/fast_snapshot.rs index 526fe2c1f..08ed522b6 100644 --- a/crates/loro-internal/src/encoding/fast_snapshot.rs +++ b/crates/loro-internal/src/encoding/fast_snapshot.rs @@ -358,11 +358,7 @@ pub(crate) fn decode_oplog(oplog: &mut OpLog, bytes: &[u8]) -> Result Result bool { ) } +/// Detached imports do not apply state. Only Move/Set element validation can +/// reject after adding history; an enclosing batch still owns its full scope. +pub(super) fn import_op_can_reject(op: &crate::op::Op, detached: bool) -> bool { + if detached { + op.container.get_type() == ContainerType::MovableList + && matches!( + op.content, + InnerContent::List( + list_op::InnerListOp::Move { .. } | list_op::InnerListOp::Set { .. } + ) + ) + } else { + state_apply_can_reject(op.container.get_type()) + } +} + impl OpLog { #[inline] pub(crate) fn new(visible_op_count: Arc) -> Self { @@ -373,7 +390,11 @@ impl OpLog { } } - pub(crate) fn preflight_import_changes(&self, changes: &[Change]) -> ImportChangesPreflight { + pub(crate) fn preflight_import_changes( + &self, + changes: &[Change], + detached: bool, + ) -> ImportChangesPreflight { let mut ans = ImportChangesPreflight::default(); for change in changes { if change.ctr_end() <= self.vv().get(&change.id.peer).copied().unwrap_or(0) { @@ -390,7 +411,7 @@ impl OpLog { if change .ops .iter() - .any(|op| state_apply_can_reject(op.container.get_type())) + .any(|op| import_op_can_reject(op, detached)) { ans.needs_state_apply_rollback = true; } @@ -411,20 +432,20 @@ impl OpLog { // small sync/import workloads. // // Scan pending last, and only when it can still change the answer: - // `has_state_apply_rollback_ops` walks every parked change and every op in + // `has_import_rollback_ops` walks every parked change and every op in // it, while a blob that only parks (deps not here yet) leaves // `applies_to_dag` false. Evaluating it eagerly made a batch of N // out-of-order blobs quadratic, since each blob re-scanned the pending set // the earlier blobs had grown. if ans.applies_to_dag && !ans.needs_state_apply_rollback - && self.pending_changes.has_state_apply_rollback_ops() + && self.pending_changes.has_import_rollback_ops(detached) { ans.needs_state_apply_rollback = true; } #[cfg(test)] - if ans.applies_to_dag { + if ans.applies_to_dag && !detached { ans.needs_state_apply_rollback = true; } diff --git a/crates/loro-internal/src/oplog/change_store.rs b/crates/loro-internal/src/oplog/change_store.rs index 08c7b28b9..2dbafb5df 100644 --- a/crates/loro-internal/src/oplog/change_store.rs +++ b/crates/loro-internal/src/oplog/change_store.rs @@ -406,37 +406,160 @@ impl ChangeStore { self.external_kv.lock().export_all() } - /// Decode the changes of a snapshot imported into a non-empty doc. - /// - /// Changes this doc already has are kept, because - /// `OpLog::check_and_trim_known_part_of_changes` compares them with the local - /// history before trimming them. When nothing is new, return no changes and - /// skip the clones. + /// Read a cold block for comparison without caching parsed ops in the arena. + /// Cached parsed history keeps using the ordinary, already cheap path. + pub(crate) fn unparsed_block_bytes(&self, id: ID) -> Option { + let external = self.external_kv.lock(); + let inner = self.inner.lock(); + if inner.retired { + return None; + } + if let Some((_, block)) = inner.mem_parsed_kv.range(..=id).next_back() { + if block.peer == id.peer && block.counter_range.1 > id.counter { + return match &block.content { + ChangesBlockContent::Bytes(b) => Some(b.bytes.clone()), + _ => None, + }; + } + } + let (key, bytes) = external + .scan(Bound::Unbounded, Bound::Included(&id.to_bytes())) + .rfind(|(key, _)| key.len() == 12)?; + if ID::from_bytes(&key).peer != id.peer + || decode_block_range(&bytes).ok()?.0 .1 <= id.counter + { + return None; + } + Some(bytes) + } + + pub(crate) fn encoded_block_counter_range(bytes: &[u8]) -> LoroResult<(Counter, Counter)> { + Ok(decode_block_range(bytes)?.0) + } + + pub(crate) fn check_text_insert_block( + &self, + id: ID, + bytes: &[u8], + on_insert: impl FnMut(Counter, &ContainerID, u32, &str, u32, &Frontiers) -> LoroResult<()>, + ) -> LoroResult { + let result = block_encode::visit_text_insert_block(bytes, on_insert); + if let Err(err) = &result { + if !matches!(err, LoroError::UsedOpID { .. }) { + self.parse_failures.record(id, err); + } + } + result + } + + /// Byte-identical, fully known blocks need neither decoding nor comparison. + /// An unflushed cached block shadows its older KV copy; never compare against + /// that stale copy. Different encodings fall back to semantic comparison. + pub(crate) fn contains_encoded_block(&self, id: ID, bytes: &[u8]) -> bool { + let external = self.external_kv.lock(); + let inner = self.inner.lock(); + if let Some(block) = inner.mem_parsed_kv.get(&id) { + return match &block.content { + ChangesBlockContent::Bytes(b) | ChangesBlockContent::Both(_, b) => { + b.bytes.as_ref() == bytes + } + ChangesBlockContent::Changes(_) => false, + }; + } + external + .get(&id.to_bytes()) + .is_some_and(|b| b.as_ref() == bytes) + } + + /// Decode only unmatched snapshot blocks in a temporary arena. Known content + /// never gets copied into the document's arena; move only the checked suffix. pub(crate) fn decode_snapshot_for_updates( bytes: Bytes, - arena: &SharedArena, - self_vv: &VersionVector, + oplog: &crate::OpLog, ) -> Result, LoroError> { - let change_store = ChangeStore::new_mem(arena, Arc::new(AtomicI64::new(0))); - let _ = change_store.import_all(bytes)?; - let mut has_new = false; - change_store.visit_all_changes(&mut |c| { - has_new |= c.ctr_end() > self_vv.get(&c.id.peer).copied().unwrap_or(0); - }); + let arena = SharedArena::new(); + let store = ChangeStore::new_mem(&arena, Arc::new(AtomicI64::new(0))); + let _ = store.import_all(bytes)?; + let external = store.external_kv.lock(); + let mut inner = store.inner.lock(); let mut changes = Vec::new(); - if has_new { - change_store.visit_all_changes(&mut |c| changes.push(c.clone())); + for (key, bytes) in external.scan(Bound::Unbounded, Bound::Unbounded) { + if key.len() != 12 { + continue; + } + let id = ID::from_bytes(&key); + if oplog.change_store.contains_encoded_block(id, &bytes) { + continue; + } + // import_all parsed the frontier blocks. Take their changes instead + // of parsing them again or cloning changes that will be dropped. + if let Some(block) = inner.mem_parsed_kv.remove(&id) { + let block = Arc::try_unwrap(block).expect("temporary store owns its blocks"); + match block.content { + ChangesBlockContent::Changes(c) | ChangesBlockContent::Both(c, _) => { + changes + .extend(Arc::try_unwrap(c).expect("temporary store owns its changes")); + } + ChangesBlockContent::Bytes(b) => changes.extend(b.parse(&arena)?), + } + } else { + changes.extend(Self::decode_block_bytes(bytes, &arena)?); + } } - - Ok(changes) + drop(inner); + drop(external); + changes.sort_unstable_by_key(|c| c.lamport); + let changes = oplog.check_and_trim_known_part_in_arena( + changes, + super::ImportedValues::Exact, + &arena, + )?; + Ok(changes + .into_iter() + .map(|mut change| { + let mut ops = RleVec::new(); + for op in change.ops.iter() { + for remote in super::local_op_to_remote(&arena, op) { + ops.push(oplog.arena.convert_single_op( + &remote.container, + change.id.peer, + remote.counter, + change.lamport + (remote.counter - change.id.counter) as Lamport, + remote.content, + )); + } + } + change.ops = ops; + register_container_and_parent_link(&oplog.arena, &change); + change + }) + .collect()) } - /// Decode an update block. Changes the doc already has are kept; see - /// [`Self::decode_snapshot_for_updates`]. pub(crate) fn decode_block_bytes(bytes: Bytes, arena: &SharedArena) -> LoroResult> { ChangesBlockBytes::new(bytes).parse(arena) } + /// Reuse the header on fallback so a new block is decoded only once. + pub(crate) fn decode_update_block( + &self, + bytes: Bytes, + vv: &VersionVector, + ) -> LoroResult> { + let block = ChangesBlockBytes::new(bytes); + block.ensure_header()?; + let header = block.header.get().unwrap(); + let id = ID::new(header.peer, header.counter); + if header.counters.last().copied().unwrap_or(0) + <= vv.get(&header.peer).copied().unwrap_or(0) + && self.contains_encoded_block(id, &block.bytes) + { + Ok(Vec::new()) + } else { + block.parse(&self.arena) + } + } + /// Rolls back the store and the arena (to `arena`, the checkpoint taken when the import /// began). See [`Self::rollback_arena`]. pub(crate) fn rollback_import( @@ -2251,6 +2374,99 @@ mod test { assert_eq!(changes_parsed, changes); } + #[test] + fn identical_encoded_block_stays_lazy_and_dirty_cache_shadows_kv() { + let source = LoroDoc::new_auto_commit(); + source.set_peer_id(1).unwrap(); + source.set_change_merge_interval(-1); + for _ in 0..40 { + let text = source.get_text("t"); + text.insert(text.len_unicode(), "abcd", PosType::Unicode) + .unwrap(); + source.commit_then_renew(); + } + let bytes = { + let oplog = source.oplog().lock(); + oplog.change_store.encode_all(oplog.vv(), oplog.frontiers()) + }; + let store = ChangeStore::new_for_test(); + store.import_all(bytes).unwrap(); + let (id, block_bytes) = store + .external_kv + .lock() + .scan(Bound::Unbounded, Bound::Unbounded) + .find(|(key, _)| key.len() == 12) + .map(|(key, bytes)| (ID::from_bytes(&key), bytes)) + .unwrap(); + assert!(!store.inner.lock().mem_parsed_kv.contains_key(&id)); + let before = store.arena.utf16_len(); + assert!(store + .decode_update_block(block_bytes.clone(), &source.oplog_vv()) + .unwrap() + .is_empty()); + assert_eq!(store.arena.utf16_len(), before); + assert!(!store.inner.lock().mem_parsed_kv.contains_key(&id)); + + let original = store.get_change(id).unwrap(); + let mut dirty = original.block.as_ref().clone(); + // Model an unflushed mutation. Its stale KV copy must not authorize a + // shortcut, even if the imported bytes match that old copy exactly. + let mut changes = dirty.content.try_changes().unwrap().clone(); + changes[0].deps = Frontiers::from_id(ID::new(9, 0)); + dirty.content = ChangesBlockContent::Changes(Arc::new(changes)); + dirty.flushed = false; + store.inner.lock().mem_parsed_kv.insert(id, Arc::new(dirty)); + assert!(!store.contains_encoded_block(id, &block_bytes)); + assert!(!store + .decode_update_block(block_bytes, &source.oplog_vv()) + .unwrap() + .is_empty()); + } + + #[test] + fn repeated_snapshot_updates_allocate_only_the_new_text_and_list_values() { + let source = LoroDoc::new_auto_commit(); + source.set_peer_id(1).unwrap(); + source.set_change_merge_interval(-1); + for i in 0..40 { + let text = source.get_text("t"); + text.insert(text.len_unicode(), "wรถrld ๐Ÿ˜€", PosType::Unicode) + .unwrap(); + source.get_list("l").push(i).unwrap(); + source.commit_then_renew(); + } + let target = LoroDoc::new_auto_commit(); + target + .import(&source.export(ExportMode::Snapshot).unwrap()) + .unwrap(); + // Once local history is parsed, repeated snapshots must not allocate + // another copy of any known strings or list values, even with Unicode. + target + .oplog() + .lock() + .change_store + .visit_all_changes(&mut |_| {}); + let text_before = target.oplog().lock().arena.utf16_len(); + let value_count = |doc: &LoroDoc| doc.oplog().lock().arena.extent().values_for_test(); + let values_before = value_count(&target); + for i in 0..4 { + let text = source.get_text("t"); + text.insert(text.len_unicode(), "Z", PosType::Unicode) + .unwrap(); + source.get_list("l").push(100 + i).unwrap(); + source.commit_then_renew(); + target + .import(&source.export(ExportMode::Snapshot).unwrap()) + .unwrap(); + assert_eq!( + target.oplog().lock().arena.utf16_len(), + text_before + i as usize + 1 + ); + assert_eq!(value_count(&target), values_before + i as usize + 1); + } + assert_eq!(target.get_deep_value(), source.get_deep_value()); + } + #[test] fn decoded_block_lamport_range_matches_counter_range() { // Regression test for a checkout hang after snapshot import. diff --git a/crates/loro-internal/src/oplog/change_store/block_encode.rs b/crates/loro-internal/src/oplog/change_store/block_encode.rs index 7d99ece62..e1abe38eb 100644 --- a/crates/loro-internal/src/oplog/change_store/block_encode.rs +++ b/crates/loro-internal/src/oplog/change_store/block_encode.rs @@ -473,19 +473,37 @@ pub fn decode_block_range( ) -> LoroResult<((Counter, Counter), (Lamport, Lamport))> { let counter_start = leb128::read::unsigned(&mut bytes).map_err(|e| { LoroError::DecodeError(format!("Failed to read counter start: {e}").into_boxed_str()) - })? as Counter; + })?; + let counter_start = + Counter::try_from(counter_start).map_err(|_| LoroError::DecodeDataCorruptionError)?; let counter_len = leb128::read::unsigned(&mut bytes).map_err(|e| { LoroError::DecodeError(format!("Failed to read counter length: {e}").into_boxed_str()) - })? as Counter; + })?; + let counter_len = + Counter::try_from(counter_len).map_err(|_| LoroError::DecodeDataCorruptionError)?; let lamport_start = leb128::read::unsigned(&mut bytes).map_err(|e| { LoroError::DecodeError(format!("Failed to read lamport start: {e}").into_boxed_str()) - })? as Lamport; + })?; + let lamport_start = + Lamport::try_from(lamport_start).map_err(|_| LoroError::DecodeDataCorruptionError)?; let lamport_len = leb128::read::unsigned(&mut bytes).map_err(|e| { LoroError::DecodeError(format!("Failed to read lamport length: {e}").into_boxed_str()) - })? as Lamport; + })?; + let lamport_len = + Lamport::try_from(lamport_len).map_err(|_| LoroError::DecodeDataCorruptionError)?; Ok(( - (counter_start, counter_start + counter_len), - (lamport_start, lamport_start + lamport_len), + ( + counter_start, + counter_start + .checked_add(counter_len) + .ok_or(LoroError::DecodeDataCorruptionError)?, + ), + ( + lamport_start, + lamport_start + .checked_add(lamport_len) + .ok_or(LoroError::DecodeDataCorruptionError)?, + ), )) } @@ -526,6 +544,70 @@ pub fn decode_cids( Ok(header) } +/// Visit a Text-insert-only block without building changes or allocating arena +/// strings. Other op kinds use the general decoder and semantic comparison. +/// Timestamps/messages are intentionally irrelevant to known-history equality. +pub(super) fn visit_text_insert_block( + bytes: &[u8], + mut on_insert: impl FnMut(Counter, &ContainerID, u32, &str, u32, &Frontiers) -> LoroResult<()>, +) -> LoroResult { + use crate::encoding::value::{ValueKind, ValueReader}; + let doc: EncodedBlock<'_> = + postcard::from_bytes(bytes).map_err(|_| LoroError::DecodeDataCorruptionError)?; + let header = decode_cids(bytes, None)?; + let cids = header.cids.get().unwrap(); + let mut reader = ValueReader::new(&doc.values); + let ops = serde_columnar::iter_from_bytes::(&doc.ops)?.ops; + let mut counter = header.counter; + let mut change_index = 0; + for op in ops { + let op = op?; + let cid = cids + .get(op.container_index as usize) + .ok_or(LoroError::DecodeDataCorruptionError)?; + if cid.container_type() != loro_common::ContainerType::Text + || op.value_type != ValueKind::Str.to_u8() + { + return Ok(false); + } + let text = reader.read_str()?; + if text.chars().count() != op.len as usize { + return Err(LoroError::DecodeDataCorruptionError); + } + let implicit_deps = Frontiers::from_id(ID::new( + header.peer, + counter + .checked_sub(1) + .ok_or(LoroError::DecodeDataCorruptionError)?, + )); + let deps = if header.counters.get(change_index) == Some(&counter) { + header + .deps_groups + .get(change_index) + .ok_or(LoroError::DecodeDataCorruptionError)? + } else { + &implicit_deps + }; + on_insert(counter, cid, op.prop as u32, text, op.len, deps)?; + counter = counter + .checked_add(op.len as Counter) + .ok_or(LoroError::DecodeDataCorruptionError)?; + if header.counters.get(change_index + 1) == Some(&counter) { + change_index += 1; + } else if header + .counters + .get(change_index + 1) + .is_none_or(|&end| counter > end) + { + return Err(LoroError::DecodeDataCorruptionError); + } + } + if Some(&counter) != header.counters.last() { + return Err(LoroError::DecodeDataCorruptionError); + } + Ok(true) +} + // MARK: decode_block pub fn decode_block( m_bytes: &[u8], @@ -711,6 +793,22 @@ mod test { use super::*; + #[test] + fn block_range_rejects_overflow_instead_of_panicking_on_external_bytes() { + for fields in [ + [u64::MAX, 0, 0, 0], + [Counter::MAX as u64, 1, 0, 0], + [0, u64::MAX, 0, 0], + [0, 0, Lamport::MAX as u64, 1], + ] { + let mut bytes = Vec::new(); + for field in fields { + leb128::write::unsigned(&mut bytes, field).unwrap(); + } + assert!(decode_block_range(&bytes).is_err()); + } + } + #[test] fn decode_block_rejects_corrupt_payload_without_panic() { let doc = LoroDoc::new_auto_commit(); diff --git a/crates/loro-internal/src/oplog/known_history.rs b/crates/loro-internal/src/oplog/known_history.rs index a5ad2fd09..c49b19719 100644 --- a/crates/loro-internal/src/oplog/known_history.rs +++ b/crates/loro-internal/src/oplog/known_history.rs @@ -48,16 +48,32 @@ impl OpLog { /// would be applied on top of a conflicting prefix. An import made only of /// known changes is a no-op either way and stays as cheap as before. pub(crate) fn check_and_trim_known_part_of_changes( + &self, + changes: Vec, + values: ImportedValues, + ) -> LoroResult> { + self.check_and_trim_known_part_in_arena(changes, values, &self.arena) + } + + /// Snapshot overlap is decoded into a temporary arena. Compare before moving + /// only the retained suffix into the document's arena. + pub(crate) fn check_and_trim_known_part_in_arena( &self, mut changes: Vec, values: ImportedValues, + imported_arena: &SharedArena, ) -> LoroResult> { let vv = self.vv(); let known_end = |c: &Change| vv.get(&c.id.peer).copied().unwrap_or(0); if changes.iter().any(|c| c.ctr_end() > known_end(c)) { for change in changes.iter() { if change.id.counter < known_end(change) { - self.check_known_part_of_change(change, known_end(change), values)?; + self.check_known_part_of_change( + change, + known_end(change), + values, + imported_arena, + )?; } } } @@ -134,6 +150,7 @@ impl OpLog { change: &Change, known_end: Counter, values: ImportedValues, + imported_arena: &SharedArena, ) -> LoroResult<()> { let peer = change.id.peer; // History before the shallow root is not stored, so it cannot be compared. @@ -141,6 +158,12 @@ impl OpLog { let end = change.ctr_end().min(known_end); let mut ctr = change.id.counter.max(shallow_start); while ctr < end { + if let Some(checked_end) = + self.check_known_text_in_cold_block(change, imported_arena, ctr, end)? + { + ctr = checked_end; + continue; + } // Missing local history cannot be compared; keep the old behavior for it. let Some(local) = self.change_store.get_change(ID::new(peer, ctr)) else { break; @@ -156,7 +179,15 @@ impl OpLog { }); } - if let Some(bad) = first_mismatch(&self.arena, values, change, &local, ctr, seg_end) { + if let Some(bad) = first_mismatch( + imported_arena, + &self.arena, + values, + change, + &local, + ctr, + seg_end, + ) { return Err(LoroError::UsedOpID { id: ID::new(peer, bad), }); @@ -167,6 +198,103 @@ impl OpLog { Ok(()) } + /// Compare cold Text insert history directly from its encoded block. This + /// preserves every dependency boundary and op atom without retaining the + /// block's parsed changes or strings. Unsupported blocks use the full check. + fn check_known_text_in_cold_block( + &self, + change: &Change, + arena: &SharedArena, + start: Counter, + end: Counter, + ) -> LoroResult> { + let Some(bytes) = self + .change_store + .unparsed_block_bytes(ID::new(change.id.peer, start)) + else { + return Ok(None); + }; + let range = crate::oplog::ChangeStore::encoded_block_counter_range(&bytes)?; + // A short imported change must not rescan a large local block for each + // atom/change. Parse it once and reuse the general comparison instead. + if end < range.1 { + return Ok(None); + } + let block_end = end.min(range.1); + let mut cursor = OpCursor::new(change, start); + let checked = self.change_store.check_text_insert_block( + ID::new(change.id.peer, range.0), + &bytes, + |counter, cid, pos, text, len, deps| { + let local_end = counter + len as Counter; + let mut ctr = counter.max(start); + let limit = local_end.min(end); + while ctr < limit { + let bad = || LoroError::UsedOpID { + id: ID::new(change.id.peer, ctr), + }; + let local_deps = if ctr == counter { + deps.clone() + } else { + Frontiers::from_id(ID::new(change.id.peer, ctr - 1)) + }; + if deps_at(change, ctr) != local_deps { + return Err(bad()); + } + let Some(op) = cursor.op_at(ctr) else { + return Err(bad()); + }; + let n = (limit - ctr).min(op.ctr_end() - ctr); + let slice = slice_op(op, ctr, n); + let InnerContent::List(InnerListOp::InsertText { + slice: imported, + pos: imported_pos, + unicode_len, + .. + }) = &slice.content + else { + return Err(bad()); + }; + if arena.idx_to_id(op.container).as_ref() != Some(cid) + || pos + .checked_add((ctr - counter) as u32) + .is_none_or(|p| *imported_pos != p) + || *unicode_len != n as u32 + { + return Err(bad()); + } + let imported = std::str::from_utf8(imported).unwrap(); + let offset = (ctr - counter) as usize; + let local = if offset == 0 && n as u32 == len { + text + } else { + let byte_start = text + .char_indices() + .nth(offset) + .map_or(text.len(), |(i, _)| i); + let byte_end = text[byte_start..] + .char_indices() + .nth(n as usize) + .map_or(text.len(), |(i, _)| byte_start + i); + &text[byte_start..byte_end] + }; + if local != imported { + let offset = local + .chars() + .zip(imported.chars()) + .position(|(a, b)| a != b) + .unwrap_or(0); + return Err(LoroError::UsedOpID { + id: ID::new(change.id.peer, ctr + offset as Counter), + }); + } + ctr += n; + } + Ok(()) + }, + )?; + Ok(checked.then_some(block_end)) + } } /// The deps of the op at `ctr`. Inside a change every op depends on the previous @@ -184,7 +312,8 @@ fn deps_at(change: &Change, ctr: Counter) -> Frontiers { /// which they differ. Each side may split or merge the ops differently, so both /// are cut at every boundary of either side before comparing. fn first_mismatch( - arena: &SharedArena, + a_arena: &SharedArena, + b_arena: &SharedArena, values: ImportedValues, a: &Change, b: &Change, @@ -203,12 +332,13 @@ fn first_mismatch( .min(b_op.ctr_end() - ctr); let a_slice = slice_op(a_op, ctr, len); let b_slice = slice_op(b_op, ctr, len); - if !op_eq(arena, values, &a_slice, &b_slice) { + if !op_eq(a_arena, b_arena, values, &a_slice, &b_slice) { // Error path only: find the exact atom. let offset = (0..len) .find(|&i| { !op_eq( - arena, + a_arena, + b_arena, values, &slice_op(a_op, ctr + i, 1), &slice_op(b_op, ctr + i, 1), @@ -260,15 +390,26 @@ fn slice_op(op: &Op, ctr: Counter, len: Counter) -> std::borrow::Cow<'_, Op> { /// Whether two ops with the same id range are the same op. Compares what the op /// means, not how it is stored: arena offsets and the direction of a one-atom /// delete are representation details. -fn op_eq(arena: &SharedArena, values: ImportedValues, a: &Op, b: &Op) -> bool { - if a.container != b.container || a.atom_len() != b.atom_len() { +fn op_eq( + a_arena: &SharedArena, + b_arena: &SharedArena, + values: ImportedValues, + a: &Op, + b: &Op, +) -> bool { + let same_container = if std::ptr::eq(a_arena, b_arena) { + a.container == b.container + } else { + a_arena.idx_to_id(a.container) == b_arena.idx_to_id(b.container) + }; + if !same_container || a.atom_len() != b.atom_len() { return false; } let exact = values == ImportedValues::Exact; match (&a.content, &b.content) { - (InnerContent::List(a), InnerContent::List(b)) => list_op_eq(arena, exact, a, b), + (InnerContent::List(a), InnerContent::List(b)) => list_op_eq(a_arena, b_arena, exact, a, b), (InnerContent::Map(a), InnerContent::Map(b)) => { a.key == b.key && a.value.is_some() == b.value.is_some() @@ -297,7 +438,13 @@ fn op_eq(arena: &SharedArena, values: ImportedValues, a: &Op, b: &Op) -> bool { } } -fn list_op_eq(arena: &SharedArena, exact: bool, a: &InnerListOp, b: &InnerListOp) -> bool { +fn list_op_eq( + a_arena: &SharedArena, + b_arena: &SharedArena, + exact: bool, + a: &InnerListOp, + b: &InnerListOp, +) -> bool { match (a, b) { ( InnerListOp::Insert { @@ -310,7 +457,14 @@ fn list_op_eq(arena: &SharedArena, exact: bool, a: &InnerListOp, b: &InnerListOp }, ) => { a_pos == b_pos - && (!exact || arena.value_slices_eq(a_slice.to_range(), b_slice.to_range())) + && (!exact + || if std::ptr::eq(a_arena, b_arena) { + a_arena.value_slices_eq(a_slice.to_range(), b_slice.to_range()) + } else { + a_arena.with_values(a_slice.to_range(), |a| { + b_arena.with_values(b_slice.to_range(), |b| a == b) + }) + }) } ( InnerListOp::InsertText { @@ -382,3 +536,92 @@ fn list_op_eq(arena: &SharedArena, exact: bool, a: &InnerListOp, b: &InnerListOp _ => false, } } + +#[cfg(test)] +mod tests { + use super::*; + use crate::{ + container::list::list_op::ListOp, + op::{ListSlice, RawOpContent}, + }; + use loro_common::{ContainerID, ContainerType}; + + fn text_change(arena: &SharedArena, text: &str) -> Change { + let op = arena.convert_single_op( + &ContainerID::new_root("t", ContainerType::Text), + 1, + 0, + 0, + RawOpContent::List(ListOp::Insert { + slice: ListSlice::RawStr { + str: text.into(), + unicode_len: text.chars().count(), + }, + pos: 0, + }), + ); + let mut ops = rle::RleVec::new(); + ops.push(op); + Change { + id: ID::new(1, 0), + lamport: 0, + timestamp: 0, + commit_msg: None, + deps: Frontiers::default(), + ops, + } + } + + #[test] + fn different_arenas_compare_container_ids_and_unicode_atoms() { + let a = SharedArena::new(); + let b = SharedArena::new(); + b.register_container(&ContainerID::new_root("other", ContainerType::Map)); + b.alloc_str("unrelated"); + let left = text_change(&a, "a๐Ÿ˜€bc"); + let same = text_change(&b, "a๐Ÿ˜€bc"); + let different = text_change(&b, "a๐Ÿ˜€Bc"); + assert_eq!( + first_mismatch(&a, &b, ImportedValues::Exact, &left, &same, 0, 4), + None + ); + assert_eq!( + first_mismatch(&a, &b, ImportedValues::Exact, &left, &different, 1, 4), + Some(2) + ); + } + + #[test] + fn cold_text_check_compares_dependencies_inside_a_merged_import() { + use crate::{cursor::PosType, LoroDoc}; + use std::sync::{atomic::AtomicI64, Arc}; + let source = LoroDoc::new_auto_commit(); + source.set_peer_id(1).unwrap(); + source.set_change_merge_interval(-1); + for i in 0..40 { + source + .get_text("t") + .insert(i, "a", PosType::Unicode) + .unwrap(); + source.commit_then_renew(); + } + let oplog = source.oplog().lock(); + let forged = crate::oplog::ChangeStore::new_mem(&oplog.arena, Arc::new(AtomicI64::new(-1))); + for i in 0..40 { + let mut c = (*oplog.get_change_at(ID::new(1, i)).unwrap()).clone(); + if i == 1 { + c.deps = Frontiers::default(); + } + forged.insert_change(c, false, true); + } + let bytes = forged.encode_all(oplog.vv(), oplog.frontiers()); + let local = crate::OpLog::new(Default::default()); + local.change_store.import_all(bytes).unwrap(); + let arena = SharedArena::new(); + let incoming = text_change(&arena, &"a".repeat(40)); + let err = local + .check_known_text_in_cold_block(&incoming, &arena, 0, 40) + .unwrap_err(); + assert!(matches!(err, LoroError::UsedOpID { id } if id == ID::new(1, 1))); + } +} diff --git a/crates/loro-internal/src/oplog/pending_changes.rs b/crates/loro-internal/src/oplog/pending_changes.rs index 78e65da0b..268f72a99 100644 --- a/crates/loro-internal/src/oplog/pending_changes.rs +++ b/crates/loro-internal/src/oplog/pending_changes.rs @@ -37,14 +37,14 @@ pub(crate) struct PendingChanges { } impl PendingChanges { - pub(crate) fn has_state_apply_rollback_ops(&self) -> bool { + pub(crate) fn has_import_rollback_ops(&self, detached: bool) -> bool { self.changes.values().any(|tree| { tree.values().any(|changes| { changes.iter().any(|change| { change .ops .iter() - .any(|op| super::state_apply_can_reject(op.container.get_type())) + .any(|op| super::import_op_can_reject(op, detached)) }) }) }) diff --git a/crates/loro-internal/src/tests/import_atomicity.rs b/crates/loro-internal/src/tests/import_atomicity.rs index bcd8cccf0..9252a56ef 100644 --- a/crates/loro-internal/src/tests/import_atomicity.rs +++ b/crates/loro-internal/src/tests/import_atomicity.rs @@ -863,3 +863,47 @@ fn import_batch_with_move_of_unknown_movable_list_elem_rolls_back() { doc.commit_then_renew(); assert_eq!(doc.state_frontiers(), doc.oplog_frontiers()); } + +#[test] +fn detached_text_import_that_unlocks_invalid_movable_list_ops_rolls_back() { + let (doc, _) = binary_update_moving_unknown_movable_list_elem(); + let text_peer = LoroDoc::new_auto_commit(); + text_peer.set_peer_id(2).unwrap(); + text_peer + .import(&doc.export(ExportMode::all_updates()).unwrap()) + .unwrap(); + text_peer + .get_text("t") + .insert(0, "unlock", PosType::Unicode) + .unwrap(); + text_peer.commit_then_renew(); + let text_update = text_peer + .export(ExportMode::updates(&doc.oplog_vv())) + .unwrap(); + let bad = binary_update_bypassing_validation( + &text_peer, + serde_json::json!({ + "schema_version": 1, "start_version": {}, "peers": ["2", "3"], + "changes": [{ + "id": "0@1", "timestamp": 0, "deps": ["0@0"], "lamport": 8, "msg": null, + "ops": [{ + "container": "cid:root-list:MovableList", "counter": 0, + "content": {"type": "move", "from": 0, "to": 1, "elem_id": "L99@0"} + }] + }] + }), + ); + let vv = doc.oplog_vv(); + let frontiers = doc.oplog_frontiers(); + let state = doc.get_deep_value(); + doc.detach(); + doc.import(&bad).unwrap(); + let pending_before = doc.oplog().lock().pending_changes_len(); + assert!(pending_before > 0); + let err = doc.import(&text_update).unwrap_err(); + assert!(matches!(err, LoroError::DecodeError(_)), "{err:?}"); + assert!(doc.is_detached()); + assert_eq!(doc.oplog().lock().pending_changes_len(), pending_before); + doc.attach(); + assert_doc_unchanged(&doc, &vv, &frontiers, &state); +} diff --git a/crates/loro/tests/import_reused_peer_id.rs b/crates/loro/tests/import_reused_peer_id.rs index 796492a70..de05c5893 100644 --- a/crates/loro/tests/import_reused_peer_id.rs +++ b/crates/loro/tests/import_reused_peer_id.rs @@ -461,3 +461,45 @@ fn reimporting_into_a_shallow_doc_still_succeeds() -> LoroResult<()> { ); Ok(()) } + +/// Exercise many cold blocks, a merged update, and the snapshot scratch arena. +/// The first mismatch must still name the exact atom far inside the prefix. +#[test] +fn cold_history_conflict_inside_a_large_known_prefix_is_rejected() -> LoroResult<()> { + let history = LoroDoc::new(); + let other = LoroDoc::new(); + for doc in [&history, &other] { + doc.set_peer_id(7)?; + doc.set_change_merge_interval(-1); + } + for i in 0..3000 { + history.get_text("t").insert(i, "a")?; + history.commit(); + other + .get_text("t") + .insert(i, if i == 2000 { "O" } else { "a" })?; + other.commit(); + } + other.get_text("t").insert(3000, "Z")?; + other.commit(); + let base = history.export(ExportMode::Snapshot)?; + for blob in [ + other.export(ExportMode::all_updates())?, + other.export(ExportMode::Snapshot)?, + ] { + for batch in [false, true] { + let target = LoroDoc::new(); + target.import(&base)?; + let before = snapshot_of(&target); + let err = if batch { + target.import_batch(&[blob.clone()]).unwrap_err() + } else { + target.import(&blob).unwrap_err() + }; + assert_used_op_id(err, ID::new(7, 2000)); + assert_unchanged(&target, &before); + assert_still_usable(&target); + } + } + Ok(()) +} From 86e5efdba322edc37e2ea0e07d421b9d3e773867 Mon Sep 17 00:00:00 2001 From: Zixuan Chen Date: Thu, 1 Oct 2026 05:24:47 +0800 Subject: [PATCH 2/2] fix: preserve import rollback and parsing decisions Co-Authored-By: GPT-6 (Codex) --- .changeset/fix-import-overlap-cost.md | 2 + context/arena-parent-links.md | 8 +- context/import-batch-atomicity.md | 9 +++ context/import-peer-id-reuse.md | 21 ++++- crates/loro-internal/src/loro.rs | 11 +-- .../loro-internal/src/oplog/change_store.rs | 76 +++++++++++++++++-- .../src/oplog/change_store/block_encode.rs | 62 +++++++++++++++ .../loro-internal/src/oplog/known_history.rs | 5 +- crates/loro/tests/import_reused_peer_id.rs | 29 +++++++ 9 files changed, 202 insertions(+), 21 deletions(-) diff --git a/.changeset/fix-import-overlap-cost.md b/.changeset/fix-import-overlap-cost.md index 27ee3a0c4..76142bd88 100644 --- a/.changeset/fix-import-overlap-cost.md +++ b/.changeset/fix-import-overlap-cost.md @@ -5,3 +5,5 @@ Reduce known-history comparison costs for snapshot-loaded documents while preserving peer-id reuse detection. Snapshot updates allocate only new history in the document arena and avoid cloning discarded changes. Avoid unnecessary rollback scopes for detached Text imports while retaining validation of pending movable-list operations. + +Keep import rollback decisions consistent across concurrent attach/detach, and fall back to ordinary history parsing when a cold-block comparison is ineligible. diff --git a/context/arena-parent-links.md b/context/arena-parent-links.md index d3eee0769..92750c3ca 100644 --- a/context/arena-parent-links.md +++ b/context/arena-parent-links.md @@ -112,7 +112,7 @@ that return `bool` or `Option` (`is_deleted`, `has_container`, `get_path`), so it has no `Err` to return, and a panic there unwinds under the locks and traps the WASM instance. Until 2026-09-30 it panicked. -Every reader of the change store that fails to decode or parse a block records +Every ordinary history reader that fails to decode or parse a block records it in `ChangeStore::parse_failures` (a leaf lock; the first failure is kept) and answers "no such change", as the readers other than the resolver always did. The resolver answers `CreatorOp::Corrupt`, which the arena treats like `Absent` @@ -144,8 +144,10 @@ Limits: - A read that returns before any failure was recorded is not undone: an `import` or `checkout` using a general history reader can finish on partial history when it is the first to read the block. Only `export` checks again at - the end. The direct cold Text comparison in `known_history.rs` returns its - decode error immediately and also records it in `parse_failures`. + the end. The direct cold Text comparison in `known_history.rs` falls back to + the ordinary reader on a decoding/eligibility error. Its stricter op-length + and change-boundary checks never record a `parse_failures` entry themselves; + only a failure of the ordinary parser declares local history unparsable. - Local edits still succeed after a failure was recorded, but they cannot be exported from this document any more. What can be salvaged is the current state (`get_deep_value`). diff --git a/context/import-batch-atomicity.md b/context/import-batch-atomicity.md index 83017d7d1..284854b45 100644 --- a/context/import-batch-atomicity.md +++ b/context/import-batch-atomicity.md @@ -104,6 +104,15 @@ MovableList pending change, so preflight includes pending `Move`/`Set` ops when replaced or nested by this optimization. Regression: `detached_text_import_that_unlocks_invalid_movable_list_ops_rolls_back`. +`import_changes_and_apply_delta_to_state_if_needed` reads the detached flag once +and uses that value for both preflight and the execution branches. Attach/detach +can change the flag without holding the op log lock; reading it again could +choose attached state application after preflight skipped its rollback scope. +The batch keeps the txn mutex across force-detach, blob imports and reattach, +and keeps its outer rollback scope; the per-blob decision does not replace it. +Under `cfg(test)`, preflight forces rollback only for +`applies_to_dag && !detached`, so detached tests still exercise this distinction. + `crates/examples/examples/import_batch_perf.rs` is the ad-hoc probe for these shapes (not part of CI; run it on two revisions and compare). diff --git a/context/import-peer-id-reuse.md b/context/import-peer-id-reuse.md index ebe8cd479..1ded5c612 100644 --- a/context/import-peer-id-reuse.md +++ b/context/import-peer-id-reuse.md @@ -93,9 +93,16 @@ encoded container IDs, positions, UTF-8 strings and dependency boundaries (`OpLog::check_known_text_in_cold_block`, `block_encode::visit_text_insert_block`). This does not parse changes or allocate strings in the local arena. Both change and op boundaries remain significant to the comparison. Unsupported op kinds -fall back to the general check; already parsed blocks keep the existing path. +or any reader error fall back to the general check; already parsed blocks keep +the existing path. A mismatch from this shortcut is returned only after the +whole block passes its eligibility checks. The shortcut never records a +`parse_failures` entry: its length/boundary requirements are stricter than +`decode_block`, so only the ordinary parser can declare history unparsable. A short imported change that ends inside a cold block uses the general path so -many short changes cannot repeatedly scan the same large block. +many short changes cannot repeatedly scan the same large block. The cold path +runs only when the imported overlap reaches the local block's end. Short +changes with different encoded bytes still parse their cold block, which +explains the remaining roughly 2.2x `overlap_snap` cost at `FAT=1`. The comparison still costs O(overlap); it is not a content-hashed version vector. The release probe is `crates/examples/examples/import_scaling_stress.rs`. @@ -117,10 +124,18 @@ universal ratio. - `oplog::known_history::tests`: cross-arena container identity and Unicode atom comparisons, and a dependency mismatch inside a merged import's cold prefix. - `oplog::change_store::test`: identical blocks stay lazy, dirty caches shadow - old KV bytes, and repeated snapshots allocate only new text/list values. + old KV bytes, repeated snapshots allocate only new text/list values, and + `merged_cold_text_overlap_is_accepted_without_parsing_the_known_block` proves + that an equal merged overlap stays cold before a successful import. +- `block_encode::test::cold_text_reader_errors_fall_back_without_recording_parse_failures`: + strings whose lengths differ from encoded op lengths and ops crossing change + boundaries remain readable through `get_change`, as with `decode_block`. - `cold_history_conflict_inside_a_large_known_prefix_is_rejected` in the public reused-peer tests: exact mismatch ID across multiple cold blocks, through updates, snapshots, and batches, with subsequent usability checks. +- `cold_text_overlap_with_merged_updates_is_accepted` in the same public tests: + merged updates with different block bytes accept a cold snapshot prefix and + leave both documents at equal state and version. ## Measurements diff --git a/crates/loro-internal/src/loro.rs b/crates/loro-internal/src/loro.rs index dd77047ea..a92afd381 100644 --- a/crates/loro-internal/src/loro.rs +++ b/crates/loro-internal/src/loro.rs @@ -782,15 +782,16 @@ impl LoroDoc { } }; - let preflight = oplog.preflight_import_changes(&changes, self.is_detached()); - if preflight.has_deps_before_shallow_root - && (self.is_detached() || !preflight.applies_to_dag) - { + // Read once: attach/detach can change the flag without the oplog lock. + // The apply branch must use the same mode as preflight's rollback decision. + let detached = self.is_detached(); + let preflight = oplog.preflight_import_changes(&changes, detached); + if preflight.has_deps_before_shallow_root && (detached || !preflight.applies_to_dag) { oplog.rollback_arena(arena_checkpoint, &mut self.state.lock()); return Err(LoroError::ImportUpdatesThatDependsOnOutdatedVersion); } - if self.is_detached() { + if detached { // An enclosing `import_batch` scope validates the whole batch before it // reattaches (`BatchImportGuard::finish`). let owns_rollback = diff --git a/crates/loro-internal/src/oplog/change_store.rs b/crates/loro-internal/src/oplog/change_store.rs index 2dbafb5df..1644097cf 100644 --- a/crates/loro-internal/src/oplog/change_store.rs +++ b/crates/loro-internal/src/oplog/change_store.rs @@ -439,17 +439,24 @@ impl ChangeStore { pub(crate) fn check_text_insert_block( &self, - id: ID, bytes: &[u8], - on_insert: impl FnMut(Counter, &ContainerID, u32, &str, u32, &Frontiers) -> LoroResult<()>, + mut on_insert: impl FnMut(Counter, &ContainerID, u32, &str, u32, &Frontiers) -> LoroResult<()>, ) -> LoroResult { - let result = block_encode::visit_text_insert_block(bytes, on_insert); - if let Err(err) = &result { - if !matches!(err, LoroError::UsedOpID { .. }) { - self.parse_failures.record(id, err); - } + // This reader has stricter eligibility checks than decode_block. Only + // the normal parser may declare local history unparsable. Delay a content + // mismatch until the entire block is eligible, otherwise fall back too. + let mut comparison = Ok(()); + let result = + block_encode::visit_text_insert_block(bytes, |counter, cid, pos, text, len, deps| { + if comparison.is_ok() { + comparison = on_insert(counter, cid, pos, text, len, deps); + } + Ok(()) + }); + match result { + Ok(true) => comparison.map(|()| true), + Ok(false) | Err(_) => Ok(false), } - result } /// Byte-identical, fully known blocks need neither decoding nor comparison. @@ -2423,6 +2430,59 @@ mod test { .is_empty()); } + #[test] + fn merged_cold_text_overlap_is_accepted_without_parsing_the_known_block() { + let history = LoroDoc::new_auto_commit(); + history.set_peer_id(1).unwrap(); + history.set_change_merge_interval(-1); + for i in 0..40 { + history.get_text("t").insert_unicode(i * 2, "a๐Ÿ˜€").unwrap(); + history.commit_then_renew(); + } + let target = LoroDoc::new_auto_commit(); + target + .import(&history.export(ExportMode::Snapshot).unwrap()) + .unwrap(); + + let merged = LoroDoc::new_auto_commit(); + merged.set_peer_id(1).unwrap(); + merged + .get_text("t") + .insert_unicode(0, &("a๐Ÿ˜€".repeat(40) + "Z")) + .unwrap(); + merged.commit_then_renew(); + { + let incoming = merged.oplog().lock(); + let change = (*incoming.get_change_at(ID::new(1, 0)).unwrap()).clone(); + let local = target.oplog().lock(); + // cfg(test) loads the frontier block at snapshot import; explicitly + // clear that cache to exercise the real cold-block comparison. + local.change_store.inner.lock().mem_parsed_kv.clear(); + let id = ID::new(1, 0); + let bytes = local.change_store.unparsed_block_bytes(id).unwrap(); + assert_ne!( + bytes.as_ref(), + block_encode::encode_block(&[change.clone()], &incoming.arena) + ); + let suffix = local + .check_and_trim_known_part_in_arena( + vec![change], + crate::oplog::known_history::ImportedValues::Exact, + &incoming.arena, + ) + .unwrap(); + assert_eq!(suffix.len(), 1); + assert_eq!(suffix[0].id, ID::new(1, 80)); + // A general comparison would have populated the parsed cache. + assert!(local.change_store.unparsed_block_bytes(id).is_some()); + } + target + .import(&merged.export(ExportMode::all_updates()).unwrap()) + .unwrap(); + assert_eq!(target.get_deep_value(), merged.get_deep_value()); + assert_eq!(target.oplog_vv(), merged.oplog_vv()); + } + #[test] fn repeated_snapshot_updates_allocate_only_the_new_text_and_list_values() { let source = LoroDoc::new_auto_commit(); diff --git a/crates/loro-internal/src/oplog/change_store/block_encode.rs b/crates/loro-internal/src/oplog/change_store/block_encode.rs index e1abe38eb..48ef273c2 100644 --- a/crates/loro-internal/src/oplog/change_store/block_encode.rs +++ b/crates/loro-internal/src/oplog/change_store/block_encode.rs @@ -832,6 +832,68 @@ mod test { assert!(result.unwrap().is_err()); } + #[test] + fn cold_text_reader_errors_fall_back_without_recording_parse_failures() { + let doc = LoroDoc::new_auto_commit(); + doc.set_peer_id(1).unwrap(); + doc.set_change_merge_interval(-1); + doc.get_text("t").insert_unicode(0, "ab").unwrap(); + doc.commit_then_renew(); + doc.get_text("t").insert_unicode(2, "c").unwrap(); + doc.commit_then_renew(); + let oplog = doc.oplog.lock(); + let mut changes = Vec::new(); + oplog + .change_store() + .visit_all_changes(&mut |change| changes.push(change.clone())); + assert_eq!(changes.len(), 2); + let block_bytes = encode_block(&changes, &oplog.arena); + + for crosses_change_boundary in [false, true] { + let mut encoded: EncodedBlock = postcard::from_bytes(&block_bytes).unwrap(); + if crosses_change_boundary { + // The header starts with a peer count, peer IDs, then the first + // change's atom length. Move its end from counter 2 to counter 1 + // while keeping the two-character insert and total range intact. + let mut header = encoded.header.as_ref(); + let peers = leb128::read::unsigned(&mut header).unwrap() as usize; + let length_offset = encoded.header.len() - header.len() + peers * 8; + assert_eq!(encoded.header[length_offset], 2); + encoded.header.to_mut()[length_offset] = 1; + } else { + // decode_block uses the actual string length for the text op, + // while the cold reader requires it to equal the encoded length. + let mut ops: EncodedOps = serde_columnar::from_bytes(&encoded.ops).unwrap(); + ops.ops[0].len = 1; + encoded.ops = Cow::Owned(serde_columnar::to_vec(&ops).unwrap()); + } + let bytes = postcard::to_allocvec(&encoded).unwrap(); + assert!(visit_text_insert_block(&bytes, |_, _, _, _, _, _| Ok(())).is_err()); + assert!(decode_block(&bytes, &SharedArena::new(), None).is_ok()); + + let store = crate::oplog::ChangeStore::new_for_test(); + let id = ID::new(1, 0); + store + .external_kv + .lock() + .set(&id.to_bytes(), bytes.clone().into()); + let mut visits = 0; + assert!(!store + .check_text_insert_block(&bytes, |_, _, _, _, _, _| { + visits += 1; + // Even a mismatch before a later reader error must fall back. + Err(LoroError::UsedOpID { id }) + }) + .unwrap()); + if crosses_change_boundary { + assert_eq!(visits, 1); + } + assert!(store.corrupt_block_error().is_ok()); + assert!(store.get_change(id).is_some()); + assert!(store.corrupt_block_error().is_ok()); + } + } + #[test] fn tree_move_payload_is_rejected_without_panic() { let arena = SharedArena::new(); diff --git a/crates/loro-internal/src/oplog/known_history.rs b/crates/loro-internal/src/oplog/known_history.rs index c49b19719..66e744d3e 100644 --- a/crates/loro-internal/src/oplog/known_history.rs +++ b/crates/loro-internal/src/oplog/known_history.rs @@ -214,7 +214,9 @@ impl OpLog { else { return Ok(None); }; - let range = crate::oplog::ChangeStore::encoded_block_counter_range(&bytes)?; + let Ok(range) = crate::oplog::ChangeStore::encoded_block_counter_range(&bytes) else { + return Ok(None); + }; // A short imported change must not rescan a large local block for each // atom/change. Parse it once and reuse the general comparison instead. if end < range.1 { @@ -223,7 +225,6 @@ impl OpLog { let block_end = end.min(range.1); let mut cursor = OpCursor::new(change, start); let checked = self.change_store.check_text_insert_block( - ID::new(change.id.peer, range.0), &bytes, |counter, cid, pos, text, len, deps| { let local_end = counter + len as Counter; diff --git a/crates/loro/tests/import_reused_peer_id.rs b/crates/loro/tests/import_reused_peer_id.rs index de05c5893..bc885d459 100644 --- a/crates/loro/tests/import_reused_peer_id.rs +++ b/crates/loro/tests/import_reused_peer_id.rs @@ -503,3 +503,32 @@ fn cold_history_conflict_inside_a_large_known_prefix_is_rejected() -> LoroResult } Ok(()) } + +/// The same history can arrive as one merged Text insert instead of thousands +/// of small changes. Its blocks differ from the snapshot's cold history, but +/// the overlap is equal and its new suffix must be accepted. +#[test] +fn cold_text_overlap_with_merged_updates_is_accepted() -> LoroResult<()> { + let history = LoroDoc::new(); + history.set_peer_id(7)?; + history.set_change_merge_interval(-1); + for i in 0..3000 { + history.get_text("t").insert(i * 2, "a๐Ÿ˜€")?; + history.commit(); + } + + let merged = LoroDoc::new(); + merged.set_peer_id(7)?; + merged + .get_text("t") + .insert(0, &("a๐Ÿ˜€".repeat(3000) + "Z"))?; + merged.commit(); + + let target = LoroDoc::new(); + target.import(&history.export(ExportMode::Snapshot)?)?; + target.import(&merged.export(ExportMode::all_updates())?)?; + assert_eq!(target.get_deep_value(), merged.get_deep_value()); + assert_eq!(target.oplog_vv(), merged.oplog_vv()); + assert_eq!(target.state_frontiers(), merged.state_frontiers()); + Ok(()) +}