Skip to content
Closed
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
33 changes: 33 additions & 0 deletions crates/mesh/src/crdt_kv/crdt.rs
Original file line number Diff line number Diff line change
Expand Up @@ -169,6 +169,39 @@ impl CrdtOrMap {
all
}

/// Live keys under `prefix`, consulting only the engine the prefix
/// routes to. Readers on request paths use this instead of
/// [`Self::keys`], which scans and clones every key of every engine.
/// A `prefix` shorter than a registered prefix could span engines and
/// falls back to the full scan.
pub fn keys_with_prefix(&self, prefix: &str) -> Vec<String> {
let engines = self.engines_snapshot();
for (registered, engine) in engines.iter() {
if prefix.starts_with(registered.as_str()) {
return engine
Comment on lines +179 to +181

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P2 Badge Fall back to all engines for overlapping prefixes

Because the engine table is longest-prefix sorted, an overlapping setup such as registered prefixes foo:bar: and foo: makes keys_with_prefix("foo:") return as soon as it reaches the foo: engine, so it never scans the foo:bar: engine even though those keys also match the requested prefix. configure_crdt_prefix only rejects exact duplicates, so this regresses CrdtNamespace::keys("") for parent namespaces that overlap child namespaces from the previous all-engine scan and silently omits live keys.

Useful? React with 👍 / 👎.

.keys()
.into_iter()
.filter(|key| key.starts_with(prefix))
.collect();
}
}
if engines
.iter()
.any(|(registered, _)| registered.starts_with(prefix))
{
return self
.keys()
.into_iter()
.filter(|key| key.starts_with(prefix))
.collect();
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
self.default_engine
.keys()
.into_iter()
.filter(|key| key.starts_with(prefix))
.collect()
}
Comment on lines +177 to +203

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.

high

Instead of adding complex routing logic in keys_with_prefix to handle overlapping or duplicate prefix registrations, we should assert against duplicate or overlapping registration of prefixes when routing keys to different engines. This ensures that each key maps to exactly one engine, making iteration order irrelevant for correctness and preventing non-monotonic resets of aggregated metrics.

Please add an assertion in the engine registration path to prevent duplicate/overlapping prefixes, which allows keys_with_prefix to remain simple and correct.

    pub fn keys_with_prefix(&self, prefix: &str) -> Vec<String> {
        let engines = self.engines_snapshot();
        if let Some((_, engine)) = engines.iter().find(|(registered, _)| prefix.starts_with(registered.as_str())) {
            engine.keys().into_iter().filter(|key| key.starts_with(prefix)).collect()
        } else {
            self.default_engine.keys().into_iter().filter(|key| key.starts_with(prefix)).collect()
        }
    }
References
  1. When routing keys to different engines or namespaces by prefix, assert against duplicate registration of prefixes. This ensures that each key maps to exactly one engine, making iteration order irrelevant for correctness.


pub fn all(&self) -> BTreeMap<String, Vec<u8>> {
let mut all = BTreeMap::new();
for engine in self.all_engines() {
Expand Down
8 changes: 3 additions & 5 deletions crates/mesh/src/kv.rs
Original file line number Diff line number Diff line change
Expand Up @@ -291,13 +291,11 @@ impl CrdtNamespace {
}

/// List all live keys matching a sub-prefix within this namespace.
/// Routed to this namespace's engine, so the cost scales with this
/// namespace's keys, not the whole store's.
pub fn keys(&self, sub_prefix: &str) -> Vec<String> {
let full_prefix = format!("{}{}", self.prefix, sub_prefix);
self.store
.keys()
.into_iter()
.filter(|k| k.starts_with(&full_prefix))
.collect()
self.store.keys_with_prefix(&full_prefix)
}

/// Subscribe to changes for keys matching a sub-prefix within this namespace.
Expand Down
23 changes: 23 additions & 0 deletions model_gateway/src/mesh/adapters/rate_limit_sync.rs
Original file line number Diff line number Diff line change
Expand Up @@ -196,6 +196,29 @@ impl RateLimitSyncAdapter {
.try_fold(0i64, |acc, s| acc.checked_add(s.count))
.unwrap_or(i64::MAX)
}

/// Cluster-wide aggregate for a counter at exactly `epoch`. Shards from
/// other epochs contribute nothing: the caller pins the window instead
/// of trusting the highest observed epoch, which after a window roll is
/// a stale window nothing has published past yet.
pub fn get_aggregate_at(&self, counter_name: &str, epoch: u64) -> i64 {
debug_assert!(
!counter_name.contains(':'),
"counter_name must not contain ':' (got {counter_name:?})",
);
let sub_prefix = format!("{counter_name}:");
let mut total = 0i64;
for key in self.rate_limits.keys(&sub_prefix) {
if let Some(bytes) = self.rate_limits.get(&key) {
if let Some(value) = decode_epoch_count(&bytes) {
if value.epoch == epoch {
total = total.checked_add(value.count).unwrap_or(i64::MAX);
}
}
}
}
total
}
}

#[cfg(test)]
Expand Down
Loading
Loading