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 @@ -146,6 +146,16 @@ TAIL_ERROR_THRESHOLD=5
# Repeat requests with the same key within this window return the original job.
# INGEST_IDEMPOTENCY_TTL_SECONDS=86400

# Data retention (G13). Duration strings (30d, 180d, 24h). 0, empty, or off
# disables that tier. Per-scope overrides: INSERT INTO scope_retention
# (scope, raw_interval, summary_interval) VALUES ('incident:prod', '7d', '90d');
# NULL column = use the env default below. The worker enqueues a purge job on
# idle poll every PURGE_INTERVAL_SECONDS (or run `raglogs purge` / `--dry-run`).
RETENTION_RAW=30d
RETENTION_SUMMARY=180d
PURGE_INTERVAL_SECONDS=3600
# PURGE_CHUNK_SIZE=1000

# Self-observability (structured JSON logs, Prometheus /metrics, optional OTLP).
# LOG_FORMAT=json|console (API defaults to json; CLI uses console on a TTY)
# OTEL_SDK_DISABLED=true skips the OpenTelemetry SDK (request ids still generated).
Expand Down
20 changes: 20 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -631,6 +631,20 @@ raglogs config llm_provider

---

### `raglogs purge`

Expire raw log rows (and cascaded log-line embeddings / cluster membership) while keeping cluster summaries and `cluster_embeddings` so `POST /v1/query/similar` still works. After `RETENTION_SUMMARY` those summaries expire too. Per-scope TTLs override env defaults via the `scope_retention` table; missing override → `RETENTION_RAW` / `RETENTION_SUMMARY`. `0`, empty, or `off` skips that tier.

The background worker (`raglogs worker`) also enqueues a purge job on idle poll about every `PURGE_INTERVAL_SECONDS` (default 3600). Purge uses `SELECT FOR UPDATE SKIP LOCKED` like ingest and deletes in `created_at`-ordered chunks (time-in-store) so ingest is not starved. Raw expiry uses `log_entries.created_at`, not the event timestamp, so historical dumps are not wiped on ingest.

```bash
raglogs purge # every scope with data
raglogs purge --scope default
raglogs purge --dry-run # count only
```

---

### `raglogs keys`

Mint, list, and revoke HTTP API keys. The API bearer token is printed **once** at create time and is never stored or logged — only an argon2 hash and a short prefix are kept. A separate **webhook signing secret** (`whsec_…`) is also printed once; it is stored server-side so ingest completion callbacks can be HMAC-signed. It is not the bearer token. `raglogs keys list` shows `whsec_****` when a signing secret exists (legacy keys minted before this column fall back to `WEBHOOK_SECRET`).
Expand Down Expand Up @@ -707,6 +721,10 @@ All settings are read from `.env`, environment variables, or CLI flags. Priority
| `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 |
| `INGEST_IDEMPOTENCY_TTL_SECONDS` | `86400` | How long `Idempotency-Key` on `POST /v1/ingestions` is remembered (batch enqueue and tail create) |
| `RETENTION_RAW` | `30d` | How long to keep raw `log_entries` measured by `created_at` (time-in-store; cascaded `log_embeddings` / `cluster_members`). `0` / empty / `off` = never purge. Per-scope override: `scope_retention.raw_interval` |
| `RETENTION_SUMMARY` | `180d` | How long to keep cluster summaries + `cluster_embeddings` after which similar-incident recall for that scope expires. Same `0` / `off` disable |
| `PURGE_INTERVAL_SECONDS` | `3600` | Idle worker poll interval between automatic purge jobs. `0` disables scheduled purge (`raglogs purge` still works) |
| `PURGE_CHUNK_SIZE` | `1000` | Max rows deleted per table per scope per purge job (worker does one chunk; CLI drains) |
| `LOG_FORMAT` | `json` | Structured log renderer: `json` or `console`. API uses this; CLI switches to console on a TTY |
| `OTEL_SDK_DISABLED` | `false` | Skip the OpenTelemetry SDK. Request ids are still generated |
| `OTEL_EXPORTER_OTLP_ENDPOINT` | _(empty)_ | Optional OTLP HTTP traces endpoint. Empty = no exporter (no collector required) |
Expand Down Expand Up @@ -1055,6 +1073,7 @@ LLM calls are separately capped by `LLM_MAX_CONCURRENCY` (default 4) so a burst
| `raglogs_llm_estimated_tokens_total` | counter | Estimated input tokens (UTF-8 chars/4, not USD) |
| `raglogs_llm_breaker_state` | gauge | `0` closed, `1` half_open, `2` open |
| `raglogs_worker_queue_depth` | gauge | Pending worker jobs (omitted until a successful scrape; left stale if DB fails) |
| `raglogs_purge_rows_total` | counter (`kind`) | Rows reclaimed by retention purge: `raw`, `summary`, or `embedding` |

OpenTelemetry spans cover ingest → cluster → explain (and the HTTP request). The default exporter is none; set `OTEL_EXPORTER_OTLP_ENDPOINT` to export, or `OTEL_SDK_DISABLED=true` to skip the SDK. Trace ids stay on response headers so `/v1/query/*` JSON (`schema_version` 1.0) is unchanged.

Expand Down Expand Up @@ -1270,6 +1289,7 @@ raglogs/
│ │ ├── normalization/ Message normalization, fingerprinting, trigger patterns
│ │ ├── parsing/ JSON and text parsers, field extractors, timestamps
│ │ ├── retrieval/ Semantic + keyword question answering
│ │ ├── retention/ Per-scope TTL, scheduled purge of raw vs summary tiers
│ │ └── timeline/ Causal timeline reconstruction
│ ├── db/ SQLAlchemy models, session management
│ └── utils/ Time window parsing, hashing helpers
Expand Down
149 changes: 149 additions & 0 deletions migrations/versions/0010_retention.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,149 @@
"""retention: CASCADE FKs, scope_retention, scope on cluster_runs/explanations

Revision ID: 0010_retention
Revises: 0009_cluster_embeddings
Create Date: 2026-08-17 00:00:00.000000

Raw purge deletes ``log_entries``; ``log_embeddings`` and ``cluster_members``
follow via ON DELETE CASCADE. Cluster summaries/embeddings stay until the
summary TTL. Per-scope overrides live in ``scope_retention``.
"""

from typing import Optional, Sequence, Union

import sqlalchemy as sa
from alembic import op

revision: str = "0010_retention"
down_revision: Union[str, None] = "0009_cluster_embeddings"
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None

_CASCADE_FKS: tuple[tuple[str, str, str, str], ...] = (
("log_embeddings", "log_entry_id", "log_entries", "id"),
("cluster_members", "log_entry_id", "log_entries", "id"),
("cluster_members", "cluster_id", "clusters", "id"),
("clusters", "cluster_run_id", "cluster_runs", "id"),
)


def _fk_name(table: str, column: str) -> Optional[str]:
bind = op.get_bind()
inspector = sa.inspect(bind)
for fk in inspector.get_foreign_keys(table):
if fk.get("constrained_columns") == [column]:
name = fk.get("name")
if isinstance(name, str) and name:
return name
return f"{table}_{column}_fkey"


def _recreate_fk(
table: str,
column: str,
referent_table: str,
referent_column: str,
*,
ondelete: Optional[str],
) -> None:
name = _fk_name(table, column)
op.drop_constraint(name, table, type_="foreignkey")
op.create_foreign_key(
name,
table,
referent_table,
[column],
[referent_column],
ondelete=ondelete,
)


def upgrade() -> None:
for table, column, referent_table, referent_column in _CASCADE_FKS:
_recreate_fk(table, column, referent_table, referent_column, ondelete="CASCADE")

op.add_column(
"cluster_runs",
sa.Column("scope", sa.String(255), nullable=False, server_default="default"),
)
op.create_index("ix_cluster_runs_scope", "cluster_runs", ["scope"])
op.add_column(
"explanations",
sa.Column("scope", sa.String(255), nullable=False, server_default="default"),
)
op.create_index("ix_explanations_scope", "explanations", ["scope"])

# Majority scope from remaining cluster membership (G8). Rows with no
# members stay at the server default ``default``.
op.execute(
sa.text(
"""
UPDATE cluster_runs AS cr
SET scope = src.scope
FROM (
SELECT DISTINCT ON (counted.cluster_run_id)
counted.cluster_run_id,
counted.scope
FROM (
SELECT c.cluster_run_id, le.scope, COUNT(*) AS n
FROM clusters c
JOIN cluster_members cm ON cm.cluster_id = c.id
JOIN log_entries le ON le.id = cm.log_entry_id
GROUP BY c.cluster_run_id, le.scope
) AS counted
ORDER BY counted.cluster_run_id, counted.n DESC, counted.scope
) AS src
WHERE cr.id = src.cluster_run_id
"""
)
)
# Best-effort: explanations inherit the majority log_entries.scope in the
# cached window (service/env filters when present).
op.execute(
sa.text(
"""
UPDATE explanations AS e
SET scope = src.scope
FROM (
SELECT DISTINCT ON (counted.explanation_id)
counted.explanation_id,
counted.scope
FROM (
SELECT e2.id AS explanation_id, le.scope, COUNT(*) AS n
FROM explanations e2
JOIN log_entries le
ON le.timestamp >= e2.window_start
AND le.timestamp <= e2.window_end
AND (e2.service_filter IS NULL OR le.service = e2.service_filter)
AND (
e2.environment_filter IS NULL
OR le.environment = e2.environment_filter
)
GROUP BY e2.id, le.scope
) AS counted
ORDER BY counted.explanation_id, counted.n DESC, counted.scope
) AS src
WHERE e.id = src.explanation_id
"""
)
)

op.create_table(
"scope_retention",
sa.Column("scope", sa.String(255), primary_key=True),
sa.Column("raw_interval", sa.String(32), nullable=True),
sa.Column("summary_interval", sa.String(32), nullable=True),
sa.Column(
"updated_at", sa.DateTime(timezone=True), server_default=sa.func.now()
),
)


def downgrade() -> None:
op.drop_table("scope_retention")
op.drop_index("ix_explanations_scope", table_name="explanations")
op.drop_column("explanations", "scope")
op.drop_index("ix_cluster_runs_scope", table_name="cluster_runs")
op.drop_column("cluster_runs", "scope")
for table, column, referent_table, referent_column in reversed(_CASCADE_FKS):
_recreate_fk(table, column, referent_table, referent_column, ondelete=None)
4 changes: 4 additions & 0 deletions src/api/routes/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
@router.get("")
def get_config():
from src.config import get_settings

settings = get_settings()
return {
"llm_provider": settings.llm_provider,
Expand All @@ -15,4 +16,7 @@ def get_config():
"default_baseline_window": settings.default_baseline_window,
"max_clusters_for_explain": settings.max_clusters_for_explain,
"max_evidence_items": settings.max_evidence_items,
"retention_raw": settings.retention_raw,
"retention_summary": settings.retention_summary,
"purge_interval_seconds": settings.purge_interval_seconds,
}
18 changes: 15 additions & 3 deletions src/api/routes/explain.py
Original file line number Diff line number Diff line change
Expand Up @@ -158,10 +158,20 @@ def _maybe_add_markdown(
)


def _save_to_cache(db, cache_hash: str, window_start: datetime, window_end: datetime,
service: Optional[str], env: Optional[str], result: dict, confidence: str, mode: str):
def _save_to_cache(
db,
cache_hash: str,
window_start: datetime,
window_end: datetime,
service: Optional[str],
env: Optional[str],
result: dict,
confidence: str,
mode: str,
scope: str = "default",
):
"""Persist an explanation result to the cache."""
from src.db.models import Explanation
from src.db.models import DEFAULT_LOG_SCOPE, Explanation

row = Explanation(
id=uuid.uuid4(),
Expand All @@ -173,6 +183,7 @@ def _save_to_cache(db, cache_hash: str, window_start: datetime, window_end: date
prompt_hash=cache_hash,
result_json=result,
confidence=confidence,
scope=scope or DEFAULT_LOG_SCOPE,
)
db.add(row)
db.flush()
Expand Down Expand Up @@ -259,6 +270,7 @@ def explain_endpoint(request: ExplainRequest, http_request: Request) -> ExplainR
_save_to_cache(
db, cache_hash, window_start, window_end,
request.service, request.env, cache_payload, result.confidence, result.mode,
scope=scope,
)

return _maybe_add_markdown(body, result, request)
Expand Down
11 changes: 9 additions & 2 deletions src/cli/commands/config_cmd.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,9 @@ def config_cmd(

if key is None:
# Show all
table = Table(title="raglogs configuration", show_header=True, header_style="bold cyan")
table = Table(
title="raglogs configuration", show_header=True, header_style="bold cyan"
)
table.add_column("Key")
table.add_column("Value")

Expand All @@ -28,13 +30,18 @@ def config_cmd(
"LLM_MODEL": settings.llm_model,
"EMBEDDINGS_PROVIDER": settings.embeddings_provider,
"EMBEDDINGS_MODEL": settings.embeddings_model,
"CLUSTER_MERGE_SIMILARITY_THRESHOLD": str(settings.cluster_merge_similarity_threshold),
"CLUSTER_MERGE_SIMILARITY_THRESHOLD": str(
settings.cluster_merge_similarity_threshold
),
"CLUSTER_MERGE_MIN_COUNT": str(settings.cluster_merge_min_count),
"ASK_SEMANTIC_TOP_K": str(settings.ask_semantic_top_k),
"ASK_SEMANTIC_MIN_SIMILARITY": str(settings.ask_semantic_min_similarity),
"DEFAULT_BASELINE_WINDOW": settings.default_baseline_window,
"MAX_CLUSTERS_FOR_EXPLAIN": str(settings.max_clusters_for_explain),
"MAX_EVIDENCE_ITEMS": str(settings.max_evidence_items),
"RETENTION_RAW": settings.retention_raw,
"RETENTION_SUMMARY": settings.retention_summary,
"PURGE_INTERVAL_SECONDS": str(settings.purge_interval_seconds),
}

for k, v in config_items.items():
Expand Down
59 changes: 59 additions & 0 deletions src/cli/commands/purge.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,59 @@
"""raglogs purge — expire raw logs and (later) cluster summaries."""

from __future__ import annotations

from typing import Optional

import typer
from rich.console import Console
from rich.table import Table

console = Console()


def purge_cmd(
scope: Optional[str] = typer.Option(
None,
"--scope",
help="Limit to one isolation scope (default: every scope with data)",
),
dry_run: bool = typer.Option(
False,
"--dry-run",
help="Count expired rows without deleting them",
),
) -> None:
"""Delete expired raw logs, then expired cluster summaries / embeddings.

Raw rows older than RETENTION_RAW (per-scope override in scope_retention)
are removed based on created_at (time-in-store); cluster_embeddings stay so
similar-incident search still works. After RETENTION_SUMMARY those summaries
go too. 0 / empty / off skips that tier. The worker also runs this on an
idle poll. --dry-run counts expired rows without deleting.
"""
from src.core.retention.purge import run_purge
from src.db.session import get_db

with get_db() as db:
counts = run_purge(
db,
scope=scope,
dry_run=dry_run,
max_chunks=None,
)

table = Table(
title="raglogs purge" + (" (dry-run)" if dry_run else ""),
show_header=True,
header_style="bold cyan",
)
table.add_column("Kind")
table.add_column("Rows", justify="right")
table.add_row("raw", f"{counts.raw:,}")
table.add_row("summary", f"{counts.summary:,}")
table.add_row("embedding", f"{counts.embedding:,}")
console.print(table)
if counts.scopes:
console.print(f"[dim]scopes:[/dim] {', '.join(counts.scopes)}")
if dry_run:
console.print("[dim]No rows were deleted (--dry-run).[/dim]")
8 changes: 6 additions & 2 deletions src/cli/commands/worker.py
Original file line number Diff line number Diff line change
@@ -1,15 +1,19 @@
"""raglogs worker — background job processor."""

import typer
from rich.console import Console

console = Console()


def worker_cmd(
poll_interval: int = typer.Option(2, "--poll-interval", help="Seconds between idle polls"),
poll_interval: int = typer.Option(
2, "--poll-interval", help="Seconds between idle polls"
),
):
"""
Start the background worker. Processes ingestion jobs enqueued via the API.
Start the background worker. Processes ingestion jobs enqueued via the API
and scheduled retention purge jobs (G13).

Safe to run multiple workers in parallel — uses SELECT FOR UPDATE SKIP LOCKED.
Graceful shutdown on SIGINT / SIGTERM (Ctrl-C).
Expand Down
2 changes: 2 additions & 0 deletions src/cli/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ def _build_app() -> typer.Typer:
from src.cli.commands.timeline import timeline_cmd
from src.cli.commands.compare import compare_cmd
from src.cli.commands.keys import app as keys_app
from src.cli.commands.purge import purge_cmd

_app = typer.Typer(
name="raglogs",
Expand All @@ -35,6 +36,7 @@ def _build_app() -> typer.Typer:
_app.command("worker")(worker_cmd)
_app.command("timeline")(timeline_cmd)
_app.command("compare")(compare_cmd)
_app.command("purge")(purge_cmd)
_app.add_typer(ask_app, name="ask")
_app.add_typer(keys_app, name="keys")
return _app
Expand Down
Loading
Loading