diff --git a/.env.example b/.env.example index 592f491..5a1da8e 100644 --- a/.env.example +++ b/.env.example @@ -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). diff --git a/README.md b/README.md index a00c441..a976a26 100644 --- a/README.md +++ b/README.md @@ -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`). @@ -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) | @@ -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. @@ -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 diff --git a/migrations/versions/0010_retention.py b/migrations/versions/0010_retention.py new file mode 100644 index 0000000..ebd0274 --- /dev/null +++ b/migrations/versions/0010_retention.py @@ -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) diff --git a/src/api/routes/config.py b/src/api/routes/config.py index 92366b5..1097abe 100644 --- a/src/api/routes/config.py +++ b/src/api/routes/config.py @@ -6,6 +6,7 @@ @router.get("") def get_config(): from src.config import get_settings + settings = get_settings() return { "llm_provider": settings.llm_provider, @@ -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, } diff --git a/src/api/routes/explain.py b/src/api/routes/explain.py index 094dd05..524b24f 100644 --- a/src/api/routes/explain.py +++ b/src/api/routes/explain.py @@ -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(), @@ -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() @@ -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) diff --git a/src/cli/commands/config_cmd.py b/src/cli/commands/config_cmd.py index 25e4bac..21b13e9 100644 --- a/src/cli/commands/config_cmd.py +++ b/src/cli/commands/config_cmd.py @@ -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") @@ -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(): diff --git a/src/cli/commands/purge.py b/src/cli/commands/purge.py new file mode 100644 index 0000000..0c8ad49 --- /dev/null +++ b/src/cli/commands/purge.py @@ -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]") diff --git a/src/cli/commands/worker.py b/src/cli/commands/worker.py index 1309c15..054914b 100644 --- a/src/cli/commands/worker.py +++ b/src/cli/commands/worker.py @@ -1,4 +1,5 @@ """raglogs worker — background job processor.""" + import typer from rich.console import Console @@ -6,10 +7,13 @@ 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). diff --git a/src/cli/main.py b/src/cli/main.py index 1b875b2..9c97db9 100644 --- a/src/cli/main.py +++ b/src/cli/main.py @@ -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", @@ -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 diff --git a/src/config/settings.py b/src/config/settings.py index 613b846..4e512ea 100644 --- a/src/config/settings.py +++ b/src/config/settings.py @@ -53,6 +53,13 @@ class Settings(BaseSettings): # Worker worker_poll_interval: int = 2 # seconds between idle polls + # Data retention (G13). Duration strings like 30d / 180d. 0, empty, or + # "off" disables that tier. Per-scope overrides live in scope_retention. + retention_raw: str = "30d" + retention_summary: str = "180d" + purge_interval_seconds: int = 3600 + purge_chunk_size: int = 1000 + # 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 diff --git a/src/core/clustering/clusterer.py b/src/core/clustering/clusterer.py index 6b1d27f..eddfbf4 100644 --- a/src/core/clustering/clusterer.py +++ b/src/core/clustering/clusterer.py @@ -117,7 +117,13 @@ def _run_clustering( if not rows: cluster_run = _create_cluster_run( - db, window_start, window_end, service, environment, save=save_to_db + db, + window_start, + window_end, + service, + environment, + save=save_to_db, + scope=scope, ) return cluster_run, [] @@ -222,6 +228,7 @@ def _run_clustering( environment, save=save_to_db, algorithm=algorithm, + scope=scope, ) if save_to_db: @@ -259,6 +266,7 @@ def _create_cluster_run( environment: Optional[str], save: bool = True, algorithm: str = "fingerprint", + scope: str = DEFAULT_LOG_SCOPE, ) -> ClusterRun: run = ClusterRun( id=uuid.uuid4(), @@ -268,6 +276,7 @@ def _create_cluster_run( environment_filter=environment, algorithm=algorithm, status="completed", + scope=scope or DEFAULT_LOG_SCOPE, ) if save: db.add(run) diff --git a/src/core/retention/__init__.py b/src/core/retention/__init__.py new file mode 100644 index 0000000..77ddaa9 --- /dev/null +++ b/src/core/retention/__init__.py @@ -0,0 +1,37 @@ +"""Data retention and purge (G13).""" + +from src.core.retention.policy import ( + RetentionPolicy, + apply_scope_override, + compute_cutoff, + is_retention_disabled, + parse_retention_interval, + resolve_scope_policy, + validate_policy_intervals, +) +from src.core.retention.purge import ( + LAST_PURGE_AT_KEY, + PURGE_JOB_TYPE, + PurgeCounts, + maybe_enqueue_purge, + run_purge, + run_purge_job, + should_enqueue_purge, +) + +__all__ = [ + "LAST_PURGE_AT_KEY", + "PURGE_JOB_TYPE", + "PurgeCounts", + "RetentionPolicy", + "apply_scope_override", + "compute_cutoff", + "is_retention_disabled", + "maybe_enqueue_purge", + "parse_retention_interval", + "resolve_scope_policy", + "validate_policy_intervals", + "run_purge", + "run_purge_job", + "should_enqueue_purge", +] diff --git a/src/core/retention/policy.py b/src/core/retention/policy.py new file mode 100644 index 0000000..b401179 --- /dev/null +++ b/src/core/retention/policy.py @@ -0,0 +1,115 @@ +"""Retention interval parsing, cutoffs, and per-scope overrides.""" + +from __future__ import annotations + +from dataclasses import dataclass +from datetime import datetime, timedelta, timezone +from typing import Any, Optional + +from src.config.settings import Settings +from src.utils.time import parse_duration + +# 0 / empty / off (and a few aliases) mean "never purge this tier". +NEVER_TOKENS = frozenset({"", "0", "off", "none", "never", "false"}) + + +@dataclass(frozen=True) +class RetentionPolicy: + scope: str + raw_interval: str + summary_interval: str + + +def is_retention_disabled(value: Optional[str]) -> bool: + """True when the interval means never purge that tier.""" + if value is None: + return True + stripped = value.strip().lower() + if stripped in NEVER_TOKENS: + return True + try: + return parse_duration(stripped).total_seconds() <= 0 + except ValueError: + return False + + +def parse_retention_interval(value: Optional[str]) -> Optional[timedelta]: + """Parse a retention string into a timedelta, or None to skip the tier. + + ``None``, empty, ``0``, and ``off`` disable expiry. Invalid strings raise + ``ValueError``. A parsed duration of zero or less also disables expiry. + """ + if value is None: + return None + stripped = value.strip() + if stripped.lower() in NEVER_TOKENS: + return None + if any(ch.isspace() for ch in stripped): + raise ValueError(f"Cannot parse retention interval: {value!r}") + delta = parse_duration(stripped) + if delta.total_seconds() <= 0: + return None + return delta + + +def compute_cutoff( + interval: Optional[str], + *, + now: Optional[datetime] = None, +) -> Optional[datetime]: + """Return the expiry cutoff, or None when that tier should not be purged.""" + delta = parse_retention_interval(interval) + if delta is None: + return None + clock = now or datetime.now(tz=timezone.utc) + if clock.tzinfo is None: + clock = clock.replace(tzinfo=timezone.utc) + return clock - delta + + +def validate_policy_intervals(policy: RetentionPolicy) -> None: + """Raise ``ValueError`` if either interval string cannot be parsed.""" + parse_retention_interval(policy.raw_interval) + parse_retention_interval(policy.summary_interval) + + +def apply_scope_override( + *, + scope: str, + default_raw: str, + default_summary: str, + override_raw: Optional[str] = None, + override_summary: Optional[str] = None, +) -> RetentionPolicy: + """Merge a ``scope_retention`` row with env defaults. + + ``None`` on an override column means "missing" → env default. An empty + string is a present override and disables that tier. + """ + raw = default_raw if override_raw is None else override_raw + summary = default_summary if override_summary is None else override_summary + return RetentionPolicy(scope=scope, raw_interval=raw, summary_interval=summary) + + +def resolve_scope_policy( + db: Any, + scope: str, + settings: Optional[Settings] = None, +) -> RetentionPolicy: + """Load per-scope overrides; fall back to ``RETENTION_RAW`` / ``RETENTION_SUMMARY``.""" + from src.config import get_settings + from src.db.models import ScopeRetention + + cfg = settings or get_settings() + row = db.get(ScopeRetention, scope) + override_raw = getattr(row, "raw_interval", None) if row is not None else None + override_summary = ( + getattr(row, "summary_interval", None) if row is not None else None + ) + return apply_scope_override( + scope=scope, + default_raw=cfg.retention_raw, + default_summary=cfg.retention_summary, + override_raw=override_raw, + override_summary=override_summary, + ) diff --git a/src/core/retention/purge.py b/src/core/retention/purge.py new file mode 100644 index 0000000..d0142ae --- /dev/null +++ b/src/core/retention/purge.py @@ -0,0 +1,659 @@ +"""Scheduled / on-demand retention purge (G13). + +Deletes expired raw log rows first (``log_embeddings`` and ``cluster_members`` +follow ON DELETE CASCADE) while leaving ``cluster_embeddings`` in place so +similar-incident search still works. After the summary TTL, cluster embeddings, +cluster runs, and explanations for that scope are removed. + +Unit tests compile the SQL statements without a database. +""" + +from __future__ import annotations + +import uuid +from dataclasses import dataclass, field +from datetime import datetime, timezone +from typing import Any, Iterable, Optional + +import structlog +from sqlalchemy import Select, delete, func, select, union +from sqlalchemy.orm import Session +from sqlalchemy.sql.dml import Delete + +from src.config.settings import Settings +from src.core.retention.policy import ( + RetentionPolicy, + compute_cutoff, + resolve_scope_policy, + validate_policy_intervals, +) +from src.db.models import ( + AppConfig, + Cluster, + ClusterEmbedding, + ClusterMember, + ClusterRun, + Explanation, + LogEmbedding, + LogEntry, + ScopeRetention, + WorkerJob, +) + +log = structlog.get_logger() + +PURGE_JOB_TYPE = "purge" +LAST_PURGE_AT_KEY = "last_purge_at" + +_RAW_EXPIRY = LogEntry.created_at +_SUMMARY_EMBEDDING_EXPIRY = func.coalesce( + ClusterEmbedding.last_seen, + ClusterEmbedding.updated_at, + ClusterEmbedding.created_at, +) +_CLUSTER_RUN_EXPIRY = func.coalesce(ClusterRun.window_end, ClusterRun.created_at) + + +@dataclass +class PurgeCounts: + raw: int = 0 + summary: int = 0 + embedding: int = 0 + scopes: list[str] = field(default_factory=list) + more_remaining: bool = False + + def add(self, other: "PurgeCounts") -> None: + self.raw += other.raw + self.summary += other.summary + self.embedding += other.embedding + self.more_remaining = self.more_remaining or other.more_remaining + for scope in other.scopes: + if scope not in self.scopes: + self.scopes.append(scope) + + def is_empty(self) -> bool: + return self.raw == 0 and self.summary == 0 and self.embedding == 0 + + def as_dict(self, *, dry_run: bool = False) -> dict[str, Any]: + return { + "raw": self.raw, + "summary": self.summary, + "embedding": self.embedding, + "scopes": list(self.scopes), + "dry_run": dry_run, + "more_remaining": self.more_remaining, + } + + +def raw_expiry_expression() -> Any: + """``log_entries.created_at`` — time-in-store, not the event timestamp.""" + return _RAW_EXPIRY + + +def summary_embedding_expiry_expression() -> Any: + """``COALESCE(last_seen, updated_at, created_at)`` for cluster embeddings.""" + return _SUMMARY_EMBEDDING_EXPIRY + + +def select_expired_log_entry_ids( + scope: str, + cutoff: datetime, + *, + limit: int, +) -> Select: + """Ids of raw rows in ``scope`` older than ``cutoff`` (chunked).""" + return ( + select(LogEntry.id) + .where(LogEntry.scope == scope) + .where(_RAW_EXPIRY < cutoff) + .order_by(_RAW_EXPIRY.asc(), LogEntry.id.asc()) + .limit(limit) + ) + + +def count_expired_log_entries_statement(scope: str, cutoff: datetime) -> Select: + return ( + select(func.count()) + .select_from(LogEntry) + .where(LogEntry.scope == scope) + .where(_RAW_EXPIRY < cutoff) + ) + + +def count_expired_log_embeddings_statement(scope: str, cutoff: datetime) -> Select: + return ( + select(func.count()) + .select_from(LogEmbedding) + .join(LogEntry, LogEmbedding.log_entry_id == LogEntry.id) + .where(LogEntry.scope == scope) + .where(_RAW_EXPIRY < cutoff) + ) + + +def count_expired_cluster_members_statement(scope: str, cutoff: datetime) -> Select: + return ( + select(func.count()) + .select_from(ClusterMember) + .join(LogEntry, ClusterMember.log_entry_id == LogEntry.id) + .where(LogEntry.scope == scope) + .where(_RAW_EXPIRY < cutoff) + ) + + +def count_log_embeddings_statement(entry_ids: Iterable[uuid.UUID]) -> Select: + ids = list(entry_ids) + return ( + select(func.count()) + .select_from(LogEmbedding) + .where(LogEmbedding.log_entry_id.in_(ids)) + ) + + +def count_cluster_members_statement(entry_ids: Iterable[uuid.UUID]) -> Select: + ids = list(entry_ids) + return ( + select(func.count()) + .select_from(ClusterMember) + .where(ClusterMember.log_entry_id.in_(ids)) + ) + + +def delete_log_entries_statement(entry_ids: Iterable[uuid.UUID]) -> Delete: + return delete(LogEntry).where(LogEntry.id.in_(list(entry_ids))) + + +def select_expired_cluster_embedding_ids( + scope: str, + cutoff: datetime, + *, + limit: int, +) -> Select: + return ( + select(ClusterEmbedding.id) + .where(ClusterEmbedding.scope == scope) + .where(_SUMMARY_EMBEDDING_EXPIRY < cutoff) + .order_by(_SUMMARY_EMBEDDING_EXPIRY.asc(), ClusterEmbedding.id.asc()) + .limit(limit) + ) + + +def count_expired_cluster_embeddings_statement(scope: str, cutoff: datetime) -> Select: + return ( + select(func.count()) + .select_from(ClusterEmbedding) + .where(ClusterEmbedding.scope == scope) + .where(_SUMMARY_EMBEDDING_EXPIRY < cutoff) + ) + + +def delete_cluster_embeddings_statement(embedding_ids: Iterable[uuid.UUID]) -> Delete: + return delete(ClusterEmbedding).where(ClusterEmbedding.id.in_(list(embedding_ids))) + + +def select_expired_explanation_ids( + scope: str, + cutoff: datetime, + *, + limit: int, +) -> Select: + return ( + select(Explanation.id) + .where(Explanation.scope == scope) + .where(Explanation.created_at < cutoff) + .order_by(Explanation.created_at.asc(), Explanation.id.asc()) + .limit(limit) + ) + + +def delete_explanations_statement(explanation_ids: Iterable[uuid.UUID]) -> Delete: + return delete(Explanation).where(Explanation.id.in_(list(explanation_ids))) + + +def count_expired_explanations_statement(scope: str, cutoff: datetime) -> Select: + return ( + select(func.count()) + .select_from(Explanation) + .where(Explanation.scope == scope) + .where(Explanation.created_at < cutoff) + ) + + +def select_expired_cluster_run_ids( + scope: str, + cutoff: datetime, + *, + limit: int, +) -> Select: + return ( + select(ClusterRun.id) + .where(ClusterRun.scope == scope) + .where(_CLUSTER_RUN_EXPIRY < cutoff) + .order_by(_CLUSTER_RUN_EXPIRY.asc(), ClusterRun.id.asc()) + .limit(limit) + ) + + +def count_clusters_for_runs_statement(run_ids: Iterable[uuid.UUID]) -> Select: + ids = list(run_ids) + return ( + select(func.count()).select_from(Cluster).where(Cluster.cluster_run_id.in_(ids)) + ) + + +def delete_cluster_runs_statement(run_ids: Iterable[uuid.UUID]) -> Delete: + return delete(ClusterRun).where(ClusterRun.id.in_(list(run_ids))) + + +def count_expired_cluster_runs_statement(scope: str, cutoff: datetime) -> Select: + return ( + select(func.count()) + .select_from(ClusterRun) + .where(ClusterRun.scope == scope) + .where(_CLUSTER_RUN_EXPIRY < cutoff) + ) + + +def count_expired_clusters_statement(scope: str, cutoff: datetime) -> Select: + return ( + select(func.count()) + .select_from(Cluster) + .join(ClusterRun, Cluster.cluster_run_id == ClusterRun.id) + .where(ClusterRun.scope == scope) + .where(_CLUSTER_RUN_EXPIRY < cutoff) + ) + + +def scopes_to_purge_statement() -> Select: + """Distinct scopes that have data or a retention override.""" + return union( + select(LogEntry.scope.label("scope")).where(LogEntry.scope.isnot(None)), + select(ClusterEmbedding.scope.label("scope")).where( + ClusterEmbedding.scope.isnot(None) + ), + select(ClusterRun.scope.label("scope")).where(ClusterRun.scope.isnot(None)), + select(Explanation.scope.label("scope")).where(Explanation.scope.isnot(None)), + select(ScopeRetention.scope.label("scope")), + ) + + +def active_purge_job_statement() -> Select: + return ( + select(WorkerJob.id) + .where(WorkerJob.job_type == PURGE_JOB_TYPE) + .where(WorkerJob.status.in_(("pending", "running"))) + .limit(1) + ) + + +def should_enqueue_purge( + *, + last_purge_at: Optional[datetime], + has_active_purge: bool, + now: datetime, + interval_seconds: int, +) -> bool: + """Idle-poll decision: enqueue when due and no purge is already claimed.""" + if has_active_purge: + return False + if interval_seconds <= 0: + return False + if last_purge_at is None: + return True + stamp = last_purge_at + if stamp.tzinfo is None: + stamp = stamp.replace(tzinfo=timezone.utc) + clock = now if now.tzinfo is not None else now.replace(tzinfo=timezone.utc) + return (clock - stamp).total_seconds() >= interval_seconds + + +def read_last_purge_at(db: Session) -> Optional[datetime]: + row = db.get(AppConfig, LAST_PURGE_AT_KEY) + if row is None: + return None + raw = row.value_json + if isinstance(raw, dict): + raw = raw.get("at") + if not isinstance(raw, str) or not raw.strip(): + return None + try: + stamp = datetime.fromisoformat(raw) + except ValueError: + return None + if stamp.tzinfo is None: + stamp = stamp.replace(tzinfo=timezone.utc) + return stamp + + +def write_last_purge_at(db: Session, now: datetime) -> None: + from sqlalchemy.dialects.postgresql import insert as pg_insert + + clock = now if now.tzinfo is not None else now.replace(tzinfo=timezone.utc) + iso = clock.isoformat() + stmt = pg_insert(AppConfig).values(key=LAST_PURGE_AT_KEY, value_json=iso) + stmt = stmt.on_conflict_do_update( + index_elements=["key"], + set_={"value_json": stmt.excluded.value_json}, + ) + db.execute(stmt) + + +def _record(kind: str, count: int, *, dry_run: bool) -> None: + if dry_run or count <= 0: + return + from src.observability.metrics import record_purge_rows + + record_purge_rows(count, kind=kind) + + +def _scalar_ids(rows: Iterable[Any]) -> list[uuid.UUID]: + out: list[uuid.UUID] = [] + for row in rows: + value = row[0] if not hasattr(row, "id") else getattr(row, "id", row[0]) + if value is None: + continue + out.append(value if isinstance(value, uuid.UUID) else uuid.UUID(str(value))) + return out + + +def _count_result(value: Any) -> int: + if value is None: + return 0 + return int(value) + + +def purge_raw_chunk( + db: Session, + scope: str, + cutoff: datetime, + *, + limit: int, + dry_run: bool = False, +) -> PurgeCounts: + """Delete up to ``limit`` expired log rows in ``scope``. Embeddings CASCADE.""" + rows = db.execute(select_expired_log_entry_ids(scope, cutoff, limit=limit)).all() + ids = _scalar_ids(rows) + counts = PurgeCounts() + if not ids: + return counts + embedding_n = _count_result( + db.execute(count_log_embeddings_statement(ids)).scalar_one() + ) + member_n = _count_result( + db.execute(count_cluster_members_statement(ids)).scalar_one() + ) + if not dry_run: + db.execute(delete_log_entries_statement(ids)) + counts.raw = len(ids) + counts.embedding = embedding_n + counts.more_remaining = len(ids) >= limit + _record("raw", counts.raw, dry_run=dry_run) + _record("raw", member_n, dry_run=dry_run) + _record("embedding", counts.embedding, dry_run=dry_run) + return counts + + +def purge_summary_chunk( + db: Session, + scope: str, + cutoff: datetime, + *, + limit: int, + dry_run: bool = False, +) -> PurgeCounts: + """Delete expired cluster embeddings / runs / explanations in ``scope``.""" + counts = PurgeCounts() + + embedding_rows = db.execute( + select_expired_cluster_embedding_ids(scope, cutoff, limit=limit) + ).all() + embedding_ids = _scalar_ids(embedding_rows) + if embedding_ids: + if not dry_run: + db.execute(delete_cluster_embeddings_statement(embedding_ids)) + counts.embedding += len(embedding_ids) + counts.more_remaining = counts.more_remaining or len(embedding_ids) >= limit + _record("embedding", len(embedding_ids), dry_run=dry_run) + + explanation_rows = db.execute( + select_expired_explanation_ids(scope, cutoff, limit=limit) + ).all() + explanation_ids = _scalar_ids(explanation_rows) + if explanation_ids: + if not dry_run: + db.execute(delete_explanations_statement(explanation_ids)) + counts.summary += len(explanation_ids) + counts.more_remaining = counts.more_remaining or len(explanation_ids) >= limit + _record("summary", len(explanation_ids), dry_run=dry_run) + + run_rows = db.execute( + select_expired_cluster_run_ids(scope, cutoff, limit=limit) + ).all() + run_ids = _scalar_ids(run_rows) + if run_ids: + cluster_n = _count_result( + db.execute(count_clusters_for_runs_statement(run_ids)).scalar_one() + ) + if not dry_run: + db.execute(delete_cluster_runs_statement(run_ids)) + summary_n = len(run_ids) + cluster_n + counts.summary += summary_n + counts.more_remaining = counts.more_remaining or len(run_ids) >= limit + _record("summary", summary_n, dry_run=dry_run) + + return counts + + +def purge_scope_batch( + db: Session, + policy: RetentionPolicy, + *, + now: datetime, + limit: int, + dry_run: bool = False, +) -> PurgeCounts: + """One raw chunk then one summary chunk for a single scope (G8).""" + counts = PurgeCounts(scopes=[policy.scope]) + raw_cutoff = compute_cutoff(policy.raw_interval, now=now) + if raw_cutoff is not None: + counts.add( + purge_raw_chunk(db, policy.scope, raw_cutoff, limit=limit, dry_run=dry_run) + ) + summary_cutoff = compute_cutoff(policy.summary_interval, now=now) + if summary_cutoff is not None: + counts.add( + purge_summary_chunk( + db, policy.scope, summary_cutoff, limit=limit, dry_run=dry_run + ) + ) + return counts + + +def count_scope_expired( + db: Session, + policy: RetentionPolicy, + *, + now: datetime, +) -> PurgeCounts: + """COUNT expired rows for a scope — used by ``--dry-run`` (no chunk loop).""" + counts = PurgeCounts(scopes=[policy.scope]) + raw_cutoff = compute_cutoff(policy.raw_interval, now=now) + if raw_cutoff is not None: + counts.raw = _count_result( + db.execute( + count_expired_log_entries_statement(policy.scope, raw_cutoff) + ).scalar_one() + ) + counts.embedding += _count_result( + db.execute( + count_expired_log_embeddings_statement(policy.scope, raw_cutoff) + ).scalar_one() + ) + member_n = _count_result( + db.execute( + count_expired_cluster_members_statement(policy.scope, raw_cutoff) + ).scalar_one() + ) + _record("raw", counts.raw, dry_run=True) + _record("raw", member_n, dry_run=True) + _record("embedding", counts.embedding, dry_run=True) + summary_cutoff = compute_cutoff(policy.summary_interval, now=now) + if summary_cutoff is not None: + cluster_emb = _count_result( + db.execute( + count_expired_cluster_embeddings_statement(policy.scope, summary_cutoff) + ).scalar_one() + ) + explanations = _count_result( + db.execute( + count_expired_explanations_statement(policy.scope, summary_cutoff) + ).scalar_one() + ) + runs = _count_result( + db.execute( + count_expired_cluster_runs_statement(policy.scope, summary_cutoff) + ).scalar_one() + ) + clusters = _count_result( + db.execute( + count_expired_clusters_statement(policy.scope, summary_cutoff) + ).scalar_one() + ) + counts.embedding += cluster_emb + counts.summary += explanations + runs + clusters + return counts + + +def list_scopes_to_purge(db: Session) -> list[str]: + rows = db.execute(scopes_to_purge_statement()).all() + scopes: list[str] = [] + seen: set[str] = set() + for row in rows: + value = row[0] if not isinstance(row, str) else row + if not value: + continue + scope = str(value) + if scope in seen: + continue + seen.add(scope) + scopes.append(scope) + scopes.sort() + return scopes + + +def run_purge( + db: Session, + *, + scope: Optional[str] = None, + dry_run: bool = False, + max_chunks: Optional[int] = 1, + now: Optional[datetime] = None, + settings: Optional[Settings] = None, +) -> PurgeCounts: + """Purge expired rows. Worker uses ``max_chunks=1``; CLI drains all chunks. + + ``--dry-run`` COUNTs expired rows once per scope (never loops). Invalid + interval strings skip that scope. A full chunk skips ``last_purge_at`` and + enqueues a follow-up job so idle workers keep draining. + """ + from src.config import get_settings + + cfg = settings or get_settings() + clock = now or datetime.now(tz=timezone.utc) + limit = max(1, int(cfg.purge_chunk_size)) + scopes = [scope] if scope else list_scopes_to_purge(db) + totals = PurgeCounts() + more_remaining = False + for item in scopes: + try: + policy = resolve_scope_policy(db, item, settings=cfg) + validate_policy_intervals(policy) + except ValueError as exc: + log.warning( + "retention_policy_invalid", + scope=item, + error=str(exc), + ) + if item not in totals.scopes: + totals.scopes.append(item) + continue + if dry_run: + totals.add(count_scope_expired(db, policy, now=clock)) + continue + chunks = 0 + while True: + batch = purge_scope_batch( + db, policy, now=clock, limit=limit, dry_run=False + ) + totals.add(batch) + chunks += 1 + db.commit() + if batch.raw == 0 and batch.summary == 0 and batch.embedding == 0: + if item not in totals.scopes: + totals.scopes.append(item) + break + if max_chunks is not None and chunks >= max_chunks: + if batch.more_remaining: + more_remaining = True + break + totals.more_remaining = more_remaining + if not dry_run: + if totals.more_remaining: + enqueue_followup_purge(db) + else: + write_last_purge_at(db, clock) + db.commit() + return totals + + +def enqueue_followup_purge(db: Session) -> None: + """Queue another purge so remaining chunks drain without waiting the interval.""" + job = WorkerJob( + job_type=PURGE_JOB_TYPE, + status="pending", + payload_json={}, + ) + db.add(job) + db.flush() + + +def run_purge_job(db: Session, worker_job: Any) -> dict[str, Any]: + """Execute a claimed ``purge`` worker job (one chunk per scope).""" + payload = getattr(worker_job, "payload_json", None) or {} + scope = payload.get("scope") if isinstance(payload, dict) else None + scope_str = str(scope).strip() if isinstance(scope, str) and scope.strip() else None + counts = run_purge(db, scope=scope_str, dry_run=False, max_chunks=1) + return counts.as_dict(dry_run=False) + + +def maybe_enqueue_purge( + db: Session, + *, + now: Optional[datetime] = None, + settings: Optional[Settings] = None, +) -> bool: + """Enqueue a purge job on idle poll when the interval has elapsed. + + Safe under SKIP LOCKED: duplicate pending rows are idempotent. Does not + enqueue when a purge job is already pending or running. + """ + from src.config import get_settings + + cfg = settings or get_settings() + clock = now or datetime.now(tz=timezone.utc) + active = db.execute(active_purge_job_statement()).scalar_one_or_none() + last = read_last_purge_at(db) + if not should_enqueue_purge( + last_purge_at=last, + has_active_purge=active is not None, + now=clock, + interval_seconds=cfg.purge_interval_seconds, + ): + return False + job = WorkerJob( + job_type=PURGE_JOB_TYPE, + status="pending", + payload_json={}, + ) + db.add(job) + db.flush() + return True diff --git a/src/db/__init__.py b/src/db/__init__.py index a8839d8..a7840e8 100644 --- a/src/db/__init__.py +++ b/src/db/__init__.py @@ -11,6 +11,7 @@ IngestionJob, LogEmbedding, LogEntry, + ScopeRetention, Source, WorkerJob, ) @@ -29,6 +30,7 @@ "Explanation", "ApiKey", "AppConfig", + "ScopeRetention", "IngestIdempotencyKey", "WorkerJob", "get_db", diff --git a/src/db/models.py b/src/db/models.py index 0906971..2ccd611 100644 --- a/src/db/models.py +++ b/src/db/models.py @@ -130,7 +130,12 @@ class LogEmbedding(Base): __tablename__ = "log_embeddings" id: Mapped[uuid.UUID] = mapped_column(UUID(as_uuid=True), primary_key=True, default=uuid.uuid4) - log_entry_id: Mapped[uuid.UUID] = mapped_column(UUID(as_uuid=True), ForeignKey("log_entries.id"), nullable=False, unique=True) + log_entry_id: Mapped[uuid.UUID] = mapped_column( + UUID(as_uuid=True), + ForeignKey("log_entries.id", ondelete="CASCADE"), + nullable=False, + unique=True, + ) embedding: Mapped[Any] = mapped_column(Vector(1536), nullable=False) model_name: Mapped[str] = mapped_column(String(255), nullable=False) created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now()) @@ -177,16 +182,23 @@ class ClusterRun(Base): environment_filter: Mapped[str | None] = mapped_column(String(100), nullable=True) algorithm: Mapped[str] = mapped_column(String(50), default="fingerprint") status: Mapped[str] = mapped_column(String(50), default="completed") + scope: Mapped[str] = mapped_column( + String(255), nullable=False, default=DEFAULT_LOG_SCOPE, server_default="default" + ) created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now()) clusters: Mapped[list["Cluster"]] = relationship("Cluster", back_populates="cluster_run") + __table_args__ = (Index("ix_cluster_runs_scope", "scope"),) + class Cluster(Base): __tablename__ = "clusters" id: Mapped[uuid.UUID] = mapped_column(UUID(as_uuid=True), primary_key=True, default=uuid.uuid4) - cluster_run_id: Mapped[uuid.UUID] = mapped_column(UUID(as_uuid=True), ForeignKey("cluster_runs.id"), nullable=False) + cluster_run_id: Mapped[uuid.UUID] = mapped_column( + UUID(as_uuid=True), ForeignKey("cluster_runs.id", ondelete="CASCADE"), nullable=False + ) cluster_key: Mapped[str] = mapped_column(String(64), nullable=False) representative_message: Mapped[str | None] = mapped_column(Text, nullable=True) fingerprint: Mapped[str | None] = mapped_column(String(64), nullable=True) @@ -209,8 +221,12 @@ class ClusterMember(Base): __tablename__ = "cluster_members" id: Mapped[uuid.UUID] = mapped_column(UUID(as_uuid=True), primary_key=True, default=uuid.uuid4) - cluster_id: Mapped[uuid.UUID] = mapped_column(UUID(as_uuid=True), ForeignKey("clusters.id"), nullable=False) - log_entry_id: Mapped[uuid.UUID] = mapped_column(UUID(as_uuid=True), ForeignKey("log_entries.id"), nullable=False) + cluster_id: Mapped[uuid.UUID] = mapped_column( + UUID(as_uuid=True), ForeignKey("clusters.id", ondelete="CASCADE"), nullable=False + ) + log_entry_id: Mapped[uuid.UUID] = mapped_column( + UUID(as_uuid=True), ForeignKey("log_entries.id", ondelete="CASCADE"), nullable=False + ) cluster: Mapped["Cluster"] = relationship("Cluster", back_populates="members") log_entry: Mapped["LogEntry"] = relationship("LogEntry", back_populates="cluster_members") @@ -229,8 +245,24 @@ class Explanation(Base): result_text: Mapped[str | None] = mapped_column(Text, nullable=True) result_json: Mapped[dict[str, Any] | None] = mapped_column(JSONB, nullable=True) confidence: Mapped[str | None] = mapped_column(String(50), nullable=True) + scope: Mapped[str] = mapped_column( + String(255), nullable=False, default=DEFAULT_LOG_SCOPE, server_default="default" + ) created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now()) + __table_args__ = (Index("ix_explanations_scope", "scope"),) + + +class ScopeRetention(Base): + """Per-scope TTL overrides. NULL interval = fall back to env default.""" + + __tablename__ = "scope_retention" + + scope: Mapped[str] = mapped_column(String(255), primary_key=True) + raw_interval: Mapped[str | None] = mapped_column(String(32), nullable=True) + summary_interval: Mapped[str | None] = mapped_column(String(32), nullable=True) + updated_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now()) + class AppConfig(Base): __tablename__ = "app_config" diff --git a/src/observability/metrics.py b/src/observability/metrics.py index 0e6eee2..7ceeef5 100644 --- a/src/observability/metrics.py +++ b/src/observability/metrics.py @@ -18,6 +18,7 @@ "ingestions", ) _INGEST_RESULTS: tuple[str, ...] = ("inserted", "deduped", "error") +_PURGE_KINDS: tuple[str, ...] = ("raw", "summary", "embedding") INGEST_DURATION = Histogram( "raglogs_ingest_duration_seconds", @@ -71,6 +72,12 @@ "LLM circuit breaker: 0=closed, 1=half_open, 2=open.", registry=REGISTRY, ) +PURGE_ROWS = Counter( + "raglogs_purge_rows_total", + "Rows reclaimed by the retention purge worker, by kind.", + ["kind"], + registry=REGISTRY, +) _QUEUE_GAUGE_NAME = "raglogs_worker_queue_depth" _QUEUE_GAUGE_HELP = ( "Pending worker jobs. Published after a successful scrape and left stale " @@ -85,6 +92,8 @@ INGEST_LINES.labels(result=_result) for _endpoint in _QUERY_ENDPOINTS: QUERY_REQUEST_DURATION.labels(endpoint=_endpoint) +for _kind in _PURGE_KINDS: + PURGE_ROWS.labels(kind=_kind) BREAKER_STATE_VALUES: dict[str, float] = { "closed": 0.0, @@ -130,6 +139,13 @@ def record_llm_estimated_tokens(count: int) -> None: LLM_ESTIMATED_TOKENS.inc(count) +def record_purge_rows(count: int, *, kind: str) -> None: + if count <= 0: + return + label = kind if kind in _PURGE_KINDS else "raw" + PURGE_ROWS.labels(kind=label).inc(count) + + def refresh_runtime_gauges() -> None: """Update breaker + queue gauges just before a /metrics scrape.""" _refresh_breaker_gauge() @@ -185,11 +201,17 @@ def _refresh_queue_depth() -> None: with get_db() as db: try: - db.execute(text(f"SET LOCAL statement_timeout = '{_QUEUE_STATEMENT_TIMEOUT_MS}'")) + db.execute( + text( + f"SET LOCAL statement_timeout = '{_QUEUE_STATEMENT_TIMEOUT_MS}'" + ) + ) except Exception: pass depth = db.execute( - select(func.count()).select_from(WorkerJob).where(WorkerJob.status == "pending") + select(func.count()) + .select_from(WorkerJob) + .where(WorkerJob.status == "pending") ).scalar_one() _set_queue_depth(int(depth or 0)) except Exception: diff --git a/src/worker/runner.py b/src/worker/runner.py index 46992bc..9e097d1 100644 --- a/src/worker/runner.py +++ b/src/worker/runner.py @@ -172,6 +172,10 @@ def process_one(db) -> bool: try: if worker_job.job_type == "ingest": result = run_ingest_job(db, worker_job) + elif worker_job.job_type == "purge": + from src.core.retention.purge import run_purge_job + + result = run_purge_job(db, worker_job) else: raise ValueError(f"Unknown job type: {worker_job.job_type!r}") @@ -194,9 +198,10 @@ def process_one(db) -> bool: # uncommitted done/failed row (poll stays pending, a crash re-claims). db.commit() - from src.core.ingestion.webhooks import maybe_deliver_ingest_callback + if worker_job.job_type == "ingest": + from src.core.ingestion.webhooks import maybe_deliver_ingest_callback - maybe_deliver_ingest_callback(db, worker_job) + maybe_deliver_ingest_callback(db, worker_job) return True # processed (even if failed) — don't sleep @@ -219,10 +224,15 @@ def _handle_signal(sig, frame): while not shutdown: try: + enqueued_purge = False with get_db() as db: processed = process_one(db) ticked = tick_tail_jobs(db) - if not processed and not ticked: + if not processed: + from src.core.retention.purge import maybe_enqueue_purge + + enqueued_purge = maybe_enqueue_purge(db) + if not processed and not ticked and not enqueued_purge: # Nothing to do — sleep before next poll for _ in range(poll_interval * 10): if shutdown: diff --git a/tests/unit/test_observability.py b/tests/unit/test_observability.py index 6b71aee..2832ae2 100644 --- a/tests/unit/test_observability.py +++ b/tests/unit/test_observability.py @@ -89,6 +89,7 @@ def test_metrics_returns_prometheus_text() -> None: "raglogs_llm_fallback_total", "raglogs_llm_breaker_state", "raglogs_cluster_count", + "raglogs_purge_rows_total", ): assert name in body, f"missing metric {name}" diff --git a/tests/unit/test_retention.py b/tests/unit/test_retention.py new file mode 100644 index 0000000..5c76d90 --- /dev/null +++ b/tests/unit/test_retention.py @@ -0,0 +1,514 @@ +"""G13 data retention: cutoffs, per-scope TTL, purge SQL, metrics, scheduler.""" + +from __future__ import annotations + +import uuid +from datetime import datetime, timedelta, timezone +from unittest.mock import MagicMock, patch + +import pytest + +from src.config.settings import Settings +from src.core.retention.policy import ( + apply_scope_override, + compute_cutoff, + is_retention_disabled, + parse_retention_interval, + resolve_scope_policy, + validate_policy_intervals, +) +from src.core.retention.purge import ( + PURGE_JOB_TYPE, + maybe_enqueue_purge, + purge_raw_chunk, + purge_scope_batch, + purge_summary_chunk, + run_purge, + select_expired_cluster_embedding_ids, + select_expired_log_entry_ids, + should_enqueue_purge, +) +from src.db.models import ClusterEmbedding, LogEntry, WorkerJob +from src.observability.metrics import REGISTRY, record_purge_rows + +NOW = datetime(2026, 8, 17, 12, 0, 0, tzinfo=timezone.utc) + + +def _compiled(stmt) -> str: + return str(stmt.compile(compile_kwargs={"literal_binds": True})).lower() + + +def _counter(kind: str) -> float: + for metric in REGISTRY.collect(): + for sample in metric.samples: + if ( + sample.name == "raglogs_purge_rows_total" + and sample.labels.get("kind") == kind + ): + return float(sample.value) + return 0.0 + + +class ScriptedSession: + """Minimal Session stand-in: queued execute() results, no database.""" + + def __init__(self, script: list[object] | None = None) -> None: + self.script = list(script or []) + self.executed: list[object] = [] + self.added: list[object] = [] + self._get_row: object | None = None + self._get_rows: dict[object, object] = {} + + def execute(self, stmt: object) -> MagicMock: + self.executed.append(stmt) + payload = self.script.pop(0) if self.script else [] + result = MagicMock() + if isinstance(payload, list): + result.all.return_value = payload + result.scalar_one.return_value = 0 + result.scalar_one_or_none.return_value = None + elif isinstance(payload, int): + result.all.return_value = [] + result.scalar_one.return_value = payload + result.scalar_one_or_none.return_value = payload + else: + result.all.return_value = [] + result.scalar_one.return_value = 0 + result.scalar_one_or_none.return_value = payload + return result + + def get(self, model: object, key: object) -> object | None: + if isinstance(getattr(self, "_get_rows", None), dict) and key in self._get_rows: + return self._get_rows[key] + return self._get_row + + def add(self, obj: object) -> None: + self.added.append(obj) + + def flush(self) -> None: + return None + + def commit(self) -> None: + return None + + +class TestParseRetentionInterval: + def test_duration_days(self) -> None: + assert parse_retention_interval("30d") == timedelta(days=30) + assert parse_retention_interval("180d") == timedelta(days=180) + + def test_zero_empty_off_mean_never(self) -> None: + for value in ("0", "", "off", "OFF", "none", "never", None): + assert parse_retention_interval(value) is None + assert is_retention_disabled(value) is True + + def test_zero_duration_means_never(self) -> None: + assert parse_retention_interval("0d") is None + assert is_retention_disabled("0h") is True + + def test_invalid_raises(self) -> None: + with pytest.raises(ValueError): + parse_retention_interval("not-a-duration") + with pytest.raises(ValueError): + parse_retention_interval("7 days") + with pytest.raises(ValueError): + parse_retention_interval("30") + with pytest.raises(ValueError): + validate_policy_intervals( + apply_scope_override( + scope="x", + default_raw="7 days", + default_summary="180d", + ) + ) + + +class TestComputeCutoff: + def test_subtracts_interval(self) -> None: + cutoff = compute_cutoff("30d", now=NOW) + assert cutoff == NOW - timedelta(days=30) + + def test_skip_when_disabled(self) -> None: + assert compute_cutoff("0", now=NOW) is None + assert compute_cutoff("off", now=NOW) is None + assert compute_cutoff("", now=NOW) is None + assert compute_cutoff(None, now=NOW) is None + + +class TestPerScopeOverride: + def test_missing_override_uses_env_default(self) -> None: + policy = apply_scope_override( + scope="incident:A", + default_raw="30d", + default_summary="180d", + override_raw=None, + override_summary=None, + ) + assert policy.raw_interval == "30d" + assert policy.summary_interval == "180d" + + def test_partial_override(self) -> None: + policy = apply_scope_override( + scope="incident:A", + default_raw="30d", + default_summary="180d", + override_raw="7d", + override_summary=None, + ) + assert policy.raw_interval == "7d" + assert policy.summary_interval == "180d" + + def test_empty_override_disables_tier(self) -> None: + policy = apply_scope_override( + scope="incident:A", + default_raw="30d", + default_summary="180d", + override_raw="0", + override_summary="off", + ) + assert compute_cutoff(policy.raw_interval, now=NOW) is None + assert compute_cutoff(policy.summary_interval, now=NOW) is None + + def test_resolve_scope_policy_reads_table_row(self) -> None: + db = ScriptedSession() + row = MagicMock() + row.raw_interval = "14d" + row.summary_interval = None + db._get_row = row + settings = Settings( + _env_file=None, retention_raw="30d", retention_summary="180d" + ) + policy = resolve_scope_policy(db, "incident:prod", settings=settings) + assert policy.raw_interval == "14d" + assert policy.summary_interval == "180d" + + def test_resolve_scope_policy_missing_row(self) -> None: + db = ScriptedSession() + settings = Settings( + _env_file=None, retention_raw="30d", retention_summary="180d" + ) + policy = resolve_scope_policy(db, "default", settings=settings) + assert policy.raw_interval == "30d" + assert policy.summary_interval == "180d" + + +class TestPurgeSql: + def test_raw_select_filters_scope_and_created_at(self) -> None: + cutoff = NOW - timedelta(days=30) + stmt = select_expired_log_entry_ids("incident:A", cutoff, limit=100) + sql = _compiled(stmt) + assert "log_entries" in sql + assert "created_at" in sql + assert "incident:a" in sql + assert "cluster_embeddings" not in sql + assert "limit" in sql + assert "coalesce" not in sql + assert "timestamp" not in sql or "created_at" in sql + + def test_raw_select_uses_created_at_not_event_timestamp(self) -> None: + stmt = select_expired_log_entry_ids("default", NOW, limit=10) + sql = _compiled(stmt) + assert "log_entries.created_at" in sql + assert "log_entries.timestamp" not in sql + + def test_raw_select_does_not_target_cluster_embeddings(self) -> None: + stmt = select_expired_log_entry_ids("default", NOW, limit=10) + sql = _compiled(stmt) + assert LogEntry.__tablename__ in sql + assert ClusterEmbedding.__tablename__ not in sql + + def test_summary_select_filters_scope(self) -> None: + stmt = select_expired_cluster_embedding_ids("incident:B", NOW, limit=50) + sql = _compiled(stmt) + assert "cluster_embeddings" in sql + assert "incident:b" in sql + assert "log_entries" not in sql + + +class TestPurgeExecution: + def test_raw_chunk_deletes_log_entries_counts_embeddings(self) -> None: + entry_id = uuid.uuid4() + db = ScriptedSession([[(entry_id,)], 2, 0, None]) + counts = purge_raw_chunk(db, "incident:A", NOW, limit=100, dry_run=False) + assert counts.raw == 1 + assert counts.embedding == 2 + assert counts.summary == 0 + sqls = [_compiled(s) for s in db.executed] + assert any("log_entries" in s and "delete" in s for s in sqls) + assert all("cluster_embeddings" not in s or "delete" not in s for s in sqls) + + def test_raw_chunk_dry_run_skips_delete(self) -> None: + entry_id = uuid.uuid4() + db = ScriptedSession([[(entry_id,)], 1, 0]) + counts = purge_raw_chunk(db, "default", NOW, limit=10, dry_run=True) + assert counts.raw == 1 + assert counts.embedding == 1 + assert not any("delete" in _compiled(s) for s in db.executed) + + def test_zero_interval_skips_raw_and_summary(self) -> None: + policy = apply_scope_override( + scope="default", + default_raw="0", + default_summary="off", + ) + db = ScriptedSession() + with ( + patch("src.core.retention.purge.purge_raw_chunk") as raw, + patch("src.core.retention.purge.purge_summary_chunk") as summary, + ): + counts = purge_scope_batch(db, policy, now=NOW, limit=100, dry_run=False) + raw.assert_not_called() + summary.assert_not_called() + assert counts.is_empty() + + def test_summary_chunk_deletes_embeddings_not_log_entries(self) -> None: + embedding_id = uuid.uuid4() + explanation_id = uuid.uuid4() + run_id = uuid.uuid4() + db = ScriptedSession( + [ + [(embedding_id,)], + None, + [(explanation_id,)], + None, + [(run_id,)], + 4, + None, + ] + ) + counts = purge_summary_chunk(db, "incident:A", NOW, limit=50, dry_run=False) + assert counts.embedding == 1 + assert counts.summary == 1 + 4 + 1 # explanation + clusters + run + sqls = [_compiled(s) for s in db.executed] + assert any("cluster_embeddings" in s and "delete" in s for s in sqls) + assert not any("log_entries" in s and "delete" in s for s in sqls) + + def test_run_purge_one_chunk_per_scope(self) -> None: + settings = Settings( + _env_file=None, + retention_raw="30d", + retention_summary="0", + purge_chunk_size=100, + ) + entry_id = uuid.uuid4() + db = ScriptedSession( + [ + [("incident:A",), ("default",)], + [(entry_id,)], + 0, + 0, + None, + ] + ) + with patch("src.core.retention.purge.write_last_purge_at") as written: + counts = run_purge( + db, + dry_run=False, + max_chunks=1, + now=NOW, + settings=settings, + ) + written.assert_called_once() + assert counts.raw >= 1 + assert "incident:A" in counts.scopes or "default" in counts.scopes + + def test_dry_run_max_chunks_none_returns_when_expired_exist(self) -> None: + """CLI ``--dry-run`` must COUNT once, not loop the same LIMIT ids.""" + settings = Settings( + _env_file=None, + retention_raw="30d", + retention_summary="0", + purge_chunk_size=1, + ) + db = ScriptedSession( + [ + [("default",)], + 5, + 2, + 1, + ] + ) + + def _fail_loop(*_a: object, **_k: object) -> None: + raise AssertionError("dry-run must not call write_last_purge_at") + + with patch( + "src.core.retention.purge.write_last_purge_at", side_effect=_fail_loop + ): + counts = run_purge( + db, + dry_run=True, + max_chunks=None, + now=NOW, + settings=settings, + ) + assert counts.raw == 5 + assert counts.embedding == 2 + assert db.script == [] + assert len(db.executed) <= 8 + + def test_invalid_policy_skips_scope_and_continues(self) -> None: + settings = Settings( + _env_file=None, + retention_raw="30d", + retention_summary="0", + purge_chunk_size=100, + ) + entry_id = uuid.uuid4() + db = ScriptedSession( + [ + [("incident:bad",), ("default",)], + [(entry_id,)], + 0, + 0, + None, + ] + ) + bad = MagicMock() + bad.raw_interval = "7 days" + bad.summary_interval = None + db._get_rows["incident:bad"] = bad + with patch("src.core.retention.purge.write_last_purge_at") as written: + counts = run_purge( + db, + dry_run=False, + max_chunks=1, + now=NOW, + settings=settings, + ) + written.assert_called_once() + assert counts.raw == 1 + assert "incident:bad" in counts.scopes + assert "default" in counts.scopes + + def test_skips_last_purge_at_when_more_chunks_remain(self) -> None: + settings = Settings( + _env_file=None, + retention_raw="30d", + retention_summary="0", + purge_chunk_size=1, + ) + entry_id = uuid.uuid4() + db = ScriptedSession( + [ + [("default",)], + [(entry_id,)], + 0, + 0, + None, + ] + ) + with ( + patch("src.core.retention.purge.write_last_purge_at") as written, + patch("src.core.retention.purge.enqueue_followup_purge") as followup, + ): + counts = run_purge( + db, + dry_run=False, + max_chunks=1, + now=NOW, + settings=settings, + ) + written.assert_not_called() + followup.assert_called_once() + assert counts.more_remaining is True + assert counts.raw == 1 + + +class TestPurgeMetrics: + def test_record_purge_rows_increments(self) -> None: + before_raw = _counter("raw") + before_emb = _counter("embedding") + record_purge_rows(3, kind="raw") + record_purge_rows(2, kind="embedding") + assert _counter("raw") == before_raw + 3 + assert _counter("embedding") == before_emb + 2 + + def test_raw_chunk_increments_metrics(self) -> None: + before_raw = _counter("raw") + before_emb = _counter("embedding") + entry_id = uuid.uuid4() + db = ScriptedSession([[(entry_id,)], 5, 3, None]) + purge_raw_chunk(db, "default", NOW, limit=10, dry_run=False) + assert _counter("raw") == before_raw + 1 + 3 + assert _counter("embedding") == before_emb + 5 + + def test_dry_run_does_not_increment_metrics(self) -> None: + before_raw = _counter("raw") + entry_id = uuid.uuid4() + db = ScriptedSession([[(entry_id,)], 1, 0]) + purge_raw_chunk(db, "default", NOW, limit=10, dry_run=True) + assert _counter("raw") == before_raw + + +class TestScheduler: + def test_due_when_never_ran(self) -> None: + assert ( + should_enqueue_purge( + last_purge_at=None, + has_active_purge=False, + now=NOW, + interval_seconds=3600, + ) + is True + ) + + def test_skip_when_recent(self) -> None: + assert ( + should_enqueue_purge( + last_purge_at=NOW - timedelta(minutes=10), + has_active_purge=False, + now=NOW, + interval_seconds=3600, + ) + is False + ) + + def test_due_when_older_than_interval(self) -> None: + assert ( + should_enqueue_purge( + last_purge_at=NOW - timedelta(hours=2), + has_active_purge=False, + now=NOW, + interval_seconds=3600, + ) + is True + ) + + def test_skip_when_active_purge(self) -> None: + assert ( + should_enqueue_purge( + last_purge_at=None, + has_active_purge=True, + now=NOW, + interval_seconds=3600, + ) + is False + ) + + def test_skip_when_interval_zero(self) -> None: + assert ( + should_enqueue_purge( + last_purge_at=None, + has_active_purge=False, + now=NOW, + interval_seconds=0, + ) + is False + ) + + def test_maybe_enqueue_adds_pending_job(self) -> None: + db = ScriptedSession([None]) + settings = Settings(_env_file=None, purge_interval_seconds=3600) + assert maybe_enqueue_purge(db, now=NOW, settings=settings) is True + assert len(db.added) == 1 + job = db.added[0] + assert isinstance(job, WorkerJob) + assert job.job_type == PURGE_JOB_TYPE + assert job.status == "pending" + + def test_maybe_enqueue_skips_when_active(self) -> None: + db = ScriptedSession([uuid.uuid4()]) + settings = Settings(_env_file=None, purge_interval_seconds=3600) + assert maybe_enqueue_purge(db, now=NOW, settings=settings) is False + assert db.added == [] diff --git a/tests/unit/test_settings.py b/tests/unit/test_settings.py index 45753a3..71ee0e1 100644 --- a/tests/unit/test_settings.py +++ b/tests/unit/test_settings.py @@ -216,3 +216,25 @@ def test_observability_settings_from_env(monkeypatch): assert settings.otel_sdk_disabled is True assert settings.otel_exporter_otlp_endpoint == "http://localhost:4318/v1/traces" assert settings.otel_service_name == "raglogs-test" + + +def test_retention_settings_defaults(): + settings = Settings(_env_file=None) + assert settings.retention_raw == "30d" + assert settings.retention_summary == "180d" + assert settings.purge_interval_seconds == 3600 + assert settings.purge_chunk_size == 1000 + + +def test_retention_settings_from_env(monkeypatch): + monkeypatch.setenv("RETENTION_RAW", "7d") + monkeypatch.setenv("RETENTION_SUMMARY", "off") + monkeypatch.setenv("PURGE_INTERVAL_SECONDS", "120") + monkeypatch.setenv("PURGE_CHUNK_SIZE", "50") + + settings = Settings(_env_file=None) + + assert settings.retention_raw == "7d" + assert settings.retention_summary == "off" + assert settings.purge_interval_seconds == 120 + assert settings.purge_chunk_size == 50 diff --git a/tests/unit/test_similar.py b/tests/unit/test_similar.py index e292db0..22882e5 100644 --- a/tests/unit/test_similar.py +++ b/tests/unit/test_similar.py @@ -325,6 +325,24 @@ def test_fingerprint_falls_through_to_log_entries(self) -> None: second_sql = _compiled(db.execute.call_args_list[1][0][0]) assert "log_entries" in second_sql + def test_fingerprint_matches_cluster_embeddings_without_log_entries(self) -> None: + """After raw purge, similar still matches fingerprints on cluster_embeddings.""" + db = MagicMock() + db.execute.return_value.all.return_value = [_row()] + matches = search_similar_fingerprint( + db, + [QueryCluster(fingerprint="abc123")], + query_scope="incident:A", + visibility=SimilarVisibility(cross_scope=True, visible_scope=None), + limit=10, + ) + assert len(matches) == 1 + assert matches[0].fingerprint == "abc123" + assert db.execute.call_count == 1 + sql = _compiled(db.execute.call_args.args[0]) + assert "cluster_embeddings" in sql + assert "log_entries" not in sql + class TestHelpers: def test_collect_query_fingerprints_dedupes(self) -> None: diff --git a/tests/unit/test_worker.py b/tests/unit/test_worker.py index 37ad1b8..b594938 100644 --- a/tests/unit/test_worker.py +++ b/tests/unit/test_worker.py @@ -445,3 +445,37 @@ def test_webhook_exception_does_not_fail_job(self): assert result is True assert job.status == "done" + + +class TestPurgeJob: + def test_purge_job_type_dispatches(self): + db = _mock_db() + job = _mock_worker_job(job_type="purge", payload={}) + db.execute.return_value.scalar_one_or_none.return_value = job + + with patch( + "src.core.retention.purge.run_purge_job", + return_value={"raw": 3, "summary": 0, "embedding": 1, "scopes": ["default"]}, + ) as mock_purge: + result = process_one(db) + + assert result is True + assert job.status == "done" + assert job.result_json["raw"] == 3 + mock_purge.assert_called_once_with(db, job) + + def test_purge_job_skips_ingest_webhook(self): + db = _mock_db() + job = _mock_worker_job(job_type="purge", payload={}) + db.execute.return_value.scalar_one_or_none.return_value = job + + with patch( + "src.core.retention.purge.run_purge_job", + return_value={"raw": 0, "summary": 0, "embedding": 0, "scopes": []}, + ), patch( + "src.core.ingestion.webhooks.maybe_deliver_ingest_callback", + ) as mock_deliver: + process_one(db) + + mock_deliver.assert_not_called() + assert job.status == "done"