From ccbd6680232f882e8d925155e67a9dc9c0bb43ec Mon Sep 17 00:00:00 2001 From: Hayden Flinner Date: Wed, 30 Sep 2026 19:30:44 -0400 Subject: [PATCH 1/4] =?UTF-8?q?feat:=20LoroDoc::undo=5Fspan=20=E2=80=94=20?= =?UTF-8?q?invert=20an=20IdSpan=20without=20rewinding=20state?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- crates/loro-common/src/lib.rs | 1 - crates/loro-internal/src/dag.rs | 8 +------- crates/loro-internal/src/undo.rs | 3 +-- crates/loro/src/lib.rs | 20 ++++++++++++++++++++ crates/loro/tests/moon_transcode.rs | 6 +++++- 5 files changed, 27 insertions(+), 11 deletions(-) diff --git a/crates/loro-common/src/lib.rs b/crates/loro-common/src/lib.rs index ea35462b5..8670d4ba3 100644 --- a/crates/loro-common/src/lib.rs +++ b/crates/loro-common/src/lib.rs @@ -18,7 +18,6 @@ mod value; pub use error::{LoroEncodeError, LoroError, LoroResult, LoroTreeError}; pub use internal_string::InternalString; -pub use logging::log::*; #[doc(hidden)] pub use rustc_hash::FxHashMap; pub use span::*; diff --git a/crates/loro-internal/src/dag.rs b/crates/loro-internal/src/dag.rs index cf4538560..0bc2f9649 100644 --- a/crates/loro-internal/src/dag.rs +++ b/crates/loro-internal/src/dag.rs @@ -1745,13 +1745,7 @@ mod tests { let root = node(1, 0, 1, 0, Frontiers::default()); let x = node(2, 0, 1, 1, ID::new(1, 0).into()); let y = node(3, 0, 1, 1, ID::new(1, 0).into()); - let merge = node( - 4, - 0, - 1, - 2, - Frontiers::from([ID::new(2, 0), ID::new(3, 0)]), - ); + let merge = node(4, 0, 1, 2, Frontiers::from([ID::new(2, 0), ID::new(3, 0)])); let dag = TestDag::new(vec![root, x, y, merge], ID::new(4, 0).into()); let left = Frontiers::from([ID::new(2, 0), ID::new(3, 0)]); diff --git a/crates/loro-internal/src/undo.rs b/crates/loro-internal/src/undo.rs index 20aea7cc4..f6e61b452 100644 --- a/crates/loro-internal/src/undo.rs +++ b/crates/loro-internal/src/undo.rs @@ -959,8 +959,7 @@ impl UndoManager { return Err(e); } Err(e) => { - get_stack(&mut self.inner.lock().borrow_mut()) - .push(span.span, span.meta); + get_stack(&mut self.inner.lock().borrow_mut()).push(span.span, span.meta); return Err(e); } }; diff --git a/crates/loro/src/lib.rs b/crates/loro/src/lib.rs index 4ff8e9779..49167a349 100644 --- a/crates/loro/src/lib.rs +++ b/crates/loro/src/lib.rs @@ -1470,6 +1470,26 @@ impl LoroDoc { self.doc.revert_to(version) } + /// Append the ops that invert `span` — the same machinery + /// [`UndoManager`] uses for undo, callable on *any* peer's ops, not + /// just the bound peer's. Unlike [`revert_to`](Self::revert_to), which + /// rewinds state to a target version (discarding whatever arrived + /// since), this inverts only the ops in the span: everything else — + /// including edits merged in after the span — is untouched. The + /// inverse ops commit under this peer's id with `origin: "undo"`. + /// + /// The spans a merge imported are `pre_vv.diff(&post_vv).forward` — + /// `undo_span` is how a merge (or any oplog range) is reverted + /// surgically rather than by rewinding the whole document. + #[inline] + pub fn undo_span(&self, span: IdSpan) -> LoroResult<()> { + let commit = self + .doc + .undo_internal(span, &mut Default::default(), None, &mut |_| {})?; + drop(commit); + Ok(()) + } + /// Apply a diff to the current document state. /// /// Internally, it will apply the diff to the current state. diff --git a/crates/loro/tests/moon_transcode.rs b/crates/loro/tests/moon_transcode.rs index c7d0f50c6..7dfa5edf1 100644 --- a/crates/loro/tests/moon_transcode.rs +++ b/crates/loro/tests/moon_transcode.rs @@ -91,7 +91,11 @@ fn run_transcode(node_bin: &str, cli_js: &Path, input: &[u8]) -> anyhow::Result< .duration_since(UNIX_EPOCH) .unwrap() .as_nanos(); - let tmp = std::env::temp_dir().join(format!("loro-moon-transcode-{}-{ts}-{}", std::process::id(), next_tmp_id())); + let tmp = std::env::temp_dir().join(format!( + "loro-moon-transcode-{}-{ts}-{}", + std::process::id(), + next_tmp_id() + )); std::fs::create_dir_all(&tmp)?; let in_path = tmp.join("in.blob"); let out_path = tmp.join("out.blob"); From 26199f1626692cdc2c1b60ffc4cdf9ae5ba66e84 Mon Sep 17 00:00:00 2001 From: Hayden Flinner Date: Sat, 3 Oct 2026 11:37:09 -0400 Subject: [PATCH 2/4] fix: keep change timestamps at millisecond precision Rounds to seconds made per-commit timing useless for tape replay. Renames merge_interval_in_s -> merge_interval_in_ms to match. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- crates/loro-internal/src/change.rs | 2 +- crates/loro-internal/src/change_meta.rs | 4 ++-- crates/loro-internal/src/configure.rs | 12 ++++++------ crates/loro-internal/src/loro.rs | 8 ++++---- crates/loro-internal/src/oplog.rs | 10 +++++++--- crates/loro/src/lib.rs | 8 ++++---- crates/loro/tests/loro_rust_test.rs | 3 ++- 7 files changed, 26 insertions(+), 21 deletions(-) diff --git a/crates/loro-internal/src/change.rs b/crates/loro-internal/src/change.rs index 72edc4a70..47c2020d7 100644 --- a/crates/loro-internal/src/change.rs +++ b/crates/loro-internal/src/change.rs @@ -32,7 +32,7 @@ pub struct Change { pub(crate) lamport: Lamport, pub(crate) deps: Frontiers, /// [Unix time](https://en.wikipedia.org/wiki/Unix_time) - /// It is the number of seconds that have elapsed since 00:00:00 UTC on 1 January 1970. + /// It is the number of milliseconds that have elapsed since 00:00:00 UTC on 1 January 1970. pub(crate) timestamp: Timestamp, pub(crate) commit_msg: Option>, pub(crate) ops: RleVec<[O; 1]>, diff --git a/crates/loro-internal/src/change_meta.rs b/crates/loro-internal/src/change_meta.rs index 93e9cbd9e..b73ebc1cd 100644 --- a/crates/loro-internal/src/change_meta.rs +++ b/crates/loro-internal/src/change_meta.rs @@ -27,7 +27,7 @@ pub struct ChangeMeta { /// The first Op id of the Change pub id: ID, /// [Unix time](https://en.wikipedia.org/wiki/Unix_time) - /// It is the number of seconds that have elapsed since 00:00:00 UTC on 1 January 1970. + /// It is the number of milliseconds that have elapsed since 00:00:00 UTC on 1 January 1970. pub timestamp: Timestamp, /// The commit message of the change pub message: Option>, @@ -91,7 +91,7 @@ impl ChangeMeta { } } - /// Get the commit timestamp in seconds since Unix epoch. + /// Get the commit timestamp in milliseconds since Unix epoch. pub fn timestamp(&self) -> crate::change::Timestamp { self.timestamp } diff --git a/crates/loro-internal/src/configure.rs b/crates/loro-internal/src/configure.rs index 90955f79c..bff5da206 100644 --- a/crates/loro-internal/src/configure.rs +++ b/crates/loro-internal/src/configure.rs @@ -11,7 +11,7 @@ use std::sync::Arc; pub struct Configure { pub(crate) text_style_config: Arc>, record_timestamp: Arc, - pub(crate) merge_interval_in_s: Arc, + pub(crate) merge_interval_in_ms: Arc, pub(crate) editable_detached_mode: Arc, pub(crate) deleted_root_containers: Arc>>, pub(crate) hide_empty_root_containers: Arc, @@ -32,7 +32,7 @@ impl Default for Configure { text_style_config: Arc::new(RwLock::new(StyleConfigMap::default_rich_text_config())), record_timestamp: Arc::new(AtomicBool::new(false)), editable_detached_mode: Arc::new(AtomicBool::new(false)), - merge_interval_in_s: Arc::new(AtomicI64::new(1000)), + merge_interval_in_ms: Arc::new(AtomicI64::new(1000)), deleted_root_containers: Arc::new(Mutex::new(Default::default())), hide_empty_root_containers: Arc::new(AtomicBool::new(false)), } @@ -47,8 +47,8 @@ impl Configure { self.record_timestamp .load(std::sync::atomic::Ordering::Relaxed), )), - merge_interval_in_s: Arc::new(AtomicI64::new( - self.merge_interval_in_s + merge_interval_in_ms: Arc::new(AtomicI64::new( + self.merge_interval_in_ms .load(std::sync::atomic::Ordering::Relaxed), )), editable_detached_mode: Arc::new(AtomicBool::new( @@ -90,12 +90,12 @@ impl Configure { } pub fn merge_interval(&self) -> i64 { - self.merge_interval_in_s + self.merge_interval_in_ms .load(std::sync::atomic::Ordering::Relaxed) } pub fn set_merge_interval(&self, interval: i64) { - self.merge_interval_in_s + self.merge_interval_in_ms .store(interval, std::sync::atomic::Ordering::Relaxed); } diff --git a/crates/loro-internal/src/loro.rs b/crates/loro-internal/src/loro.rs index 6d86ab69d..df5871f9d 100644 --- a/crates/loro-internal/src/loro.rs +++ b/crates/loro-internal/src/loro.rs @@ -487,10 +487,10 @@ impl LoroDoc { self.config.set_record_timestamp(record); } - /// Set the interval of mergeable changes, in seconds. + /// Set the interval of mergeable changes, in milliseconds. /// /// If two continuous local changes are within the interval, they will be merged into one change. - /// The default value is 1000 seconds. + /// The default value is 1000 milliseconds. #[inline] pub fn set_change_merge_interval(&self, interval: i64) { self.config.set_merge_interval(interval); @@ -3514,7 +3514,7 @@ pub struct CommitOptions { /// Defaults to true. pub immediate_renew: bool, - /// Custom timestamp for the commit in seconds since Unix epoch. + /// Custom timestamp for the commit in milliseconds since Unix epoch. /// If None, the current time will be used. pub timestamp: Option, @@ -3547,7 +3547,7 @@ impl CommitOptions { /// Set the timestamp of the commit. /// - /// The timestamp is the number of **seconds** that have elapsed since 00:00:00 UTC on January 1, 1970. + /// The timestamp is the number of **milliseconds** that have elapsed since 00:00:00 UTC on January 1, 1970. pub fn timestamp(mut self, timestamp: Timestamp) -> Self { self.timestamp = Some(timestamp); self diff --git a/crates/loro-internal/src/oplog.rs b/crates/loro-internal/src/oplog.rs index 1db21511a..850217189 100644 --- a/crates/loro-internal/src/oplog.rs +++ b/crates/loro-internal/src/oplog.rs @@ -166,7 +166,7 @@ impl OpLog { pub(crate) fn new(visible_op_count: Arc) -> Self { let arena = SharedArena::new(); let cfg = Configure::default(); - let change_store = ChangeStore::new_mem(&arena, cfg.merge_interval_in_s.clone()); + let change_store = ChangeStore::new_mem(&arena, cfg.merge_interval_in_ms.clone()); arena.set_creator_resolver(change_store.creator_resolver()); Self { visible_op_count, @@ -494,7 +494,7 @@ impl OpLog { let configure = self.configure.clone(); // Also rolls back the arena; see `ChangeStore::retire`. self.change_store.retire(arena_checkpoint); - let change_store = ChangeStore::new_mem(&arena, configure.merge_interval_in_s.clone()); + let change_store = ChangeStore::new_mem(&arena, configure.merge_interval_in_ms.clone()); arena.set_creator_resolver(change_store.creator_resolver()); self.history_cache = Mutex::new(ContainerHistoryCache::new(change_store.clone(), None)); self.dag = AppDag::new(change_store.clone()); @@ -1582,7 +1582,11 @@ pub(crate) fn local_op_to_remote( } pub(crate) fn get_timestamp_now_txn() -> Timestamp { - (get_sys_timestamp() as Timestamp + 500) / 1000 + // Milliseconds — `get_sys_timestamp` already yields ms, and every + // other consumer (awareness, undo, diff timeouts) treats Timestamp as + // ms. Rounding to seconds here made per-commit timing useless for + // replay (creation tapes pace edits on these values). + get_sys_timestamp() as Timestamp } #[cfg(test)] diff --git a/crates/loro/src/lib.rs b/crates/loro/src/lib.rs index 49167a349..a89fdc866 100644 --- a/crates/loro/src/lib.rs +++ b/crates/loro/src/lib.rs @@ -280,13 +280,13 @@ impl LoroDoc { self.doc.is_detached_editing_enabled() } - /// Set the interval of mergeable changes, **in seconds**. + /// Set the interval of mergeable changes, **in milliseconds**. /// /// If two continuous local changes are within the interval, they will be merged into one change. - /// The default value is 1000 seconds. + /// The default value is 1000 milliseconds. /// - /// By default, we record timestamps in seconds for each change. So if the merge interval is 1, and changes A and B - /// have timestamps of 3 and 4 respectively, then they will be merged into one change. + /// By default, we record timestamps in milliseconds for each change. So if the merge interval is 1000, and changes A and B + /// have timestamps of 3000 and 4000 respectively, then they will be merged into one change. #[inline] pub fn set_change_merge_interval(&self, interval: i64) { self.doc.set_change_merge_interval(interval); diff --git a/crates/loro/tests/loro_rust_test.rs b/crates/loro/tests/loro_rust_test.rs index 42df863df..6d538625d 100644 --- a/crates/loro/tests/loro_rust_test.rs +++ b/crates/loro/tests/loro_rust_test.rs @@ -196,7 +196,8 @@ fn timestamp() { doc1.commit(); doc1.with_oplog(|oplog| { let c = oplog.get_change_at(ID::new(1, 2)).unwrap(); - assert!(c.timestamp() < last_timestamp + 10); + // timestamps are milliseconds — 10s of slack around the previous commit + assert!(c.timestamp() < last_timestamp + 10_000); }); } From 8c4d1b214c5dabc8eca8cc0639cc9e958da689a6 Mon Sep 17 00:00:00 2001 From: Hayden Flinner Date: Sun, 4 Oct 2026 10:49:26 -0400 Subject: [PATCH 3/4] =?UTF-8?q?perf:=20coalesce=20id=E2=86=92cursor=20frag?= =?UTF-8?q?ments=20in=20the=20richtext=20tracker?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Fragments were hard-capped at 256 atoms, so a uniform ~1M-atom span covered ~4k fragments. Each rope leaf split then re-mapped every fragment in the span via update_insert_batch, giving O(splits x span) — ~80M fragment iterations on a 2.2MB doc import (~13s to attach, ~47s to re-export state). Bound fragments by run count instead of atom count — a single-run fragment has one cursor at any size: - Cursor::try_merge joins two adjacent Small insert sets when their combined runs fit SMALL_SET_MAX_LEN, collapsing boundary runs that share a leaf. - IdToCursor::coalesce merges adjacent compatible fragments in the range update_insert/update_insert_batch just touched. Adjacency is required (a.counter_end() == b.counter): the list can carry counter gaps left by other containers' ops, and merging across a gap would change get_insert results inside it. - insert() no longer splits a large uniform single-run cursor into MAX_FRAGMENT_LEN chunks, and merges with the previous adjacent fragment when possible. Result: the 2.2MB production doc attaches in ~9ms (was 13.4s), and perf_import_insert_split_quadratic_e2e (2M atoms, 8192 boundary splits, ~33.5M expected fragment updates before) now runs in ~15ms. State exports verified byte-identical to the pre-fix build. --- .../richtext/tracker/id_to_cursor.rs | 98 +++++++++++++------ 1 file changed, 69 insertions(+), 29 deletions(-) diff --git a/crates/loro-internal/src/container/richtext/tracker/id_to_cursor.rs b/crates/loro-internal/src/container/richtext/tracker/id_to_cursor.rs index d716069a3..c6953551f 100644 --- a/crates/loro-internal/src/container/richtext/tracker/id_to_cursor.rs +++ b/crates/loro-internal/src/container/richtext/tracker/id_to_cursor.rs @@ -41,41 +41,43 @@ impl IdToCursor { if let Some(last) = list.last_mut() { let last_end = last.counter + last.cursor.rle_len() as Counter; debug_assert!(last_end <= id.counter, "id:{}, {:#?}", id, &self); - if last_end == id.counter - && last.cursor.can_merge(&cursor) - && last.cursor.rle_len() + cursor.rle_len() < MAX_FRAGMENT_LEN - { - last.cursor.merge_right(&cursor); + if last_end == id.counter && last.cursor.try_merge(&cursor) { return; } } - if let Cursor::Insert(InsertSet::Small(set)) = cursor { - if set.len > MAX_FRAGMENT_LEN as u32 { - assert!(set.set.len() == 1); - let insert = set.set[0]; - let mut counter = id.counter; - for start in (0..set.len).step_by(MAX_FRAGMENT_LEN) { - let end = (start + MAX_FRAGMENT_LEN as u32).min(set.len); - let len = (end - start) as usize; - list.push(Fragment { - counter, - cursor: Cursor::new_insert(insert.leaf, len), - }); - counter += len as Counter; - } - } else { - list.push(Fragment { - counter: id.counter, - cursor: Cursor::Insert(InsertSet::Small(set)), - }); - } - } else { - list.push(Fragment { - counter: id.counter, - cursor, + // A fragment's cost is bounded by its run count, not its atom + // length — a single-run fragment has one cursor at any size, so + // there is no reason to split a large uniform span into + // MAX_FRAGMENT_LEN chunks. Coalescing keeps later id→leaf remaps + // proportional to leaf-boundary runs instead of atoms. + list.push(Fragment { + counter: id.counter, + cursor, + }); + } + + /// Merge adjacent fragments in `list[lo..hi]` whose cursors combine + /// within the small-set run cap. Only adjacent fragments + /// (`a.counter_end() == b.counter`) may merge: the list can carry + /// counter gaps left by other containers' ops, and merging across a + /// gap would change `get_insert` results inside it. + fn coalesce(list: &mut Vec, lo: usize, hi: usize) { + if hi - lo < 2 { + return; + } + + let removed: Vec = list.splice(lo..hi, []).collect(); + let mut merged: Vec = Vec::with_capacity(removed.len()); + for f in removed { + let absorbed = merged.last_mut().is_some_and(|last| { + last.counter_end() == f.counter && last.cursor.try_merge(&f.cursor) }); + if !absorbed { + merged.push(f); + } } + list.splice(lo..lo, merged); } /// Update the given id_span to the new_leaf @@ -91,6 +93,7 @@ impl IdToCursor { Err(index) => index.saturating_sub(1), }; + let first = index; let mut start_counter = id_span.counter.start; while start_counter < id_span.counter.end && index < list.len() @@ -107,6 +110,7 @@ impl IdToCursor { } assert_eq!(start_counter, id_span.counter.end); + Self::coalesce(list, first.saturating_sub(1), (index + 1).min(list.len())); } pub fn update_insert_batch(&mut self, updates: &mut [(IdSpan, LeafIndex)]) { @@ -133,6 +137,7 @@ impl IdToCursor { continue; }; + let mut per_fragment: FxHashMap> = FxHashMap::default(); @@ -170,10 +175,15 @@ impl IdToCursor { let mut fragment_indexes: Vec = per_fragment.keys().copied().collect(); fragment_indexes.sort_unstable(); + let (lo, hi) = ( + *fragment_indexes.first().unwrap(), + *fragment_indexes.last().unwrap(), + ); for index in fragment_indexes { let updates = per_fragment.get(&index).unwrap(); list[index].cursor.update_insert_many(updates); } + Self::coalesce(list, lo.saturating_sub(1), (hi + 2).min(list.len())); } } @@ -957,6 +967,36 @@ impl Cursor { } } + /// Merge `rhs` (the cursor immediately following `self`) into `self` + /// when both are Small insert sets whose combined run count stays + /// within the small-set cap. The caller must guarantee the two spans + /// are counter-adjacent. + fn try_merge(&mut self, rhs: &Self) -> bool { + let (Cursor::Insert(InsertSet::Small(a)), Cursor::Insert(InsertSet::Small(b))) = + (self, rhs) + else { + return false; + }; + let boundary_merge = matches!( + (a.set.last(), b.set.first()), + (Some(l), Some(f)) if l.leaf == f.leaf + ); + let runs = a.set.len() + b.set.len() - boundary_merge as usize; + if runs > SMALL_SET_MAX_LEN { + return false; + } + + if boundary_merge { + let add = b.set[0].len; + a.set.last_mut().unwrap().len += add; + a.set.extend(b.set.iter().skip(1).copied()); + } else { + a.set.extend(b.set.iter().copied()); + } + a.len += b.len; + true + } + fn get_insert(&self, pos: usize) -> Option { if pos >= self.rle_len() { return None; From 08231db78699f87def069d0c5c6eff311e5a3ef6 Mon Sep 17 00:00:00 2001 From: Hayden Flinner Date: Sun, 4 Oct 2026 10:49:31 -0400 Subject: [PATCH 4/4] =?UTF-8?q?feat:=20tracker-stats=20feature=20=E2=80=94?= =?UTF-8?q?=20diagnostic=20counters=20for=20the=20richtext=20tracker?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adds an off-by-default `tracker-stats` cargo feature (loro + loro-internal) with zero-cost `#[cfg]`-gated counters on the richtext Tracker's hot paths: - op counts: inserts, checkouts, retreat/forward elements - skip_applied forwarding, Fugue in_between scans, leaf splits - id→cursor work: per-fragment update iterations, update_many call totals, dense-array rebuilds, large-set updates, iterator yields, batch call counts, span atoms, and max fragment-list length - stats::dump(tag) prints the totals; called at the end of calc_diff_internal and the richtext tracker rebuild Written to quantify a pathological fragment-remapping case on a 2.2MB document (~80M fragment iterations); kept because it makes the inner loop magnitudes visible without recompiling instrumentation. --- crates/loro-internal/Cargo.toml | 1 + .../src/container/richtext/tracker.rs | 65 +++++++++++++++++++ .../container/richtext/tracker/crdt_rope.rs | 2 + .../richtext/tracker/id_to_cursor.rs | 27 ++++++++ crates/loro-internal/src/diff_calc.rs | 4 ++ crates/loro/Cargo.toml | 1 + 6 files changed, 100 insertions(+) diff --git a/crates/loro-internal/Cargo.toml b/crates/loro-internal/Cargo.toml index ca6bfd102..5e3771992 100644 --- a/crates/loro-internal/Cargo.toml +++ b/crates/loro-internal/Cargo.toml @@ -96,6 +96,7 @@ test_utils = ["arbitrary", "tabled"] counter = ["loro-common/counter"] logging = ["loro-common/logging"] jsonpath = [] +tracker-stats = [] [[bench]] name = "text_r" diff --git a/crates/loro-internal/src/container/richtext/tracker.rs b/crates/loro-internal/src/container/richtext/tracker.rs index 751b92c2c..abaed55bc 100644 --- a/crates/loro-internal/src/container/richtext/tracker.rs +++ b/crates/loro-internal/src/container/richtext/tracker.rs @@ -27,6 +27,54 @@ pub(crate) const UNKNOWN_SPAN_LEN: u32 = u32::MAX / 4; pub(crate) use crdt_rope::CrdtRopeDelta; +#[cfg(feature = "tracker-stats")] +pub(crate) mod stats { + use std::sync::atomic::{AtomicUsize, Ordering}; + pub static INSERTS: AtomicUsize = AtomicUsize::new(0); + pub static CHECKOUTS: AtomicUsize = AtomicUsize::new(0); + pub static RETREAT_ELEMS: AtomicUsize = AtomicUsize::new(0); + pub static FORWARD_ELEMS: AtomicUsize = AtomicUsize::new(0); + pub static SKIP_FORWARD_CALLS: AtomicUsize = AtomicUsize::new(0); + pub static SKIP_FORWARD_ELEMS: AtomicUsize = AtomicUsize::new(0); + pub static IN_BETWEEN_ELEMS: AtomicUsize = AtomicUsize::new(0); + pub static SPLIT_LEAVES: AtomicUsize = AtomicUsize::new(0); + pub static UPDATE_INSERT_FRAGS: AtomicUsize = AtomicUsize::new(0); + pub static UPDATE_MANY_DENSE: AtomicUsize = AtomicUsize::new(0); + pub static LARGE_SEQ_UPDATES: AtomicUsize = AtomicUsize::new(0); + pub static ITER_YIELDS: AtomicUsize = AtomicUsize::new(0); + pub static UPDATE_MANY_CALLS: AtomicUsize = AtomicUsize::new(0); + pub static BATCH_CALLS: AtomicUsize = AtomicUsize::new(0); + pub static SPAN_ATOMS: AtomicUsize = AtomicUsize::new(0); + pub static MAX_LIST_LEN: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0); + pub fn bump(c: &'static AtomicUsize, n: usize) { + c.fetch_add(n, Ordering::Relaxed); + } + pub fn dump(tag: &str) { + eprintln!( + "tracker-stats {tag}: inserts={} checkouts={} retreat_elems={} forward_elems={} skip_fwd_calls={} skip_fwd_elems={} in_between={} split_leaves={} upd_frag_iters={} dense_elems={} large_seq={} iter_yields={} update_many_calls={}", + INSERTS.load(Ordering::Relaxed), + CHECKOUTS.load(Ordering::Relaxed), + RETREAT_ELEMS.load(Ordering::Relaxed), + FORWARD_ELEMS.load(Ordering::Relaxed), + SKIP_FORWARD_CALLS.load(Ordering::Relaxed), + SKIP_FORWARD_ELEMS.load(Ordering::Relaxed), + IN_BETWEEN_ELEMS.load(Ordering::Relaxed), + SPLIT_LEAVES.load(Ordering::Relaxed), + UPDATE_INSERT_FRAGS.load(Ordering::Relaxed), + UPDATE_MANY_DENSE.load(Ordering::Relaxed), + LARGE_SEQ_UPDATES.load(Ordering::Relaxed), + ITER_YIELDS.load(Ordering::Relaxed), + UPDATE_MANY_CALLS.load(Ordering::Relaxed), + ); + eprintln!( + "tracker-stats {tag}: batch_calls={} span_atoms={} max_list_len={}", + BATCH_CALLS.load(Ordering::Relaxed), + SPAN_ATOMS.load(Ordering::Relaxed), + MAX_LIST_LEN.load(Ordering::Relaxed), + ); + } +} + #[derive(Debug)] pub(crate) struct Tracker { applied_vv: VersionVector, @@ -98,6 +146,8 @@ impl Tracker { // &pos, // &content // ); + #[cfg(feature = "tracker-stats")] + stats::bump(&stats::INSERTS, 1); // tracing::span!(tracing::Level::INFO, "TrackerInsert"); if let ControlFlow::Break(_) = self.skip_applied(op_id.id(), content.len(), |applied_counter_end| { @@ -166,6 +216,8 @@ impl Tracker { } fn update_insert_by_split(&mut self, split: &[LeafIndex]) { + #[cfg(feature = "tracker-stats")] + stats::bump(&stats::SPLIT_LEAVES, split.len()); match split.len() { 0 => {} 1 => { @@ -274,6 +326,11 @@ impl Tracker { IdSpan::new(op_id.peer, cnt_start, op_id.counter + len as Counter), &mut updates, ); + #[cfg(feature = "tracker-stats")] + { + stats::bump(&stats::SKIP_FORWARD_CALLS, 1); + stats::bump(&stats::SKIP_FORWARD_ELEMS, updates.len()); + } self.batch_update(updates, false); } @@ -363,11 +420,15 @@ impl Tracker { self.rope.clear_diff_status(); } + #[cfg(feature = "tracker-stats")] + stats::bump(&stats::CHECKOUTS, 1); let current_vv = std::mem::take(&mut self.current_vv); let (retreat, forward) = current_vv.diff_iter(vv); let mut updates = Vec::new(); for span in retreat { for c in self.id_to_cursor.iter(span) { + #[cfg(feature = "tracker-stats")] + stats::bump(&stats::RETREAT_ELEMS, 1); match c { id_to_cursor::IterCursor::Insert { leaf, id_span } => { updates.push(crdt_rope::LeafUpdate { @@ -453,9 +514,13 @@ impl Tracker { } } + #[cfg(feature = "tracker-stats")] + let fwd_before = updates.len(); for span in forward { self.forward(span, &mut updates); } + #[cfg(feature = "tracker-stats")] + stats::bump(&stats::FORWARD_ELEMS, updates.len() - fwd_before); if !on_diff_status { self.current_vv = vv.clone(); diff --git a/crates/loro-internal/src/container/richtext/tracker/crdt_rope.rs b/crates/loro-internal/src/container/richtext/tracker/crdt_rope.rs index 3e3576c5f..b5840926e 100644 --- a/crates/loro-internal/src/container/richtext/tracker/crdt_rope.rs +++ b/crates/loro-internal/src/container/richtext/tracker/crdt_rope.rs @@ -147,6 +147,8 @@ impl CrdtRope { (origin_right, parent_right_idx, in_between) }; + #[cfg(feature = "tracker-stats")] + super::stats::bump(&super::stats::IN_BETWEEN_ELEMS, in_between.len()); content.origin_left = origin_left.map(|x| x.try_into().unwrap()); content.origin_right = origin_right.map(|x| x.try_into().unwrap()); diff --git a/crates/loro-internal/src/container/richtext/tracker/id_to_cursor.rs b/crates/loro-internal/src/container/richtext/tracker/id_to_cursor.rs index c6953551f..d551fb2ea 100644 --- a/crates/loro-internal/src/container/richtext/tracker/id_to_cursor.rs +++ b/crates/loro-internal/src/container/richtext/tracker/id_to_cursor.rs @@ -8,6 +8,8 @@ use rustc_hash::FxHashMap; use smallvec::smallvec; use self::insert_set::InsertSet; +#[cfg(feature = "tracker-stats")] +use super::stats; // If we make this too large, we may have too many cursors inside a fragment // and trigger the worst case @@ -99,6 +101,8 @@ impl IdToCursor { && index < list.len() && start_counter < list[index].counter_end() { + #[cfg(feature = "tracker-stats")] + stats::bump(&stats::UPDATE_INSERT_FRAGS, 1); let fragment = &mut list[index]; let from = (start_counter - fragment.counter) as usize; let to = @@ -137,12 +141,23 @@ impl IdToCursor { continue; }; + #[cfg(feature = "tracker-stats")] + { + stats::bump(&stats::BATCH_CALLS, 1); + let l = list.len() as u64; + let m = stats::MAX_LIST_LEN.load(std::sync::atomic::Ordering::Relaxed); + if l > m { + stats::MAX_LIST_LEN.store(l, std::sync::atomic::Ordering::Relaxed); + } + } let mut per_fragment: FxHashMap> = FxHashMap::default(); for (id_span, new_leaf) in peer_updates { debug_assert!(!id_span.is_reversed()); + #[cfg(feature = "tracker-stats")] + stats::bump(&stats::SPAN_ATOMS, id_span.atom_len()); let mut index = match list.binary_search_by_key(&id_span.counter.start, |x| x.counter) { Ok(index) => index, @@ -237,6 +252,8 @@ impl IdToCursor { continue; }; + #[cfg(feature = "tracker-stats")] + stats::bump(&stats::ITER_YIELDS, 1); return Some(next); } @@ -272,11 +289,15 @@ impl IdToCursor { continue; } + #[cfg(feature = "tracker-stats")] + stats::bump(&stats::ITER_YIELDS, 1); return Some(IterCursor::Delete(span.slice(from as usize, to as usize))); } Cursor::Move { from, to } => { index += 1; let op_id = ID::new(iter_id_span.peer, f.counter); + #[cfg(feature = "tracker-stats")] + stats::bump(&stats::ITER_YIELDS, 1); return Some(IterCursor::Move { from_id: *from, to_leaf: *to, @@ -431,6 +452,8 @@ mod insert_set { } pub(crate) fn update_many(&mut self, updates: &[(usize, usize, LeafIndex)]) { + #[cfg(feature = "tracker-stats")] + stats::bump(&stats::UPDATE_MANY_CALLS, updates.len()); if updates.is_empty() { return; } @@ -443,12 +466,16 @@ mod insert_set { let len = self.len(); if len > MAX_FRAGMENT_LEN { + #[cfg(feature = "tracker-stats")] + stats::bump(&stats::LARGE_SEQ_UPDATES, updates.len()); for &(from, to, leaf) in updates { self.update(from, to, leaf); } return; } + #[cfg(feature = "tracker-stats")] + stats::bump(&stats::UPDATE_MANY_DENSE, len); let mut dense: SmallVec<[LeafIndex; MAX_FRAGMENT_LEN]> = SmallVec::with_capacity(len); match self { InsertSet::Small(set) => { diff --git a/crates/loro-internal/src/diff_calc.rs b/crates/loro-internal/src/diff_calc.rs index 825ddcf19..4bd4759f7 100644 --- a/crates/loro-internal/src/diff_calc.rs +++ b/crates/loro-internal/src/diff_calc.rs @@ -430,6 +430,8 @@ impl DiffCalculator { } } + #[cfg(feature = "tracker-stats")] + crate::container::richtext::tracker::stats::dump("calc_diff"); ( ans.into_values().map(|x| x.1).collect_vec(), origin_diff_mode, @@ -1503,6 +1505,8 @@ impl RichtextDiffCalculator { }; replay_container_ops_from_empty(idx, oplog, vv, &mut visitor); + #[cfg(feature = "tracker-stats")] + crate::container::richtext::tracker::stats::dump("rebuild"); (tracker, styles) } } diff --git a/crates/loro/Cargo.toml b/crates/loro/Cargo.toml index 9fa1b0579..c5efa149b 100644 --- a/crates/loro/Cargo.toml +++ b/crates/loro/Cargo.toml @@ -42,6 +42,7 @@ default = ["counter"] counter = ["loro-internal/counter"] jsonpath = ["loro-internal/jsonpath"] logging = ["loro-internal/logging"] +tracker-stats = ["loro-internal/tracker-stats"] [lints.rust] unexpected_cfgs = { level = "warn", check-cfg = ['cfg(loom)'] }