diff --git a/.env.example b/.env.example index 27c6a0e..d09aa52 100644 --- a/.env.example +++ b/.env.example @@ -94,6 +94,17 @@ MAX_CLUSTERS_FOR_EXPLAIN=10 # Seconds the worker sleeps between polling for new jobs when the queue is empty. WORKER_POLL_INTERVAL=2 +# Ingest backpressure (minimal stand-in for full API/LLM rate limiting). +# POST /v1/ingestions and POST /v1/ingestions/lines return 429 + Retry-After +# when pending worker jobs >= INGEST_QUEUE_MAX. +INGEST_QUEUE_MAX=100 +INGEST_RETRY_AFTER_SECONDS=5 +# 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. +TAIL_POLL_INTERVAL=30 +TAIL_ERROR_THRESHOLD=5 + # ── Adapters ────────────────────────────────────────────────────────────────── # AWS region for the cloudwatch source adapter (raglogs ingest --adapter cloudwatch). diff --git a/README.md b/README.md index 843b687..ad0cc80 100644 --- a/README.md +++ b/README.md @@ -923,7 +923,7 @@ curl -X POST http://localhost:8000/v1/query/explain \ | Role | Allowed | |---|---| -| `ingest` | `POST /v1/ingestions` (and the deprecated `/ingestions` alias) | +| `ingest` | `POST /v1/ingestions`, `POST /v1/ingestions/lines`, tail lifecycle `pause` / `resume` / `stop` (and the deprecated `/ingestions` aliases) | | `query` | `GET /v1/ingestions*`, `POST /v1/query/*`, web UI (`GET /`, `/static`), OpenAPI (`/docs`) | | `admin` | everything, including `GET /v1/config` | @@ -949,10 +949,14 @@ 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) | -| `POST` | `/v1/ingestions` | Enqueue an ingest job (`adapter`: `file`, `cloudwatch`, `datadog`, `loki`, or `k8s`) | +| `GET` | `/health` | Service and DB health check (unversioned). Includes `tail_jobs: {running, paused}`. | +| `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 | +| `POST` | `/v1/ingestions/{id}:resume` | Resume a paused tail job | +| `POST` | `/v1/ingestions/{id}:stop` | Stop a tail job (terminal; cannot resume) | | `GET` | `/v1/ingestions` | List recent completed ingestion jobs, newest first | -| `GET` | `/v1/ingestions/{job_id}` | Poll ingestion job status | +| `GET` | `/v1/ingestions/{job_id}` | Fetch ingestion job detail | | `GET` | `/v1/ingestions/latest` | ID of the most recently completed ingestion job, if any | | `POST` | `/v1/query/explain` | Explain a time window | | `POST` | `/v1/query/ask` | Answer a natural language question | @@ -961,6 +965,33 @@ If auth is disabled and the process binds a non-loopback address (`0.0.0.0`, `:: | `POST` | `/v1/query/compare` | Diff two time windows (same semantics as `raglogs compare`) | | `GET` | `/v1/config` | Read effective configuration | +**Push NDJSON.** `POST /v1/ingestions/lines` accepts newline-delimited lines (`Content-Type: application/x-ndjson`, `application/jsonl`, or `text/plain`). Each line is a raw log string or a JSON object with at least `message` / `raw` / `text` (optional `timestamp`, `service`, `level`, `host`, `env`). Cap is `INGEST_PUSH_MAX_LINES` (default 5000); over the cap returns 400. + +```bash +curl -X POST http://localhost:8000/v1/ingestions/lines \ + -H "Content-Type: application/x-ndjson" \ + --data-binary $'{"message":"timeout talking to payments","level":"error","service":"api"}\nplain syslog line\n' +``` + +**Tail jobs.** `POST /v1/ingestions` with `"mode": "tail"` (adapters `cloudwatch`, `datadog`, or `loki` only) creates a long-lived job. The worker re-runs the adapter from the saved cursor about every `TAIL_POLL_INTERVAL` seconds (default 30). Pause, resume, or stop with: + +```bash +curl -X POST http://localhost:8000/v1/ingestions \ + -H "Content-Type: application/json" \ + -d '{"adapter":"loki","params":{"query":"{app=\\"api\\"}"},"mode":"tail"}' +# → { "ingestion_job_id": "...", "mode": "tail", "status": "running" } + +curl -X POST http://localhost:8000/v1/ingestions/$ID:pause +curl -X POST http://localhost:8000/v1/ingestions/$ID:resume +curl -X POST http://localhost:8000/v1/ingestions/$ID:stop +``` + +`stop` is terminal. After `TAIL_ERROR_THRESHOLD` consecutive poll failures (default 5) a tail job auto-pauses; `/health` reports `tail_jobs.running` and `tail_jobs.paused`. + +**Backpressure.** When pending worker jobs ≥ `INGEST_QUEUE_MAX` (default 100), `POST /v1/ingestions` and `POST /v1/ingestions/lines` return **429** with `Retry-After` (`INGEST_RETRY_AFTER_SECONDS`, default 5) and body `{"error_code":"INGEST_QUEUE_FULL","message":"..."}`. This is a queue-depth stand-in, not full API/LLM rate limiting. + +Overlapping tail poll windows can **double-count** the same physical lines until content dedup (G6) lands. Prefer a single tail job per source and avoid also batch-ingesting the same window. + Unversioned `/ingestions`, `/query/*`, and `/config` remain as **deprecated aliases** for one release. They behave the same as the `/v1` paths and send `Deprecation: true` plus a `Link: ; rel="successor-version"` header. `/health`, the web UI (`/`), and `/static` stay unversioned. **Compatibility policy.** Additive changes stay in `v1`. Breaking path or method removals require `v2`. JSON response bodies are unchanged in this release (a stable evidence schema is tracked separately as G7). diff --git a/clients/openapi.json b/clients/openapi.json index c1d6818..db2afe0 100644 --- a/clients/openapi.json +++ b/clients/openapi.json @@ -53,7 +53,7 @@ "ingestions" ], "summary": "Create Ingestion", - "description": "Enqueue an ingest job. Returns immediately with worker_job_id.\nPoll GET /ingestions/jobs/{worker_job_id} for progress.\nWhen done, use ingestion_job_id from the result to scope /query/explain.", + "description": "Enqueue an ingest job. Returns immediately with worker_job_id.\nPoll GET /ingestions/jobs/{worker_job_id} for progress.\nWhen done, use ingestion_job_id from the result to scope /query/explain.\n\n``mode=tail`` creates a long-lived job that the worker polls; the response\nincludes ``ingestion_job_id`` and no worker_job_id.", "operationId": "v1_ingestions_create_ingestion__v1_ingestions", "requestBody": { "content": { @@ -154,6 +154,174 @@ } } }, + "/v1/ingestions/lines": { + "post": { + "tags": [ + "ingestions" + ], + "summary": "Push Ingestion Lines", + "description": "Push NDJSON of raw or pre-parsed log lines. Persisted synchronously.", + "operationId": "v1_ingestions_push_ingestion_lines__v1_ingestions_lines", + "requestBody": { + "content": { + "application/x-ndjson": { + "schema": { + "type": "string" + } + }, + "application/jsonl": { + "schema": { + "type": "string" + } + }, + "text/plain": { + "schema": { + "type": "string" + } + } + }, + "required": true + }, + "responses": { + "200": { + "description": "Successful Response", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/PushLinesResponse" + } + } + } + } + } + } + }, + "/v1/ingestions/{ingestion_job_id}:pause": { + "post": { + "tags": [ + "ingestions" + ], + "summary": "Pause Tail Job", + "operationId": "v1_ingestions_pause_tail_job__v1_ingestions__ingestion_job_id__pause", + "parameters": [ + { + "name": "ingestion_job_id", + "in": "path", + "required": true, + "schema": { + "type": "string", + "title": "Ingestion Job Id" + } + } + ], + "responses": { + "200": { + "description": "Successful Response", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/TailLifecycleResponse" + } + } + } + }, + "422": { + "description": "Validation Error", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/HTTPValidationError" + } + } + } + } + } + } + }, + "/v1/ingestions/{ingestion_job_id}:resume": { + "post": { + "tags": [ + "ingestions" + ], + "summary": "Resume Tail Job", + "operationId": "v1_ingestions_resume_tail_job__v1_ingestions__ingestion_job_id__resume", + "parameters": [ + { + "name": "ingestion_job_id", + "in": "path", + "required": true, + "schema": { + "type": "string", + "title": "Ingestion Job Id" + } + } + ], + "responses": { + "200": { + "description": "Successful Response", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/TailLifecycleResponse" + } + } + } + }, + "422": { + "description": "Validation Error", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/HTTPValidationError" + } + } + } + } + } + } + }, + "/v1/ingestions/{ingestion_job_id}:stop": { + "post": { + "tags": [ + "ingestions" + ], + "summary": "Stop Tail Job", + "operationId": "v1_ingestions_stop_tail_job__v1_ingestions__ingestion_job_id__stop", + "parameters": [ + { + "name": "ingestion_job_id", + "in": "path", + "required": true, + "schema": { + "type": "string", + "title": "Ingestion Job Id" + } + } + ], + "responses": { + "200": { + "description": "Successful Response", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/TailLifecycleResponse" + } + } + } + }, + "422": { + "description": "Validation Error", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/HTTPValidationError" + } + } + } + } + } + } + }, "/v1/ingestions/{ingestion_job_id}": { "get": { "tags": [ @@ -224,7 +392,7 @@ "ingestions" ], "summary": "Create Ingestion", - "description": "Enqueue an ingest job. Returns immediately with worker_job_id.\nPoll GET /ingestions/jobs/{worker_job_id} for progress.\nWhen done, use ingestion_job_id from the result to scope /query/explain.", + "description": "Enqueue an ingest job. Returns immediately with worker_job_id.\nPoll GET /ingestions/jobs/{worker_job_id} for progress.\nWhen done, use ingestion_job_id from the result to scope /query/explain.\n\n``mode=tail`` creates a long-lived job that the worker polls; the response\nincludes ``ingestion_job_id`` and no worker_job_id.", "operationId": "legacy_ingestions_create_ingestion__ingestions", "requestBody": { "content": { @@ -328,6 +496,178 @@ "deprecated": true } }, + "/ingestions/lines": { + "post": { + "tags": [ + "ingestions" + ], + "summary": "Push Ingestion Lines", + "description": "Push NDJSON of raw or pre-parsed log lines. Persisted synchronously.", + "operationId": "legacy_ingestions_push_ingestion_lines__ingestions_lines", + "requestBody": { + "content": { + "application/x-ndjson": { + "schema": { + "type": "string" + } + }, + "application/jsonl": { + "schema": { + "type": "string" + } + }, + "text/plain": { + "schema": { + "type": "string" + } + } + }, + "required": true + }, + "responses": { + "200": { + "description": "Successful Response", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/PushLinesResponse" + } + } + } + } + }, + "deprecated": true + } + }, + "/ingestions/{ingestion_job_id}:pause": { + "post": { + "tags": [ + "ingestions" + ], + "summary": "Pause Tail Job", + "operationId": "legacy_ingestions_pause_tail_job__ingestions__ingestion_job_id__pause", + "deprecated": true, + "parameters": [ + { + "name": "ingestion_job_id", + "in": "path", + "required": true, + "schema": { + "type": "string", + "title": "Ingestion Job Id" + } + } + ], + "responses": { + "200": { + "description": "Successful Response", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/TailLifecycleResponse" + } + } + } + }, + "422": { + "description": "Validation Error", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/HTTPValidationError" + } + } + } + } + } + } + }, + "/ingestions/{ingestion_job_id}:resume": { + "post": { + "tags": [ + "ingestions" + ], + "summary": "Resume Tail Job", + "operationId": "legacy_ingestions_resume_tail_job__ingestions__ingestion_job_id__resume", + "deprecated": true, + "parameters": [ + { + "name": "ingestion_job_id", + "in": "path", + "required": true, + "schema": { + "type": "string", + "title": "Ingestion Job Id" + } + } + ], + "responses": { + "200": { + "description": "Successful Response", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/TailLifecycleResponse" + } + } + } + }, + "422": { + "description": "Validation Error", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/HTTPValidationError" + } + } + } + } + } + } + }, + "/ingestions/{ingestion_job_id}:stop": { + "post": { + "tags": [ + "ingestions" + ], + "summary": "Stop Tail Job", + "operationId": "legacy_ingestions_stop_tail_job__ingestions__ingestion_job_id__stop", + "deprecated": true, + "parameters": [ + { + "name": "ingestion_job_id", + "in": "path", + "required": true, + "schema": { + "type": "string", + "title": "Ingestion Job Id" + } + } + ], + "responses": { + "200": { + "description": "Successful Response", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/TailLifecycleResponse" + } + } + } + }, + "422": { + "description": "Validation Error", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/HTTPValidationError" + } + } + } + } + } + } + }, "/ingestions/{ingestion_job_id}": { "get": { "tags": [ @@ -1084,17 +1424,39 @@ "EnqueuedResponse": { "properties": { "worker_job_id": { - "type": "string", + "anyOf": [ + { + "type": "string" + }, + { + "type": "null" + } + ], "title": "Worker Job Id" }, "status": { "type": "string", "title": "Status" + }, + "ingestion_job_id": { + "anyOf": [ + { + "type": "string" + }, + { + "type": "null" + } + ], + "title": "Ingestion Job Id" + }, + "mode": { + "type": "string", + "title": "Mode", + "default": "batch" } }, "type": "object", "required": [ - "worker_job_id", "status" ], "title": "EnqueuedResponse" @@ -1248,6 +1610,20 @@ }, "type": "object", "title": "Adapters" + }, + "tail_jobs": { + "anyOf": [ + { + "additionalProperties": { + "type": "integer" + }, + "type": "object" + }, + { + "type": "null" + } + ], + "title": "Tail Jobs" } }, "type": "object", @@ -1373,6 +1749,15 @@ "type": "boolean", "title": "With Embeddings", "default": false + }, + "mode": { + "type": "string", + "enum": [ + "batch", + "tail" + ], + "title": "Mode", + "default": "batch" } }, "type": "object", @@ -1455,6 +1840,27 @@ } ], "title": "Finished At" + }, + "mode": { + "type": "string", + "title": "Mode", + "default": "batch" + }, + "last_polled_at": { + "anyOf": [ + { + "type": "string" + }, + { + "type": "null" + } + ], + "title": "Last Polled At" + }, + "consecutive_errors": { + "type": "integer", + "title": "Consecutive Errors", + "default": 0 } }, "type": "object", @@ -1545,6 +1951,67 @@ ], "title": "LatestIngestionResponse" }, + "PushLinesResponse": { + "properties": { + "ingestion_job_id": { + "type": "string", + "title": "Ingestion Job Id" + }, + "status": { + "type": "string", + "title": "Status" + }, + "line_count": { + "type": "integer", + "title": "Line Count" + }, + "parsed_count": { + "type": "integer", + "title": "Parsed Count" + }, + "error_count": { + "type": "integer", + "title": "Error Count" + }, + "mode": { + "type": "string", + "title": "Mode", + "default": "push" + } + }, + "type": "object", + "required": [ + "ingestion_job_id", + "status", + "line_count", + "parsed_count", + "error_count" + ], + "title": "PushLinesResponse" + }, + "TailLifecycleResponse": { + "properties": { + "ingestion_job_id": { + "type": "string", + "title": "Ingestion Job Id" + }, + "mode": { + "type": "string", + "title": "Mode" + }, + "status": { + "type": "string", + "title": "Status" + } + }, + "type": "object", + "required": [ + "ingestion_job_id", + "mode", + "status" + ], + "title": "TailLifecycleResponse" + }, "TimelineRequest": { "properties": { "since": { diff --git a/migrations/versions/0005_ingest_modes.py b/migrations/versions/0005_ingest_modes.py new file mode 100644 index 0000000..a1d42a7 --- /dev/null +++ b/migrations/versions/0005_ingest_modes.py @@ -0,0 +1,42 @@ +"""ingestion job modes for batch / push / tail + +Revision ID: 0005_ingest_modes +Revises: 0004_api_keys +Create Date: 2026-08-17 00:00:00.000000 +""" + +from typing import Sequence, Union + +import sqlalchemy as sa + +from alembic import op + +revision: str = "0005_ingest_modes" +down_revision: Union[str, None] = "0004_api_keys" +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + op.add_column( + "ingestion_jobs", + sa.Column("mode", sa.String(20), nullable=False, server_default="batch"), + ) + op.add_column("ingestion_jobs", sa.Column("cursor", sa.Text(), nullable=True)) + op.add_column( + "ingestion_jobs", + sa.Column("last_polled_at", sa.DateTime(timezone=True), nullable=True), + ) + op.add_column( + "ingestion_jobs", + sa.Column( + "consecutive_errors", sa.Integer(), nullable=False, server_default="0" + ), + ) + + +def downgrade() -> None: + op.drop_column("ingestion_jobs", "consecutive_errors") + op.drop_column("ingestion_jobs", "last_polled_at") + op.drop_column("ingestion_jobs", "cursor") + op.drop_column("ingestion_jobs", "mode") diff --git a/src/api/routes/health.py b/src/api/routes/health.py index e86270b..45fe0f1 100644 --- a/src/api/routes/health.py +++ b/src/api/routes/health.py @@ -17,6 +17,7 @@ class HealthResponse(BaseModel): db: str # connected | disconnected worker_queue_depth: Optional[int] # pending worker jobs; None if DB unreachable adapters: dict[str, str] # adapter name -> "ok" | "unavailable: " + tail_jobs: Optional[dict[str, int]] = None # {running, paused}; None if DB unreachable def _adapter_health() -> dict[str, str]: @@ -59,21 +60,32 @@ def health_check(): adapters = _adapter_health() if not db_ok: - return HealthResponse(status="degraded", db="disconnected", worker_queue_depth=None, adapters=adapters) + return HealthResponse( + status="degraded", + db="disconnected", + worker_queue_depth=None, + adapters=adapters, + tail_jobs=None, + ) try: - from src.db.models import WorkerJob from sqlalchemy import func, select + + from src.core.ingestion.tail import tail_job_counts + from src.db.models import WorkerJob with get_db() as db: depth = db.execute( select(func.count()).select_from(WorkerJob).where(WorkerJob.status == "pending") ).scalar_one() + tail_jobs = tail_job_counts(db) except Exception: depth = None + tail_jobs = None return HealthResponse( status="ok" if db_ok else "degraded", db="connected", worker_queue_depth=depth, adapters=adapters, + tail_jobs=tail_jobs, ) diff --git a/src/api/routes/ingestions.py b/src/api/routes/ingestions.py index f3b0b9f..dfa713b 100644 --- a/src/api/routes/ingestions.py +++ b/src/api/routes/ingestions.py @@ -2,23 +2,30 @@ Ingestion API routes. POST /ingestions — enqueue async ingest, return worker_job_id immediately +POST /ingestions/lines — push NDJSON lines (sync persist) +POST /ingestions/{id}:pause|resume|stop — tail job lifecycle GET /ingestions — list recent completed ingestion jobs (newest first) GET /ingestions/jobs/{id} — poll worker job status (pending|running|done|failed) GET /ingestions/latest — most recently completed ingestion job, if any -GET /ingestions/{id} — fetch completed IngestionJob detail by ingestion_job_id +GET /ingestions/{id} — fetch IngestionJob detail by ingestion_job_id """ + import uuid -from datetime import datetime -from typing import Any, Optional +from datetime import datetime, timezone +from typing import Any, Literal, Optional -from fastapi import APIRouter, HTTPException +from fastapi import APIRouter, HTTPException, Request +from fastapi.responses import JSONResponse from pydantic import BaseModel, model_validator +from src.core.ingestion.tail import TAIL_ADAPTERS + router = APIRouter() # ── Schemas ────────────────────────────────────────────────────────────────── + class IngestRequest(BaseModel): paths: list[str] = [] recursive: bool = False @@ -37,9 +44,16 @@ class IngestRequest(BaseModel): to_time: Optional[datetime] = None resume_ingestion_job_id: Optional[str] = None with_embeddings: bool = False + mode: Literal["batch", "tail"] = "batch" @model_validator(mode="after") def _require_paths_for_file_adapter(self): + if self.mode == "tail": + if self.adapter not in TAIL_ADAPTERS: + raise ValueError( + "mode='tail' requires adapter cloudwatch, datadog, or loki" + ) + return self if self.adapter == "file" and not self.paths: raise ValueError("paths is required when adapter='file'") if self.adapter in ("k8s", "kubernetes"): @@ -50,14 +64,16 @@ def _require_paths_for_file_adapter(self): class EnqueuedResponse(BaseModel): - worker_job_id: str - status: str # always "pending" at enqueue time + worker_job_id: Optional[str] = None # set for batch enqueue; null for tail + status: str # pending (batch) | running (tail) + ingestion_job_id: Optional[str] = None + mode: str = "batch" class WorkerJobStatus(BaseModel): worker_job_id: str - status: str # pending | running | done | failed - ingestion_job_id: Optional[str] # set once worker completes ingestion + status: str # pending | running | done | failed + ingestion_job_id: Optional[str] # set once worker completes ingestion error: Optional[str] created_at: str started_at: Optional[str] @@ -78,6 +94,9 @@ class IngestionJobDetail(BaseModel): created_at: str started_at: Optional[str] finished_at: Optional[str] + mode: str = "batch" + last_polled_at: Optional[str] = None + consecutive_errors: int = 0 class LatestIngestionResponse(BaseModel): @@ -95,65 +114,100 @@ class IngestionListResponse(BaseModel): ingestions: list[IngestionSummary] -# ── Routes ─────────────────────────────────────────────────────────────────── +class PushLinesResponse(BaseModel): + ingestion_job_id: str + status: str + line_count: int + parsed_count: int + error_count: int + mode: str = "push" + + +class TailLifecycleResponse(BaseModel): + ingestion_job_id: str + mode: str + status: str + + +# ── Helpers ────────────────────────────────────────────────────────────────── + + +def _queue_full_response() -> JSONResponse: + from src.config import get_settings + + settings = get_settings() + return JSONResponse( + status_code=429, + content={ + "error_code": "INGEST_QUEUE_FULL", + "message": "Ingest worker queue is full; retry after the Retry-After delay.", + }, + headers={"Retry-After": str(settings.ingest_retry_after_seconds)}, + ) -@router.post("", response_model=EnqueuedResponse, status_code=202) -def create_ingestion(request: IngestRequest): - """ - Enqueue an ingest job. Returns immediately with worker_job_id. - Poll GET /ingestions/jobs/{worker_job_id} for progress. - When done, use ingestion_job_id from the result to scope /query/explain. - """ - from src.db.models import WorkerJob - from src.db.session import get_db + +def _reject_if_queue_full(db: Any) -> Optional[JSONResponse]: + from src.config import get_settings + from src.core.ingestion.backpressure import ( + ingest_queue_is_full, + pending_worker_job_count, + ) + + settings = get_settings() + pending = pending_worker_job_count(db) + if ingest_queue_is_full(pending, settings.ingest_queue_max): + return _queue_full_response() + return None + + +def _validate_non_file_adapter(request: IngestRequest) -> dict[str, Any]: + """Discover + window-validate a pull adapter. Returns possibly rewritten params.""" + from src.adapters.base import SourceSpec + from src.adapters.registry import get_adapter + from src.config import get_settings + from src.core.errors import AdapterUnavailableError params = request.params - if request.adapter == "file": - from src.adapters.file.adapter import discover_files + if request.adapter in ("k8s", "kubernetes"): + from src.adapters.k8s.adapter import build_k8s_params - # Validate paths before enqueuing — fail fast - files = discover_files(request.paths, recursive=request.recursive) - if not files: - raise HTTPException(status_code=400, detail="No log files found at the given paths") - else: - from src.adapters.base import SourceSpec - from src.adapters.registry import get_adapter - from src.config import get_settings - from src.core.errors import AdapterUnavailableError + params = build_k8s_params( + request.params, paths=request.paths, recursive=request.recursive + ) - if request.adapter in ("k8s", "kubernetes"): - from src.adapters.k8s.adapter import build_k8s_params + try: + adapter = get_adapter(request.adapter, get_settings()) + refs = list( + adapter.discover(SourceSpec(adapter=request.adapter, params=params)) + ) + except AdapterUnavailableError as e: + raise HTTPException( + status_code=400, + detail={"error_code": e.error_code, "message": str(e)}, + ) - params = build_k8s_params( - request.params, paths=request.paths, recursive=request.recursive - ) + if request.adapter in ("k8s", "kubernetes") and not refs: + raise HTTPException( + status_code=400, detail="No log files found at the given paths" + ) + + if request.since or request.from_time or request.to_time: + from src.utils.time import resolve_window try: - adapter = get_adapter(request.adapter, get_settings()) - # Fail fast on missing/invalid params (e.g. no log_group / query) — - # discover() for pull adapters is local-only validation, no network call. - refs = list(adapter.discover(SourceSpec(adapter=request.adapter, params=params))) - except AdapterUnavailableError as e: - raise HTTPException( - status_code=400, - detail={"error_code": e.error_code, "message": str(e)}, + resolve_window( + since=request.since, + from_time=request.from_time, + to_time=request.to_time, ) + except ValueError as e: + raise HTTPException(status_code=400, detail=str(e)) - if request.adapter in ("k8s", "kubernetes") and not refs: - raise HTTPException(status_code=400, detail="No log files found at the given paths") - - if request.since or request.from_time or request.to_time: - # Same validation POST /query/explain does at the route level — otherwise - # an unparseable `since` (e.g. "not-a-duration") 202s here and only fails - # once the worker picks it up. - from src.utils.time import resolve_window + return params - try: - resolve_window(since=request.since, from_time=request.from_time, to_time=request.to_time) - except ValueError as e: - raise HTTPException(status_code=400, detail=str(e)) - payload = { +def _ingest_payload(request: IngestRequest, params: dict[str, Any]) -> dict[str, Any]: + return { "paths": request.paths, "recursive": request.recursive, "format": request.format, @@ -167,9 +221,110 @@ def create_ingestion(request: IngestRequest): "to_time": request.to_time.isoformat() if request.to_time else None, "resume_ingestion_job_id": request.resume_ingestion_job_id, "with_embeddings": request.with_embeddings, + "mode": request.mode, } + +def _create_tail_job( + request: IngestRequest, params: dict[str, Any], db: Any +) -> EnqueuedResponse: + from src.core.ingestion.service import _get_or_create_source + from src.db.models import IngestionJob + + source_name = request.source_name or f"tail:{request.adapter}" + source = _get_or_create_source(db, name=source_name, type_=request.adapter) + now = datetime.now(tz=timezone.utc) + job = IngestionJob( + id=uuid.uuid4(), + source_id=source.id, + status="running", + started_at=now, + metadata_json=_ingest_payload(request, params), + source_adapter=request.adapter, + source_ref=source_name, + mode="tail", + consecutive_errors=0, + ) + db.add(job) + db.flush() + return EnqueuedResponse( + worker_job_id=None, + status="running", + ingestion_job_id=str(job.id), + mode="tail", + ) + + +def _parse_ingestion_job_id(ingestion_job_id: str) -> uuid.UUID: + try: + return uuid.UUID(ingestion_job_id) + except ValueError: + raise HTTPException(status_code=400, detail="Invalid ingestion_job_id") + + +def _apply_tail_action( + ingestion_job_id: str, action: Literal["pause", "resume", "stop"] +) -> TailLifecycleResponse: + from src.core.ingestion.tail import TailLifecycleError, apply_tail_lifecycle + from src.db.models import IngestionJob + from src.db.session import get_db + + ij_uuid = _parse_ingestion_job_id(ingestion_job_id) with get_db() as db: + job = db.query(IngestionJob).filter(IngestionJob.id == ij_uuid).first() + if not job: + raise HTTPException(status_code=404, detail="Ingestion job not found") + try: + job.status = apply_tail_lifecycle(job.mode, job.status, action) + except TailLifecycleError as e: + raise HTTPException(status_code=409, detail=str(e)) + if action == "stop": + job.finished_at = datetime.now(tz=timezone.utc) + db.flush() + return TailLifecycleResponse( + ingestion_job_id=str(job.id), + mode=job.mode, + status=job.status, + ) + + +# ── Routes ─────────────────────────────────────────────────────────────────── + + +@router.post("", response_model=EnqueuedResponse, status_code=202) +def create_ingestion(request: IngestRequest): + """ + Enqueue an ingest job. Returns immediately with worker_job_id. + Poll GET /ingestions/jobs/{worker_job_id} for progress. + When done, use ingestion_job_id from the result to scope /query/explain. + + ``mode=tail`` creates a long-lived job that the worker polls; the response + includes ``ingestion_job_id`` and no worker_job_id. + """ + from src.db.models import WorkerJob + from src.db.session import get_db + + params = request.params + if request.adapter == "file": + from src.adapters.file.adapter import discover_files + + files = discover_files(request.paths, recursive=request.recursive) + if not files: + raise HTTPException( + status_code=400, detail="No log files found at the given paths" + ) + else: + params = _validate_non_file_adapter(request) + + with get_db() as db: + rejected = _reject_if_queue_full(db) + if rejected is not None: + return rejected + + if request.mode == "tail": + return _create_tail_job(request, params, db) + + payload = _ingest_payload(request, params) job = WorkerJob( id=uuid.uuid4(), job_type="ingest", @@ -180,7 +335,7 @@ def create_ingestion(request: IngestRequest): db.flush() worker_job_id = str(job.id) - return EnqueuedResponse(worker_job_id=worker_job_id, status="pending") + return EnqueuedResponse(worker_job_id=worker_job_id, status="pending", mode="batch") @router.get("", response_model=IngestionListResponse) @@ -206,7 +361,9 @@ def list_ingestions(): ingestion_job_id=str(job.id), source_name=source_name, parsed_count=job.parsed_count, - finished_at=job.finished_at.isoformat() if job.finished_at else None, + finished_at=job.finished_at.isoformat() + if job.finished_at + else None, ) for job, source_name in rows ] @@ -235,7 +392,9 @@ def get_worker_job_status(worker_job_id: str): return WorkerJobStatus( worker_job_id=str(job.id), status=job.status, - ingestion_job_id=str(job.ingestion_job_id) if job.ingestion_job_id else None, + ingestion_job_id=str(job.ingestion_job_id) + if job.ingestion_job_id + else None, error=job.error, created_at=job.created_at.isoformat(), started_at=job.started_at.isoformat() if job.started_at else None, @@ -265,22 +424,86 @@ def get_latest_ingestion(): return LatestIngestionResponse(ingestion_job_id=str(job_id) if job_id else None) +@router.post( + "/lines", + response_model=PushLinesResponse, + status_code=200, + openapi_extra={ + "requestBody": { + "required": True, + "content": { + "application/x-ndjson": {"schema": {"type": "string"}}, + "application/jsonl": {"schema": {"type": "string"}}, + "text/plain": {"schema": {"type": "string"}}, + }, + } + }, +) +async def push_ingestion_lines(request: Request) -> PushLinesResponse: + """Push NDJSON of raw or pre-parsed log lines. Persisted synchronously.""" + from src.config import get_settings + from src.core.ingestion.push import NdjsonParseError, parse_ndjson_payload + from src.core.ingestion.service import ingest_push_lines + from src.db.session import get_db + + settings = get_settings() + body = (await request.body()).decode("utf-8", errors="replace") + try: + raw_lines = parse_ndjson_payload(body, settings.ingest_push_max_lines) + except NdjsonParseError as e: + raise HTTPException(status_code=400, detail=str(e)) + + with get_db() as db: + rejected = _reject_if_queue_full(db) + if rejected is not None: + return rejected # type: ignore[return-value] + job, stats = ingest_push_lines(db, raw_lines) + + return PushLinesResponse( + ingestion_job_id=str(job.id), + status=job.status, + line_count=stats.lines_read, + parsed_count=stats.parsed_count, + error_count=stats.error_count, + mode="push", + ) + + +@router.get("/lines", include_in_schema=False) +def push_lines_get_not_allowed() -> None: + """Static /lines must not fall through to /{ingestion_job_id} UUID parsing.""" + raise HTTPException(status_code=405, detail="Method Not Allowed") + + +@router.post("/{ingestion_job_id}:pause", response_model=TailLifecycleResponse) +def pause_tail_job(ingestion_job_id: str) -> TailLifecycleResponse: + return _apply_tail_action(ingestion_job_id, "pause") + + +@router.post("/{ingestion_job_id}:resume", response_model=TailLifecycleResponse) +def resume_tail_job(ingestion_job_id: str) -> TailLifecycleResponse: + return _apply_tail_action(ingestion_job_id, "resume") + + +@router.post("/{ingestion_job_id}:stop", response_model=TailLifecycleResponse) +def stop_tail_job(ingestion_job_id: str) -> TailLifecycleResponse: + return _apply_tail_action(ingestion_job_id, "stop") + + @router.get("/{ingestion_job_id}", response_model=IngestionJobDetail) def get_ingestion_detail(ingestion_job_id: str): """Fetch an IngestionJob record by its ID.""" from src.db.models import IngestionJob from src.db.session import get_db - try: - ij_uuid = uuid.UUID(ingestion_job_id) - except ValueError: - raise HTTPException(status_code=400, detail="Invalid ingestion_job_id") + ij_uuid = _parse_ingestion_job_id(ingestion_job_id) with get_db() as db: job = db.query(IngestionJob).filter(IngestionJob.id == ij_uuid).first() if not job: raise HTTPException(status_code=404, detail="Ingestion job not found") + polled = getattr(job, "last_polled_at", None) return IngestionJobDetail( ingestion_job_id=str(job.id), status=job.status, @@ -294,4 +517,7 @@ def get_ingestion_detail(ingestion_job_id: str): created_at=job.created_at.isoformat(), started_at=job.started_at.isoformat() if job.started_at else None, finished_at=job.finished_at.isoformat() if job.finished_at else None, + mode=getattr(job, "mode", None) or "batch", + last_polled_at=polled.isoformat() if isinstance(polled, datetime) else None, + consecutive_errors=int(getattr(job, "consecutive_errors", 0) or 0), ) diff --git a/src/clients/v1.py b/src/clients/v1.py index 559392a..7fda40e 100644 --- a/src/clients/v1.py +++ b/src/clients/v1.py @@ -111,6 +111,40 @@ def get_ingestion(self, ingestion_job_id: str) -> dict[str, Any]: """GET /v1/ingestions/{ingestion_job_id}.""" return self._get(f"/v1/ingestions/{ingestion_job_id}") + def push_lines( + self, + body: str, + content_type: str = "application/x-ndjson", + ) -> dict[str, Any]: + """POST /v1/ingestions/lines — NDJSON of raw or pre-parsed log lines.""" + headers = self._headers() + headers["Content-Type"] = content_type + response = self._client.request( + "POST", + self._url("/v1/ingestions/lines"), + headers=headers, + content=body, + ) + if response.status_code >= 400: + try: + err_body: Any = response.json() + except ValueError: + err_body = response.text + raise RaglogsAPIError(response.status_code, err_body) + return response.json() + + def pause_ingestion(self, ingestion_job_id: str) -> dict[str, Any]: + """POST /v1/ingestions/{id}:pause — pause a tail job.""" + return self._post(f"/v1/ingestions/{ingestion_job_id}:pause", {}) + + def resume_ingestion(self, ingestion_job_id: str) -> dict[str, Any]: + """POST /v1/ingestions/{id}:resume — resume a paused tail job.""" + return self._post(f"/v1/ingestions/{ingestion_job_id}:resume", {}) + + def stop_ingestion(self, ingestion_job_id: str) -> dict[str, Any]: + """POST /v1/ingestions/{id}:stop — stop a tail job (terminal).""" + return self._post(f"/v1/ingestions/{ingestion_job_id}:stop", {}) + def explain( self, *, diff --git a/src/config/settings.py b/src/config/settings.py index c84c8de..8f80768 100644 --- a/src/config/settings.py +++ b/src/config/settings.py @@ -50,6 +50,13 @@ class Settings(BaseSettings): # Worker worker_poll_interval: int = 2 # seconds between idle polls + # Ingest backpressure (minimal G9 stand-in) and push/tail (G4) + ingest_queue_max: int = 100 + ingest_retry_after_seconds: int = 5 + ingest_push_max_lines: int = 5000 + tail_poll_interval: int = 30 + tail_error_threshold: int = 5 + # Adapters adapter_cloudwatch_region: str = "us-east-1" datadog_api_key: str = "" diff --git a/src/core/ingestion/backpressure.py b/src/core/ingestion/backpressure.py new file mode 100644 index 0000000..45b9efb --- /dev/null +++ b/src/core/ingestion/backpressure.py @@ -0,0 +1,19 @@ +"""Ingest worker-queue backpressure (minimal G9 stand-in).""" + +from sqlalchemy import func, select +from sqlalchemy.orm import Session + +from src.db.models import WorkerJob + + +def pending_worker_job_count(db: Session) -> int: + """Return the number of WorkerJobs waiting to be claimed.""" + count = db.execute( + select(func.count()).select_from(WorkerJob).where(WorkerJob.status == "pending") + ).scalar_one() + return int(count or 0) + + +def ingest_queue_is_full(pending_count: int, max_depth: int) -> bool: + """True when enqueue/push should reject with 429.""" + return pending_count >= max_depth diff --git a/src/core/ingestion/push.py b/src/core/ingestion/push.py new file mode 100644 index 0000000..5287c62 --- /dev/null +++ b/src/core/ingestion/push.py @@ -0,0 +1,64 @@ +"""NDJSON parsing for POST /ingestions/lines. + +Each non-empty line is either a raw log string or a JSON value: +- JSON object with ``message`` / ``raw`` / ``text`` (plus optional timestamp, + service, level, host, env) → serialized as a JSON log line +- JSON string → the unquoted string is the raw line +- anything else → treated as a raw text log line +""" + +from __future__ import annotations + +from typing import Any + +import orjson + + +class NdjsonParseError(ValueError): + """Invalid or oversized NDJSON payload.""" + + +def parse_ndjson_payload(body: str, max_lines: int) -> list[str]: + """Split an NDJSON body into raw lines ready for ``_process_line``. + + Raises ``NdjsonParseError`` when the payload is empty or exceeds ``max_lines``. + Blank lines are ignored and do not count toward the cap. + """ + if max_lines < 1: + raise NdjsonParseError("max_lines must be >= 1") + + physical: list[str] = [line for line in body.splitlines() if line.strip()] + if not physical: + raise NdjsonParseError("NDJSON body is empty") + if len(physical) > max_lines: + raise NdjsonParseError(f"payload exceeds INGEST_PUSH_MAX_LINES ({max_lines})") + + return [_coerce_ndjson_line(line) for line in physical] + + +def _coerce_ndjson_line(line: str) -> str: + stripped = line.strip() + try: + parsed: Any = orjson.loads(stripped) + except orjson.JSONDecodeError: + return stripped + + if isinstance(parsed, str): + return parsed + + if isinstance(parsed, dict): + return _object_to_json_line(parsed) + + # Arrays / numbers / bools / null — keep the original text so the + # existing parsers can reject or stringify them. + return stripped + + +def _object_to_json_line(obj: dict[str, Any]) -> str: + """Ensure a JSON object has a ``message`` field the json parser can use.""" + if "message" not in obj: + for key in ("raw", "text", "msg"): + if key in obj and obj[key] is not None: + obj = {**obj, "message": obj[key]} + break + return orjson.dumps(obj).decode() diff --git a/src/core/ingestion/service.py b/src/core/ingestion/service.py index f3fc6fb..a01f7d8 100644 --- a/src/core/ingestion/service.py +++ b/src/core/ingestion/service.py @@ -175,7 +175,9 @@ def ingest_files( raise ValueError(f"No files found for paths: {paths}") # Create or find source - source = _get_or_create_source(db, name=source_name or ", ".join(paths[:2]), type_="file") + source = _get_or_create_source( + db, name=source_name or ", ".join(paths[:2]), type_="file" + ) # Create ingestion job job = IngestionJob( @@ -187,6 +189,7 @@ def ingest_files( metadata_json={"paths": paths}, source_adapter="file", source_ref=", ".join(str(p) for p in files[:5]), + mode="batch", ) db.add(job) db.flush() @@ -208,8 +211,15 @@ def ingest_files( for line in read_lines(file_path): entry = _process_line( - line, file_fmt, file_service, default_env, - source, job, "file", str(file_path), stats, + line, + file_fmt, + file_service, + default_env, + source, + job, + "file", + str(file_path), + stats, ) if entry is None: continue @@ -256,6 +266,8 @@ def ingest_from_source( resume_completed_streams: Optional[Iterable[str]] = None, progress_callback: Optional[Callable[[int, int], None]] = None, with_embeddings: bool = False, + existing_job: Optional[IngestionJob] = None, + finalize: bool = True, ) -> tuple[IngestionJob, IngestionStats]: """ Adapter-driven ingestion entry point (e.g. CloudWatch, Loki). Discovers streams via the @@ -277,6 +289,11 @@ def ingest_from_source( entirely rather than re-read from the start of the window — a completed stream's saved cursor is None (exhausted), and applying that as a "resume from" token would otherwise restart it and duplicate every row it already produced. + + ``existing_job`` appends lines onto a long-lived job (tail mode) instead of + creating a new IngestionJob. When ``finalize`` is False the job is left + running: status is not set to completed/failed, streams are not marked + completed, and counts are additive. Tail ticks pass both. """ import time @@ -300,33 +317,52 @@ def ingest_from_source( stats.files_processed = len(refs) if not refs: - raise ValueError(f"No streams discovered for adapter={spec.adapter!r} params={spec.params!r}") + raise ValueError( + f"No streams discovered for adapter={spec.adapter!r} params={spec.params!r}" + ) completed_at_start = set(resume_completed_streams or ()) if resume_cursors: for ref in refs: - if ref.stream_id in resume_cursors and ref.stream_id not in completed_at_start: + if ( + ref.stream_id in resume_cursors + and ref.stream_id not in completed_at_start + ): ref.cursor = resume_cursors[ref.stream_id] - source = _get_or_create_source( - db, - name=source_name or ", ".join(r.stream_id for r in refs[:2]), - type_=spec.adapter, - ) - - job = IngestionJob( - id=uuid.uuid4(), - source_id=source.id, - status="running", - started_at=datetime.now(tz=timezone.utc), - file_count=len(refs), - metadata_json={"adapter": spec.adapter, "params": spec.params}, - source_adapter=spec.adapter, - source_ref=", ".join(r.stream_id for r in refs[:5]), - ) - db.add(job) - db.flush() + if existing_job is not None: + job = existing_job + source = db.query(Source).filter(Source.id == job.source_id).first() + if source is None: + source = _get_or_create_source( + db, + name=source_name or ", ".join(r.stream_id for r in refs[:2]), + type_=spec.adapter, + ) + job.source_id = source.id + if not job.source_ref: + job.source_ref = ", ".join(r.stream_id for r in refs[:5]) + job.file_count = max(job.file_count or 0, len(refs)) + else: + source = _get_or_create_source( + db, + name=source_name or ", ".join(r.stream_id for r in refs[:2]), + type_=spec.adapter, + ) + job = IngestionJob( + id=uuid.uuid4(), + source_id=source.id, + status="running", + started_at=datetime.now(tz=timezone.utc), + file_count=len(refs), + metadata_json={"adapter": spec.adapter, "params": spec.params}, + source_adapter=spec.adapter, + source_ref=", ".join(r.stream_id for r in refs[:5]), + mode="batch", + ) + db.add(job) + db.flush() batch: list[LogEntry] = [] adapter_errors: list[str] = [] @@ -369,7 +405,9 @@ def ingest_from_source( except AdapterUnavailableError as e: adapter_errors.append(str(e)) else: - if ref.cursor is None: + # Tail ticks re-read the same streams forever; cursor=None means + # "caught up for this window", not "never read again". + if finalize and ref.cursor is None: completed_streams.add(ref.stream_id) finally: cursors[ref.stream_id] = ref.cursor @@ -382,18 +420,98 @@ def ingest_from_source( _flush_log_batch(db, batch, embedder) stats.duration_seconds = time.time() - start_time + _apply_ingest_counts(job, stats, additive=existing_job is not None) + _store_job_cursors( + job, + cursors, + completed_streams, + adapter_errors, + finalize=finalize, + ) + if finalize: + job.status = "completed" + job.finished_at = datetime.now(tz=timezone.utc) + db.flush() + except Exception as exc: + if finalize: + job.status = "failed" + job.error_message = str(exc) + job.finished_at = datetime.now(tz=timezone.utc) + _apply_ingest_counts(job, stats, additive=existing_job is not None) + _store_job_cursors( + job, + cursors, + completed_streams, + adapter_errors, + finalize=True, + ) + db.flush() + raise + + return job, stats + + +def ingest_push_lines( + db: Session, + raw_lines: list[str], + source_name: Optional[str] = None, + default_service: Optional[str] = None, + default_env: Optional[str] = None, + fmt: str = "auto", + with_embeddings: bool = False, +) -> tuple[IngestionJob, IngestionStats]: + """Persist caller-pushed raw lines through parse → fingerprint → LogEntry.""" + import time + + start_time = time.time() + stats = IngestionStats() + embedder = ingest_embeddings_provider(with_embeddings) + + source = _get_or_create_source(db, name=source_name or "push", type_="push") + job = IngestionJob( + id=uuid.uuid4(), + source_id=source.id, + status="running", + started_at=datetime.now(tz=timezone.utc), + file_count=0, + metadata_json={"source": "push"}, + source_adapter="push", + source_ref="push", + mode="push", + ) + db.add(job) + db.flush() + + batch: list[LogEntry] = [] + try: + for line in raw_lines: + effective_fmt = _resolve_fmt(line, fmt) + entry = _process_line( + line, + effective_fmt, + default_service, + default_env, + source, + job, + "push", + "push", + stats, + ) + if entry is None: + continue + batch.append(entry) + if len(batch) >= BATCH_SIZE: + _flush_log_batch(db, batch, embedder) + + if batch: + _flush_log_batch(db, batch, embedder) + stats.duration_seconds = time.time() - start_time job.status = "completed" job.finished_at = datetime.now(tz=timezone.utc) job.line_count = stats.lines_read job.error_count = stats.error_count job.parsed_count = stats.parsed_count - job.metadata_json = { - **(job.metadata_json or {}), - "cursors": cursors, - "completed_streams": sorted(completed_streams), - **({"partial": True} if adapter_errors else {}), - } db.flush() except Exception as exc: job.status = "failed" @@ -402,17 +520,48 @@ def ingest_from_source( job.line_count = stats.lines_read job.error_count = stats.error_count job.parsed_count = stats.parsed_count - job.metadata_json = { - **(job.metadata_json or {}), - "cursors": cursors, - "completed_streams": sorted(completed_streams), - } db.flush() raise return job, stats +def _apply_ingest_counts( + job: IngestionJob, + stats: IngestionStats, + *, + additive: bool, +) -> None: + if additive: + job.line_count = (job.line_count or 0) + stats.lines_read + job.error_count = (job.error_count or 0) + stats.error_count + job.parsed_count = (job.parsed_count or 0) + stats.parsed_count + else: + job.line_count = stats.lines_read + job.error_count = stats.error_count + job.parsed_count = stats.parsed_count + + +def _store_job_cursors( + job: IngestionJob, + cursors: dict[str, Optional[str]], + completed_streams: set[str], + adapter_errors: list[str], + *, + finalize: bool, +) -> None: + import json + + job.cursor = json.dumps(cursors) + meta = dict(job.metadata_json or {}) + meta["cursors"] = cursors + if finalize: + meta["completed_streams"] = sorted(completed_streams) + if adapter_errors: + meta["partial"] = True + job.metadata_json = meta + + def _get_or_create_source(db: Session, name: str, type_: str = "file") -> Source: existing = db.query(Source).filter(Source.name == name).first() if existing: @@ -428,6 +577,7 @@ def _infer_service_from_filename(path: Path) -> Optional[str]: name = path.stem # filename without extension # Remove common suffixes like .log, dates, numbers import re + name = re.sub(r"[-_]?\d{4}[-_]\d{2}[-_]\d{2}.*$", "", name) name = re.sub(r"[-_]?\d+$", "", name) name = name.strip("-_") diff --git a/src/core/ingestion/tail.py b/src/core/ingestion/tail.py new file mode 100644 index 0000000..0940bab --- /dev/null +++ b/src/core/ingestion/tail.py @@ -0,0 +1,311 @@ +"""Tail-mode ingestion: lifecycle, auto-pause, and worker ticks.""" + +from __future__ import annotations + +import json +from datetime import datetime, timedelta, timezone +from typing import Any, Literal, Optional + +import structlog +from sqlalchemy import or_, select +from sqlalchemy.orm import Session + +from src.adapters.base import SourceSpec, TimeWindow +from src.db.models import IngestionJob + +log = structlog.get_logger() + +TAIL_ADAPTERS: frozenset[str] = frozenset({"cloudwatch", "datadog", "loki"}) +TAIL_LIFECYCLE_ACTIONS: frozenset[str] = frozenset({"pause", "resume", "stop"}) + + +class TailLifecycleError(ValueError): + """Invalid tail lifecycle transition. ``code`` is ``not_tail`` or ``conflict``.""" + + def __init__(self, code: str, message: str) -> None: + self.code = code + super().__init__(message) + + +def apply_tail_lifecycle( + mode: str, + status: str, + action: Literal["pause", "resume", "stop"], +) -> str: + """Return the next status for a tail job. + + Raises ``TailLifecycleError`` when the job is not tail-mode or the + transition is illegal (resume/pause of a stopped job). + """ + if mode != "tail": + raise TailLifecycleError("not_tail", "lifecycle actions require mode=tail") + if action not in TAIL_LIFECYCLE_ACTIONS: + raise TailLifecycleError("conflict", f"unknown action {action!r}") + + if action == "stop": + return "stopped" + + if status == "stopped": + raise TailLifecycleError( + "conflict", + "stopped tail jobs cannot be paused or resumed", + ) + + if action == "pause": + return "paused" + return "running" + + +def consecutive_errors_after_failure( + current: int, + threshold: int, +) -> tuple[int, bool]: + """Return ``(new_count, should_pause)`` after one failed tail tick.""" + new_count = current + 1 + return new_count, new_count >= threshold + + +def cursors_from_job(job: IngestionJob) -> dict[str, Optional[str]]: + """Load stream cursors from the dedicated column, falling back to metadata.""" + if job.cursor: + try: + data = json.loads(job.cursor) + if isinstance(data, dict): + return data + except (ValueError, TypeError): + pass + meta = job.metadata_json or {} + stored = meta.get("cursors") + if isinstance(stored, dict): + return stored + return {} + + +def persist_cursors(job: IngestionJob, cursors: dict[str, Optional[str]]) -> None: + """Write cursors to both ``job.cursor`` (JSON text) and metadata_json.""" + job.cursor = json.dumps(cursors) + meta = dict(job.metadata_json or {}) + meta["cursors"] = cursors + job.metadata_json = meta + + +def has_open_cursors(cursors: dict[str, Optional[str]]) -> bool: + """True when any stream still has a pagination token (not caught up).""" + return any(value is not None and str(value) != "" for value in cursors.values()) + + +def is_tail_job_due(job: IngestionJob, now: datetime, poll_interval: int) -> bool: + """Whether a tail job should be ticked at ``now`` given ``TAIL_POLL_INTERVAL``.""" + if job.mode != "tail" or job.status != "running": + return False + if job.last_polled_at is None: + return True + return job.last_polled_at <= now - timedelta(seconds=poll_interval) + + +def _parse_stored_datetime(value: Any) -> Optional[datetime]: + if value is None: + return None + if isinstance(value, datetime): + return value + if isinstance(value, str): + try: + return datetime.fromisoformat(value) + except ValueError: + return None + return None + + +def _held_paging_window(job: IngestionJob) -> Optional[TimeWindow]: + """Return the stored paging window if an open cursor still needs it.""" + if not has_open_cursors(cursors_from_job(job)): + return None + meta = job.metadata_json or {} + start = _parse_stored_datetime(meta.get("tail_window_start")) + end = _parse_stored_datetime(meta.get("tail_window_end")) + if start is None or end is None: + return None + return TimeWindow(start=start, end=end) + + +def tick_window_for_job( + job: IngestionJob, + now: datetime, +) -> TimeWindow: + """Window for one tail poll: last poll → now, else first-window from spec / 1m. + + ``last_polled_at`` is only the poll clock. While any stream still has a + pagination cursor, reuse ``metadata_json['tail_window_start'|'tail_window_end']`` + — Datadog/CloudWatch page tokens are only valid for that original from/to. + """ + held = _held_paging_window(job) + if held is not None: + return held + + if job.last_polled_at is not None: + return TimeWindow(start=job.last_polled_at, end=now) + + from src.utils.time import resolve_window + + meta = job.metadata_json or {} + since = meta.get("since") + from_raw = meta.get("from_time") + to_raw = meta.get("to_time") + if since or from_raw or to_raw: + from_dt = datetime.fromisoformat(from_raw) if from_raw else None + to_dt = datetime.fromisoformat(to_raw) if to_raw else None + start, end = resolve_window(since=since, from_time=from_dt, to_time=to_dt) + return TimeWindow(start=start, end=end) + + start, end = resolve_window(since="1m") + return TimeWindow(start=start, end=end) + + +def _apply_tail_window_progress( + job: IngestionJob, window: TimeWindow, now: datetime +) -> None: + """Record poll clock and either hold or release the paging window. + + ``last_polled_at`` is always ``now``. Open cursors persist + ``tail_window_start`` / ``tail_window_end``; exhausted streams drop them. + """ + job.last_polled_at = now + meta = dict(job.metadata_json or {}) + if has_open_cursors(cursors_from_job(job)): + meta["tail_window_start"] = window.start.isoformat() + meta["tail_window_end"] = window.end.isoformat() + job.metadata_json = meta + return + meta.pop("tail_window_start", None) + meta.pop("tail_window_end", None) + job.metadata_json = meta + + +def spec_from_tail_job(job: IngestionJob) -> SourceSpec: + meta = job.metadata_json or {} + return SourceSpec( + adapter=str(meta.get("adapter") or job.source_adapter), + params=dict(meta.get("params") or {}), + service=meta.get("service"), + env=meta.get("env"), + ) + + +def _due_tail_jobs( + db: Session, now: datetime, poll_interval: int +) -> list[IngestionJob]: + cutoff = now - timedelta(seconds=poll_interval) + stmt = ( + select(IngestionJob) + .where( + IngestionJob.mode == "tail", + IngestionJob.status == "running", + or_( + IngestionJob.last_polled_at.is_(None), + IngestionJob.last_polled_at <= cutoff, + ), + ) + .with_for_update(skip_locked=True) + ) + jobs = list(db.execute(stmt).scalars().all()) + # Re-filter in Python so unit tests with mocked execute() still honor + # last_polled_at / paused status (SQL WHERE is not applied on MagicMock). + return [job for job in jobs if is_tail_job_due(job, now, poll_interval)] + + +def tick_one_tail_job( + db: Session, job: IngestionJob, now: datetime, error_threshold: int +) -> None: + """Run one adapter poll against an existing tail IngestionJob.""" + from src.core.ingestion.service import ingest_from_source + + meta = job.metadata_json or {} + spec = spec_from_tail_job(job) + window = tick_window_for_job(job, now) + fmt = str(meta.get("format") or "auto") + source_name = meta.get("source_name") + with_embeddings = bool(meta.get("with_embeddings") or False) + + try: + _, _stats = ingest_from_source( + db=db, + spec=spec, + window=window, + source_name=source_name, + fmt=fmt, + resume_cursors=cursors_from_job(job), + resume_completed_streams=None, + with_embeddings=with_embeddings, + existing_job=job, + finalize=False, + ) + except Exception as exc: + new_count, should_pause = consecutive_errors_after_failure( + job.consecutive_errors or 0, + error_threshold, + ) + job.consecutive_errors = new_count + job.error_message = str(exc) + if should_pause: + job.status = "paused" + log.error( + "tail_job_auto_paused", + ingestion_job_id=str(job.id), + consecutive_errors=new_count, + error=str(exc), + ) + else: + log.warning( + "tail_job_tick_failed", + ingestion_job_id=str(job.id), + consecutive_errors=new_count, + error=str(exc), + ) + # Stamp last_polled_at so the job is not due again until TAIL_POLL_INTERVAL + # (otherwise the worker busy-loops and re-reads already-flushed lines). + job.last_polled_at = now + db.flush() + return + + job.consecutive_errors = 0 + job.error_message = None + _apply_tail_window_progress(job, window, now) + db.flush() + + +def tick_tail_jobs(db: Session) -> int: + """Poll due running tail jobs. Returns the number of jobs ticked.""" + from src.config import get_settings + from src.core.ingestion.backpressure import ( + ingest_queue_is_full, + pending_worker_job_count, + ) + + settings = get_settings() + pending = pending_worker_job_count(db) + if ingest_queue_is_full(pending, settings.ingest_queue_max): + log.info("tail_ticks_skipped_queue_full", pending=pending) + return 0 + + now = datetime.now(tz=timezone.utc) + jobs = _due_tail_jobs(db, now, settings.tail_poll_interval) + for job in jobs: + tick_one_tail_job(db, job, now, settings.tail_error_threshold) + return len(jobs) + + +def tail_job_counts(db: Session) -> dict[str, Any]: + """``{running, paused}`` counts for /health.""" + from sqlalchemy import func + + running = db.execute( + select(func.count()) + .select_from(IngestionJob) + .where(IngestionJob.mode == "tail", IngestionJob.status == "running") + ).scalar_one() + paused = db.execute( + select(func.count()) + .select_from(IngestionJob) + .where(IngestionJob.mode == "tail", IngestionJob.status == "paused") + ).scalar_one() + return {"running": int(running or 0), "paused": int(paused or 0)} diff --git a/src/db/models.py b/src/db/models.py index 876674d..3c9cba9 100644 --- a/src/db/models.py +++ b/src/db/models.py @@ -53,6 +53,10 @@ class IngestionJob(Base): error_message: Mapped[str | None] = mapped_column(Text, nullable=True) source_adapter: Mapped[str] = mapped_column(String(50), nullable=False, default="file") source_ref: Mapped[str | None] = mapped_column(Text, nullable=True) + mode: Mapped[str] = mapped_column(String(20), nullable=False, default="batch") + cursor: Mapped[str | None] = mapped_column(Text, nullable=True) + last_polled_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) + consecutive_errors: Mapped[int] = mapped_column(Integer, default=0) created_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), server_default=func.now()) source: Mapped["Source"] = relationship("Source", back_populates="jobs") diff --git a/src/worker/runner.py b/src/worker/runner.py index 8d22d3d..f875534 100644 --- a/src/worker/runner.py +++ b/src/worker/runner.py @@ -165,6 +165,7 @@ def process_one(db) -> bool: def run_worker(poll_interval: int = POLL_INTERVAL): """Main worker loop. Runs until SIGINT/SIGTERM.""" + from src.core.ingestion.tail import tick_tail_jobs from src.db.session import get_db shutdown = False @@ -183,7 +184,8 @@ def _handle_signal(sig, frame): try: with get_db() as db: processed = process_one(db) - if not processed: + ticked = tick_tail_jobs(db) + if not processed and not ticked: # Nothing to do — sleep before next poll for _ in range(poll_interval * 10): if shutdown: diff --git a/tests/unit/test_api.py b/tests/unit/test_api.py index ae47887..7b0d40d 100644 --- a/tests/unit/test_api.py +++ b/tests/unit/test_api.py @@ -36,6 +36,10 @@ def _ctx_db(query_result=_UNSET, execute_scalar=_UNSET): ) if execute_scalar is not _UNSET: mock_db.execute.return_value.scalar_one.return_value = execute_scalar + else: + # Default 0 pending worker jobs so ingest backpressure does not 429 + # every mocked POST /ingestions. + mock_db.execute.return_value.scalar_one.return_value = 0 return mock_db @@ -99,6 +103,15 @@ def test_health_includes_worker_queue_depth(self): assert resp.json()["worker_queue_depth"] == 7 + def test_health_includes_tail_job_counts(self): + with patch("src.db.session.check_connection", return_value=True), \ + _patch_get_db(execute_scalar=2): + resp = client.get("/health") + + tail = resp.json()["tail_jobs"] + assert tail["running"] == 2 + assert tail["paused"] == 2 + def test_health_includes_file_adapter_ok(self): with patch("src.db.session.check_connection", return_value=True), \ _patch_get_db(execute_scalar=3): diff --git a/tests/unit/test_auth.py b/tests/unit/test_auth.py index aaf2d3f..fc5d1a4 100644 --- a/tests/unit/test_auth.py +++ b/tests/unit/test_auth.py @@ -131,6 +131,11 @@ def test_ingest_post_vs_get(self) -> None: assert required_roles("POST", "/ingestions") == frozenset({"ingest", "admin"}) assert required_roles("GET", "/ingestions") == frozenset({"query", "admin"}) assert required_roles("GET", "/ingestions/latest") == frozenset({"query", "admin"}) + assert required_roles("POST", "/ingestions/lines") == frozenset({"ingest", "admin"}) + assert required_roles("POST", "/ingestions/abc:pause") == frozenset({"ingest", "admin"}) + assert required_roles("POST", "/v1/ingestions/lines") == frozenset({"ingest", "admin"}) + assert required_roles("POST", "/v1/ingestions/abc:resume") == frozenset({"ingest", "admin"}) + assert required_roles("POST", "/v1/ingestions/abc:stop") == frozenset({"ingest", "admin"}) def test_query_and_config_and_ui(self) -> None: assert required_roles("POST", "/query/explain") == frozenset({"query", "admin"}) @@ -341,17 +346,28 @@ def test_query_role_cannot_post_v1_ingestions(self) -> None: assert resp.status_code == 403 assert resp.json()["error_code"] == "AUTH_FORBIDDEN" - def test_ingest_role_can_post_v1_ingestions(self) -> None: + def test_ingest_role_can_post_v1_ingestions_lines(self) -> None: with patch("src.config.get_settings", return_value=_auth_settings()), \ patch("src.api.auth.keys.lookup_api_key", return_value=_key("ingest")): resp = client.post( - "/v1/ingestions", - json={}, + "/v1/ingestions/lines", + content="", headers={"Authorization": "Bearer rlk_ingestrole1"}, ) - assert resp.status_code == 422 + assert resp.status_code == 400 assert resp.status_code != 403 + def test_query_role_cannot_post_v1_ingestions_lines(self) -> None: + with patch("src.config.get_settings", return_value=_auth_settings()), \ + patch("src.api.auth.keys.lookup_api_key", return_value=_key("query")): + resp = client.post( + "/v1/ingestions/lines", + content="hello\n", + headers={"Authorization": "Bearer rlk_queryrolexx"}, + ) + assert resp.status_code == 403 + assert resp.json()["error_code"] == "AUTH_FORBIDDEN" + def test_docs_not_exempt(self) -> None: with patch("src.config.get_settings", return_value=_auth_settings()): resp = client.get("/docs") diff --git a/tests/unit/test_ingestion_service.py b/tests/unit/test_ingestion_service.py index 580df20..29a2428 100644 --- a/tests/unit/test_ingestion_service.py +++ b/tests/unit/test_ingestion_service.py @@ -9,6 +9,7 @@ from datetime import datetime, timezone from unittest.mock import MagicMock +import uuid import pytest @@ -246,6 +247,7 @@ def test_timestamp_falls_back_to_received_at(self, monkeypatch): def test_raw_line_defaults_fill_service_env_host(self, monkeypatch): """Adapters (Loki labels, k8s fields) may set RawLogLine defaults; ingestion should persist them when the parsed line has no service/env/host.""" + class LabelAdapter: name = "labels" @@ -267,7 +269,9 @@ def read(self, ref, window): _patch_get_adapter(monkeypatch, LabelAdapter()) job, stats = ingest_from_source( - db=db, spec=SourceSpec(adapter="labels", params={}), window=_WINDOW, + db=db, + spec=SourceSpec(adapter="labels", params={}), + window=_WINDOW, ) assert job.status == "completed" @@ -369,6 +373,58 @@ def test_partial_when_one_stream_unavailable(self, monkeypatch): assert stats.lines_read == 1 assert stats.parsed_count == 1 + def test_existing_job_stays_running_and_counts_add(self, monkeypatch): + db = _mock_db(existing_source=MagicMock()) + existing = IngestionJob( + id=uuid.uuid4(), + source_id=uuid.uuid4(), + status="running", + mode="tail", + line_count=4, + parsed_count=4, + error_count=0, + file_count=1, + metadata_json={"adapter": "fake", "params": {}}, + source_adapter="fake", + source_ref="stream-0", + ) + _patch_get_adapter(monkeypatch, FakeAdapter(["plain extra line"])) + + job, stats = ingest_from_source( + db=db, + spec=SourceSpec(adapter="fake", params={}), + window=_WINDOW, + existing_job=existing, + finalize=False, + ) + + assert job is existing + assert job.status == "running" + assert job.finished_at is None + assert stats.parsed_count == 1 + assert job.parsed_count == 5 + assert job.line_count == 5 + assert not any(isinstance(o, IngestionJob) for o in db.added) + + +class TestIngestPushLines: + def test_persists_raw_lines(self): + from src.core.ingestion.service import ingest_push_lines + + db = _mock_db() + job, stats = ingest_push_lines( + db, + raw_lines=[ + '{"message": "pushed", "level": "error", "service": "api"}', + "text line", + ], + ) + assert job.mode == "push" + assert job.status == "completed" + assert stats.parsed_count == 2 + assert stats.lines_read == 2 + assert "api" in stats.services_detected + class TestIngestFiles: def test_happy_path_sets_source_adapter(self, tmp_path): @@ -399,7 +455,9 @@ def fake_persist(db, entries, provider=None, settings=None): monkeypatch.setattr( "src.core.ingestion.service.ingest_embeddings_provider", - lambda with_embeddings, settings=None: object() if with_embeddings else None, + lambda with_embeddings, settings=None: ( + object() if with_embeddings else None + ), ) monkeypatch.setattr( "src.core.ingestion.service.persist_log_embeddings", fake_persist diff --git a/tests/unit/test_openapi_contract.py b/tests/unit/test_openapi_contract.py index 133ffef..8ec62ca 100644 --- a/tests/unit/test_openapi_contract.py +++ b/tests/unit/test_openapi_contract.py @@ -7,7 +7,11 @@ from __future__ import annotations from src.api.app import app -from src.api.deprecation import is_deprecated_alias, is_versioned_api_path, successor_path +from src.api.deprecation import ( + is_deprecated_alias, + is_versioned_api_path, + successor_path, +) REQUIRED_GET: tuple[str, ...] = ( "/health", @@ -22,6 +26,10 @@ "/v1/query/compare", "/v1/query/clusters", "/v1/ingestions", + "/v1/ingestions/lines", + "/v1/ingestions/{ingestion_job_id}:pause", + "/v1/ingestions/{ingestion_job_id}:resume", + "/v1/ingestions/{ingestion_job_id}:stop", ) diff --git a/tests/unit/test_push_tail_ingest.py b/tests/unit/test_push_tail_ingest.py new file mode 100644 index 0000000..0242320 --- /dev/null +++ b/tests/unit/test_push_tail_ingest.py @@ -0,0 +1,530 @@ +"""Unit tests for push NDJSON ingest, queue backpressure, and tail jobs.""" + +from __future__ import annotations + +import uuid +from datetime import datetime, timedelta, timezone +from unittest.mock import MagicMock, patch + +import pytest +from fastapi.testclient import TestClient + +from src.adapters.base import TimeWindow +from src.api.app import app +from src.core.ingestion.backpressure import ingest_queue_is_full +from src.core.ingestion.push import NdjsonParseError, parse_ndjson_payload +from src.core.ingestion.tail import ( + TailLifecycleError, + _due_tail_jobs, + apply_tail_lifecycle, + consecutive_errors_after_failure, + persist_cursors, + tick_one_tail_job, + tick_window_for_job, +) +from src.db.models import IngestionJob + +client = TestClient(app, raise_server_exceptions=False) + + +def _ctx_db(*, execute_scalar: int = 0, query_result=None): + mock_db = MagicMock() + mock_db.__enter__ = MagicMock(return_value=mock_db) + mock_db.__exit__ = MagicMock(return_value=False) + mock_db.query.return_value.filter.return_value.first.return_value = query_result + mock_db.execute.return_value.scalar_one.return_value = execute_scalar + return mock_db + + +# ── NDJSON parsing ──────────────────────────────────────────────────────────── + + +class TestParseNdjson: + def test_raw_and_json_objects(self) -> None: + body = ( + "plain text line\n" + '{"message": "boom", "level": "error", "service": "api"}\n' + '{"raw": "from raw field", "timestamp": "2026-01-01T00:00:00Z"}\n' + '"quoted string line"\n' + ) + lines = parse_ndjson_payload(body, max_lines=5000) + assert lines[0] == "plain text line" + assert "boom" in lines[1] + assert "from raw field" in lines[2] + assert lines[3] == "quoted string line" + + def test_blank_lines_ignored(self) -> None: + lines = parse_ndjson_payload("\n\na\n \nb\n", max_lines=10) + assert lines == ["a", "b"] + + def test_empty_body_raises(self) -> None: + with pytest.raises(NdjsonParseError, match="empty"): + parse_ndjson_payload("", max_lines=10) + with pytest.raises(NdjsonParseError, match="empty"): + parse_ndjson_payload(" \n\n", max_lines=10) + + def test_over_max_lines_raises(self) -> None: + body = "a\nb\nc\n" + with pytest.raises(NdjsonParseError, match="INGEST_PUSH_MAX_LINES"): + parse_ndjson_payload(body, max_lines=2) + + +# ── Backpressure helper ─────────────────────────────────────────────────────── + + +class TestQueueBackpressure: + def test_full_at_max(self) -> None: + assert ingest_queue_is_full(100, 100) is True + assert ingest_queue_is_full(101, 100) is True + assert ingest_queue_is_full(99, 100) is False + assert ingest_queue_is_full(0, 100) is False + + +# ── Tail lifecycle state machine ────────────────────────────────────────────── + + +class TestTailLifecycle: + def test_pause_resume_stop(self) -> None: + assert apply_tail_lifecycle("tail", "running", "pause") == "paused" + assert apply_tail_lifecycle("tail", "paused", "resume") == "running" + assert apply_tail_lifecycle("tail", "running", "stop") == "stopped" + assert apply_tail_lifecycle("tail", "paused", "stop") == "stopped" + assert apply_tail_lifecycle("tail", "stopped", "stop") == "stopped" + + def test_not_tail_is_error(self) -> None: + with pytest.raises(TailLifecycleError) as exc: + apply_tail_lifecycle("batch", "running", "pause") + assert exc.value.code == "not_tail" + + def test_cannot_resume_or_pause_stopped(self) -> None: + with pytest.raises(TailLifecycleError) as exc: + apply_tail_lifecycle("tail", "stopped", "resume") + assert exc.value.code == "conflict" + with pytest.raises(TailLifecycleError): + apply_tail_lifecycle("tail", "stopped", "pause") + + def test_idempotent_pause_and_resume(self) -> None: + assert apply_tail_lifecycle("tail", "paused", "pause") == "paused" + assert apply_tail_lifecycle("tail", "running", "resume") == "running" + + +class TestAutoPauseCounter: + def test_pauses_at_threshold(self) -> None: + count, pause = consecutive_errors_after_failure(4, 5) + assert count == 5 + assert pause is True + + def test_does_not_pause_before_threshold(self) -> None: + count, pause = consecutive_errors_after_failure(0, 5) + assert count == 1 + assert pause is False + + +# ── Routes ──────────────────────────────────────────────────────────────────── + + +class TestPushLinesRoute: + def test_returns_counts_when_persist_mocked(self) -> None: + mock_job = MagicMock() + mock_job.id = uuid.uuid4() + mock_job.status = "completed" + mock_stats = MagicMock() + mock_stats.lines_read = 2 + mock_stats.parsed_count = 2 + mock_stats.error_count = 0 + mock_db = _ctx_db() + + with ( + patch("src.db.session.get_db", side_effect=lambda: mock_db), + patch( + "src.core.ingestion.service.ingest_push_lines", + return_value=(mock_job, mock_stats), + ), + ): + resp = client.post( + "/v1/ingestions/lines", + content='{"message":"a"}\nplain line\n', + headers={"Content-Type": "application/x-ndjson"}, + ) + + assert resp.status_code == 200 + data = resp.json() + assert data["ingestion_job_id"] == str(mock_job.id) + assert data["parsed_count"] == 2 + assert data["line_count"] == 2 + assert data["mode"] == "push" + + def test_unversioned_alias_also_works(self) -> None: + mock_job = MagicMock() + mock_job.id = uuid.uuid4() + mock_job.status = "completed" + mock_stats = MagicMock() + mock_stats.lines_read = 1 + mock_stats.parsed_count = 1 + mock_stats.error_count = 0 + mock_db = _ctx_db() + + with ( + patch("src.db.session.get_db", side_effect=lambda: mock_db), + patch( + "src.core.ingestion.service.ingest_push_lines", + return_value=(mock_job, mock_stats), + ), + ): + resp = client.post( + "/ingestions/lines", + content="hello\n", + headers={"Content-Type": "application/jsonl"}, + ) + + assert resp.status_code == 200 + assert resp.headers.get("deprecation") == "true" + + def test_empty_body_400(self) -> None: + resp = client.post( + "/v1/ingestions/lines", + content="", + headers={"Content-Type": "application/x-ndjson"}, + ) + assert resp.status_code == 400 + + def test_over_max_lines_400(self) -> None: + from src.config.settings import Settings + + settings = Settings(_env_file=None, ingest_push_max_lines=1) + with patch("src.config.get_settings", return_value=settings): + resp = client.post( + "/v1/ingestions/lines", + content="a\nb\n", + headers={"Content-Type": "text/plain"}, + ) + assert resp.status_code == 400 + assert "INGEST_PUSH_MAX_LINES" in resp.json()["detail"] + + def test_get_lines_is_not_invalid_uuid(self) -> None: + resp = client.get("/ingestions/lines") + assert resp.status_code == 405 + assert resp.status_code != 422 + assert resp.status_code != 400 + + resp_v1 = client.get("/v1/ingestions/lines") + assert resp_v1.status_code == 405 + + +class TestQueueFull429: + def test_post_ingestions_429_when_pending_at_max(self) -> None: + mock_db = _ctx_db(execute_scalar=100) + with ( + patch("src.adapters.file.adapter.discover_files", return_value=["f.log"]), + patch("src.db.session.get_db", side_effect=lambda: mock_db), + ): + resp = client.post("/v1/ingestions", json={"paths": ["/logs"]}) + + assert resp.status_code == 429 + body = resp.json() + assert body["error_code"] == "INGEST_QUEUE_FULL" + assert "message" in body + assert resp.headers.get("retry-after") == "5" + + def test_push_lines_429_when_queue_full(self) -> None: + mock_db = _ctx_db(execute_scalar=100) + with ( + patch("src.db.session.get_db", side_effect=lambda: mock_db), + patch("src.core.ingestion.service.ingest_push_lines") as mock_ingest, + ): + resp = client.post( + "/v1/ingestions/lines", + content="hello\n", + headers={"Content-Type": "application/x-ndjson"}, + ) + + assert resp.status_code == 429 + assert resp.json()["error_code"] == "INGEST_QUEUE_FULL" + mock_ingest.assert_not_called() + + +class TestTailCreateAndLifecycleRoutes: + def test_create_tail_job(self) -> None: + mock_db = _ctx_db() + added: list[object] = [] + + def capture_add(obj: object) -> None: + added.append(obj) + if getattr(obj, "id", None) is None: + obj.id = uuid.uuid4() # type: ignore[attr-defined] + + mock_db.add.side_effect = capture_add + mock_adapter = MagicMock() + mock_adapter.discover.return_value = [MagicMock(stream_id="/aws/lambda/x")] + + with ( + patch("src.adapters.registry.get_adapter", return_value=mock_adapter), + patch("src.db.session.get_db", side_effect=lambda: mock_db), + ): + resp = client.post( + "/v1/ingestions", + json={ + "adapter": "cloudwatch", + "params": {"log_group": "/aws/lambda/x"}, + "mode": "tail", + "paths": [], + }, + ) + + assert resp.status_code == 202 + data = resp.json() + assert data["mode"] == "tail" + assert data["status"] == "running" + assert data["worker_job_id"] is None + assert data["ingestion_job_id"] + jobs = [o for o in added if isinstance(o, IngestionJob)] + assert len(jobs) == 1 + assert jobs[0].mode == "tail" + assert jobs[0].status == "running" + + def test_tail_rejects_file_adapter(self) -> None: + resp = client.post( + "/v1/ingestions", + json={"adapter": "file", "paths": ["/logs"], "mode": "tail"}, + ) + assert resp.status_code == 422 + + def test_pause_resume_stop(self) -> None: + job = MagicMock() + job.id = uuid.uuid4() + job.mode = "tail" + job.status = "running" + job.finished_at = None + mock_db = _ctx_db(query_result=job) + + with patch("src.db.session.get_db", side_effect=lambda: mock_db): + paused = client.post(f"/v1/ingestions/{job.id}:pause") + assert paused.status_code == 200 + assert paused.json()["status"] == "paused" + assert job.status == "paused" + + resumed = client.post(f"/v1/ingestions/{job.id}:resume") + assert resumed.status_code == 200 + assert resumed.json()["status"] == "running" + + stopped = client.post(f"/v1/ingestions/{job.id}:stop") + assert stopped.status_code == 200 + assert stopped.json()["status"] == "stopped" + assert job.finished_at is not None + + def test_pause_missing_job_404(self) -> None: + mock_db = _ctx_db(query_result=None) + with patch("src.db.session.get_db", side_effect=lambda: mock_db): + resp = client.post(f"/v1/ingestions/{uuid.uuid4()}:pause") + assert resp.status_code == 404 + + def test_pause_batch_job_409(self) -> None: + job = MagicMock() + job.id = uuid.uuid4() + job.mode = "batch" + job.status = "completed" + mock_db = _ctx_db(query_result=job) + with patch("src.db.session.get_db", side_effect=lambda: mock_db): + resp = client.post(f"/v1/ingestions/{job.id}:pause") + assert resp.status_code == 409 + + +class TestTickAutoPause: + def test_auto_pauses_after_threshold(self) -> None: + job = MagicMock() + job.id = uuid.uuid4() + job.mode = "tail" + job.status = "running" + job.consecutive_errors = 4 + job.metadata_json = {"adapter": "cloudwatch", "params": {"log_group": "g"}} + job.cursor = None + job.last_polled_at = None + job.source_adapter = "cloudwatch" + job.error_message = None + + db = MagicMock() + now = datetime(2026, 8, 17, tzinfo=timezone.utc) + with patch( + "src.core.ingestion.service.ingest_from_source", + side_effect=RuntimeError("adapter down"), + ): + tick_one_tail_job(db, job, now, error_threshold=5) + + assert job.status == "paused" + assert job.consecutive_errors == 5 + assert "adapter down" in job.error_message + assert job.last_polled_at == now + + def test_failed_tick_is_not_due_until_poll_interval(self) -> None: + job = MagicMock() + job.id = uuid.uuid4() + job.mode = "tail" + job.status = "running" + job.consecutive_errors = 0 + job.metadata_json = {"adapter": "cloudwatch", "params": {"log_group": "g"}} + job.cursor = None + job.last_polled_at = None + job.source_adapter = "cloudwatch" + job.error_message = None + + db = MagicMock() + now = datetime(2026, 8, 17, 12, 0, tzinfo=timezone.utc) + with patch( + "src.core.ingestion.service.ingest_from_source", + side_effect=RuntimeError("boom"), + ): + tick_one_tail_job(db, job, now, error_threshold=5) + + assert job.last_polled_at == now + assert job.consecutive_errors == 1 + assert job.status == "running" + + db.execute.return_value.scalars.return_value.all.return_value = [job] + assert _due_tail_jobs(db, now + timedelta(seconds=1), poll_interval=30) == [] + assert _due_tail_jobs(db, now + timedelta(seconds=31), poll_interval=30) == [ + job + ] + + job.status = "paused" + assert _due_tail_jobs(db, now + timedelta(seconds=31), poll_interval=30) == [] + + def test_success_resets_counter(self) -> None: + job = MagicMock() + job.id = uuid.uuid4() + job.mode = "tail" + job.status = "running" + job.consecutive_errors = 3 + job.metadata_json = {"adapter": "loki", "params": {"query": '{app="api"}'}} + job.cursor = "{}" + job.last_polled_at = None + job.source_adapter = "loki" + job.error_message = "old" + + db = MagicMock() + now = datetime(2026, 8, 17, 12, 0, tzinfo=timezone.utc) + mock_stats = MagicMock() + with patch( + "src.core.ingestion.service.ingest_from_source", + return_value=(job, mock_stats), + ): + tick_one_tail_job(db, job, now, error_threshold=5) + + assert job.status == "running" + assert job.consecutive_errors == 0 + assert job.error_message is None + assert job.last_polled_at == now + + +class TestTailWindowPagination: + def test_open_cursor_holds_window_until_exhausted(self) -> None: + window_start = datetime(2026, 8, 17, 12, 0, tzinfo=timezone.utc) + window_end = datetime(2026, 8, 17, 12, 1, tzinfo=timezone.utc) + now = datetime(2026, 8, 17, 12, 1, 30, tzinfo=timezone.utc) + job = MagicMock() + job.id = uuid.uuid4() + job.mode = "tail" + job.status = "running" + job.consecutive_errors = 0 + job.metadata_json = {"adapter": "datadog", "params": {"query": "*"}} + job.cursor = None + job.last_polled_at = None + job.source_adapter = "datadog" + job.error_message = None + + def ingest_with_open_cursor( + *, existing_job: MagicMock, **_kwargs: object + ) -> tuple: + persist_cursors(existing_job, {"stream-0": "page-2"}) + return existing_job, MagicMock() + + db = MagicMock() + with ( + patch( + "src.core.ingestion.service.ingest_from_source", + side_effect=ingest_with_open_cursor, + ), + patch( + "src.core.ingestion.tail.tick_window_for_job", + return_value=TimeWindow(start=window_start, end=window_end), + ), + ): + tick_one_tail_job(db, job, now, error_threshold=5) + + assert job.last_polled_at == now + assert job.metadata_json["tail_window_start"] == window_start.isoformat() + assert job.metadata_json["tail_window_end"] == window_end.isoformat() + + held = tick_window_for_job(job, now + timedelta(minutes=5)) + assert held.start == window_start + assert held.end == window_end + + later = now + timedelta(seconds=30) + + def ingest_exhausted(*, existing_job: MagicMock, **_kwargs: object) -> tuple: + persist_cursors(existing_job, {"stream-0": None}) + return existing_job, MagicMock() + + with patch( + "src.core.ingestion.service.ingest_from_source", + side_effect=ingest_exhausted, + ): + tick_one_tail_job(db, job, later, error_threshold=5) + + assert job.last_polled_at == later + assert "tail_window_start" not in job.metadata_json + assert "tail_window_end" not in job.metadata_json + + def test_failed_continuation_keeps_held_window(self) -> None: + window_start = datetime(2026, 8, 17, 12, 0, tzinfo=timezone.utc) + window_end = datetime(2026, 8, 17, 12, 1, tzinfo=timezone.utc) + hold_now = datetime(2026, 8, 17, 12, 1, 30, tzinfo=timezone.utc) + job = MagicMock() + job.id = uuid.uuid4() + job.mode = "tail" + job.status = "running" + job.consecutive_errors = 0 + job.metadata_json = {"adapter": "datadog", "params": {"query": "*"}} + job.cursor = None + job.last_polled_at = None + job.source_adapter = "datadog" + job.error_message = None + + def ingest_with_open_cursor( + *, existing_job: MagicMock, **_kwargs: object + ) -> tuple: + persist_cursors(existing_job, {"stream-0": "page-2"}) + return existing_job, MagicMock() + + db = MagicMock() + with ( + patch( + "src.core.ingestion.service.ingest_from_source", + side_effect=ingest_with_open_cursor, + ), + patch( + "src.core.ingestion.tail.tick_window_for_job", + return_value=TimeWindow(start=window_start, end=window_end), + ), + ): + tick_one_tail_job(db, job, hold_now, error_threshold=5) + + fail_now = hold_now + timedelta(minutes=5) + with patch( + "src.core.ingestion.service.ingest_from_source", + side_effect=RuntimeError("page failed"), + ): + tick_one_tail_job(db, job, fail_now, error_threshold=5) + + assert job.last_polled_at == fail_now + assert job.status == "running" + assert job.consecutive_errors == 1 + db.execute.return_value.scalars.return_value.all.return_value = [job] + assert ( + _due_tail_jobs(db, fail_now + timedelta(seconds=1), poll_interval=30) == [] + ) + + later = fail_now + timedelta(minutes=10) + held = tick_window_for_job(job, later) + assert held.start == window_start + assert held.end == window_end + assert held != TimeWindow(start=fail_now, end=later) diff --git a/tests/unit/test_settings.py b/tests/unit/test_settings.py index b96cdec..f048d8c 100644 --- a/tests/unit/test_settings.py +++ b/tests/unit/test_settings.py @@ -1,4 +1,5 @@ """Tests for unprefixed environment variable names on Settings.""" + import pytest from src.config.settings import Settings @@ -51,7 +52,10 @@ def test_ignores_legacy_raglogs_prefix(monkeypatch): settings = Settings(_env_file=None) assert settings.llm_provider == "disabled" - assert settings.db_url == "postgresql+psycopg://postgres:postgres@localhost:5432/raglogs" + assert ( + settings.db_url + == "postgresql+psycopg://postgres:postgres@localhost:5432/raglogs" + ) def test_auth_settings_defaults_disabled(): @@ -81,3 +85,19 @@ def test_auth_settings_from_env(monkeypatch): assert settings.oidc_jwks_url == "https://idp.example/jwks.json" assert settings.api_bind_host == "0.0.0.0" assert settings.auth_refuse_insecure_bind is True + + +def test_ingest_backpressure_and_tail_settings_from_env(monkeypatch): + monkeypatch.setenv("INGEST_QUEUE_MAX", "10") + monkeypatch.setenv("INGEST_RETRY_AFTER_SECONDS", "8") + monkeypatch.setenv("INGEST_PUSH_MAX_LINES", "100") + monkeypatch.setenv("TAIL_POLL_INTERVAL", "15") + monkeypatch.setenv("TAIL_ERROR_THRESHOLD", "3") + + settings = Settings(_env_file=None) + + assert settings.ingest_queue_max == 10 + assert settings.ingest_retry_after_seconds == 8 + assert settings.ingest_push_max_lines == 100 + assert settings.tail_poll_interval == 15 + assert settings.tail_error_threshold == 3