Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -113,6 +113,10 @@ TAIL_ERROR_THRESHOLD=5
# WEBHOOK_MAX_RETRIES=5
# WEBHOOK_TIMEOUT=10

# Idempotency-Key TTL for POST /v1/ingestions (batch enqueue and tail create).
# Repeat requests with the same key within this window return the original job.
# INGEST_IDEMPOTENCY_TTL_SECONDS=86400


# ── Adapters ──────────────────────────────────────────────────────────────────
# AWS region for the cloudwatch source adapter (raglogs ingest --adapter cloudwatch).
Expand Down
12 changes: 12 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -688,6 +688,7 @@ All settings are read from `.env`, environment variables, or CLI flags. Priority
| `WEBHOOK_SECRET` | _(empty)_ | Fallback HMAC secret for ingest completion callbacks when auth is off or the API key has no per-key `whsec_` |
| `WEBHOOK_MAX_RETRIES` | `5` | Extra webhook POST attempts after the first (6 POSTs by default) on 5xx / 429 / connect errors |
| `WEBHOOK_TIMEOUT` | `10` | Per-attempt HTTP timeout in seconds for completion callbacks |
| `INGEST_IDEMPOTENCY_TTL_SECONDS` | `86400` | How long `Idempotency-Key` on `POST /v1/ingestions` is remembered (batch enqueue and tail create) |

---

Expand Down Expand Up @@ -993,6 +994,17 @@ curl -X POST http://localhost:8000/v1/ingestions/$ID:stop

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

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

```bash
curl -X POST http://localhost:8000/v1/ingestions \
-H "Content-Type: application/json" \
-H "Idempotency-Key: incident-123-retry" \
-d '{"paths":["/var/log/app"]}'
```

**Content dedup.** Every persist path (`ingest_files`, `ingest_from_source` including tail ticks, and push `/lines`) stores `original_line_hash` (SHA-256 of the **raw** line, distinct from the normalized fingerprint) and upserts on `(scope, source_ref, original_line_hash, timestamp)`. Re-reading the same physical lines is a no-op, so cluster counts stay stable across overlapping windows and tail/push retries. Missing `source_ref` is stored as `""` so uniqueness works (Postgres NULLs are distinct). `scope` defaults to `"default"`, or the API key's scope when a principal is present; queries are **not** filtered by scope yet. Duplicate lines are skipped, not errors.

**Completion callbacks.** Optional `callback_url` on `POST /v1/ingestions` (http or https only; `file:` and empty hosts are rejected). When a **batch** worker job reaches a terminal state (`done` / `failed`), raglogs POSTs an HMAC-SHA256-signed JSON body to that URL. Delivery is fail-open: retries with jittered exponential backoff (`WEBHOOK_MAX_RETRIES`, default 5 extra attempts) on 5xx, 429, and connect errors; 4xx other than 429 are not retried. Failures are logged and **do not** change ingest status — poll `GET /v1/ingestions/jobs/{worker_job_id}` still works.

`job_id` in the payload is the **ingestion_job_id** when ingest created a row; if the worker failed before that, it is the `worker_job_id`. `scope` is the authenticated key's scope, or `"default"` when auth is off. `counts.clusters` is `0` (clustering is not part of ingest). Worker `done` maps to `"succeeded"` (or `"partial"` when `error_count > 0`); `failed` maps to `"failed"`.
Expand Down
138 changes: 89 additions & 49 deletions clients/openapi.json
Original file line number Diff line number Diff line change
Expand Up @@ -28,42 +28,42 @@
}
},
"/v1/ingestions": {
"get": {
"tags": [
"ingestions"
],
"summary": "List Ingestions",
"description": "List recent completed ingestion jobs, newest first. Used by the web UI's ingestion picker.",
"operationId": "v1_ingestions_list_ingestions__v1_ingestions",
"responses": {
"200": {
"description": "Successful Response",
"content": {
"application/json": {
"schema": {
"$ref": "#/components/schemas/IngestionListResponse"
}
}
}
}
}
},
"post": {
"tags": [
"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.\n\n``mode=tail`` creates a long-lived job that the worker polls; the response\nincludes ``ingestion_job_id`` and no worker_job_id.\n\nOptional ``callback_url`` (http/https) is stored on batch worker jobs. The\nworker POSTs an HMAC-signed completion payload on terminal state. Tail jobs\ndo not fire callbacks (long-lived; poll or ``:stop`` instead).",
"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.\n\nOptional ``callback_url`` (http/https) is stored on batch worker jobs. The\nworker POSTs an HMAC-signed completion payload on terminal state. Tail jobs\ndo not fire callbacks (long-lived; poll or ``:stop`` instead).\n\nOptional ``Idempotency-Key`` (also accepted as ``idempotency-key``): a\nrepeat within the TTL returns the original 202 job rather than a new one.",
"operationId": "v1_ingestions_create_ingestion__v1_ingestions",
"parameters": [
{
"name": "Idempotency-Key",
"in": "header",
"required": false,
"schema": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"description": "If set, a repeat POST within INGEST_IDEMPOTENCY_TTL_SECONDS (default 86400) returns the original job instead of enqueueing a new one. Empty values are rejected with 400. Applies to batch enqueue and tail create.",
"title": "Idempotency-Key"
},
"description": "If set, a repeat POST within INGEST_IDEMPOTENCY_TTL_SECONDS (default 86400) returns the original job instead of enqueueing a new one. Empty values are rejected with 400. Applies to batch enqueue and tail create."
}
],
"requestBody": {
"required": true,
"content": {
"application/json": {
"schema": {
"$ref": "#/components/schemas/IngestRequest"
}
}
},
"required": true
}
},
"responses": {
"202": {
Expand All @@ -87,6 +87,26 @@
}
}
}
},
"get": {
"tags": [
"ingestions"
],
"summary": "List Ingestions",
"description": "List recent completed ingestion jobs, newest first. Used by the web UI's ingestion picker.",
"operationId": "v1_ingestions_list_ingestions__v1_ingestions",
"responses": {
"200": {
"description": "Successful Response",
"content": {
"application/json": {
"schema": {
"$ref": "#/components/schemas/IngestionListResponse"
}
}
}
}
}
}
},
"/v1/ingestions/jobs/{worker_job_id}": {
Expand Down Expand Up @@ -366,43 +386,43 @@
}
},
"/ingestions": {
"get": {
"tags": [
"ingestions"
],
"summary": "List Ingestions",
"description": "List recent completed ingestion jobs, newest first. Used by the web UI's ingestion picker.",
"operationId": "legacy_ingestions_list_ingestions__ingestions",
"responses": {
"200": {
"description": "Successful Response",
"content": {
"application/json": {
"schema": {
"$ref": "#/components/schemas/IngestionListResponse"
}
}
}
}
},
"deprecated": true
},
"post": {
"tags": [
"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.\n\n``mode=tail`` creates a long-lived job that the worker polls; the response\nincludes ``ingestion_job_id`` and no worker_job_id.\n\nOptional ``callback_url`` (http/https) is stored on batch worker jobs. The\nworker POSTs an HMAC-signed completion payload on terminal state. Tail jobs\ndo not fire callbacks (long-lived; poll or ``:stop`` instead).",
"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.\n\nOptional ``callback_url`` (http/https) is stored on batch worker jobs. The\nworker POSTs an HMAC-signed completion payload on terminal state. Tail jobs\ndo not fire callbacks (long-lived; poll or ``:stop`` instead).\n\nOptional ``Idempotency-Key`` (also accepted as ``idempotency-key``): a\nrepeat within the TTL returns the original 202 job rather than a new one.",
"operationId": "legacy_ingestions_create_ingestion__ingestions",
"deprecated": true,
"parameters": [
{
"name": "Idempotency-Key",
"in": "header",
"required": false,
"schema": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"description": "If set, a repeat POST within INGEST_IDEMPOTENCY_TTL_SECONDS (default 86400) returns the original job instead of enqueueing a new one. Empty values are rejected with 400. Applies to batch enqueue and tail create.",
"title": "Idempotency-Key"
},
"description": "If set, a repeat POST within INGEST_IDEMPOTENCY_TTL_SECONDS (default 86400) returns the original job instead of enqueueing a new one. Empty values are rejected with 400. Applies to batch enqueue and tail create."
}
],
"requestBody": {
"required": true,
"content": {
"application/json": {
"schema": {
"$ref": "#/components/schemas/IngestRequest"
}
}
},
"required": true
}
},
"responses": {
"202": {
Expand All @@ -425,8 +445,28 @@
}
}
}
},
"deprecated": true
}
},
"get": {
"tags": [
"ingestions"
],
"summary": "List Ingestions",
"description": "List recent completed ingestion jobs, newest first. Used by the web UI's ingestion picker.",
"operationId": "legacy_ingestions_list_ingestions__ingestions",
"deprecated": true,
"responses": {
"200": {
"description": "Successful Response",
"content": {
"application/json": {
"schema": {
"$ref": "#/components/schemas/IngestionListResponse"
}
}
}
}
}
}
},
"/ingestions/jobs/{worker_job_id}": {
Expand Down
78 changes: 78 additions & 0 deletions migrations/versions/0007_ingest_idempotency.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,78 @@
"""ingest idempotency keys and log-entry content dedup

Revision ID: 0007_ingest_idempotency
Revises: 0006_api_key_webhook_secret
Create Date: 2026-08-17 00:00:00.000000
"""

from typing import Sequence, Union

import sqlalchemy as sa
from sqlalchemy.dialects import postgresql

from alembic import op

revision: str = "0007_ingest_idempotency"
down_revision: Union[str, None] = "0006_api_key_webhook_secret"
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None


def upgrade() -> None:
op.add_column(
"log_entries",
sa.Column("original_line_hash", sa.String(64), nullable=True),
)
op.add_column(
"log_entries",
sa.Column("scope", sa.String(255), nullable=False, server_default="default"),
)
op.create_index(
"ux_log_entries_dedup",
"log_entries",
["scope", "source_ref", "original_line_hash", "timestamp"],
unique=True,
postgresql_where=sa.text(
"original_line_hash IS NOT NULL AND timestamp IS NOT NULL "
"AND source_ref IS NOT NULL"
),
)
op.create_table(
"ingest_idempotency_keys",
sa.Column("key", sa.String(256), primary_key=True),
sa.Column(
"worker_job_id",
postgresql.UUID(as_uuid=True),
sa.ForeignKey("worker_jobs.id"),
nullable=True,
),
sa.Column(
"ingestion_job_id",
postgresql.UUID(as_uuid=True),
sa.ForeignKey("ingestion_jobs.id"),
nullable=True,
),
sa.Column("mode", sa.String(20), nullable=False, server_default="batch"),
sa.Column(
"created_at",
sa.DateTime(timezone=True),
server_default=sa.func.now(),
),
sa.Column("expires_at", sa.DateTime(timezone=True), nullable=False),
)
op.create_index(
"ix_ingest_idempotency_keys_expires_at",
"ingest_idempotency_keys",
["expires_at"],
)


def downgrade() -> None:
op.drop_index(
"ix_ingest_idempotency_keys_expires_at",
table_name="ingest_idempotency_keys",
)
op.drop_table("ingest_idempotency_keys")
op.drop_index("ux_log_entries_dedup", table_name="log_entries")
op.drop_column("log_entries", "scope")
op.drop_column("log_entries", "original_line_hash")
Loading
Loading