Skip to content
Closed
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
7 changes: 7 additions & 0 deletions chart/interloper/templates/NOTES.txt
Original file line number Diff line number Diff line change
Expand Up @@ -32,5 +32,12 @@ HTTPRoute hostnames: {{ join ", " .Values.httpRoute.hostnames }}
Generate one and re-apply: openssl rand -base64 32
{{- end }}

{{- if and (eq (default "" (dig "runner" "type" "" .Values.config)) "k8s") (not .Values.secrets.eventIngestToken) }}

⚠ No secrets.eventIngestToken set — per-asset workers will persist events via
log-scraping (best-effort) instead of the durable ingest endpoint.
Generate one and re-apply: openssl rand -base64 32
{{- end }}

Tail the scheduler logs:
kubectl -n {{ .Release.Namespace }} logs -f deploy/{{ include "interloper.fullname" . }}-scheduler
20 changes: 20 additions & 0 deletions chart/interloper/templates/_helpers.tpl
Original file line number Diff line number Diff line change
Expand Up @@ -147,6 +147,18 @@ Postgres host: either the bundled subchart service or external config.
{{- end -}}
{{- end -}}

{{/*
Base API URL that worker pods POST events to. Defaults to the in-cluster API
service; override via events.ingestUrl.
*/}}
{{- define "interloper.eventIngestUrl" -}}
{{- if .Values.events.ingestUrl -}}
{{ .Values.events.ingestUrl }}
{{- else -}}
http://{{ include "interloper.fullname" . }}-api.{{ .Release.Namespace }}.svc:{{ .Values.api.service.port }}
{{- end -}}
{{- end -}}

{{/*
Environment variables shared by all interloper pods (scheduler, api, frontend).
Contains Postgres connection info + the encryption key secret.
Expand Down Expand Up @@ -177,4 +189,12 @@ Contains Postgres connection info + the encryption key secret.
name: {{ include "interloper.secretName" . }}
key: INTERLOPER_AUTH_GOOGLE_CLIENT_SECRET
optional: true
- name: INTERLOPER_EVENTS_INGEST_URL
value: {{ include "interloper.eventIngestUrl" . | quote }}
- name: INTERLOPER_EVENTS_INGEST_TOKEN
valueFrom:
secretKeyRef:
name: {{ include "interloper.secretName" . }}
key: INTERLOPER_EVENTS_INGEST_TOKEN
optional: true
{{- end -}}
3 changes: 3 additions & 0 deletions chart/interloper/templates/secret.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -17,4 +17,7 @@ stringData:
{{- if .Values.secrets.googleClientSecret }}
INTERLOPER_AUTH_GOOGLE_CLIENT_SECRET: {{ .Values.secrets.googleClientSecret | quote }}
{{- end }}
{{- if .Values.secrets.eventIngestToken }}
INTERLOPER_EVENTS_INGEST_TOKEN: {{ .Values.secrets.eventIngestToken | quote }}
{{- end }}
{{- end }}
14 changes: 14 additions & 0 deletions chart/interloper/values.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -187,6 +187,20 @@ secrets:
encryptionKey: "" # recommended to generate: openssl rand -base64 32
smtpPassword: ""
googleClientSecret: "" # OAuth client secret (paired with config.auth.google_client_id)
# Shared service token for the internal event-ingest endpoint. When set,
# per-asset worker pods POST their events to the API instead of relying on
# log-scraping. When empty, the ingest endpoint is disabled (workers fall
# back to the stderr path). Generate: openssl rand -base64 32
eventIngestToken: ""

# =============================================================================
# Events
# =============================================================================

events:
# Base API URL that worker pods POST events to. Defaults to the in-cluster
# API service; override only if workers must reach the API by another route.
ingestUrl: ""

# =============================================================================
# Postgres
Expand Down
12 changes: 11 additions & 1 deletion packages/interloper-api/src/interloper_api/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,14 +9,21 @@
from interloper.catalog.base import Catalog
from interloper_db import Store

from interloper_api.dependencies import set_auth_config, set_catalog, set_smtp_config, set_store
from interloper_api.dependencies import (
set_auth_config,
set_catalog,
set_ingest_token,
set_smtp_config,
set_store,
)
from interloper_api.routes import (
admin,
assets,
auth,
backfills,
destinations,
external,
internal,
jobs,
oauth,
organisations,
Expand All @@ -33,6 +40,7 @@ def create_app(
catalog: Catalog | None = None,
auth_config: Any | None = None,
smtp_config: Any | None = None,
event_ingest_token: str | None = None,
cors_origins: list[str] | None = None,
**kwargs: Any,
) -> FastAPI:
Expand Down Expand Up @@ -69,6 +77,7 @@ def create_app(
set_auth_config(auth_config)
if smtp_config:
set_smtp_config(smtp_config)
set_ingest_token(event_ingest_token)

api = APIRouter(prefix="/api")
api.include_router(auth.router, tags=["auth"])
Expand All @@ -80,6 +89,7 @@ def create_app(
api.include_router(destinations.router, prefix="/destinations", tags=["destinations"])
api.include_router(jobs.router, prefix="/jobs", tags=["jobs"])
api.include_router(runs.router, prefix="/runs", tags=["runs"])
api.include_router(internal.router, prefix="/internal", tags=["internal"])
api.include_router(backfills.router, prefix="/backfills", tags=["backfills"])
api.include_router(assets.router, prefix="/assets", tags=["assets"])
api.include_router(oauth.router, tags=["oauth"])
Expand Down
52 changes: 51 additions & 1 deletion packages/interloper-api/src/interloper_api/dependencies.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,10 +2,11 @@

from __future__ import annotations

import hmac
from typing import Any
from uuid import UUID

from fastapi import Cookie, Depends, HTTPException
from fastapi import Cookie, Depends, Header, HTTPException
from interloper.catalog.base import Catalog
from interloper_db import Organisation, Profile, Store
from interloper_db.models import Session as SessionModel
Expand All @@ -14,6 +15,7 @@
_catalog: Catalog | None = None
_auth_config: Any | None = None
_smtp_config: Any | None = None
_ingest_token: str | None = None

# Role hierarchy: admin > editor > viewer
_ROLE_RANK = {"viewer": 0, "editor": 1, "admin": 2}
Expand Down Expand Up @@ -110,6 +112,21 @@ def get_smtp_config() -> Any:
return _smtp_config


def set_ingest_token(token: str | None) -> None:
"""Set the shared service token that authenticates internal event ingest.

Args:
token: The secret, or ``None``/empty to disable the ingest endpoint.
"""
global _ingest_token # noqa: PLW0603
_ingest_token = token or None


def get_ingest_token() -> str | None:
"""Return the configured event-ingest token, or ``None`` if disabled."""
return _ingest_token


# -- Auth dependencies -------------------------------------------------------


Expand Down Expand Up @@ -310,3 +327,36 @@ def require_super_admin(
if not user.is_super_admin:
raise HTTPException(status_code=403, detail="Requires super-admin privileges")
return user


# -- Service-token auth (internal machine-to-machine) ------------------------


def _bearer_token(authorization: str | None) -> str | None:
"""Extract the token from an ``Authorization: Bearer <token>`` header."""
if not authorization:
return None
scheme, _, token = authorization.partition(" ")
if scheme.lower() != "bearer" or not token:
return None
return token


def require_ingest_token(authorization: str | None = Header(default=None)) -> None:
"""Authenticate an internal event-ingest request via a bearer service token.

This is machine-to-machine auth for trusted in-cluster callers (per-asset
worker processes), deliberately separate from the cookie-session user auth.
The token is a shared secret set via :func:`set_ingest_token`; when unset,
the ingest surface is disabled entirely.

Raises:
HTTPException: 503 if ingest is not configured, 401 if the token is
missing or wrong.
"""
expected = get_ingest_token()
if not expected:
raise HTTPException(status_code=503, detail="Event ingest is not configured")
provided = _bearer_token(authorization)
if not provided or not hmac.compare_digest(provided, expected):
raise HTTPException(status_code=401, detail="Invalid or missing ingest token")
77 changes: 77 additions & 0 deletions packages/interloper-api/src/interloper_api/routes/internal.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
"""Internal (machine-to-machine) API: reliable event ingest.

Per-asset worker processes run in slim, DB-less pods, so instead of routing
their events through stdout/log-scraping (lossy), they POST them here. The
endpoint is authenticated with a shared service token (not the user session)
and writes each event with :meth:`Store.save_event`, which is idempotent — so
the worker can retry a batch freely without creating duplicates.
"""

from __future__ import annotations

import logging
from typing import Any
from uuid import UUID

import interloper as il
from fastapi import APIRouter, Depends, HTTPException
from interloper.errors import NotFoundError
from interloper_db import Store
from pydantic import BaseModel

from interloper_api.dependencies import get_store, require_ingest_token

logger = logging.getLogger(__name__)

router = APIRouter()


class EventIngestRequest(BaseModel):
"""A batch of serialized events for a single run.

Each entry is the flat dict produced by ``Event.to_dict()`` (``event_id``,
``type``, ``timestamp`` plus inlined metadata).
"""

events: list[dict[str, Any]]


class EventIngestResponse(BaseModel):
"""Result of an ingest batch."""

accepted: int
rejected: int


@router.post("/runs/{run_id}/events")
def ingest_run_events(
run_id: UUID,
body: EventIngestRequest,
_: None = Depends(require_ingest_token),
store: Store = Depends(get_store),
) -> EventIngestResponse:
"""Persist a batch of events for *run_id*.

The ``org_id`` is resolved server-side from the run, so callers only need
the run id. Malformed events are skipped (and counted in ``rejected``)
rather than failing the batch; a persistence error is allowed to surface as
a 5xx so the caller retries — idempotent writes make that safe.
"""
try:
run = store.get_run(run_id)
except NotFoundError:
raise HTTPException(status_code=404, detail=f"Run {run_id} not found")

accepted = 0
rejected = 0
for raw in body.events:
try:
event = il.Event.from_dict(raw)
except Exception: # noqa: BLE001 - malformed event: skip, don't retry
rejected += 1
logger.warning("Rejected malformed event for run %s", run_id)
continue
store.save_event(event, org_id=run.org_id, run_id=run_id)
accepted += 1

return EventIngestResponse(accepted=accepted, rejected=rejected)
Loading
Loading