Skip to content

feat(smg): cluster-wide rate limiting via epoch-windowed mesh counters - #1715

Closed
CatherineSue wants to merge 5 commits into
mainfrom
chang/mesh-d3b-middleware
Closed

CatherineSue wants to merge 5 commits into
mainfrom
chang/mesh-d3b-middleware

Conversation

@CatherineSue

@CatherineSue CatherineSue commented Jun 13, 2026 •

Copy link
Copy Markdown
Member

Description

Problem

The v1 cluster-wide rate limit (MeshSyncManager::check_global_rate_limit) was removed in the worker-sync PR, leaving only per-node token-bucket limiting. A deployment configured with a cluster-wide limit_per_second had no enforcement of it on the v2 mesh — the hook the concurrency middleware comment promised would "return through the v2 RateLimitSyncAdapter."

Solution

Each gateway counts the requests it admits into a wall-clock-second window and publishes the cumulative count as its rl:global:{node} shard. Counts are monotone within a window — the EpochMaxWins merge pins the per-shard maximum, so a published count must never regress inside an epoch; a window roll advances the epoch so a fresh zero replaces the old peak instead of fighting it. The concurrency middleware consults the limiter before the local token bucket and 429s once the epoch-pinned cluster aggregate plus this node's unflushed delta reaches the limit.

  • GlobalRateLimiter (owned by MeshAdapters) holds the window. It deliberately is not part of RateLimitSyncAdapter — the adapter is a generic, policy-free counter bridge (its doc: "the caller owns the epoch clock"), while this is one specific policy (the global counter, config:rate_limit, window/grace/flush). This transport-vs-policy seam matches WorkerSyncAdapter (bridge) vs WorkerRegistry (policy).
  • check() / record_admit() split: check consumes no budget; only requests that also clear local admission are recorded, so a saturated node whose local bucket rejects a flood never burns the cluster's budget for its peers.
  • Inline, traffic-driven publishing — no background task. An admit advances the count and publishes the grown shard at most once per flush interval; once admits stop, the published value is frozen (the unpublished tail is bounded by one flush interval of admits and is enforcement-harmless). A design-review panel confirmed a background ticker buys nothing here: gossip ships ops on its own ~1s round regardless, so nothing delivers shards to peers faster than that floor — and it would cost a linted tokio::spawn + idle wakeups.
  • get_aggregate_at(epoch) pins the budget read to the current window's epoch. Trusting the highest observed epoch would let a stale exhausted window gate admission after a roll — and since only admitted requests advance the counter, no node would ever publish the new epoch and the cluster would reject forever.
  • Regression grace: a backward NTP step or a spurious far-future reading re-anchors the epoch after a bounded grace, so a clock step costs at most ~one grace period, not the step's magnitude.
  • keys_with_prefix (mesh crate): the per-request aggregate read no longer scans and clones every key of every engine in the shared store — it consults only the rl: engine.
  • Config is fail-open: a missing key disables cluster limiting; an undecodable value also fails open but warns once per failure streak (silently-disabled limiting must not look like unconfigured).

Known limitations (disclosed)

  • Coordination-free, not an exact global counter. A peer's spend reaches a node after that peer's publish (≤ flush interval) and the next gossip round (≤ ~1s) — the gossip term dominates. The effective guarantee is between limit and roughly limit + nodes × (flush-interval + gossip-period) admits per window; under sustained load it trends toward the per-node-island worst case.
  • Wall-clock windows assume NTP discipline. A node skewed ≥1s enforces its own window island (worst case islands × limit). No skew metric yet.
  • In-flight requests can all pass check before any record_admit (separate locks) — per-node overshoot bounded by in-flight concurrency; deliberate, dwarfed by the gossip-cadence slack.
  • The frozen unpublished tail under-reads each window's last sub-interval for any future rl: shard-history consumer (none today); wire flush-on-roll when one appears.

The abstraction and inline-flush mechanism were validated by a 4-agent design-review panel (architecture / correctness / codebase-precedent / adversarial-challenger), unanimous on keeping GlobalRateLimiter and 3–1 (challenger concurring) on inline flush; the four doc corrections it surfaced are folded in.

Changes

  • model_gateway/src/mesh/global_rate_limit.rs (new) — GlobalRateLimiter, RateLimitConfig, RATE_LIMIT_CONFIG_KEY
  • model_gateway/src/mesh/wiring.rs, mesh/mod.rs — construct + expose the limiter on MeshAdapters
  • model_gateway/src/mesh/adapters/rate_limit_sync.rs — get_aggregate_at
  • model_gateway/src/middleware/concurrency.rs — check-before-bucket / record-after-admit
  • crates/mesh/src/crdt_kv/crdt.rs, kv.rs — keys_with_prefix

Test Plan

  • cargo test -p smg and cargo test -p smg-mesh green; cargo clippy --workspace --all-targets --all-features -- -D warnings clean.
  • New global_rate_limit tests: no-config passthrough, undecodable-config fail-open, admit-to-limit-then-429, check-consumes-no-budget, window-roll reset, monotone cumulative flush, aggregate+unflushed enforcement, stale-epoch pinning, backward-step and forward-glitch recovery after the grace.

Independent of the worker/op-log mesh stacks; targets main directly.

Checklist
  • cargo +nightly fmt passes
  • cargo clippy --all-targets --all-features -- -D warnings passes
  • (Optional) Documentation updated
  • (Optional) Please join us on Slack #sig-smg to discuss, review, and merge PRs

🤖 Generated with Claude Code

Summary by CodeRabbit

  • New Features

    • Added cluster-wide, per-second rate limiting backed by shared state; requests that exceed the shared budget are rejected with HTTP 429.
    • Wired the global limiter into both scheduler admission and request middleware, including correct admit/budget accounting.
  • Enhancements

    • Improved prefix-based key enumeration to resolve sub-prefix queries more efficiently.
    • Added the ability to compute rate-limit aggregates for a specific epoch.
  • Documentation / Metrics

    • Added a global_rejected metric label to distinguish cluster-wide rejections from local ones.

@coderabbitai

coderabbitai Bot commented Jun 13, 2026 •

Copy link
Copy Markdown

Review Change Stack

Note

Reviews paused

It looks like this branch is under active development. To avoid overwhelming you with review comments due to an influx of new commits, CodeRabbit has automatically paused this review. You can configure this behavior by changing the reviews.auto_review.auto_pause_after_reviewed_commits setting.

Use the following commands to manage reviews:

  • @coderabbitai resume to resume automatic reviews.
  • @coderabbitai review to trigger a single review.

Use the checkboxes below for quick actions:

  • ▶️ Resume reviews
  • 🔍 Trigger review
📝 Walkthrough

Walkthrough

Adds a cluster-wide epoch-based global rate limiter backed by CRDT shard tracking, wires it through mesh startup and scheduler state, and enforces it in concurrency and scheduler admission paths with updated rate-limit metrics.

Changes

Global rate limiting feature

Layer / File(s) Summary
CRDT prefix-based key enumeration
crates/mesh/src/crdt_kv/crdt.rs, crates/mesh/src/kv.rs
Adds keys_with_prefix() to route prefix queries to a single matching engine when possible; CrdtNamespace::keys now delegates to store.keys_with_prefix(full_prefix).
RateLimitSyncAdapter epoch-pinned aggregation
model_gateway/src/mesh/adapters/rate_limit_sync.rs
Adds get_aggregate_at(counter_name, epoch) to compute the cluster aggregate for a specific epoch, summing only shards whose epoch equals the pinned epoch with overflow saturation.
GlobalRateLimiter core with window management and tests
model_gateway/src/mesh/global_rate_limit.rs
Adds GlobalRateLimiter, RateLimitConfig, RATE_LIMIT_CONFIG_KEY, unix_now_secs(), check(now_secs), record_admit(now_secs), and test-only flush(now_secs). Includes epoch window rotation, fail-open config deserialization with single-per-streak warning, paced monotone publishes to CRDT, and enforcement combining CRDT aggregate with unflushed local delta. Comprehensive unit tests validate unlimited/fail-open behavior, limit enforcement, window rotation, monotone publishing, and clock anomaly recovery.
Mesh module exports and adapter wiring
model_gateway/src/mesh/mod.rs, model_gateway/src/mesh/wiring.rs
Exports global_rate_limit submodule and types from mesh module; extends MeshAdapters with a global_rate_limit: Arc<GlobalRateLimiter> field constructed in start() alongside existing adapters, and provides accessor method.
Scheduler state threading and server wiring
model_gateway/src/middleware/scheduler/state.rs, model_gateway/src/server.rs
Threads optional GlobalRateLimiter into SchedulerState and AdmissionMode::from_config; updates server startup to derive and pass global_rate_limit to admission configuration when mesh adapters are present.
Concurrency middleware global rate limit enforcement
model_gateway/src/middleware/concurrency.rs
Integrates global limiter: performs early check(unix_now_secs()) (returns HTTP 429 on deny with RATE_LIMIT_GLOBAL_REJECTED), and records global admits via record_admit(now_secs) when requests are actually admitted locally (no-local-limiter passthrough, immediate token acquisition, and queued-path completion).
Scheduler admission global rate limit enforcement
model_gateway/src/middleware/scheduler/admission.rs, model_gateway/src/observability/metrics.rs
Adds cluster-wide global rate-limit gate to priority admission: pre-checks and denies with RATE_LIMIT_GLOBAL_REJECTED before RPS/scheduler logic, and records cluster-budget consumption after successful admission. Updates the rate-limit metric label set to include global_rejected.

Estimated code review effort: 4 (Complex) | ~50 minutes

Possibly related PRs

  • lightseekorg/smg#534: Introduces the CrdtOrMap::keys_with_prefix API that this PR now uses for namespace key listing.
  • lightseekorg/smg#1327: Adds the epoch-aware RateLimitSyncAdapter::get_aggregate_at(...) path extended here for cluster-wide limiting.
  • lightseekorg/smg#1577: Modifies the same scheduler admission flow that this PR extends with global rate-limit checks and admit accounting.

Suggested reviewers: tonyluj, llfl, slin1237

Poem

🐰 I hopped through shards and counted each beat,
Epochs stayed in line, neat as a treat.
A global gate now watches the lane,
Hop, admit, publish — balance in the rain.

🚥 Pre-merge checks | ✅ 5
✅ Passed checks (5 passed)
Check name Status Explanation
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The title accurately summarizes the main change: adding cluster-wide rate limiting using epoch-windowed mesh counters.
Docstring Coverage ✅ Passed Docstring coverage is 100.00% which is sufficient. The required threshold is 80.00%.
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.
✨ Finishing Touches
📝 Generate docstrings
  • Create stacked PR
  • Commit on current branch
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Commit unit tests in branch chang/mesh-d3b-middleware

Comment @coderabbitai help to get the list of available commands.

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: 58ddd19d5d

ℹ️ About Codex in GitHub

Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you

  • Open a pull request for review
  • Mark a draft as ready
  • Comment "@codex review".

If Codex has suggestions, it will comment; otherwise it will react with 👍.

When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".

match bincode::deserialize::<RateLimitConfig>(&bytes) {
Ok(config) => {
self.config_decode_failed.store(false, Ordering::Release);
Some(config.limit_per_second)

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 Preserve the documented zero-as-disabled rate limit

The mesh rate-limit docs state that limit_per_second: 0 disables rate limiting, but this path returns Some(0) as an active limit; check() then compares aggregate + unflushed < limit, which can never be true for a nonnegative aggregate. If an operator uses the documented zero value to disable the cluster limiter, every request is rejected with 429 instead of failing open.

Useful? React with 👍 / 👎.

Comment on lines +271 to +272
if let Some(limiter) = &global_limiter {
limiter.record_admit(unix_now_secs());

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 Recheck the global budget after queued requests get a token

When the local bucket is empty and queuing is enabled, the global check only happens before enqueue, but this record_admit(unix_now_secs()) runs after the request may have waited into a later second. Because check() intentionally reserves no budget, a burst that fills the queue can all pass the earlier check and then drain in a new window without ever checking that window's cluster budget, admitting up to the queue capacity beyond the configured global limit. Recheck immediately after the permit is granted (or reserve budget at enqueue) before running the request.

Useful? React with 👍 / 👎.

@gemini-code-assist gemini-code-assist Bot left a comment

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.

Code Review

This pull request introduces a cluster-wide rate limiter (GlobalRateLimiter) that aggregates request counts across nodes using an epoch-pinned CRDT store, integrated directly into the concurrency middleware. It also optimizes prefix-based key retrieval in the CRDT store with a new keys_with_prefix method. The review feedback highlights several critical improvement opportunities: asserting against duplicate prefix registrations to simplify routing logic, releasing the window lock before calling the CRDT store in check() to prevent severe lock contention, and caching the system clock reading in the middleware to avoid redundant calls and potential epoch inconsistencies.

Important

The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.

Comment on lines +177 to +203
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
.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();
}
self.default_engine
.keys()
.into_iter()
.filter(|key| key.starts_with(prefix))
.collect()
}

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.

Comment on lines +201 to +210
let mut window = self.window.lock();
window.rotate(now_secs, self.regression_grace);
// Pinned to this window's epoch: after a roll, the stale window's
// shards must not gate admission (nothing would ever publish the
// new epoch and the cluster would reject forever).
let aggregate = self
.rate_limit
.get_aggregate_at(GLOBAL_COUNTER, window.epoch);
let unflushed = (window.count - window.flushed).max(0);
aggregate.saturating_add(unflushed) < limit as i64

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

In check(), calling self.rate_limit.get_aggregate_at() while holding the window mutex lock can cause severe lock contention under high concurrency. get_aggregate_at performs a prefix scan and multiple lookups on the CRDT store, which is relatively slow.

Since we only need the epoch and unflushed delta from the window state to perform the check, we can rotate the window and extract these values under a brief lock acquisition, and then perform the CRDT store read outside of the lock. This significantly improves concurrency and throughput.

Additionally, casting limit directly to i64 (limit as i64) can lead to signed overflow/wrap-around bugs if limit is configured to a very large value. We should perform the comparison using safe u64 casting instead. Note that since the values are physically constrained to be far below u64::MAX, we do not need defensive overflow checks like saturating_add and can use standard addition.

        let (epoch, unflushed) = {
            let mut window = self.window.lock();
            window.rotate(now_secs, self.regression_grace);
            (window.epoch, (window.count - window.flushed).max(0))
        };
        // Pinned to this window's epoch: after a roll, the stale window's
        // shards must not gate admission (nothing would ever publish the
        // new epoch and the cluster would reject forever).
        // Read the aggregate outside the window lock to avoid holding it
        // during CRDT store prefix scans and lookups, preventing lock contention.
        let aggregate = self
            .rate_limit
            .get_aggregate_at(GLOBAL_COUNTER, epoch);
        let total = (aggregate.max(0) as u64) + (unflushed as u64);
        total < limit
References
  1. Do not add defensive overflow checks (e.g., checked_add or saturating_add) when the domain logic guarantees that the values are physically constrained to be far below the type's maximum limit.

Comment on lines +213 to +241
let global_limiter = app_state
.mesh_adapters
.as_ref()
.map(|adapters| adapters.global_rate_limit().clone());
if let Some(limiter) = &global_limiter {
if !limiter.check(unix_now_secs()) {
debug!("Cluster-wide rate limit exceeded, returning 429");
Metrics::record_http_rate_limit(metrics_labels::RATE_LIMIT_REJECTED);
return StatusCode::TOO_MANY_REQUESTS.into_response();
}
}

let token_bucket = match &app_state.context.rate_limiter {
Some(bucket) => bucket.clone(),
None => {
// Rate limiting disabled, pass through immediately
// No local limiter: the global check above was the admission.
if let Some(limiter) = &global_limiter {
limiter.record_admit(unix_now_secs());
}
return next.run(request).await;
}
};

// Try to acquire token immediately
if token_bucket.try_acquire(1.0).is_ok() {
debug!("Acquired token immediately");
if let Some(limiter) = &global_limiter {
limiter.record_admit(unix_now_secs());
}

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

Reading the system clock multiple times via unix_now_secs() throughout the middleware execution path is redundant and can lead to inconsistent epoch views if the clock ticks over a second boundary between the check and the record phases.

We should read the timestamp once at the start of the middleware and reuse it consistently. Please also update the queue acquisition path (line 272) to use this now_secs variable.

    let now_secs = unix_now_secs();
    let global_limiter = app_state
        .mesh_adapters
        .as_ref()
        .map(|adapters| adapters.global_rate_limit().clone());
    if let Some(limiter) = &global_limiter {
        if !limiter.check(now_secs) {
            debug!("Cluster-wide rate limit exceeded, returning 429");
            Metrics::record_http_rate_limit(metrics_labels::RATE_LIMIT_REJECTED);
            return StatusCode::TOO_MANY_REQUESTS.into_response();
        }
    }

    let token_bucket = match &app_state.context.rate_limiter {
        Some(bucket) => bucket.clone(),
        None => {
            // No local limiter: the global check above was the admission.
            if let Some(limiter) = &global_limiter {
                limiter.record_admit(now_secs);
            }
            return next.run(request).await;
        }
    };

    // Try to acquire token immediately
    if token_bucket.try_acquire(1.0).is_ok() {
        debug!("Acquired token immediately");
        if let Some(limiter) = &global_limiter {
            limiter.record_admit(now_secs);
        }

@github-actions github-actions Bot added model-gateway Model gateway crate changes mesh Mesh crate changes labels Jun 13, 2026

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Actionable comments posted: 2

🤖 Prompt for all review comments with 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.

Inline comments:
In `@model_gateway/src/mesh/global_rate_limit.rs`:
- Around line 253-258: The unix_now_secs() function currently maps a
SystemTime-before-UNIX_EPOCH error to 0 which can freeze the window rotation;
change the error handling so it returns a sentinel (e.g., u64::MAX) instead of
0: replace .unwrap_or(0) with .unwrap_or(u64::MAX) and add a short comment
noting this is to signal an anomalous system clock rather than silently using
epoch 0 so windows continue to rotate; alternatively, if you prefer fail-fast
behavior, replace it with a panic including the system time error—keep the
change localized to unix_now_secs().
- Around line 46-50: RateLimitConfig.limit_per_second is u64 but is cast to i64
in RateLimitConfig::check(), which can overflow and produce negative values; fix
by bounding the value before casting: in RateLimitConfig::limit() (or at
construction/validation in RateLimitConfig::check()) clamp or cap
limit_per_second to i64::MAX (e.g., use std::cmp::min(limit_per_second, i64::MAX
as u64)) and then cast to i64, or return an error/validation failure if the
configured value exceeds i64::MAX so downstream comparisons in check() remain
safe.
🪄 Autofix (Beta)

Fix all unresolved CodeRabbit comments on this PR:

  • Push a commit to this branch (recommended)
  • Create a new PR with the fixes

ℹ️ Review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: ASSERTIVE

Plan: Pro

Run ID: a6d0b1b6-b6df-4c0a-90a9-81adf29d8801

📥 Commits

Reviewing files that changed from the base of the PR and between 195db92 and 58ddd19.

📒 Files selected for processing (7)
  • crates/mesh/src/crdt_kv/crdt.rs
  • crates/mesh/src/kv.rs
  • model_gateway/src/mesh/adapters/rate_limit_sync.rs
  • model_gateway/src/mesh/global_rate_limit.rs
  • model_gateway/src/mesh/mod.rs
  • model_gateway/src/mesh/wiring.rs
  • model_gateway/src/middleware/concurrency.rs

Comment on lines +46 to +50
#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
pub struct RateLimitConfig {
/// Maximum admitted requests per second across the whole cluster.
pub limit_per_second: u64,
}

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

🧹 Nitpick | 🔵 Trivial | 💤 Low value

Consider bounding limit_per_second to prevent cast overflow.

The limit_per_second field is u64, but at line 210 in check() it's cast to i64 for comparison. If a limit > i64::MAX were configured (unrealistic but possible), the cast would wrap to negative in release builds, breaking enforcement.

🛡️ Defensive fix to cap the limit
 #[derive(Debug, Clone, Copy, Serialize, Deserialize)]
 pub struct RateLimitConfig {
     /// Maximum admitted requests per second across the whole cluster.
     pub limit_per_second: u64,
 }

Then in limit() method (around line 170):

         match bincode::deserialize::<RateLimitConfig>(&bytes) {
             Ok(config) => {
                 self.config_decode_failed.store(false, Ordering::Release);
-                Some(config.limit_per_second)
+                Some(config.limit_per_second.min(i64::MAX as u64))
             }
🤖 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 `@model_gateway/src/mesh/global_rate_limit.rs` around lines 46 - 50,
RateLimitConfig.limit_per_second is u64 but is cast to i64 in
RateLimitConfig::check(), which can overflow and produce negative values; fix by
bounding the value before casting: in RateLimitConfig::limit() (or at
construction/validation in RateLimitConfig::check()) clamp or cap
limit_per_second to i64::MAX (e.g., use std::cmp::min(limit_per_second, i64::MAX
as u64)) and then cast to i64, or return an error/validation failure if the
configured value exceeds i64::MAX so downstream comparisons in check() remain
safe.

Comment on lines +253 to +258
pub(crate) fn unix_now_secs() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|elapsed| elapsed.as_secs())
.unwrap_or(0)
}

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

🧹 Nitpick | 🔵 Trivial | ⚖️ Poor tradeoff

Verify handling when SystemTime precedes UNIX_EPOCH.

If the system clock is set before 1970-01-01 (unrealistic but possible), duration_since fails and the function returns 0. This would cause the window to never rotate (since now_secs=0 matches window.epoch=0), enforcing the limit as a single eternal window rather than per-second windows.

A more defensive approach would panic or use a sentinel value like u64::MAX to signal the anomaly, but given modern systems won't have pre-1970 clocks, this is likely acceptable as-is.

🤖 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 `@model_gateway/src/mesh/global_rate_limit.rs` around lines 253 - 258, The
unix_now_secs() function currently maps a SystemTime-before-UNIX_EPOCH error to
0 which can freeze the window rotation; change the error handling so it returns
a sentinel (e.g., u64::MAX) instead of 0: replace .unwrap_or(0) with
.unwrap_or(u64::MAX) and add a short comment noting this is to signal an
anomalous system clock rather than silently using epoch 0 so windows continue to
rotate; alternatively, if you prefer fail-fast behavior, replace it with a panic
including the system time error—keep the change localized to unix_now_secs().

if let Some(limiter) = &global_limiter {
if !limiter.check(unix_now_secs()) {
debug!("Cluster-wide rate limit exceeded, returning 429");
Metrics::record_http_rate_limit(metrics_labels::RATE_LIMIT_REJECTED);

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

🟡 Nit: This records RATE_LIMIT_REJECTED — the same label the local token-bucket rejections use at lines 285, 296, and 302. An operator triaging a 429 spike can't tell whether the cluster-wide limiter or the per-node bucket is responsible.

Consider a distinct label (e.g., RATE_LIMIT_GLOBAL_REJECTED) so dashboards and alerts can differentiate the two rejection sources.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Good call — fixed in 6c5bee9. Cluster-limiter rejections now record a distinct RATE_LIMIT_GLOBAL_REJECTED label (alongside the local bucket's RATE_LIMIT_REJECTED), so a 429 spike is attributable to the cluster limiter vs this node's token bucket on dashboards/alerts.

@claude claude Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Clean, well-documented design. The epoch-pinned aggregation, check/record split, and regression grace are all well-reasoned. One minor nit on metric label distinguishability.

Summary: 0 🔴 Important · 1 🟡 Nit · 0 🟣 Pre-existing

CatherineSue added a commit that referenced this pull request Jun 13, 2026
Cluster-limiter 429s reused RATE_LIMIT_REJECTED, the same label the
local token bucket uses, so a 429 spike couldn't be attributed to the
cluster limiter vs this node's bucket. Cluster rejections now record
RATE_LIMIT_GLOBAL_REJECTED. Reported by review on #1715.

Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: 6c5bee9c18

ℹ️ About Codex in GitHub

Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you

  • Open a pull request for review
  • Mark a draft as ready
  • Comment "@codex review".

If Codex has suggestions, it will comment; otherwise it will react with 👍.

When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".

Comment on lines +213 to +216
let global_limiter = app_state
.mesh_adapters
.as_ref()
.map(|adapters| adapters.global_rate_limit().clone());

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 Enforce global rate limits in priority mode

When priority_scheduler_enabled is true, with_admission_layer installs priority_admission_middleware instead of this middleware (model_gateway/src/server.rs lines 780-784), and that path only checks the local RPS bucket before scheduler admission. Because the new GlobalRateLimiter check/record is wired only here, any mesh deployment using the priority scheduler ignores config:rate_limit entirely, so the advertised cluster-wide limit is not enforced for those protected routes.

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Real gap — fixed in 9023596. SchedulerState now carries the GlobalRateLimiter (threaded through AdmissionMode::from_config from app_state.mesh_adapters), and priority_admission_middleware checks it before the local RPS bucket and scheduler (so a cluster-rejected request consumes neither slot), recording the admission only once the request clears both — preserving the no-burn-on-local-reject invariant the legacy path has. Cluster rejections record RATE_LIMIT_GLOBAL_REJECTED. config:rate_limit is now enforced under the priority scheduler too.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Actionable comments posted: 1

🤖 Prompt for all review comments with 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.

Inline comments:
In `@model_gateway/src/observability/metrics.rs`:
- Around line 435-437: The Prometheus counter description string that currently
says "Rate limiting decisions by result (allowed/rejected)" must be updated to
reflect the new label value introduced by RATE_LIMIT_GLOBAL_REJECTED; locate the
metric creation where the counter description is set (the counter for rate
limiting decisions) and change the description to include "global_rejected" (for
example: "Rate limiting decisions by result (allowed/rejected/global_rejected)")
so the metric docs match the actual label cardinality and operators can see all
possible values.
🪄 Autofix (Beta)

Fix all unresolved CodeRabbit comments on this PR:

  • Push a commit to this branch (recommended)
  • Create a new PR with the fixes

ℹ️ Review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: ASSERTIVE

Plan: Pro

Run ID: af955aef-aa82-4404-9342-2e7b2a2c3b99

📥 Commits

Reviewing files that changed from the base of the PR and between 58ddd19 and 6c5bee9.

📒 Files selected for processing (2)
  • model_gateway/src/middleware/concurrency.rs
  • model_gateway/src/observability/metrics.rs

Comment thread model_gateway/src/observability/metrics.rs
CatherineSue added a commit that referenced this pull request Jun 13, 2026
The smg_http_rate_limit_total describe still read (allowed/rejected)
after the global_rejected label was added. Reported by review on #1715.

Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
CatherineSue added a commit that referenced this pull request Jun 13, 2026
The GlobalRateLimiter was wired only into the legacy concurrency
middleware, so a mesh deployment running the priority scheduler ignored
config:rate_limit entirely — the advertised cluster-wide limit silently
did nothing. SchedulerState now carries the limiter (threaded through
AdmissionMode::from_config from app_state.mesh_adapters); the priority
admission path checks it before the local RPS bucket and scheduler (so a
cluster-rejected request consumes neither) and records the admission
only once the request clears both, preserving the no-burn-on-local-reject
invariant. Cluster rejections record RATE_LIMIT_GLOBAL_REJECTED, matching
the legacy path. Reported by review on #1715.

Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: 9023596242

ℹ️ About Codex in GitHub

Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you

  • Open a pull request for review
  • Mark a draft as ready
  • Comment "@codex review".

If Codex has suggestions, it will comment; otherwise it will react with 👍.

When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".

// count this admission toward the cluster window now (only
// admitted requests consume cluster budget).
if let Some(limiter) = &state.global_rate_limit {
limiter.record_admit(unix_now_secs());

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 Recheck the global limit after priority admission waits

When the priority scheduler is enabled and saturated, state.scheduler.admit(...).await can queue this request until a later wall-clock second, but the global check was done only before the wait and check() does not reserve budget. In that case every queued request that passed the old-window check can be recorded into the new window here without rechecking the new cluster aggregate, allowing a burst up to the scheduler queue capacity over config:rate_limit.

Useful? React with 👍 / 👎.

Comment on lines +179 to +181
for (registered, engine) in engines.iter() {
if prefix.starts_with(registered.as_str()) {
return engine

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 👍 / 👎.

if let Some(limiter) = &state.global_rate_limit {
if !limiter.check(unix_now_secs()) {
Metrics::record_http_rate_limit(metrics_labels::RATE_LIMIT_GLOBAL_REJECTED);
return SchedulerError::QueueFull.into_response();

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

🟡 Nit: SchedulerError::QueueFull produces code: "scheduler_queue_full" and message: "request queue is full for this priority class" in the 429 body. For a cluster-wide rate-limit rejection that's misleading — an operator examining client-side errors would investigate per-class queue depths (empty) instead of the cluster budget (exhausted).

The metric label (RATE_LIMIT_GLOBAL_REJECTED) correctly differentiates, so monitoring is fine. But a dedicated variant (e.g. SchedulerError::RateLimited) with an accurate message would save debugging time. The local RPS bucket at line 135 has the same issue, so this could be a follow-up that covers both.

CatherineSue added a commit that referenced this pull request Jun 17, 2026
Cluster-limiter 429s reused RATE_LIMIT_REJECTED, the same label the
local token bucket uses, so a 429 spike couldn't be attributed to the
cluster limiter vs this node's bucket. Cluster rejections now record
RATE_LIMIT_GLOBAL_REJECTED. Reported by review on #1715.

Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
CatherineSue added a commit that referenced this pull request Jun 17, 2026
The smg_http_rate_limit_total describe still read (allowed/rejected)
after the global_rejected label was added. Reported by review on #1715.

Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
CatherineSue added a commit that referenced this pull request Jun 17, 2026
The GlobalRateLimiter was wired only into the legacy concurrency
middleware, so a mesh deployment running the priority scheduler ignored
config:rate_limit entirely — the advertised cluster-wide limit silently
did nothing. SchedulerState now carries the limiter (threaded through
AdmissionMode::from_config from app_state.mesh_adapters); the priority
admission path checks it before the local RPS bucket and scheduler (so a
cluster-rejected request consumes neither) and records the admission
only once the request clears both, preserving the no-burn-on-local-reject
invariant. Cluster rejections record RATE_LIMIT_GLOBAL_REJECTED, matching
the legacy path. Reported by review on #1715.

Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
@CatherineSue
CatherineSue force-pushed the chang/mesh-d3b-middleware branch from 9023596 to dba4dd8 Compare June 17, 2026 19:37

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Actionable comments posted: 2

♻️ Duplicate comments (1)
model_gateway/src/mesh/global_rate_limit.rs (1)

166-172: ⚠️ Potential issue | 🟠 Major | ⚡ Quick win

Avoid lossy u64 -> i64 cast in limit enforcement.

On Line 210, limit as i64 can wrap when config sets limit_per_second > i64::MAX, which can turn the threshold negative and reject essentially all traffic.

💡 Suggested fix
-    fn limit(&self) -> Option<u64> {
+    fn limit(&self) -> Option<i64> {
         let bytes = self.configs.get(RATE_LIMIT_CONFIG_KEY)?;
         match bincode::deserialize::<RateLimitConfig>(&bytes) {
             Ok(config) => {
                 self.config_decode_failed.store(false, Ordering::Release);
-                Some(config.limit_per_second)
+                Some(i64::try_from(config.limit_per_second).unwrap_or(i64::MAX))
             }
             Err(err) => {
                 if !self.config_decode_failed.swap(true, Ordering::AcqRel) {
                     warn!(
                         %err,
@@
-        aggregate.saturating_add(unflushed) < limit as i64
+        aggregate.saturating_add(unflushed) < limit

Also applies to: 210-210

🤖 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 `@model_gateway/src/mesh/global_rate_limit.rs` around lines 166 - 172, The
limit() method returns config.limit_per_second as u64, but the code at line 210
casts this limit to i64 without handling overflow. When limit_per_second exceeds
i64::MAX, the cast will wrap to a negative value, causing all traffic to be
rejected. Replace the unsafe cast with a saturating or checked conversion that
handles the overflow case properly, ensuring that values larger than i64::MAX
are either clamped to i64::MAX or handled explicitly to prevent negative
thresholds in the rate limiting logic.
🤖 Prompt for all review comments with 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.

Inline comments:
In `@crates/mesh/src/crdt_kv/crdt.rs`:
- Around line 177-197: The keys_with_prefix function has an early return on the
first registered prefix match that prevents checking other engines with
overlapping prefixes. When both foo: and foo:bar: are registered engines and
keys_with_prefix("foo:") is called, the function returns keys from only the
first matching engine instead of collecting keys from all relevant engines.
Remove the early return after finding the first prefix match in the first loop
iteration and instead accumulate keys from all registered prefixes that match,
then return the complete collection of keys from all matching engines.

In `@model_gateway/src/middleware/scheduler/admission.rs`:
- Around line 116-126: The global rate limit rejection in the check method is
incorrectly returning SchedulerError::QueueFull, which conflates global rate
limiting with scheduler queue saturation. Add a new SchedulerError variant
called GlobalRateLimited in error.rs and provide the necessary mappings for it
(StatusCode::TOO_MANY_REQUESTS and error code string "global_rate_limited").
Then update the check in admission.rs to return
SchedulerError::GlobalRateLimited instead of SchedulerError::QueueFull when the
global rate limiter rejects the request.

---

Duplicate comments:
In `@model_gateway/src/mesh/global_rate_limit.rs`:
- Around line 166-172: The limit() method returns config.limit_per_second as
u64, but the code at line 210 casts this limit to i64 without handling overflow.
When limit_per_second exceeds i64::MAX, the cast will wrap to a negative value,
causing all traffic to be rejected. Replace the unsafe cast with a saturating or
checked conversion that handles the overflow case properly, ensuring that values
larger than i64::MAX are either clamped to i64::MAX or handled explicitly to
prevent negative thresholds in the rate limiting logic.
🪄 Autofix (Beta)

Fix all unresolved CodeRabbit comments on this PR:

  • Push a commit to this branch (recommended)
  • Create a new PR with the fixes

ℹ️ Review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: ASSERTIVE

Plan: Pro

Run ID: 38d7ec66-10d0-44a2-aa2a-0b71e1df5fd9

📥 Commits

Reviewing files that changed from the base of the PR and between 9023596 and dba4dd8.

📒 Files selected for processing (11)
  • crates/mesh/src/crdt_kv/crdt.rs
  • crates/mesh/src/kv.rs
  • model_gateway/src/mesh/adapters/rate_limit_sync.rs
  • model_gateway/src/mesh/global_rate_limit.rs
  • model_gateway/src/mesh/mod.rs
  • model_gateway/src/mesh/wiring.rs
  • model_gateway/src/middleware/concurrency.rs
  • model_gateway/src/middleware/scheduler/admission.rs
  • model_gateway/src/middleware/scheduler/state.rs
  • model_gateway/src/observability/metrics.rs
  • model_gateway/src/server.rs

Comment thread crates/mesh/src/crdt_kv/crdt.rs
Comment on lines +116 to +126
// Cluster-wide rate limit, checked before the local RPS bucket and
// scheduler so a cluster-rejected request consumes neither. Mirrors the
// legacy concurrency path; the check consumes no budget (recorded only
// once the request is actually admitted below), so a rejection here
// never burns the cluster's budget.
if let Some(limiter) = &state.global_rate_limit {
if !limiter.check(unix_now_secs()) {
Metrics::record_http_rate_limit(metrics_labels::RATE_LIMIT_GLOBAL_REJECTED);
return SchedulerError::QueueFull.into_response();
}
}

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

🧹 Nitpick | 🔵 Trivial | ⚡ Quick win

Consider adding a dedicated SchedulerError variant for global rate-limit rejections.

Returning SchedulerError::QueueFull for global rate-limit denials conflates two distinct rejection reasons. The error code in responses will be "scheduler_queue_full" (per error.rs), making it harder for operators to distinguish cluster-wide throttling from actual scheduler queue saturation in logs and client diagnostics.

The metric (RATE_LIMIT_GLOBAL_REJECTED) is correctly distinct, but the HTTP error body will misattribute the cause.

💡 Suggested approach

Add a new variant to SchedulerError:

// In error.rs
GlobalRateLimited,

With mappings:

Self::GlobalRateLimited => StatusCode::TOO_MANY_REQUESTS,
// ...
Self::GlobalRateLimited => "global_rate_limited",

Then in this file:

 if let Some(limiter) = &state.global_rate_limit {
     if !limiter.check(unix_now_secs()) {
         Metrics::record_http_rate_limit(metrics_labels::RATE_LIMIT_GLOBAL_REJECTED);
-        return SchedulerError::QueueFull.into_response();
+        return SchedulerError::GlobalRateLimited.into_response();
     }
 }
🤖 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 `@model_gateway/src/middleware/scheduler/admission.rs` around lines 116 - 126,
The global rate limit rejection in the check method is incorrectly returning
SchedulerError::QueueFull, which conflates global rate limiting with scheduler
queue saturation. Add a new SchedulerError variant called GlobalRateLimited in
error.rs and provide the necessary mappings for it
(StatusCode::TOO_MANY_REQUESTS and error code string "global_rate_limited").
Then update the check in admission.rs to return
SchedulerError::GlobalRateLimited instead of SchedulerError::QueueFull when the
global rate limiter rejects the request.

@github-actions

github-actions Bot commented Jul 2, 2026

Copy link
Copy Markdown

This pull request has been automatically marked as stale because it has not had any activity within 14 days. It will be automatically closed if no further activity occurs within 16 days. Leave a comment if you feel this pull request should remain open. Thank you!

@github-actions github-actions Bot added the stale PR has been inactive for 14+ days label Jul 2, 2026
CrdtNamespace::keys scanned and cloned every key of every engine in the
shared store before filtering. keys_with_prefix consults only the engine
the prefix routes to, falling back to the full scan for a prefix shorter
than any registered one — keeps per-request rl: aggregate reads bounded by
the rl: engine's key count instead of the whole store.

Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
Restores the cluster-wide limit the v1 check_global_rate_limit hook
enforced, on the v2 mesh: each gateway counts the requests it admits into
a wall-clock-second window and publishes the cumulative count as its
rl:global:{node} shard (monotone within the window, as EpochMaxWins
requires; a window roll advances the epoch so a fresh zero replaces the
old peak). The concurrency middleware checks the limiter before the local
token bucket and records only locally-admitted requests, so a saturated
node's rejected flood never burns cluster budget; 429 once the
epoch-pinned cluster aggregate plus this node's unflushed delta reaches
config:rate_limit's limit_per_second.

GlobalRateLimiter (owned by MeshAdapters) holds the window; publishing is
traffic-driven and inline (no background task), paced by a flush interval.
A regression grace re-anchors the epoch after a wall-clock step so it
costs at most one grace period, and an undecodable config fails open with
a warn-once. The budget read is pinned to the current window's epoch via
RateLimitSyncAdapter::get_aggregate_at: trusting the highest observed
epoch would let a stale exhausted window wedge the cluster into rejecting
forever, since only admitted requests advance the counter.

Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
Cluster-limiter 429s reused RATE_LIMIT_REJECTED, the same label the
local token bucket uses, so a 429 spike couldn't be attributed to the
cluster limiter vs this node's bucket. Cluster rejections now record
RATE_LIMIT_GLOBAL_REJECTED. Reported by review on #1715.

Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
The smg_http_rate_limit_total describe still read (allowed/rejected)
after the global_rejected label was added. Reported by review on #1715.

Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
The GlobalRateLimiter was wired only into the legacy concurrency
middleware, so a mesh deployment running the priority scheduler ignored
config:rate_limit entirely — the advertised cluster-wide limit silently
did nothing. SchedulerState now carries the limiter (threaded through
AdmissionMode::from_config from app_state.mesh_adapters); the priority
admission path checks it before the local RPS bucket and scheduler (so a
cluster-rejected request consumes neither) and records the admission
only once the request clears both, preserving the no-burn-on-local-reject
invariant. Cluster rejections record RATE_LIMIT_GLOBAL_REJECTED, matching
the legacy path. Reported by review on #1715.

Signed-off-by: Chang Su <8605658+CatherineSue@users.noreply.github.com>
@CatherineSue
CatherineSue force-pushed the chang/mesh-d3b-middleware branch from dba4dd8 to 3602bdb Compare July 2, 2026 16:22

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: 3602bdbbc6

ℹ️ About Codex in GitHub

Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you

  • Open a pull request for review
  • Mark a draft as ready
  • Comment "@codex review".

If Codex has suggestions, it will comment; otherwise it will react with 👍.

When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".

Comment on lines +106 to +108
let regressed = now_secs < self.epoch && self.rotated_at.elapsed() >= grace;
if now_secs > self.epoch || regressed {
self.epoch = now_secs;

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 Preserve lower-epoch counts after clock recovery

When the wall clock moves backward, or recovers from a published far-future glitch, this branch re-anchors the local window to a lower epoch even though the rl: CRDT keeps the previously published higher epoch as the shard's current value. Subsequent lower-epoch sync_counter writes are therefore not returned by decode_epoch_count/get_aggregate_at, but publish still marks them flushed, so both this node and its peers omit those admitted requests until real time catches up to the old epoch; under load that can repeatedly admit another budget every flush interval during the clock-recovery period.

Useful? React with 👍 / 👎.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Actionable comments posted: 1

Caution

Some comments are outside the diff and can’t be posted inline due to platform limitations.

⚠️ Outside diff range comments (1)
model_gateway/src/mesh/wiring.rs (1)

94-153: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick win

Consider a wiring-level integration test for global_rate_limit().

GlobalRateLimiter's internals are thoroughly unit-tested in global_rate_limit.rs, but those tests build their own limiter directly and never go through MeshAdapters::start. A wiring bug (e.g., the limiter accidentally consulting the wrong CRDT namespace for mesh_kv.configs()) wouldn't be caught by either suite today. A small test here — put RATE_LIMIT_CONFIG_KEY via mesh.configs(), then assert on adapters.global_rate_limit().check(...) — would close the gap, mirroring the existing rl_namespace_uses_epoch_max_wins test style.

🤖 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 `@model_gateway/src/mesh/wiring.rs` around lines 94 - 153, Add a wiring-level
integration test around MeshAdapters::start and global_rate_limit() to cover the
CRDT namespace path. Use the existing test patterns in this module, set
RATE_LIMIT_CONFIG_KEY through mesh.configs(), then verify behavior via
adapters.global_rate_limit().check(...) instead of constructing
GlobalRateLimiter directly. This should live alongside
rl_namespace_uses_epoch_max_wins and exercise the real wiring so a namespace
mismatch in MeshAdapters or GlobalRateLimiter is caught.
♻️ Duplicate comments (1)
crates/mesh/src/crdt_kv/crdt.rs (1)

172-204: 🗄️ Data Integrity & Integration | 🟠 Major | ⚡ Quick win

Early-return still drops keys when registered prefixes overlap.

This is the same defect flagged in a previous review: the first loop returns as soon as it finds any registered prefix that the query prefix starts with, before checking whether a longer registered prefix (a descendant, e.g. foo:bar:) also matches. Because engines are sorted longest-prefix-first, a query like keys_with_prefix("foo:") will still miss keys routed to foo:bar:'s engine if foo: matches first in this loop before the multi-engine fallback is reached.

Currently dormant only because the two prefixes actually registered (worker:, rl:) don't overlap — but this is a general-purpose crate API, so a future third overlapping registration would silently under-return keys.

💡 Suggested fix (same as prior review)
 pub fn keys_with_prefix(&self, prefix: &str) -> Vec<String> {
     let engines = self.engines_snapshot();
+    // Descendant registered prefixes mean this query can span multiple engines.
+    if engines.iter().any(|(registered, _)| {
+        registered.len() > prefix.len() && registered.starts_with(prefix)
+    }) {
+        return self
+            .keys()
+            .into_iter()
+            .filter(|key| key.starts_with(prefix))
+            .collect();
+    }
     for (registered, engine) in engines.iter() {
         if prefix.starts_with(registered.as_str()) {
             return engine
                 .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();
-    }
     self.default_engine
         .keys()
         .into_iter()
         .filter(|key| key.starts_with(prefix))
         .collect()
 }
🤖 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/mesh/src/crdt_kv/crdt.rs` around lines 172 - 204, The early return in
keys_with_prefix causes missed keys when registered prefixes overlap, because it
stops on the first matching engine before considering a longer descendant prefix
that may also match. Update keys_with_prefix to preserve the longest-prefix
routing behavior used by engines_snapshot, either by selecting the most specific
matching registered prefix before returning or by reusing the same routing logic
used elsewhere in crdt.rs so overlapping registrations are handled correctly.
Ensure the fallback full-scan path still only runs when the query prefix can
span multiple engines.
🤖 Prompt for all review comments with 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.

Inline comments:
In `@model_gateway/src/middleware/concurrency.rs`:
- Around line 227-236: The passthrough branch in concurrency.rs is missing the
local allowed metric update, so `smg_http_rate_limit_total` undercounts admitted
requests when `app_state.context.rate_limiter` is absent. Update the
no-local-limiter path in `concurrency` to record the same `RATE_LIMIT_ALLOWED`
outcome used by the immediate-acquire and queued-acquire paths before calling
`record_admit` and `next.run`, keeping metric parity with the other admission
flows.

---

Outside diff comments:
In `@model_gateway/src/mesh/wiring.rs`:
- Around line 94-153: Add a wiring-level integration test around
MeshAdapters::start and global_rate_limit() to cover the CRDT namespace path.
Use the existing test patterns in this module, set RATE_LIMIT_CONFIG_KEY through
mesh.configs(), then verify behavior via adapters.global_rate_limit().check(...)
instead of constructing GlobalRateLimiter directly. This should live alongside
rl_namespace_uses_epoch_max_wins and exercise the real wiring so a namespace
mismatch in MeshAdapters or GlobalRateLimiter is caught.

---

Duplicate comments:
In `@crates/mesh/src/crdt_kv/crdt.rs`:
- Around line 172-204: The early return in keys_with_prefix causes missed keys
when registered prefixes overlap, because it stops on the first matching engine
before considering a longer descendant prefix that may also match. Update
keys_with_prefix to preserve the longest-prefix routing behavior used by
engines_snapshot, either by selecting the most specific matching registered
prefix before returning or by reusing the same routing logic used elsewhere in
crdt.rs so overlapping registrations are handled correctly. Ensure the fallback
full-scan path still only runs when the query prefix can span multiple engines.
🪄 Autofix (Beta)

Fix all unresolved CodeRabbit comments on this PR:

  • Push a commit to this branch (recommended)
  • Create a new PR with the fixes

ℹ️ Review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: ASSERTIVE

Plan: Pro

Run ID: a00166c7-e29f-4a3b-9dae-9d3dbdc4c981

📥 Commits

Reviewing files that changed from the base of the PR and between dba4dd8 and 3602bdb.

📒 Files selected for processing (11)
  • crates/mesh/src/crdt_kv/crdt.rs
  • crates/mesh/src/kv.rs
  • model_gateway/src/mesh/adapters/rate_limit_sync.rs
  • model_gateway/src/mesh/global_rate_limit.rs
  • model_gateway/src/mesh/mod.rs
  • model_gateway/src/mesh/wiring.rs
  • model_gateway/src/middleware/concurrency.rs
  • model_gateway/src/middleware/scheduler/admission.rs
  • model_gateway/src/middleware/scheduler/state.rs
  • model_gateway/src/observability/metrics.rs
  • model_gateway/src/server.rs

Comment on lines 227 to 236
let token_bucket = match &app_state.context.rate_limiter {
Some(bucket) => bucket.clone(),
None => {
// Rate limiting disabled, pass through immediately
// No local limiter: the global check above was the admission.
if let Some(limiter) = &global_limiter {
limiter.record_admit(unix_now_secs());
}
return next.run(request).await;
}
};

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

🎯 Functional Correctness | 🔵 Trivial | ⚡ Quick win

Passthrough admits aren't counted as allowed in smg_http_rate_limit_total.

The immediate-acquire and queued-acquire paths both record RATE_LIMIT_ALLOWED alongside record_admit, but this no-local-limiter passthrough branch only records the admit to the cluster limiter, not the local metric. Since the global check is now a real gating decision here, dashboards computing rejection ratio from smg_http_rate_limit_total will under-count total decisions on nodes without a local token bucket, skewing the ratio whenever the global limiter denies.

📊 Proposed fix for metric parity
         None => {
             // No local limiter: the global check above was the admission.
             if let Some(limiter) = &global_limiter {
                 limiter.record_admit(unix_now_secs());
+                Metrics::record_http_rate_limit(metrics_labels::RATE_LIMIT_ALLOWED);
             }
             return next.run(request).await;
         }
📝 Committable suggestion

‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.

Suggested change
let token_bucket = match &app_state.context.rate_limiter {
Some(bucket) => bucket.clone(),
None => {
// Rate limiting disabled, pass through immediately
// No local limiter: the global check above was the admission.
if let Some(limiter) = &global_limiter {
limiter.record_admit(unix_now_secs());
}
return next.run(request).await;
}
};
let token_bucket = match &app_state.context.rate_limiter {
Some(bucket) => bucket.clone(),
None => {
// No local limiter: the global check above was the admission.
if let Some(limiter) = &global_limiter {
limiter.record_admit(unix_now_secs());
Metrics::record_http_rate_limit(metrics_labels::RATE_LIMIT_ALLOWED);
}
return next.run(request).await;
}
};
🤖 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 `@model_gateway/src/middleware/concurrency.rs` around lines 227 - 236, The
passthrough branch in concurrency.rs is missing the local allowed metric update,
so `smg_http_rate_limit_total` undercounts admitted requests when
`app_state.context.rate_limiter` is absent. Update the no-local-limiter path in
`concurrency` to record the same `RATE_LIMIT_ALLOWED` outcome used by the
immediate-acquire and queued-acquire paths before calling `record_admit` and
`next.run`, keeping metric parity with the other admission flows.

@github-actions github-actions Bot removed the stale PR has been inactive for 14+ days label Jul 3, 2026
@github-actions

Copy link
Copy Markdown

This pull request has been automatically marked as stale because it has not had any activity within 14 days. It will be automatically closed if no further activity occurs within 16 days. Leave a comment if you feel this pull request should remain open. Thank you!

@github-actions github-actions Bot added the stale PR has been inactive for 14+ days label Jul 18, 2026
@github-actions

github-actions Bot commented Aug 3, 2026

Copy link
Copy Markdown

This pull request has been automatically closed due to inactivity. Please feel free to reopen if you intend to continue working on it. Thank you!

@github-actions github-actions Bot closed this Aug 3, 2026
@lightseek-bot
lightseek-bot deleted the chang/mesh-d3b-middleware branch August 4, 2026 08:57
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

mesh Mesh crate changes model-gateway Model gateway crate changes stale PR has been inactive for 14+ days

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant