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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
345 changes: 340 additions & 5 deletions crates/kv_index/src/string_tree.rs
Original file line number Diff line number Diff line change
Expand Up @@ -407,9 +407,23 @@ impl Tree {
.entry(Arc::clone(&tenant_id))
.or_insert(0);

// Track remaining text as a slice - no allocation needed
let mut remaining = text;
let mut prev = Arc::clone(&self.root);
// Descend from the root inserting the whole text.
self.insert_from(Arc::clone(&self.root), text, tenant_id);
}

/// Insert `remaining` for `tenant_id` starting the descent at `start`
/// (which is treated as the parent of the first edge), reusing the exact
/// node-split / tenant-attach / leaf-timestamp logic of [`Self::insert_text`].
///
/// `insert_text` is `insert_from(root, text, tenant_id)` after the root
/// bookkeeping. [`Self::match_and_insert_with`] reuses it to splice only the
/// *unmatched suffix* at the fall-off node, so the already-matched prefix is
/// never re-walked. The caller is responsible for the root bookkeeping
/// (`tenant_last_access_time` / `tenant_char_count` entries) and, when
/// resuming below the root, for attaching `tenant_id` to the ancestor nodes
/// on the matched path (which this method does not touch).
fn insert_from(&self, start: NodeRef, mut remaining: &str, tenant_id: TenantId) {
let mut prev = start;

// Result type to carry state out of the match block
// This allows the entry guard to be dropped before we update prev
Expand Down Expand Up @@ -608,7 +622,7 @@ impl Tree {
.iter()
.next()
.map(|kv| Arc::clone(kv.key()))
.unwrap_or_else(|| Arc::from("empty"));
.unwrap_or_else(|| intern_tenant("empty"));
*curr.last_tenant.write() = Some(Arc::clone(&t));
t
}
Expand All @@ -620,7 +634,7 @@ impl Tree {
.iter()
.next()
.map(|kv| Arc::clone(kv.key()))
.unwrap_or_else(|| Arc::from("empty"));
.unwrap_or_else(|| intern_tenant("empty"));
*curr.last_tenant.write() = Some(Arc::clone(&t));
t
}
Expand Down Expand Up @@ -648,6 +662,231 @@ impl Tree {
}
}

/// Read-only resolution of a node's "owning" tenant, mirroring the tenant
/// pick in [`Self::match_prefix_with_counts`] **without** the cache-populate
/// or probabilistic timestamp side effects. Returns the tenant plus whether
/// the slow (iteration) path was taken — the caller replays the original's
/// single deferred side effect on the final node.
///
/// Used by [`Self::match_and_insert`] so the match tenant is resolved from a
/// node's state **before** the fused insert adds the inserting tenant to it,
/// reproducing the original "match runs fully before insert" ordering for the
/// (otherwise non-deterministic) slow-path tenant pick.
#[expect(
clippy::unused_self,
reason = "method logically belongs to the tree; mirrors match_prefix_with_counts resolution"
)]
fn resolve_tenant_readonly(&self, node: &NodeRef) -> (TenantId, bool) {
let cached = node.last_tenant.read();
if let Some(ref t) = *cached {
if node.tenant_last_access_time.contains_key(t.as_ref()) {
return (Arc::clone(t), false);
}
}
drop(cached);
let t = node
.tenant_last_access_time
.iter()
.next()
.map(|kv| Arc::clone(kv.key()))
.unwrap_or_else(|| intern_tenant("empty"));
(t, true)
}

/// Combined match + insert for a known `tenant` in a single descent.
///
/// Thin wrapper over [`Self::match_and_insert_with`] for callers that already
/// know the tenant to insert for before matching (e.g. the imbalanced
/// min-load path, which routes by worker load, not cache affinity). Returns
/// the match result for any cache bookkeeping the caller needs.
pub fn match_and_insert(&self, text: &str, tenant: &str) -> PrefixMatchResult {
self.match_and_insert_with(text, move |_| Some(tenant))
}

/// Match in a single descent, then choose the insert tenant from the match
/// result and insert for it — still a SINGLE walk of the matched prefix.
///
/// String analogue of `TokenTree::match_and_insert_with`. The cache-aware
/// router picks the worker it inserts for *from* the match outcome, so the
/// tenant is unknown until the match finishes. `select` runs once, after the
/// match, with the [`PrefixMatchResult`]; `Some(tenant)` inserts `text` for
/// that tenant, `None` skips the insert (the router's "no worker selected"
/// branch).
///
/// # How the single descent is achieved
///
/// The match phase walks the prefix exactly like
/// [`Self::match_prefix_with_counts`] (including its single deferred,
/// probabilistic timestamp touch on the resolved node — done via
/// [`Self::finish_match_and_insert`]), while recording the chain of
/// full-match nodes it descended through. After `select` yields the tenant we
/// re-attach it to those recorded ancestor nodes directly (the epoch-0
/// intermediate attach `insert_text` performs while descending) and splice
/// only the *unmatched suffix* at the fall-off node via
/// [`Self::insert_from`]. The matched prefix is therefore compared once, not
/// twice.
///
/// # Preserved semantics
///
/// Identical to `match_prefix_with_counts` followed by
/// `insert_text(text, tenant)`: same matched/-input char counts, same
/// resolved tenant and probabilistic touch, same node splits, same epoch-0
/// ancestor attaches, same final-leaf real timestamp, and same
/// `tenant_char_count` accounting. Because the match phase runs fully before
/// any insert mutation (just like the two separate calls), the tenant pick
/// observes the un-polluted tree; only the exact monotonic timestamp values
/// of insert's writes shift by a few ticks, which is immaterial to LRU.
pub fn match_and_insert_with<'t, F>(&self, text: &str, select: F) -> PrefixMatchResult
where
F: FnOnce(&PrefixMatchResult) -> Option<&'t str>,
{
// ---- Phase 1: MATCH descent (mirrors match_prefix_with_counts) ----
// Record the full-match nodes (for ancestor re-attach) and the fall-off
// point (node + remaining slice) for the suffix splice. No mutation here
// beyond match's own deferred touch below, so the tenant pick sees the
// un-polluted tree exactly like the standalone match.
let mut remaining = text;
let mut matched_chars = 0usize;
let mut current = Arc::clone(&self.root);
// The node match resolves its tenant on (its final `curr`): the deepest
// full-match node, the partial child, or the root if nothing matched.
let mut match_curr = Arc::clone(&self.root);
// (node, char_count) for each full-match edge, in order.
// Pre-allocated; most matched paths are well under this depth.
let mut path: Vec<(NodeRef, usize)> = Vec::with_capacity(16);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

On the request hot path, we can avoid re-allocating the Vec for path on every call by pre-allocating it with a small capacity (e.g., 16), which is sufficient for most tree depths.

Suggested change
let mut path: Vec<(NodeRef, usize)> = Vec::with_capacity(16);

while let Some(first_char) = remaining.chars().next() {
let child_node = current.children.get(&first_char).map(|e| e.value().clone());

let Some(matched_node) = child_node else {
// No child for this char: match stops at `current`.
break;
};

let matched_text_guard = matched_node.text.read();
let matched_node_text_count = matched_text_guard.char_count();
let shared_count = shared_prefix_count(remaining, matched_text_guard.as_str());
drop(matched_text_guard);

if shared_count == matched_node_text_count {
// Full match -> continue. Record for ancestor re-attach.
matched_chars += shared_count;
path.push((Arc::clone(&matched_node), matched_node_text_count));
remaining = advance_by_chars(remaining, shared_count);
current = Arc::clone(&matched_node);
match_curr = matched_node;
} else {
// Partial match: match stops, resolving on this child node.
matched_chars += shared_count;
match_curr = matched_node;
// `current` stays the parent — the splice re-probes it and
// splits the partial child exactly like `insert_text`.
break;
}
}

// ---- Match side effect + result (verbatim match_prefix_with_counts) ----
let (match_tenant, match_slow) = self.resolve_tenant_readonly(&match_curr);
let result = self.finish_match_and_insert(
&match_curr,
&match_tenant,
match_slow,
matched_chars,
text,
);

// ---- Decide the insert tenant from the match result ----
let Some(tenant) = select(&result) else {
return result;
};
// Intern through the shared pool so the stored Arc dedups (matches
// insert_text's interning).
let tenant_id = intern_tenant(tenant);

// ---- Phase 2: INSERT for `tenant_id` without re-walking the prefix ----
// Root bookkeeping (mirrors insert_text).
self.root
.tenant_last_access_time
.entry(Arc::clone(&tenant_id))
.or_insert(0);
self.tenant_char_count
.entry(Arc::clone(&tenant_id))
.or_insert(0);

// Re-attach the inserting tenant to every full-match ancestor node, the
// epoch-0 intermediate attach `insert_text` performs while descending.
for (node, char_count) in &path {
if !node
.tenant_last_access_time
.contains_key(tenant_id.as_ref())
{
self.tenant_char_count
.entry(Arc::clone(&tenant_id))
.and_modify(|count| *count += *char_count)
.or_insert(*char_count);
node.tenant_last_access_time
.insert(Arc::clone(&tenant_id), 0);
}
}

if remaining.is_empty() {
// Loop-end: `current` is the final leaf; give it the real timestamp
// (insert_text's tail). It already received the epoch-0 attach above
// if it is on the path.
let epoch = get_epoch();
current
.tenant_last_access_time
.insert(Arc::clone(&tenant_id), epoch);
} else {
// Fall-off: splice only the unmatched suffix below `current`,
// reusing insert_text's exact loop (handles split-continue + the
// final-leaf real timestamp). The matched prefix is not re-walked.
self.insert_from(current, remaining, tenant_id);
Comment on lines +816 to +844

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🟠 Major | 🏗️ Heavy lift

Revalidate replayed nodes before folding their cached char_count into tenant_char_count.

This replay assumes every recorded NodeRef still represents the same matched edge it did in Phase 1. If another insert splits one of those nodes before Lines 815-843 run, that ref can now be the suffix child while char_count still reflects the old unsplit edge. The replay then attaches the tenant to the stale node, increments tenant_char_count by the stale full-edge length, and skips the new intermediate. The route may “self-heal” on a later request, but the inflated size accounting does not, so size-based eviction will act on corrupted data. Please validate that each recorded node is still on the matched path or fall back to a fresh insert when the path has changed.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@crates/kv_index/src/string_tree.rs` around lines 815 - 843, Recorded ancestor
nodes in path must be revalidated before folding their cached char_count into
tenant_char_count because intervening splits can change which node represents
the matched edge; update the loop that iterates path (the block that checks
node.tenant_last_access_time and updates self.tenant_char_count) to first
confirm each NodeRef still corresponds to the same matched edge/length (compare
the node's current edge length or parent/child linkage against the recorded
char_count or expected child on the path), and if any mismatch is detected abort
the replay and fall back to a fresh insertion (call self.insert_from(current,
remaining, tenant_id) or otherwise re-run insert_text path discovery) instead of
applying stale char_count or attaching the tenant to a stale node; ensure
tenant_last_access_time updates only occur for validated nodes.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Avoid splicing from stale matched nodes

When another request splits the matched path between Phase 1 and this splice, current and the recorded path can refer to nodes whose prefix changed. For example, one request can match abcdef... and pause with current at the abcdef node; a concurrent insert of abc... splits it into abc + def; resuming here inserts the suffix under the old def node and never attaches the tenant to the new abc intermediate, while the old match_prefix_with_counts + insert_text path re-walked from the root and would attach both nodes. This leaves tenant-specific lookups and eviction counts inconsistent under concurrent same-prefix inserts.

Useful? React with 👍 / 👎.

}

result
}

/// Replay `match_prefix_with_counts`'s single deferred side effect (populate
/// the `last_tenant` cache on the slow path, then a probabilistic 1/8
/// timestamp touch) on the resolved final node, and build the match result.
/// Factored out so [`Self::match_and_insert_with`] applies it identically to
/// the standalone `match_prefix_with_counts` tail.
#[expect(
clippy::unused_self,
reason = "method logically belongs to the tree; mirrors match_prefix_with_counts tail"
)]
fn finish_match_and_insert(
&self,
match_node: &NodeRef,
tenant: &TenantId,
took_slow_path: bool,
matched_chars: usize,
text: &str,
) -> PrefixMatchResult {
// On the slow path the original resolution populates the cache.
if took_slow_path {
*match_node.last_tenant.write() = Some(Arc::clone(tenant));
}

// Probabilistic (1 in 8) timestamp touch on the resolved tenant, skipping
// the synthetic "empty" tenant — identical to match_prefix_with_counts.
let epoch = get_epoch();
if epoch & 0x7 == 0 && tenant.as_ref() != "empty" {
match_node
.tenant_last_access_time
.insert(Arc::clone(tenant), epoch);
}

let input_char_count = text.chars().count();

PrefixMatchResult {
tenant: Arc::clone(tenant),
matched_char_count: matched_chars,
input_char_count,
}
}

/// Legacy prefix_match API for backward compatibility.
/// Note: This computes matched_text which has allocation overhead.
pub fn prefix_match_legacy(&self, text: &str) -> (String, String) {
Expand Down Expand Up @@ -3015,4 +3254,100 @@ mod tests {
paths.sort();
assert_eq!(paths, vec!["", "你好", "你好世界", "你好朋友"]);
}

/// Run the same op sequence two ways — `match_prefix_with_counts` +
/// `insert_text` (legacy pair) vs the fused `match_and_insert` — and assert
/// the match counts and resulting per-tenant char counts agree. Timestamps
/// and the probabilistic tenant cache are not compared.
fn assert_fused_matches_pair(ops: &[(&str, &str)]) {
let pair = Tree::new();
let fused = Tree::new();
for (text, tenant) in ops {
let r_pair = pair.match_prefix_with_counts(text);
pair.insert_text(text, tenant);

let r_fused = fused.match_and_insert(text, tenant);

assert_eq!(
r_pair.matched_char_count, r_fused.matched_char_count,
"matched_char_count mismatch for tenant {tenant}"
);
assert_eq!(
r_pair.input_char_count, r_fused.input_char_count,
"input_char_count mismatch"
);
assert_eq!(
r_pair.tenant.as_ref(),
r_fused.tenant.as_ref(),
"matched tenant mismatch for input {text:?}"
);
}
assert_eq!(
pair.get_tenant_char_count(),
fused.get_tenant_char_count(),
"tenant char counts diverged between pair and fused"
);
}

#[test]
fn test_string_match_and_insert_equiv_basic() {
// Note: every query here either misses on an empty tree (resolves to the
// synthetic "empty") or falls off on a node whose `last_tenant` is
// single-valued, so the (otherwise approximate) tenant pick is
// deterministic and comparable between the two construction paths.
assert_fused_matches_pair(&[
("hello world", "w1"), // fresh leaf
("hello world", "w1"), // exact re-insert (full match, loop-end)
("hello there", "w2"), // shared "hello " prefix -> split
("hello", "w3"), // prefix of "hello " split node
]);
}

#[test]
fn test_string_match_and_insert_equiv_deep_chain() {
// Build a multi-node chain ("abc" -> {"def","xyz"}), then a query that
// FULL-matches several nodes before falling off — exercising the fused
// path's ancestor re-attach (`path`) plus a deep `insert_from` splice.
assert_fused_matches_pair(&[
("abcdef", "w1"), // leaf "abcdef"
("abcxyz", "w2"), // split -> "abc" + "def"/"xyz"
("abcdefghi", "w1"), // full-match "abc"+"def", then append "ghi"
("abcdefghi", "w3"), // full-match whole chain for a new tenant
]);
}

#[test]
fn test_string_match_and_insert_equiv_empty_and_unicode() {
assert_fused_matches_pair(&[
("", "w1"), // empty text -> tenant attached at root
("你好世界", "w1"), // multi-byte fresh
("你好朋友", "w2"), // multi-byte shared-prefix split
]);
}

#[test]
fn test_string_match_and_insert_with_select_and_skip() {
let text = "the quick brown fox";

let via_with = Tree::new();
via_with.match_and_insert_with(text, |_| Some("w1"));
let via_plain = Tree::new();
via_plain.match_and_insert(text, "w1");
assert_eq!(
via_with.get_tenant_char_count(),
via_plain.get_tenant_char_count()
);

// None selection must leave the tree untouched.
let tree = Tree::new();
tree.insert_text(text, "w1");
let before = tree.get_tenant_char_count();
let r = tree.match_and_insert_with(text, |_| None);
assert_eq!(r.matched_char_count, text.chars().count());
assert_eq!(
tree.get_tenant_char_count(),
before,
"None selection must not insert"
);
}
}
Loading
Loading