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
2 changes: 2 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,8 @@ DB_PASSWORD=__SECRET_HEX_16__
REDIS_PASSWORD=__SECRET_HEX_16__
# dev: REDIS_URL=redis://:${REDIS_PASSWORD}@localhost:6379
# prod: REDIS_URL=redis://:${REDIS_PASSWORD}@redis:6379
# Redis Cluster: REDIS_URL=redis-cluster://:${REDIS_PASSWORD}@redis-0:6379?node=redis-1:6379
# (see deploy/helm/think-watch/README.md, "Redis Cluster")

# --- Application ---
JWT_SECRET=__SECRET_HEX_32__
Expand Down
54 changes: 54 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,60 @@ target.

## [Unreleased]

### Read before upgrading

- **Rate-limit windows start empty.** Rate-limit counters move to new Redis
keys (one hash per counter, tagged so that Redis Cluster can run them), and
the counts from before the upgrade are not carried over: every window starts
empty and fills from the first request after the upgrade. The old keys
expire by themselves within two window lengths. Budget counters are kept.
- **Route health starts fresh.** A route's samples, circuit breaker and
lifetime request count move to new keys for the same reason, so every route
starts closed with nothing counted. The old lifetime counters never expire;
`redis-cli --scan --pattern 'route_health:[0-9a-f]*' | xargs redis-cli del`
removes them (the new keys start `route_health:{`).
- **A key's limits no longer replace its owner's.** A rate limit or budget set
on an API key used to take the place of the owner's limit for the same
window or period. Both now apply, each on its own counter: a key's limits can
narrow what its owner may do through that key, never widen it. A key given a
higher limit than its owner to give it more room needs the owner's limit
raised instead.

### Fixed

- **Token limits refuse requests.** A `tokens` rate limit never refused
anything, and stopped counting once a request would have taken it past its
limit. A request is now refused once the window's recorded usage reaches the
limit, and every request's tokens are recorded after it, even past the limit
— a window can overshoot by what was in flight when it filled.
- **Several request limits at once.** With two or more `requests` limits on a
user (per minute and per hour, say), every request that passed was counted
twice, and only one of the limits could refuse; with
`security.rate_limit_fail_closed` on, every request was refused as
`rate_limiter_unavailable`.
- **An API key's limits count on that key.** They were counted on its owner's
counter, which every key of the owner shared, and the usage the console
reads for a key (`/api/admin/limits/api_key/{id}/usage`) was always 0. Each
key now has counters of its own, rate limits and budgets alike, for the
gateway and the MCP gateway, and its usage shows what it used.
- **A refused request counts against nothing.** A request refused by a spent
budget, or by one rate limit after another had passed, was still counted
against the request limits. Budgets are now checked first and every rate
limit in one step, so a refused request leaves every counter as it was.
- **`Retry-After` says when to retry.** A `429` from the gateway's own limits
said `Retry-After: 30` whatever the limit. It now gives the seconds until the
window has room for another request, or until a spent budget's period ends
(the next midnight, Monday or 1st of the month, UTC). A spent budget also
sends `x-should-retry: false`, so the OpenAI and Anthropic SDKs don't retry it
by themselves. The body stays in the caller's API format.
- **Redis Cluster.** The rate-limit, route-health and quota scripts touched
keys of several hash slots, which a Redis Cluster refuses: on a cluster, rate
limits silently stopped applying (or refused every request with
`security.rate_limit_fail_closed`), circuit breakers never tripped, and cache
invalidation reached one node only. Every key a script touches now shares a
hash tag, pattern deletes scan every node, and the Helm chart's README
describes a `redis-cluster://` URL.

## [3.1.0] — 2026-10-05

The thinkwatch-core crates move from v0.59.0 to v0.62.0. Two changes reach
Expand Down
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

4 changes: 3 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -92,7 +92,9 @@ The gateway (port `3000`) is the only part that clients need to reach. The conso
- A model's maximum output tokens, set on the Models page, caps `max_tokens` on every request to that model; it replaces the old output length guardrail.

**Limits and budgets**
- Request-count limits are checked before the request; token limits and budgets are counted after the response, so one request can cross a budget before the next is refused.
- Every limit and budget is checked before the request, against what earlier requests used; tokens are counted after the response, so one request can cross a token limit or a budget before the next is refused. A refused request counts against nothing, and its `429` says in `Retry-After` when the limit frees.
- A limit set on an API key applies on top of its owner's, on a counter of its own.
- Redis can be a single node or a Redis Cluster.
- If Redis is unavailable, limits fail open by default. Setting `security.rate_limit_fail_closed` refuses requests instead.
- Budget alerts fire once per period at 50%, 80%, 95% and 100%.

Expand Down
4 changes: 3 additions & 1 deletion README.zh-CN.md
Original file line number Diff line number Diff line change
Expand Up @@ -92,7 +92,9 @@ cd web && pnpm install && pnpm dev
- 模型的「最大输出 token」在模型页设置,限制发往该模型的每个请求的 `max_tokens`,取代原来的输出长度护栏。

**限流与预算**
- 请求数限制在请求发出前检查;Token 限制与预算在响应返回后计入,因此单个请求可能越过预算,此后的请求才会被拒绝。
- 所有限流与预算都在请求发出前按此前的用量检查;Token 在响应返回后计入,因此单个请求可能越过 Token 限制或预算,此后的请求才会被拒绝。被拒绝的请求不计入任何限制,其 `429` 响应以 `Retry-After` 说明限制何时解除。
- 设置在 API Key 上的限制叠加在其所属用户的限制之上,单独计数。
- Redis 可以是单节点,也可以是 Redis Cluster。
- Redis 不可用时,限流默认放行;设置 `security.rate_limit_fail_closed` 后改为拒绝请求。
- 预算提醒在每个周期内于 50%、80%、95% 和 100% 各触发一次。

Expand Down
61 changes: 13 additions & 48 deletions crates/auth/src/rbac.rs
Original file line number Diff line number Diff line change
Expand Up @@ -373,61 +373,26 @@ pub async fn compute_user_surface_constraints(
Ok(apply_user_overrides(role_merged, override_constraints))
}

/// Like [`compute_user_surface_constraints`] but also folds in any
/// active rate-limit / budget overrides keyed on a specific
/// `api_key_id`. Per-key overrides REPLACE the user-derived value
/// in the matching `(surface, metric, window)` / `(surface, period)`
/// slot — the same merge semantics as user-scope overrides.
///
/// Use this in the gateway hot path (where the auth middleware
/// already knows the api_key id). Other callers (analytics
/// dashboards, admin "what does this user see today" views) keep
/// using the user-only variant since they have no key context.
pub async fn compute_effective_surface_constraints(
/// The limits attached to one API key — its lineage's active
/// `rate_limit_rules` / `budget_caps` rows — on their own. They are not
/// merged into the owner's: the gateway counts them on the lineage's
/// counters and checks them on top of the owner's limits, so a key's
/// limits can narrow what its owner may do through it but never widen
/// it. Keyed on the lineage so they survive rotation.
pub async fn compute_key_surface_constraints(
pool: &PgPool,
user_id: Uuid,
api_key_id: Uuid,
lineage_id: Uuid,
) -> Result<think_watch_common::limits::SurfaceConstraints, sqlx::Error> {
use think_watch_common::limits::{
self, BudgetSubject, RateLimitSubject, apply_user_overrides,
list_enabled_caps_for_subjects, list_enabled_rules_for_subjects, side_table_as_constraints,
BudgetSubject, RateLimitSubject, list_enabled_caps_for_subjects,
list_enabled_rules_for_subjects, side_table_as_constraints,
};

let user_merged = compute_user_surface_constraints(pool, user_id).await?;

// Per-key overrides are stored against the key's `lineage_id`
// (subject_kind = 'api_key_lineage') so they survive rotation.
// Resolve api_key_id → lineage_id once and bind every lookup on
// the lineage. A non-existent api_key_id (never happens at
// runtime — the auth middleware just authenticated this id —
// but defend anyway) maps to an empty override set.
let lineage_id: Option<uuid::Uuid> =
sqlx::query_scalar("SELECT lineage_id FROM api_keys WHERE id = $1")
.bind(api_key_id)
.fetch_optional(pool)
.await?;
let Some(lineage_id) = lineage_id else {
return Ok(user_merged);
};

// The api_key-side overrides are loaded separately so the
// existing `compute_user_*` helper stays a pure function of
// user_id (used by analytics + admin views). Loading both kinds
// here is two queries instead of one — that's bounded and the
// hot path already does enough DB round-trips that it doesn't
// dominate latency.
let key_rules =
let rules =
list_enabled_rules_for_subjects(pool, &[(RateLimitSubject::ApiKeyLineage, lineage_id)])
.await?;
let key_caps =
let caps =
list_enabled_caps_for_subjects(pool, &[(BudgetSubject::ApiKeyLineage, lineage_id)]).await?;
let key_overrides = side_table_as_constraints(&key_rules, &key_caps);

if key_overrides == limits::SurfaceConstraints::default() {
return Ok(user_merged);
}

Ok(apply_user_overrides(user_merged, key_overrides))
Ok(side_table_as_constraints(&rules, &caps))
}

/// Compute the set of permissions that are explicitly denied to `user_id`
Expand Down
1 change: 1 addition & 0 deletions crates/common/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ hmac = { workspace = true }
metrics = { workspace = true }
clickhouse = { workspace = true }
bytes = { workspace = true }
futures = { workspace = true }
url = { workspace = true }
# S3-compatible body offload (matches the same SigV4 + reqwest pattern
# the Bedrock provider uses — no aws-sdk-s3 dependency, so the build
Expand Down
1 change: 1 addition & 0 deletions crates/common/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -49,5 +49,6 @@ pub mod crypto; // AES-256-GCM envelope for secrets at rest
pub mod fixed_window;
pub mod json_secret; // `{"$enc": ...}` — a secret nested inside a JSONB column
pub mod pii; // BlobRedactor — at-rest body redaction shared by gateway + mcp-gateway
pub mod redis_keys; // pattern deletes that work on one node and on a Redis Cluster
pub mod tasks; // supervised_spawn — panic-isolated background tasks
pub mod validation;
4 changes: 2 additions & 2 deletions crates/common/src/lifecycle/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,8 +8,8 @@
//!
//! ```text
//! Raw<S>
//! → check_limits (rate-limit gate)
//! → check_budget (pre-call budget peek)
//! → check_budget (pre-call budget peek; charges nothing)
//! → check_limits (rate-limit gate; charges only when it passes)
//! → check_access (allowed_models / allowed_tools)
//! → surface-specific (cache lookup, breaker, credential resolution,
//! invoke_upstream → Invocation<S>)
Expand Down
90 changes: 42 additions & 48 deletions crates/common/src/lifecycle/stages/check_budget.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,26 +16,27 @@
//! (request allowed) unless the caller passes `fail_closed = true`
//! — matching the [`super::check_limits`] semantics.
//!
//! Ordering note: this stage runs AFTER `check_limits`. A request
//! that hits the per-window requests counter via `check_limits` and
//! then gets rejected here for being over-budget will have its
//! `requests` counter incremented anyway — rate-limit counters
//! measure "attempts", not "successes". Operators querying the
//! rate-limit metric will see budget-rejected requests reflected
//! there, by design.
//! Ordering note: this stage runs BEFORE `check_limits`, on the
//! [`Raw`] state. `check_limits` charges the `requests` counters, so
//! running it first would charge a request this stage then refuses.
//! Since this stage charges nothing, a request refused by either
//! charges nothing.

use chrono::Utc;
use fred::clients::Client;

use crate::audit::AuditLogger;
use crate::limits::{BudgetCap, budget};

use super::super::Surface;
use super::super::state::LimitsChecked;
use super::super::state::Raw;

/// Run a read-only spend check against every supplied cap. On
/// allow, returns the input state unchanged (passthrough — no
/// extra type narrowing). On deny, emits a `"budget_exceeded"`
/// audit row + short-circuits with `S::budget_exceeded_response`.
/// audit row + short-circuits with `S::budget_exceeded_response`,
/// naming the spent cap whose period ends last — the request can't
/// succeed before then — and the seconds until it does.
///
/// `fail_closed` controls behaviour on a Redis read error:
/// - `false` (default): bumps `lifecycle_budget_fail_open_total`
Expand All @@ -47,12 +48,12 @@ use super::super::state::LimitsChecked;
fields(trace_id = %state.trace_id, cap_count = caps.len()),
)]
pub async fn check_budget<S: Surface>(
state: LimitsChecked<S>,
state: Raw<S>,
caps: &[BudgetCap],
redis: &Client,
fail_closed: bool,
audit: &AuditLogger,
) -> Result<LimitsChecked<S>, S::Response> {
) -> Result<Raw<S>, S::Response> {
if caps.is_empty() {
return Ok(state);
}
Expand Down Expand Up @@ -84,27 +85,34 @@ pub async fn check_budget<S: Surface>(

// Walk caps + statuses in lockstep — `current_spend` preserves
// input order so the pairing is positional.
for (cap, status) in caps.iter().zip(statuses.iter()) {
if status.current >= status.limit {
let label = budget_label(cap);
metrics::counter!("lifecycle_budget_exceeded_total").increment(1);
tracing::warn!(
trace_id = %state.trace_id,
cap = %label,
current = status.current,
limit = status.limit,
"budget cap exceeded"
);
let entry = S::audit_entry(&state.identity, "budget_exceeded")
.trace_id(state.trace_id.clone())
.detail(serde_json::json!({
"limit": label,
"current": status.current,
"max": status.limit,
}));
audit.log(entry);
return Err(S::budget_exceeded_response(&label));
}
let now = Utc::now();
let spent = caps
.iter()
.zip(statuses.iter())
.filter(|(_, status)| status.current >= status.limit)
.map(|(cap, status)| (cap, status, budget::secs_until_period_end(cap.period, now)))
.max_by_key(|(_, _, wait)| *wait);
if let Some((cap, status, retry_after_secs)) = spent {
let label = budget_label(cap);
metrics::counter!("lifecycle_budget_exceeded_total").increment(1);
tracing::warn!(
trace_id = %state.trace_id,
cap = %label,
current = status.current,
limit = status.limit,
retry_after_secs,
"budget cap exceeded"
);
let entry = S::audit_entry(&state.identity, "budget_exceeded")
.trace_id(state.trace_id.clone())
.detail(serde_json::json!({
"limit": label,
"current": status.current,
"max": status.limit,
"retry_after_secs": retry_after_secs,
}));
audit.log(entry);
return Err(S::budget_exceeded_response(&label, retry_after_secs));
}

Ok(state)
Expand All @@ -125,7 +133,6 @@ fn budget_label(cap: &BudgetCap) -> String {

#[cfg(test)]
mod tests {
use super::super::super::state::{LimitCheckRecord, LimitsChecked};
use super::super::super::test_surface::{TestResponse, TestSurface, make_raw};
use super::*;
use crate::limits::{BudgetPeriod, BudgetSubject};
Expand All @@ -142,26 +149,13 @@ mod tests {
crate::audit::AuditLogger::test_drain()
}

fn make_limits_checked(user_id: Uuid) -> LimitsChecked<TestSurface> {
let raw = make_raw(user_id);
LimitsChecked {
identity: raw.identity,
trace_id: raw.trace_id,
started_at: raw.started_at,
client_ip: raw.client_ip,
limit_check: LimitCheckRecord {
currents: Vec::new(),
},
}
}

/// No caps configured → trivially pass through without touching
/// Redis. The disconnected dummy client confirms the function
/// never reached out.
#[tokio::test]
async fn passes_through_with_no_caps() {
let user_id = Uuid::new_v4();
let state = make_limits_checked(user_id);
let state = make_raw(user_id);
let trace_id = state.trace_id.clone();
let started_at = state.started_at;

Expand Down Expand Up @@ -194,7 +188,7 @@ mod tests {
/// short-circuit Response type to be propagated via `?`-style
/// match — same shape `check_limits` uses.
#[allow(dead_code)]
async fn type_state_compiles(state: LimitsChecked<TestSurface>) {
async fn type_state_compiles(state: Raw<TestSurface>) {
let _ = check_budget::<TestSurface>(state, &[], &dummy_redis(), false, &dummy_audit())
.await
.map_err(|r| match r {
Expand Down
Loading
Loading