diff --git a/.changeset/fix-import-overlap-cost.md b/.changeset/fix-import-overlap-cost.md new file mode 100644 index 000000000..76142bd88 --- /dev/null +++ b/.changeset/fix-import-overlap-cost.md @@ -0,0 +1,9 @@ +--- +"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. + +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 0d0d4665c..92750c3ca 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`, @@ -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` @@ -142,8 +142,12 @@ 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` 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`). @@ -182,6 +186,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..284854b45 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,24 @@ 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`. + +`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 638f8f612..1ded5c612 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,33 @@ 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 +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. 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`. +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 +120,130 @@ 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, 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 + +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..1644097cf 100644 --- a/crates/loro-internal/src/oplog/change_store.rs +++ b/crates/loro-internal/src/oplog/change_store.rs @@ -406,37 +406,167 @@ 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, + bytes: &[u8], + mut on_insert: impl FnMut(Counter, &ContainerID, u32, &str, u32, &Frontiers) -> LoroResult<()>, + ) -> LoroResult { + // 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), + } + } + + /// 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 +2381,152 @@ 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 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(); + 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..48ef273c2 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(); @@ -734,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 a5ad2fd09..66e744d3e 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,104 @@ 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 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 { + 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( + &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 +313,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 +333,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 +391,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 +439,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 +458,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 +537,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..bc885d459 100644 --- a/crates/loro/tests/import_reused_peer_id.rs +++ b/crates/loro/tests/import_reused_peer_id.rs @@ -461,3 +461,74 @@ 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(()) +} + +/// 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(()) +}