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
6 changes: 3 additions & 3 deletions openapi-spec/openapi.json
Original file line number Diff line number Diff line change
Expand Up @@ -8581,7 +8581,7 @@
"System"
],
"summary": "Get Schema Status",
"description": "What the last schema bootstrap pass did, object by object.\n\nThis is the operator's record of the schema apply, replacing the completed\nArgoCD Job the engine took over from. ``state`` is what readiness follows:\nconverged or observed is ready, failed leaves the pod up and NotReady with\nthe cause here.",
"description": "What the last schema bootstrap pass did, object by object.\n\nThis is the operator's record of the schema apply, replacing the completed\nArgoCD Job the engine took over from. ``state`` is what readiness follows:\nconverged or observed is ready, failed leaves the pod up and NotReady with\nthe stage that failed here and the cause in the engine log.",
"operationId": "get_schema_status_api_v1_system_schema_get",
"responses": {
"200": {
Expand Down Expand Up @@ -35628,7 +35628,7 @@
"error": {
"type": "string",
"title": "Error",
"description": "Why the pass failed; empty when it did not."
"description": "The stage that failed, its cause in the engine log; empty when none did."
},
"counts": {
"additionalProperties": {
Expand Down Expand Up @@ -39342,7 +39342,7 @@
"error"
],
"title": "TopicFailure",
"description": "One topic the broker would not do the thing to, and what it said."
"description": "One topic the broker would not do the thing to; the broker's own text is in the log."
},
"TopicRemoveResponse": {
"properties": {
Expand Down
87 changes: 83 additions & 4 deletions src/dfe_engine/api/errors.py
Original file line number Diff line number Diff line change
Expand Up @@ -126,6 +126,71 @@ class ErrorCode:
UNRESOLVED_REFERENCE = "unresolved_reference"


SERVICE_UNAVAILABLE_MESSAGE = "Backing service is unreachable, retry shortly"
GITOPS_UNAVAILABLE_MESSAGE = "The deploy repo is unreachable, retry after the Retry-After interval"


def hide_backend_text(message: str, exc: BaseException, *, event: str, **fields: Any) -> str:
"""Log ``exc``'s text under ``event`` and return ``message`` for the caller in its place.

A backend's own error text can carry statement fragments, user names and password
hashes, so it reaches the engine log and never a response.

Args:
message: What the caller is told instead.
exc: The backend failure whose text is logged.
event: The log line's message.
**fields: Further structured fields for the log line.

Returns:
``message``, unchanged.
"""
logger.warning(event, error=str(exc), **fields)
return message


def backend_failure(
status_code: int, code: str, message: str, exc: BaseException, *, event: str, **fields: Any
) -> HTTPException:
"""The HTTPException a route raises for a backend failure, its text logged instead of sent.

Args:
status_code: The HTTP status to answer.
code: The machine-readable error code.
message: What the caller is told.
exc: The backend failure whose text is logged.
event: The log line's message.
**fields: Further structured fields for the log line.

Returns:
The exception to raise, chained from ``exc`` by the caller.
"""
return HTTPException(
status_code=status_code,
detail={"code": code, "message": hide_backend_text(message, exc, event=event, **fields)},
)


def engine_message(exc: BaseException, fallback: str, *, event: str, **fields: Any) -> str:
"""``exc``'s own message when the engine wrote it, else ``fallback`` with its text logged.

The error types this serves are raised ``from`` a backend exception whenever they
carry that backend's text, and with no cause when the engine composed the message.

Args:
exc: The error a route is about to report.
fallback: What the caller is told when ``exc`` carries a backend's text.
event: The log line's message.
**fields: Further structured fields for the log line.

Returns:
The message that is safe to send.
"""
if exc.__cause__ is None:
return str(exc)
return hide_backend_text(fallback, exc, event=event, **fields)


def raise_exchange_http(exc: Exception) -> NoReturn:
"""Map an import/export failure to its HTTP answer.

Expand Down Expand Up @@ -243,9 +308,18 @@ async def protected_account_handler(request: Request, exc: ProtectedAccountError
from scalo.resilience import ServiceUnavailable

@app.exception_handler(ServiceUnavailable)
async def service_unavailable_handler(_request: Request, exc: ServiceUnavailable):
async def service_unavailable_handler(request: Request, exc: ServiceUnavailable):
waking = bool(getattr(exc, "waking", False))
message = "Backing service is warming up, retry shortly" if waking else str(exc)
if waking:
message = "Backing service is warming up, retry shortly"
else:
# scalo's message ends with the backend's last error, verbatim.
message = hide_backend_text(
SERVICE_UNAVAILABLE_MESSAGE,
exc,
event="backing service unavailable",
path=request.url.path,
)
body = ErrorResponse(
code=ErrorCode.SERVICE_UNAVAILABLE,
message=message,
Expand All @@ -258,10 +332,15 @@ async def service_unavailable_handler(_request: Request, exc: ServiceUnavailable
from dfe_engine.gitops.repo import GitopsUnavailableError

@app.exception_handler(GitopsUnavailableError)
async def gitops_unavailable_handler(_request: Request, exc: GitopsUnavailableError):
async def gitops_unavailable_handler(request: Request, exc: GitopsUnavailableError):
body = ErrorResponse(
code=ErrorCode.SERVICE_UNAVAILABLE,
message=str(exc),
message=hide_backend_text(
GITOPS_UNAVAILABLE_MESSAGE,
exc,
event="deploy repo unavailable",
path=request.url.path,
),
context={"service": "gitops", "remote": exc.remote},
)
return JSONResponse(
Expand Down
23 changes: 20 additions & 3 deletions src/dfe_engine/api/task_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,9 @@
from typing import Any

from pydantic import BaseModel, Field
from scalo.logger import logger

TASK_FAILED_MESSAGE = "The task failed; the engine log has the reason"


def _json_safe(value: Any) -> Any:
Expand Down Expand Up @@ -159,20 +162,33 @@ def __init__(self, max_completed: int = 1000) -> None:
self._max_completed = max_completed

def submit(
self, kind: str, coro_fn, *args, held_to: Iterable[str] | None = None, **kwargs
self,
kind: str,
coro_fn,
*args,
held_to: Iterable[str] | None = None,
reportable: tuple[type[Exception], ...] = (),
**kwargs,
) -> TaskInfo:
"""Submit an async callable for background execution.

The coroutine function receives a ``task`` keyword argument -- a
``_Task`` instance whose ``set_progress()`` method can be called
to report progress.

Every failure is logged. Its message reaches the task's readers only when
it is one of ``reportable``; any other failure reads as
:data:`TASK_FAILED_MESSAGE`, because a backend's own error text can carry
statement fragments, user names and password hashes.

Args:
kind: The task kind, for ``list(kind=...)``.
coro_fn: The coroutine function to run, called with ``*args`` and ``**kwargs``.
*args: Positional arguments for ``coro_fn``.
held_to: The orgs the task's result is held to. None means it may hold
any org's data, so only a reader of every org sees it.
reportable: Exception types whose message the engine composes about the
request, and which a reader may therefore see.
**kwargs: Keyword arguments for ``coro_fn``.

Returns:
Expand All @@ -196,9 +212,10 @@ async def _run() -> None:
task.status = TaskStatus.CANCELLED
task.message = "Cancelled"
except Exception as exc:
logger.warning("background task failed", task_id=task.id, kind=kind, error=str(exc))
task.status = TaskStatus.FAILED
task.error = str(exc)
task.message = f"Failed: {exc}"
task.error = str(exc) if isinstance(exc, reportable) else TASK_FAILED_MESSAGE
task.message = f"Failed: {task.error}"
finally:
task.completed_at = datetime.now(UTC).isoformat()
task._progress_event.set()
Expand Down
25 changes: 15 additions & 10 deletions src/dfe_engine/api/v1/apps.py
Original file line number Diff line number Diff line change
Expand Up @@ -65,7 +65,7 @@
get_source_registry,
require_action,
)
from dfe_engine.api.errors import ErrorResponse
from dfe_engine.api.errors import ErrorResponse, backend_failure
from dfe_engine.api.v1.app_contracts import read_contract
from dfe_engine.api.v1.helm import conflict_error
from dfe_engine.api.v1.sampler import check_sample_scope, sample_org_scope
Expand Down Expand Up @@ -2149,6 +2149,17 @@ def _reader(request: Request, client: Any) -> OperationalReader:
return OperationalReader(client, settings.clickhouse.effective_data_database)


def metrics_unavailable(exc: MetricsUnavailableError) -> HTTPException:
"""The 503 for an otel read ClickHouse refused, its text logged rather than sent."""
return backend_failure(
503,
"metrics_unavailable",
"the telemetry tables could not be read; the engine log has ClickHouse's reason",
exc,
event="otel query failed",
)


@router.get("/{service}/{instance}/status")
def get_status(
service: str, instance: str, user: CurrentUser, request: Request, client: ClickHouseClient
Expand All @@ -2159,9 +2170,7 @@ def get_status(
try:
status = _reader(request, client).status(app.telemetry_name)
except MetricsUnavailableError as exc:
raise HTTPException(
503, detail={"code": "metrics_unavailable", "message": str(exc)}
) from exc
raise metrics_unavailable(exc) from exc
return StatusResponse(
telemetry_name=status.telemetry_name,
reporting=status.reporting,
Expand All @@ -2181,9 +2190,7 @@ def get_metrics(
try:
found = _reader(request, client).metrics(app.telemetry_name)
except MetricsUnavailableError as exc:
raise HTTPException(
503, detail={"code": "metrics_unavailable", "message": str(exc)}
) from exc
raise metrics_unavailable(exc) from exc
return MetricsResponse(
telemetry_name=found.telemetry_name,
window_seconds=found.window_seconds,
Expand Down Expand Up @@ -2217,9 +2224,7 @@ def get_resource_series(
app.telemetry_name, window_seconds=window, bucket_seconds=bucket
)
except MetricsUnavailableError as exc:
raise HTTPException(
503, detail={"code": "metrics_unavailable", "message": str(exc)}
) from exc
raise metrics_unavailable(exc) from exc
return ResourceSeriesResponse(
telemetry_name=app.telemetry_name,
window_seconds=window,
Expand Down
27 changes: 17 additions & 10 deletions src/dfe_engine/api/v1/discovery.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,10 +21,17 @@
from pydantic import BaseModel

from dfe_engine.api.deps import CurrentUser, require_action
from dfe_engine.api.errors import backend_failure
from dfe_engine.auth.rbac_scopes import scopes_dict

router = APIRouter(prefix="/discovery", tags=["discovery"])

_QUERY_FAILED = "The ClickHouse metadata query failed; the engine log has ClickHouse's reason"


def _query_failed(exc: Exception) -> HTTPException:
return backend_failure(500, "query_error", _QUERY_FAILED, exc, event="discovery query failed")


# -- Response models -----------------------------------------

Expand Down Expand Up @@ -75,13 +82,13 @@ def _get_ch_client(request: Request):
# Discovery is an admin-level cross-database listing, so it reads over "default".
return conn_registry.get_client("default")
except Exception as exc:
raise HTTPException(
status_code=503,
detail={
"code": "connection_error",
"message": f"Cannot connect to ClickHouse: {exc}",
},
)
raise backend_failure(
503,
"connection_error",
"Cannot connect to ClickHouse; the engine log has the reason",
exc,
event="discovery: ClickHouse connection failed",
) from exc


# -- Endpoints -----------------------------------------------
Expand All @@ -107,7 +114,7 @@ async def list_databases(
if row[0] not in ("system", "information_schema", "INFORMATION_SCHEMA")
]
except Exception as exc:
raise HTTPException(status_code=500, detail={"code": "query_error", "message": str(exc)})
raise _query_failed(exc) from exc


@router.get(
Expand Down Expand Up @@ -149,7 +156,7 @@ async def list_tables(
for row in result.result_rows
]
except Exception as exc:
raise HTTPException(status_code=500, detail={"code": "query_error", "message": str(exc)})
raise _query_failed(exc) from exc


@router.get(
Expand Down Expand Up @@ -195,4 +202,4 @@ async def list_columns(
except HTTPException:
raise
except Exception as exc:
raise HTTPException(status_code=500, detail={"code": "query_error", "message": str(exc)})
raise _query_failed(exc) from exc
22 changes: 13 additions & 9 deletions src/dfe_engine/api/v1/governance.py
Original file line number Diff line number Diff line change
Expand Up @@ -24,9 +24,9 @@

from fastapi import APIRouter, Depends, HTTPException, Query, Request, Response
from pydantic import BaseModel
from scalo.logger import logger

from dfe_engine.api.deps import CurrentUser, live_role_config, require_action
from dfe_engine.api.errors import backend_failure
from dfe_engine.api.review import apply_review_headers
from dfe_engine.api.write_turn import WRITE_TURN
from dfe_engine.appmgmt import contract
Expand Down Expand Up @@ -446,10 +446,12 @@ def reconcile_ch_rbac_endpoint(user: CurrentUser, request: Request) -> dict[str,
try:
admin_client = ch_admin_client(settings)
except Exception as exc:
logger.warning("CH RBAC reconcile: ClickHouse unreachable", error=str(exc))
raise HTTPException(
status_code=503,
detail={"code": "clickhouse_unavailable", "message": _CH_UNAVAILABLE},
raise backend_failure(
503,
"clickhouse_unavailable",
_CH_UNAVAILABLE,
exc,
event="CH RBAC reconcile: ClickHouse unreachable",
) from exc

try:
Expand All @@ -463,10 +465,12 @@ def reconcile_ch_rbac_endpoint(user: CurrentUser, request: Request) -> dict[str,
except Exception as exc:
if not is_connection_error(exc):
raise
logger.warning("CH RBAC reconcile: ClickHouse connection lost", error=str(exc))
raise HTTPException(
status_code=503,
detail={"code": "clickhouse_unavailable", "message": _CH_UNAVAILABLE},
raise backend_failure(
503,
"clickhouse_unavailable",
_CH_UNAVAILABLE,
exc,
event="CH RBAC reconcile: ClickHouse connection lost",
) from exc
# A full run covers every identity, so a startup retry still waiting stands down.
note_ch_rbac_reconciled(request.app.state)
Expand Down
2 changes: 1 addition & 1 deletion src/dfe_engine/api/v1/kafka_topics.py
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,7 @@


class TopicFailure(BaseModel):
"""One topic the broker would not do the thing to, and what it said."""
"""One topic the broker would not do the thing to; the broker's own text is in the log."""

name: str
error: str
Expand Down
8 changes: 6 additions & 2 deletions src/dfe_engine/api/v1/oidc_login.py
Original file line number Diff line number Diff line change
Expand Up @@ -239,10 +239,14 @@ async def oidc_callback(
logger.warning("OIDC callback failed", provider=provider, error=str(exc))
# A refused credential is an audit event, not just an operational log line.
audit_login_denied("unknown", "oidc", _get_client_ip(request), str(exc))
# The IdP's own answer stays in the log: this route takes unauthenticated callers.
raise HTTPException(
status_code=401,
detail={"code": "unauthorized", "message": f"OIDC login failed: {exc}"},
)
detail={
"code": "unauthorized",
"message": "OIDC login failed; the engine log has the reason",
},
) from exc

if not identity.subject:
audit_login_denied("unknown", "oidc", _get_client_ip(request), "no_subject")
Expand Down
Loading
Loading