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
19 changes: 16 additions & 3 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -94,11 +94,24 @@ MAX_CLUSTERS_FOR_EXPLAIN=10
# Seconds the worker sleeps between polling for new jobs when the queue is empty.
WORKER_POLL_INTERVAL=2

# Ingest backpressure (minimal stand-in for full API/LLM rate limiting).
# POST /v1/ingestions and POST /v1/ingestions/lines return 429 + Retry-After
# when pending worker jobs >= INGEST_QUEUE_MAX.
# Ingest queue backpressure. POST /v1/ingestions and POST /v1/ingestions/lines
# return 429 INGEST_QUEUE_FULL + Retry-After when pending worker jobs >= max.
# Tail ticks skip the same ceiling. Distinct from RATE_LIMITED (token bucket).
INGEST_QUEUE_MAX=100
INGEST_RETRY_AFTER_SECONDS=5

# API token-bucket rate limiting (ingest writes + query routes). In-memory per
# process. 0 rps = unlimited for that category. Defaults are high so local demo
# and tests are not 429'd. Identity is API key id, or "anonymous" when auth is off.
# RATELIMIT_ENABLED=true
# RATELIMIT_INGEST_RPS=100
# RATELIMIT_QUERY_RPS=100
# RATELIMIT_BURST=100
# RATELIMIT_RETRY_AFTER_SECONDS=1

# Max in-flight LLM provider calls (process-wide; independent of API concurrency).
# 0 = unlimited. Noop provider skips the wait.
# LLM_MAX_CONCURRENCY=4
# Max NDJSON lines accepted by POST /v1/ingestions/lines (over cap → 400).
INGEST_PUSH_MAX_LINES=5000
# Tail-mode poll cadence and consecutive-error auto-pause.
Expand Down
17 changes: 15 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -688,6 +688,14 @@ All settings are read from `.env`, environment variables, or CLI flags. Priority
| `OIDC_JWKS_URL` | _(empty)_ | Optional JWKS URL; default is `{issuer}/.well-known/openid-configuration` then `{issuer}/.well-known/jwks.json` |
| `API_BIND_HOST` | `127.0.0.1` | Host the startup guard treats as the bind address. `make api` sets this to `0.0.0.0` to match uvicorn |
| `AUTH_REFUSE_INSECURE_BIND` | `false` | If `true`, refuse to start when auth is off and the bind host is not loopback; if `false`, log a warning only |
| `RATELIMIT_ENABLED` | `true` | Token-bucket rate limiting on ingest writes and query routes. In-memory per process (not shared across workers) |
| `RATELIMIT_INGEST_RPS` | `100` | Steady-state tokens/sec for `POST /v1/ingestions*`. `0` = unlimited |
| `RATELIMIT_QUERY_RPS` | `100` | Steady-state tokens/sec for `/v1/query*`. `0` = unlimited |
| `RATELIMIT_BURST` | `100` | Bucket size (max tokens) per API key (or `anonymous` when auth is off) |
| `RATELIMIT_RETRY_AFTER_SECONDS` | `1` | `Retry-After` value on `429 RATE_LIMITED` |
| `INGEST_QUEUE_MAX` | `100` | Pending worker-job ceiling; over this, ingest returns `429 INGEST_QUEUE_FULL` |
| `INGEST_RETRY_AFTER_SECONDS` | `5` | `Retry-After` value on `429 INGEST_QUEUE_FULL` |
| `LLM_MAX_CONCURRENCY` | `4` | Max in-flight LLM provider calls process-wide. `0` = unlimited. Noop does not wait |
| `WEBHOOK_SECRET` | _(empty)_ | Fallback HMAC secret for ingest completion callbacks when auth is off or the API key has no per-key `whsec_` |
| `WEBHOOK_MAX_RETRIES` | `5` | Extra webhook POST attempts after the first (6 POSTs by default) on 5xx / 429 / connect errors |
| `WEBHOOK_TIMEOUT` | `10` | Per-attempt HTTP timeout in seconds for completion callbacks |
Expand All @@ -699,7 +707,7 @@ All settings are read from `.env`, environment variables, or CLI flags. Priority

raglogs is fully useful without any LLM. The `--no-llm` flag (or `LLM_PROVIDER=disabled`) activates deterministic template-based summaries.

When an LLM is configured, it receives only a small curated evidence packet — not raw logs. The prompt enforces fixed output structure, prohibits fabrication, and requires explicit uncertainty statements when evidence is insufficient.
When an LLM is configured, it receives only a small curated evidence packet — not raw logs. The prompt enforces fixed output structure, prohibits fabrication, and requires explicit uncertainty statements when evidence is insufficient. In-flight provider calls are capped by `LLM_MAX_CONCURRENCY` (CLI and API share the process semaphore; the noop provider does not block).

### OpenAI

Expand Down Expand Up @@ -1001,7 +1009,12 @@ curl -X POST http://localhost:8000/v1/ingestions/$ID:stop

`stop` is terminal. After `TAIL_ERROR_THRESHOLD` consecutive poll failures (default 5) a tail job auto-pauses; `/health` reports `tail_jobs.running` and `tail_jobs.paused`.

**Backpressure.** When pending worker jobs ≥ `INGEST_QUEUE_MAX` (default 100), `POST /v1/ingestions` and `POST /v1/ingestions/lines` return **429** with `Retry-After` (`INGEST_RETRY_AFTER_SECONDS`, default 5) and body `{"error_code":"INGEST_QUEUE_FULL","message":"..."}`. This is a queue-depth stand-in, not full API/LLM rate limiting.
**Backpressure and rate limiting.** Two independent 429s:

- **Queue depth.** When pending worker jobs ≥ `INGEST_QUEUE_MAX` (default 100), `POST /v1/ingestions` and `POST /v1/ingestions/lines` return **429** with `Retry-After` (`INGEST_RETRY_AFTER_SECONDS`, default 5) and body `{"error_code":"INGEST_QUEUE_FULL","message":"..."}`. Tail ticks skip the same ceiling.
- **API token bucket.** `POST /v1/ingestions*` (writes) and `/v1/query*` (plus unversioned aliases) are limited per API key (`request.state.auth_principal.key_id`, or a single `anonymous` bucket when `AUTH_ENABLED=false`). Exceeding the bucket returns **429** with `Retry-After` (`RATELIMIT_RETRY_AFTER_SECONDS`, default 1) and body `{"error_code":"RATE_LIMITED","message":"..."}`. Defaults (`RATELIMIT_INGEST_RPS` / `RATELIMIT_QUERY_RPS` / `RATELIMIT_BURST` = 100) are high enough for local demo and tests; `0` rps means unlimited for that category. `/health`, `/docs`, static UI, and `/config` are not limited. Buckets are in-memory per process.

LLM calls are separately capped by `LLM_MAX_CONCURRENCY` (default 4) so a burst of `explain` cannot fan out unbounded provider requests.

**Idempotency-Key.** `POST /v1/ingestions` (batch enqueue and tail create; also the deprecated `/ingestions` alias) honors an `Idempotency-Key` header (max 256 characters). A repeat **in the same isolation scope** within `INGEST_IDEMPOTENCY_TTL_SECONDS` (default 86400) returns the original **202** job — the same `worker_job_id` for batch, the same `ingestion_job_id` for tail — instead of starting a new one. Reusing another scope's key returns **409** `IDEMPOTENCY_SCOPE_CONFLICT`. Empty keys return **400**. GET routes ignore the header. `POST /v1/ingestions/lines` does not use the header; duplicate push/tail lines are handled by content dedup instead.

Expand Down
5 changes: 4 additions & 1 deletion src/api/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
from src.api.auth.middleware import AuthMiddleware
from src.api.auth.scope import ScopeResolutionError, scope_error_response
from src.api.deprecation import DeprecationHeaderMiddleware
from src.api.ratelimit import RateLimitMiddleware
from src.api.routes import ask, clusters, compare_windows, config, explain, health, ingestions, timeline, ui

_OPENAPI_DESCRIPTION = """Incident explanation API — ask your logs what happened.
Expand Down Expand Up @@ -60,7 +61,9 @@ def generate(route: APIRoute) -> str:
lifespan=lifespan,
)

# Last added middleware is outermost: deprecation headers apply even to auth errors.
# Last added middleware is outermost: deprecation headers apply even to auth
# errors. Rate limiting sits inside auth so request.state.auth_principal is set.
app.add_middleware(RateLimitMiddleware)
app.add_middleware(AuthMiddleware)
app.add_middleware(DeprecationHeaderMiddleware)

Expand Down
160 changes: 160 additions & 0 deletions src/api/ratelimit.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,160 @@
"""Per-API-key token-bucket rate limiting for ingest and query routes (G9).

Buckets are in-memory and per process — they are not shared across
uvicorn workers or hosts. Identity is ``request.state.auth_principal.key_id``
when present; otherwise ``\"anonymous\"`` (including AUTH_ENABLED=false).

``RATELIMIT_INGEST_RPS`` / ``RATELIMIT_QUERY_RPS`` of 0 means unlimited for
that category. ``RATELIMIT_ENABLED=false`` disables limiting entirely.
Health, docs, static UI, config, and key-admin routes are not limited.
"""

from __future__ import annotations

import re
import threading
import time
from typing import Awaitable, Callable, Literal

from starlette.middleware.base import BaseHTTPMiddleware
from starlette.requests import Request
from starlette.responses import JSONResponse, Response

ERROR_RATE_LIMITED = "RATE_LIMITED"

RateLimitKind = Literal["ingest", "query"]

_API_VERSION_PREFIX = re.compile(r"^/v\d+(?=/|$)")

_store_lock = threading.Lock()
_buckets: dict[tuple[str, str], "_TokenBucket"] = {}


class _TokenBucket:
"""Thread-safe token bucket. Rate/burst changes reset the fill level."""

def __init__(self) -> None:
self._lock = threading.Lock()
self._tokens: float = 0.0
self._updated: float = time.monotonic()
self._rate: float = 0.0
self._burst: float = 0.0

def consume(self, rate: float, burst: float, tokens: float = 1.0) -> bool:
if rate <= 0:
return True
capacity = max(float(burst), tokens)
with self._lock:
now = time.monotonic()
if rate != self._rate or capacity != self._burst:
self._rate = rate
self._burst = capacity
self._tokens = capacity
self._updated = now
else:
elapsed = now - self._updated
self._tokens = min(capacity, self._tokens + elapsed * rate)
self._updated = now
if self._tokens >= tokens:
self._tokens -= tokens
return True
return False


def reset_rate_limiter() -> None:
"""Drop all buckets. Used by tests so cases do not leak fill state."""
with _store_lock:
_buckets.clear()


def _normalize_path(path: str) -> str:
if not path:
return "/"
if len(path) > 1 and path.endswith("/"):
return path.rstrip("/")
return path


def rate_limit_kind(method: str, path: str) -> RateLimitKind | None:
"""Return ingest/query for limited routes, else None (exempt)."""
normalized = _API_VERSION_PREFIX.sub("", _normalize_path(path), count=1)
if not normalized:
normalized = "/"
verb = method.upper()
if normalized == "/ingestions" or normalized.startswith("/ingestions/"):
if verb == "POST":
return "ingest"
return None
if normalized == "/query" or normalized.startswith("/query/"):
return "query"
return None


def bucket_identity(request: Request) -> str:
"""Per-key identity, or anonymous when auth is off / principal missing."""
principal = getattr(request.state, "auth_principal", None)
if principal is None:
return "anonymous"
key_id = getattr(principal, "key_id", None)
if key_id:
return str(key_id)
subject = getattr(principal, "subject", None)
if subject:
return f"oidc:{subject}"
return "anonymous"


def allow_request(kind: RateLimitKind, identity: str, rate: float, burst: float) -> bool:
"""Consume one token from the (kind, identity) bucket. True if allowed."""
if rate <= 0:
return True
key = (kind, identity)
with _store_lock:
bucket = _buckets.get(key)
if bucket is None:
bucket = _TokenBucket()
_buckets[key] = bucket
return bucket.consume(rate, burst)


def _limited_response(retry_after: int) -> JSONResponse:
return JSONResponse(
status_code=429,
content={
"error_code": ERROR_RATE_LIMITED,
"message": "Rate limit exceeded; retry after the Retry-After delay.",
},
headers={"Retry-After": str(int(retry_after))},
)


class RateLimitMiddleware(BaseHTTPMiddleware):
"""Token-bucket limiter applied after auth so key_id is available."""

async def dispatch(
self,
request: Request,
call_next: Callable[[Request], Awaitable[Response]],
) -> Response:
from src.config import get_settings

settings = get_settings()
if not settings.ratelimit_enabled:
return await call_next(request)

kind = rate_limit_kind(request.method, request.url.path)
if kind is None:
return await call_next(request)

rate = (
settings.ratelimit_ingest_rps
if kind == "ingest"
else settings.ratelimit_query_rps
)
if rate <= 0:
return await call_next(request)

identity = bucket_identity(request)
if not allow_request(kind, identity, rate, settings.ratelimit_burst):
return _limited_response(settings.ratelimit_retry_after_seconds)
return await call_next(request)
17 changes: 16 additions & 1 deletion src/config/settings.py
Original file line number Diff line number Diff line change
Expand Up @@ -50,13 +50,28 @@ class Settings(BaseSettings):
# Worker
worker_poll_interval: int = 2 # seconds between idle polls

# Ingest backpressure (minimal G9 stand-in) and push/tail (G4)
# Ingest backpressure (G9) and push/tail (G4). Tail ticks skip when the
# pending worker queue is at ingest_queue_max.
ingest_queue_max: int = 100
ingest_retry_after_seconds: int = 5
ingest_push_max_lines: int = 5000
tail_poll_interval: int = 30
tail_error_threshold: int = 5

# API token-bucket rate limiting (G9). In-memory per process — not shared
# across workers. Defaults are high so local demo and TestClient suites
# are not 429'd. 0 rps = unlimited for that category. Identity is
# auth_principal.key_id, else "anonymous" when AUTH_ENABLED=false.
ratelimit_enabled: bool = True
ratelimit_ingest_rps: float = 100.0
ratelimit_query_rps: float = 100.0
ratelimit_burst: float = 100.0
ratelimit_retry_after_seconds: int = 1

# Process-wide cap on in-flight LLM provider calls (G9). Independent of
# API concurrency. 0 = unlimited. Noop provider skips the wait.
llm_max_concurrency: int = 4

# HMAC-signed ingest completion webhooks (G5). WEBHOOK_SECRET is the
# fallback when AUTH_ENABLED=false or the API key has no per-key secret.
webhook_secret: str = ""
Expand Down
13 changes: 6 additions & 7 deletions src/core/explain/summarizer.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@
from src.core.explain.confidence import compute_confidence
from src.core.explain.evidence import EvidencePacket, assemble_evidence
from src.core.explain.templates import render_insufficient_evidence, render_text_summary
from src.core.llm.provider import NoopLLMProvider, build_llm_provider
from src.core.llm.provider import build_llm_provider
from src.db.models import DEFAULT_LOG_SCOPE, IngestionJob
from src.db.scope_filter import filter_ingestion_jobs_by_scope

Expand Down Expand Up @@ -116,12 +116,11 @@ def explain_window(
if not no_llm and settings.llm_provider != "disabled":
try:
llm = build_llm_provider(settings)
if not isinstance(llm, NoopLLMProvider):
evidence_dict = _packet_to_dict(packet)
llm_text = llm.generate_summary(evidence_dict)
if llm_text:
summary_text = llm_text
mode = "llm"
evidence_dict = _packet_to_dict(packet)
llm_text = llm.generate_summary(evidence_dict)
if llm_text:
summary_text = llm_text
mode = "llm"
except Exception:
# Degrade gracefully
pass
Expand Down
7 changes: 6 additions & 1 deletion src/core/ingestion/backpressure.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,9 @@
"""Ingest worker-queue backpressure (minimal G9 stand-in)."""
"""Ingest worker-queue backpressure (G9).

When pending WorkerJobs >= INGEST_QUEUE_MAX, POST /ingestions and
POST /ingestions/lines return 429 INGEST_QUEUE_FULL. Tail ticks skip
the same ceiling. Distinct from API RATE_LIMITED (token bucket).
"""

from sqlalchemy import func, select
from sqlalchemy.orm import Session
Expand Down
Loading
Loading