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
10 changes: 10 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,16 @@ INGEST_RETRY_AFTER_SECONDS=5
# Max in-flight LLM provider calls (process-wide; independent of API concurrency).
# 0 = unlimited. Noop provider skips the wait.
# LLM_MAX_CONCURRENCY=4

# LLM resilience (timeout, retries, token ceiling, circuit breaker). Unused when
# LLM_PROVIDER=disabled (noop). On timeout/error/open breaker, explain/ask fall
# back to deterministic templates and set llm.fell_back=true; the request succeeds.
# LLM_TIMEOUT=30
# LLM_MAX_RETRIES=2 # extra attempts after the first (3 total)
# LLM_MAX_TOKENS=600 # completion cap (OpenAI max_tokens / Ollama num_predict)
# LLM_MAX_INPUT_TOKENS=0 # 0 = derive from LLM_MAX_TOKENS (chars/4 estimate)
# LLM_BREAKER_THRESHOLD=5 # consecutive failures before the breaker opens
# LLM_BREAKER_COOLDOWN_SECONDS=60
# 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
12 changes: 10 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -696,6 +696,12 @@ All settings are read from `.env`, environment variables, or CLI flags. Priority
| `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 |
| `LLM_TIMEOUT` | `30` | Per-attempt HTTP timeout in seconds for OpenAI/Ollama calls |
| `LLM_MAX_RETRIES` | `2` | Extra attempts after the first (3 total) with jittered exponential backoff |
| `LLM_MAX_TOKENS` | `600` | Completion cap (`max_tokens` / Ollama `num_predict`) |
| `LLM_MAX_INPUT_TOKENS` | `0` | Estimated input-token budget (`chars/4`). `0` derives from `LLM_MAX_TOKENS`. Over budget: trim evidence (respecting `MAX_EVIDENCE_ITEMS`) or fall back |
| `LLM_BREAKER_THRESHOLD` | `5` | Consecutive LLM failures before the process-local breaker opens |
| `LLM_BREAKER_COOLDOWN_SECONDS` | `60` | Seconds the breaker stays open before a half-open probe |
| `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 @@ -709,6 +715,8 @@ raglogs is fully useful without any LLM. The `--no-llm` flag (or `LLM_PROVIDER=d

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).

Every provider call has a timeout (`LLM_TIMEOUT`) and bounded jittered retries (`LLM_MAX_RETRIES`). On timeout, HTTP error, exhausted retries, or an over-budget evidence payload, explain/ask **fall back to the same deterministic templates** and set `llm.fell_back=true` — the request still succeeds. After `LLM_BREAKER_THRESHOLD` consecutive failures a process-local circuit breaker opens for `LLM_BREAKER_COOLDOWN_SECONDS`; while open, raglogs skips the provider entirely and serves templates. `GET /health` exposes `llm_breaker: {state, consecutive_failures, cooldown_remaining_seconds}` (`closed` / `open` / `half_open`). An open breaker marks `status` as `degraded` but still returns HTTP 200 so probes do not fail. Fallback never invents: it only renders the curated evidence packet.

### OpenAI

```env
Expand Down Expand Up @@ -970,7 +978,7 @@ If auth is disabled and the process binds a non-loopback address (`0.0.0.0`, `::

| Method | Endpoint | Description |
|---|---|---|
| `GET` | `/health` | Service and DB health check (unversioned). Includes `tail_jobs: {running, paused}`. |
| `GET` | `/health` | Service and DB health check (unversioned). Includes `tail_jobs: {running, paused}` and `llm_breaker` (`closed` / `open` / `half_open`). Open breaker → `status: degraded`, still HTTP 200. |
| `POST` | `/v1/ingestions` | Enqueue a batch ingest job (`adapter`: `file`, `cloudwatch`, `datadog`, `loki`, or `k8s`). Set `"mode": "tail"` for pull adapters to start a long-lived tail job. |
| `POST` | `/v1/ingestions/lines` | Push NDJSON of raw or pre-parsed log lines (sync persist) |
| `POST` | `/v1/ingestions/{id}:pause` | Pause a tail job |
Expand Down Expand Up @@ -1014,7 +1022,7 @@ curl -X POST http://localhost:8000/v1/ingestions/$ID:stop
- **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.
LLM calls are separately capped by `LLM_MAX_CONCURRENCY` (default 4) so a burst of `explain` cannot fan out unbounded provider requests. Timeouts, retries, automatic template fallback (`llm.fell_back`), and the process-local circuit breaker are described under [LLM integration](#llm-integration).

**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
26 changes: 24 additions & 2 deletions src/api/routes/health.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,12 +12,19 @@
_adapter_health_cache: dict[str, tuple[float, str]] = {}


class LlmBreakerHealth(BaseModel):
state: str # closed | open | half_open
consecutive_failures: int
cooldown_remaining_seconds: float


class HealthResponse(BaseModel):
status: str # ok | degraded
db: str # connected | disconnected
worker_queue_depth: Optional[int] # pending worker jobs; None if DB unreachable
adapters: dict[str, str] # adapter name -> "ok" | "unavailable: <reason>"
tail_jobs: Optional[dict[str, int]] = None # {running, paused}; None if DB unreachable
llm_breaker: Optional[LlmBreakerHealth] = None


def _adapter_health() -> dict[str, str]:
Expand Down Expand Up @@ -52,12 +59,25 @@ def _adapter_health() -> dict[str, str]:
return statuses


def _llm_breaker_health() -> LlmBreakerHealth:
from src.core.llm.resilience import breaker_health

snap = breaker_health()
return LlmBreakerHealth(
state=str(snap["state"]),
consecutive_failures=int(snap["consecutive_failures"]),
cooldown_remaining_seconds=float(snap["cooldown_remaining_seconds"]),
)


@router.get("/health", response_model=HealthResponse)
def health_check():
def health_check() -> HealthResponse:
from src.db.session import check_connection, get_db

db_ok = check_connection()
adapters = _adapter_health()
llm_breaker = _llm_breaker_health()
breaker_open = llm_breaker.state == "open"

if not db_ok:
return HealthResponse(
Expand All @@ -66,6 +86,7 @@ def health_check():
worker_queue_depth=None,
adapters=adapters,
tail_jobs=None,
llm_breaker=llm_breaker,
)

try:
Expand All @@ -83,9 +104,10 @@ def health_check():
tail_jobs = None

return HealthResponse(
status="ok" if db_ok else "degraded",
status="degraded" if breaker_open else "ok",
db="connected",
worker_queue_depth=depth,
adapters=adapters,
tail_jobs=tail_jobs,
llm_breaker=llm_breaker,
)
10 changes: 10 additions & 0 deletions src/config/settings.py
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,16 @@ class Settings(BaseSettings):
# API concurrency. 0 = unlimited. Noop provider skips the wait.
llm_max_concurrency: int = 4

# LLM resilience (G10). Timeouts/retries wrap every provider HTTP call.
# On failure the pipeline falls back to deterministic templates.
# 0 input-token budget means "derive from LLM_MAX_TOKENS".
llm_timeout: float = 30.0
llm_max_retries: int = 2
llm_max_tokens: int = 600
llm_max_input_tokens: int = 0
llm_breaker_threshold: int = 5
llm_breaker_cooldown_seconds: float = 60.0

# 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
8 changes: 6 additions & 2 deletions src/core/explain/summarizer.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
from datetime import datetime
from typing import Optional

import structlog
from sqlalchemy.orm import Session

from src.config import get_settings
Expand All @@ -14,6 +15,8 @@
from src.db.models import DEFAULT_LOG_SCOPE, IngestionJob
from src.db.scope_filter import filter_ingestion_jobs_by_scope

log = structlog.get_logger()


@dataclass
class ExplainResult:
Expand Down Expand Up @@ -122,8 +125,9 @@ def explain_window(
summary_text = llm_text
mode = "llm"
except Exception:
# Degrade gracefully
pass
# Timeout, retries exhausted, open breaker, or budget: keep mode
# "rules" so llm.fell_back is true when an LLM was requested.
log.warning("llm_explain_failed", exc_info=True)

if not summary_text:
summary_text = render_text_summary(packet, confidence)
Expand Down
97 changes: 76 additions & 21 deletions src/core/llm/provider.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,8 @@
from typing import Any, Protocol, runtime_checkable

import httpx
from tenacity import retry, stop_after_attempt, wait_exponential

from src.core.llm.resilience import invoke_llm, prepare_llm_packet


@runtime_checkable
Expand Down Expand Up @@ -39,6 +40,18 @@ def generate_summary(self, evidence_packet: dict) -> str:
If evidence is insufficient, say so clearly. Keep the entire output under 300 words. No markdown formatting."""


def _llm_timeout() -> float:
from src.config import get_settings

return float(get_settings().llm_timeout)


def _llm_max_tokens() -> int:
from src.config import get_settings

return int(get_settings().llm_max_tokens)


class OpenAILLMProvider:
def __init__(
self,
Expand All @@ -50,12 +63,9 @@ def __init__(
self.model = model
self.base_url = base_url.rstrip("/")

@retry(stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, min=2, max=10))
def generate_summary(self, evidence_packet: dict) -> str:
payload = json.dumps(evidence_packet, default=str, indent=2)
user_message = f"Analyze this incident evidence and produce a summary:\n\n{payload}"

with httpx.Client(timeout=60) as client:
def complete(self, system_prompt: str, user_message: str) -> str:
"""Single OpenAI chat-completions attempt (timeout/max_tokens from settings)."""
with httpx.Client(timeout=_llm_timeout()) as client:
response = client.post(
f"{self.base_url}/chat/completions",
headers={
Expand All @@ -64,10 +74,10 @@ def generate_summary(self, evidence_packet: dict) -> str:
},
json={
"model": self.model,
"max_tokens": 600,
"max_tokens": _llm_max_tokens(),
"temperature": 0,
"messages": [
{"role": "system", "content": SYSTEM_PROMPT},
{"role": "system", "content": system_prompt},
{"role": "user", "content": user_message},
],
},
Expand All @@ -76,6 +86,11 @@ def generate_summary(self, evidence_packet: dict) -> str:
data = response.json()
return data["choices"][0]["message"]["content"].strip()

def generate_summary(self, evidence_packet: dict) -> str:
payload = json.dumps(evidence_packet, default=str, indent=2)
user_message = f"Analyze this incident evidence and produce a summary:\n\n{payload}"
return self.complete(SYSTEM_PROMPT, user_message)


class OllamaLLMProvider:
def __init__(
Expand All @@ -86,25 +101,55 @@ def __init__(
self.model = model
self.base_url = base_url.rstrip("/")

@retry(stop=stop_after_attempt(2), wait=wait_exponential(multiplier=1, min=1, max=5))
def generate_summary(self, evidence_packet: dict) -> str:
payload = json.dumps(evidence_packet, default=str, indent=2)
prompt = f"{SYSTEM_PROMPT}\n\nIncident evidence:\n{payload}\n\nSummary:"

with httpx.Client(timeout=120) as client:
def complete(self, system_prompt: str, user_message: str) -> str:
"""Single Ollama /api/generate attempt (timeout/num_predict from settings)."""
prompt = f"{system_prompt}\n\n{user_message}"
with httpx.Client(timeout=_llm_timeout()) as client:
response = client.post(
f"{self.base_url}/api/generate",
json={
"model": self.model,
"prompt": prompt,
"stream": False,
"options": {"temperature": 0},
"options": {
"temperature": 0,
"num_predict": _llm_max_tokens(),
},
},
)
response.raise_for_status()
data = response.json()
return data.get("response", "").strip()

def generate_summary(self, evidence_packet: dict) -> str:
payload = json.dumps(evidence_packet, default=str, indent=2)
user_message = f"Incident evidence:\n{payload}\n\nSummary:"
return self.complete(SYSTEM_PROMPT, user_message)


class ResilientLLMProvider:
"""G10 breaker + token budget + retries around an inner provider.

Lives *inside* ``CappedLLMProvider`` so in-flight slots cover the whole
retry sequence, but the breaker check still skips HTTP when open.
"""

def __init__(self, inner: LLMProvider) -> None:
self.inner = inner

def generate_summary(self, evidence_packet: dict) -> str:
from src.config import get_settings

prepared = prepare_llm_packet(evidence_packet, get_settings())
return invoke_llm(lambda: self.inner.generate_summary(prepared))

def complete(self, system_prompt: str, user_message: str) -> str:
inner = self.inner
complete = getattr(inner, "complete", None)
if complete is None:
return ""
return invoke_llm(lambda: complete(system_prompt, user_message))


_llm_sem_lock = threading.Lock()
_llm_semaphore: threading.Semaphore | None = None
Expand Down Expand Up @@ -153,21 +198,24 @@ class CappedLLMProvider:

Noop inner providers skip the wait so deterministic mode never blocks,
but still go through this entrypoint (CLI and API share the semaphore).

Stack (outer → inner): CappedLLMProvider → ResilientLLMProvider → OpenAI/Ollama.
Noop skips ResilientLLMProvider entirely.
"""

def __init__(self, inner: LLMProvider) -> None:
self.inner = inner

def generate_summary(self, evidence_packet: dict) -> str:
skip = isinstance(self.inner, NoopLLMProvider)
skip = isinstance(unwrap_llm_provider(self), NoopLLMProvider)
with llm_concurrency_slot(skip=skip):
return self.inner.generate_summary(evidence_packet)


def unwrap_llm_provider(provider: LLMProvider) -> LLMProvider:
"""Return the inner provider if ``provider`` is concurrency-capped."""
"""Return the inner provider if ``provider`` is concurrency-capped or resilient."""
inner: LLMProvider = provider
while isinstance(inner, CappedLLMProvider):
while isinstance(inner, (CappedLLMProvider, ResilientLLMProvider)):
inner = inner.inner
return inner

Expand All @@ -190,5 +238,12 @@ def _build_inner_llm_provider(settings: Any) -> LLMProvider:


def build_llm_provider(settings: Any) -> LLMProvider:
"""Factory: build the configured LLM provider (concurrency-capped)."""
return CappedLLMProvider(_build_inner_llm_provider(settings))
"""Factory: build the configured LLM provider (resilient + concurrency-capped).

Order: CappedLLMProvider (G9) wraps ResilientLLMProvider (G10) wraps the
HTTP provider. Noop skips the resilience wrapper so disabled mode is unchanged.
"""
inner = _build_inner_llm_provider(settings)
if not isinstance(inner, NoopLLMProvider):
inner = ResilientLLMProvider(inner)
return CappedLLMProvider(inner)
Loading
Loading