diff --git a/docs/agent-tools/claude-tools.json b/docs/agent-tools/claude-tools.json index 1897227..46fb112 100644 --- a/docs/agent-tools/claude-tools.json +++ b/docs/agent-tools/claude-tools.json @@ -882,6 +882,151 @@ } } }, + { + "name": "list_exceptions_v1_ops_exceptions_get", + "description": "List Exceptions", + "input_schema": { + "type": "object", + "properties": { + "source": { + "anyOf": [ + { + "enum": [ + "deadletter", + "webhook_delivery", + "reconciliation" + ], + "type": "string" + }, + { + "type": "null" + } + ], + "title": "Source" + }, + "status": { + "anyOf": [ + { + "enum": [ + "open", + "in_progress", + "acknowledged", + "resolved" + ], + "type": "string" + }, + { + "type": "null" + } + ], + "title": "Status" + }, + "page": { + "type": "integer", + "minimum": 1, + "default": 1, + "title": "Page" + }, + "page_size": { + "type": "integer", + "maximum": 100, + "minimum": 1, + "default": 50, + "title": "Page Size" + } + } + } + }, + { + "name": "exceptions_stats_v1_ops_exceptions_stats_get", + "description": "Exceptions Stats", + "input_schema": { + "type": "object", + "properties": {} + } + }, + { + "name": "acknowledge_exception_v1_ops_exceptions__item_id__acknowledge_post", + "description": "Acknowledge Exception", + "input_schema": { + "type": "object", + "properties": { + "item_id": { + "type": "string", + "title": "Item Id" + }, + "body": { + "anyOf": [ + { + "properties": { + "note": { + "anyOf": [ + { + "type": "string" + }, + { + "type": "null" + } + ], + "title": "Note" + } + }, + "type": "object", + "title": "TriageActionRequest" + }, + { + "type": "null" + } + ], + "title": "Payload" + } + }, + "required": [ + "item_id" + ] + } + }, + { + "name": "resolve_exception_v1_ops_exceptions__item_id__resolve_post", + "description": "Resolve Exception", + "input_schema": { + "type": "object", + "properties": { + "item_id": { + "type": "string", + "title": "Item Id" + }, + "body": { + "anyOf": [ + { + "properties": { + "note": { + "anyOf": [ + { + "type": "string" + }, + { + "type": "null" + } + ], + "title": "Note" + } + }, + "type": "object", + "title": "TriageActionRequest" + }, + { + "type": "null" + } + ], + "title": "Payload" + } + }, + "required": [ + "item_id" + ] + } + }, { "name": "search_v1_search_get", "description": "Search", diff --git a/docs/agent-tools/openai-tools.json b/docs/agent-tools/openai-tools.json index 3f68f2f..bfb69e4 100644 --- a/docs/agent-tools/openai-tools.json +++ b/docs/agent-tools/openai-tools.json @@ -981,6 +981,163 @@ } } }, + { + "type": "function", + "function": { + "name": "list_exceptions_v1_ops_exceptions_get", + "description": "List Exceptions", + "parameters": { + "type": "object", + "properties": { + "source": { + "anyOf": [ + { + "enum": [ + "deadletter", + "webhook_delivery", + "reconciliation" + ], + "type": "string" + }, + { + "type": "null" + } + ], + "title": "Source" + }, + "status": { + "anyOf": [ + { + "enum": [ + "open", + "in_progress", + "acknowledged", + "resolved" + ], + "type": "string" + }, + { + "type": "null" + } + ], + "title": "Status" + }, + "page": { + "type": "integer", + "minimum": 1, + "default": 1, + "title": "Page" + }, + "page_size": { + "type": "integer", + "maximum": 100, + "minimum": 1, + "default": 50, + "title": "Page Size" + } + } + } + } + }, + { + "type": "function", + "function": { + "name": "exceptions_stats_v1_ops_exceptions_stats_get", + "description": "Exceptions Stats", + "parameters": { + "type": "object", + "properties": {} + } + } + }, + { + "type": "function", + "function": { + "name": "acknowledge_exception_v1_ops_exceptions__item_id__acknowledge_post", + "description": "Acknowledge Exception", + "parameters": { + "type": "object", + "properties": { + "item_id": { + "type": "string", + "title": "Item Id" + }, + "body": { + "anyOf": [ + { + "properties": { + "note": { + "anyOf": [ + { + "type": "string" + }, + { + "type": "null" + } + ], + "title": "Note" + } + }, + "type": "object", + "title": "TriageActionRequest" + }, + { + "type": "null" + } + ], + "title": "Payload" + } + }, + "required": [ + "item_id" + ] + } + } + }, + { + "type": "function", + "function": { + "name": "resolve_exception_v1_ops_exceptions__item_id__resolve_post", + "description": "Resolve Exception", + "parameters": { + "type": "object", + "properties": { + "item_id": { + "type": "string", + "title": "Item Id" + }, + "body": { + "anyOf": [ + { + "properties": { + "note": { + "anyOf": [ + { + "type": "string" + }, + { + "type": "null" + } + ], + "title": "Note" + } + }, + "type": "object", + "title": "TriageActionRequest" + }, + { + "type": "null" + } + ], + "title": "Payload" + } + }, + "required": [ + "item_id" + ] + } + } + }, { "type": "function", "function": { diff --git a/docs/openapi.json b/docs/openapi.json index f6a91ce..8a8cbf4 100644 --- a/docs/openapi.json +++ b/docs/openapi.json @@ -1857,6 +1857,244 @@ } } }, + "/v1/ops/exceptions": { + "get": { + "tags": [ + "ops" + ], + "summary": "List Exceptions", + "operationId": "list_exceptions_v1_ops_exceptions_get", + "parameters": [ + { + "name": "source", + "in": "query", + "required": false, + "schema": { + "anyOf": [ + { + "enum": [ + "deadletter", + "webhook_delivery", + "reconciliation" + ], + "type": "string" + }, + { + "type": "null" + } + ], + "title": "Source" + } + }, + { + "name": "status", + "in": "query", + "required": false, + "schema": { + "anyOf": [ + { + "enum": [ + "open", + "in_progress", + "acknowledged", + "resolved" + ], + "type": "string" + }, + { + "type": "null" + } + ], + "title": "Status" + } + }, + { + "name": "page", + "in": "query", + "required": false, + "schema": { + "type": "integer", + "minimum": 1, + "default": 1, + "title": "Page" + } + }, + { + "name": "page_size", + "in": "query", + "required": false, + "schema": { + "type": "integer", + "maximum": 100, + "minimum": 1, + "default": 50, + "title": "Page Size" + } + } + ], + "responses": { + "200": { + "description": "Successful Response", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/ExceptionsListResponse" + } + } + } + }, + "422": { + "description": "Validation Error", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/HTTPValidationError" + } + } + } + } + } + } + }, + "/v1/ops/exceptions/stats": { + "get": { + "tags": [ + "ops" + ], + "summary": "Exceptions Stats", + "operationId": "exceptions_stats_v1_ops_exceptions_stats_get", + "responses": { + "200": { + "description": "Successful Response", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/ExceptionsStatsResponse" + } + } + } + } + } + } + }, + "/v1/ops/exceptions/{item_id}/acknowledge": { + "post": { + "tags": [ + "ops" + ], + "summary": "Acknowledge Exception", + "operationId": "acknowledge_exception_v1_ops_exceptions__item_id__acknowledge_post", + "parameters": [ + { + "name": "item_id", + "in": "path", + "required": true, + "schema": { + "type": "string", + "title": "Item Id" + } + } + ], + "requestBody": { + "content": { + "application/json": { + "schema": { + "anyOf": [ + { + "$ref": "#/components/schemas/TriageActionRequest" + }, + { + "type": "null" + } + ], + "title": "Payload" + } + } + } + }, + "responses": { + "200": { + "description": "Successful Response", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/TriageActionResponse" + } + } + } + }, + "422": { + "description": "Validation Error", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/HTTPValidationError" + } + } + } + } + } + } + }, + "/v1/ops/exceptions/{item_id}/resolve": { + "post": { + "tags": [ + "ops" + ], + "summary": "Resolve Exception", + "operationId": "resolve_exception_v1_ops_exceptions__item_id__resolve_post", + "parameters": [ + { + "name": "item_id", + "in": "path", + "required": true, + "schema": { + "type": "string", + "title": "Item Id" + } + } + ], + "requestBody": { + "content": { + "application/json": { + "schema": { + "anyOf": [ + { + "$ref": "#/components/schemas/TriageActionRequest" + }, + { + "type": "null" + } + ], + "title": "Payload" + } + } + } + }, + "responses": { + "200": { + "description": "Successful Response", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/TriageActionResponse" + } + } + } + }, + "422": { + "description": "Validation Error", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/HTTPValidationError" + } + } + } + } + } + } + }, "/v1/search": { "get": { "tags": [ @@ -2775,6 +3013,29 @@ ], "title": "DismissResponse" }, + "EntityRef": { + "properties": { + "kind": { + "type": "string", + "enum": [ + "event", + "order", + "webhook" + ], + "title": "Kind" + }, + "id": { + "type": "string", + "title": "Id" + } + }, + "type": "object", + "required": [ + "kind", + "id" + ], + "title": "EntityRef" + }, "EntityResponse": { "properties": { "entity_type": { @@ -2827,6 +3088,154 @@ ], "title": "EntityResponse" }, + "ExceptionAction": { + "properties": { + "action": { + "type": "string", + "title": "Action" + }, + "href": { + "type": "string", + "title": "Href" + } + }, + "type": "object", + "required": [ + "action", + "href" + ], + "title": "ExceptionAction" + }, + "ExceptionItem": { + "properties": { + "item_id": { + "type": "string", + "title": "Item Id" + }, + "source": { + "type": "string", + "enum": [ + "deadletter", + "webhook_delivery", + "reconciliation" + ], + "title": "Source" + }, + "severity": { + "type": "string", + "enum": [ + "high", + "medium", + "low" + ], + "title": "Severity" + }, + "occurred_at": { + "type": "string", + "format": "date-time", + "title": "Occurred At" + }, + "last_seen_at": { + "type": "string", + "format": "date-time", + "title": "Last Seen At" + }, + "entity_ref": { + "$ref": "#/components/schemas/EntityRef" + }, + "title": { + "type": "string", + "title": "Title" + }, + "detail": { + "type": "string", + "title": "Detail" + }, + "status": { + "type": "string", + "enum": [ + "open", + "in_progress", + "acknowledged", + "resolved" + ], + "title": "Status" + }, + "actions": { + "items": { + "$ref": "#/components/schemas/ExceptionAction" + }, + "type": "array", + "title": "Actions" + } + }, + "type": "object", + "required": [ + "item_id", + "source", + "severity", + "occurred_at", + "last_seen_at", + "entity_ref", + "title", + "detail", + "status", + "actions" + ], + "title": "ExceptionItem" + }, + "ExceptionsListResponse": { + "properties": { + "items": { + "items": { + "$ref": "#/components/schemas/ExceptionItem" + }, + "type": "array", + "title": "Items" + }, + "pagination": { + "additionalProperties": { + "type": "integer" + }, + "type": "object", + "title": "Pagination" + } + }, + "type": "object", + "required": [ + "items", + "pagination" + ], + "title": "ExceptionsListResponse" + }, + "ExceptionsStatsResponse": { + "properties": { + "by_source": { + "additionalProperties": { + "additionalProperties": { + "type": "integer" + }, + "type": "object" + }, + "type": "object", + "title": "By Source" + }, + "last_24h": { + "type": "integer", + "title": "Last 24H" + }, + "manual_resolutions": { + "type": "integer", + "title": "Manual Resolutions" + } + }, + "type": "object", + "required": [ + "last_24h", + "manual_resolutions" + ], + "title": "ExceptionsStatsResponse" + }, "ExplainRequest": { "properties": { "question": { @@ -3959,6 +4368,41 @@ "type": "object", "title": "StuckOrdersSummary" }, + "TriageActionRequest": { + "properties": { + "note": { + "anyOf": [ + { + "type": "string" + }, + { + "type": "null" + } + ], + "title": "Note" + } + }, + "type": "object", + "title": "TriageActionRequest" + }, + "TriageActionResponse": { + "properties": { + "item_id": { + "type": "string", + "title": "Item Id" + }, + "status": { + "type": "string", + "title": "Status" + } + }, + "type": "object", + "required": [ + "item_id", + "status" + ], + "title": "TriageActionResponse" + }, "ValidationError": { "properties": { "loc": { diff --git a/src/serving/api/routers/ops.py b/src/serving/api/routers/ops.py index b87814f..fb6d461 100644 --- a/src/serving/api/routers/ops.py +++ b/src/serving/api/routers/ops.py @@ -1,22 +1,35 @@ """Ops surfaces — GET /v1/ops/stuck-orders, the stuck-orders worklist -(ops-surfaces-spec.md §3, D3). - -Composes exactly the QueryEngine port: the open-orders read -(``fetch_orders_by_status``) and the stage-clock journal read -(``fetch_pipeline_events``), the same two reads the Order 360 timeline uses. -No raw engine connection reach, no vault DSN (invariant I1). +(ops-surfaces-spec.md §3, D3); GET/POST /v1/ops/exceptions*, the exception +inbox (ops-surfaces-spec.md §4, D4). + +Composes exactly the QueryEngine port (the open-orders read +``fetch_orders_by_status`` and the stage-clock journal read +``fetch_pipeline_events``, the same two reads the Order 360 timeline uses) +and the ControlPlaneStore port (dead-letter reads, the triage overlay, +webhook dead-delivery reads). No raw engine connection reach, no vault DSN +(invariant I1). """ from __future__ import annotations import math -from datetime import datetime +from datetime import UTC, datetime, timedelta from typing import Any, Literal, cast -from fastapi import APIRouter, Query, Request +from fastapi import APIRouter, HTTPException, Query, Request from pydantic import BaseModel, Field from starlette.concurrency import run_in_threadpool +from src.serving.control_plane import ( + TriageState, + get_control_plane_store, + stuck_replay_threshold_seconds, +) +from src.serving.semantic_layer.reconciliation import ( + ReconciliationFinding, + check_journal_vs_store, + check_stuck_replay, +) from src.serving.semantic_layer.stage_clock import ( coerce_dt, ladder_stage_names, @@ -26,6 +39,14 @@ router = APIRouter(prefix="/v1/ops", tags=["ops"]) +_DEADLETTER_STATUS_MAP = { + "failed": "open", + "replay_pending": "in_progress", + "replayed": "resolved", + "dismissed": "resolved", +} +_SEVERITY_RANK = {"high": 0, "medium": 1, "low": 2} + class StuckOrderItem(BaseModel): order_id: str @@ -188,3 +209,346 @@ async def get_stuck_orders( _build_stuck_orders_payload, request, stage, include_within_sla, page, page_size ) return StuckOrdersResponse.model_validate(payload) + + +# --------------------------------------------------------------------------- +# Exception inbox (D4, ops-surfaces-spec.md §4) +# --------------------------------------------------------------------------- + + +class EntityRef(BaseModel): + kind: Literal["event", "order", "webhook"] + id: str + + +class ExceptionAction(BaseModel): + action: str + href: str + + +class ExceptionItem(BaseModel): + item_id: str + source: Literal["deadletter", "webhook_delivery", "reconciliation"] + severity: Literal["high", "medium", "low"] + occurred_at: datetime + last_seen_at: datetime + entity_ref: EntityRef + title: str + detail: str + status: Literal["open", "in_progress", "acknowledged", "resolved"] + actions: list[ExceptionAction] + + +class ExceptionsListResponse(BaseModel): + items: list[ExceptionItem] + pagination: dict[str, int] + + +class ExceptionsStatsResponse(BaseModel): + by_source: dict[str, dict[str, int]] = Field(default_factory=dict) + last_24h: int + manual_resolutions: int + + +class TriageActionRequest(BaseModel): + note: str | None = None + + +class TriageActionResponse(BaseModel): + item_id: str + status: str + + +def _tenant_id(request: Request) -> str: + tenant_key = getattr(request.state, "tenant_key", None) + tenant_id = getattr(request.state, "tenant_id", None) + return str(tenant_id or getattr(tenant_key, "tenant", "default")) + + +def _deadletter_row_to_item(row: dict[str, Any]) -> dict[str, Any]: + status = _DEADLETTER_STATUS_MAP.get(row["status"], "open") + occurred_at = coerce_dt(row.get("received_at")) or datetime.now(UTC) + last_seen_at = coerce_dt(row.get("last_retried_at")) or occurred_at + event_id = row["event_id"] + actions = ( + [] + if status == "resolved" + else [ + {"action": "replay", "href": f"/v1/deadletter/{event_id}/replay"}, + {"action": "dismiss", "href": f"/v1/deadletter/{event_id}/dismiss"}, + ] + ) + return { + "item_id": f"dl:{event_id}", + "source": "deadletter", + "severity": "high", + "occurred_at": occurred_at, + "last_seen_at": last_seen_at, + "entity_ref": {"kind": "event", "id": event_id}, + "title": f"Dead-letter: {row.get('failure_reason') or 'unknown reason'}", + "detail": row.get("failure_detail") or "", + "status": status, + "actions": actions, + } + + +def _webhook_delivery_row_to_item( + row: dict[str, Any], state: TriageState | None, now: datetime +) -> dict[str, Any]: + item_id = f"wh:{row['webhook_id']}:{row['event_id']}" + status = state.status if state is not None else "open" + fallback_at = coerce_dt(row.get("updated_at")) or now + # Overlay timestamps come back naive from DuckDB, aware from PostgreSQL + # (TIMESTAMPTZ) — coerce_dt normalizes either to aware UTC. + occurred_at = fallback_at + last_seen_at = fallback_at + if state is not None: + occurred_at = coerce_dt(state.first_seen_at) or fallback_at + last_seen_at = coerce_dt(state.last_seen_at) or fallback_at + actions = ( + [] + if status == "resolved" + else [ + {"action": "acknowledge", "href": f"/v1/ops/exceptions/{item_id}/acknowledge"}, + {"action": "resolve", "href": f"/v1/ops/exceptions/{item_id}/resolve"}, + ] + ) + return { + "item_id": item_id, + "source": "webhook_delivery", + "severity": "medium", + "occurred_at": occurred_at, + "last_seen_at": last_seen_at, + "entity_ref": {"kind": "webhook", "id": row["webhook_id"]}, + "title": f"Webhook delivery dead for {row['webhook_id']}", + "detail": row.get("last_error") or "", + "status": status, + "actions": actions, + } + + +def _reconciliation_finding_to_item( + finding: ReconciliationFinding, state: TriageState | None +) -> dict[str, Any]: + item_id = f"rc:{finding.dedupe_key}" + status = state.status if state is not None else "open" + occurred_at = finding.occurred_at + last_seen_at = finding.occurred_at + if state is not None: + occurred_at = coerce_dt(state.first_seen_at) or finding.occurred_at + last_seen_at = coerce_dt(state.last_seen_at) or finding.occurred_at + actions = ( + [] + if status == "resolved" + else [ + {"action": "acknowledge", "href": f"/v1/ops/exceptions/{item_id}/acknowledge"}, + {"action": "resolve", "href": f"/v1/ops/exceptions/{item_id}/resolve"}, + ] + ) + return { + "item_id": item_id, + "source": "reconciliation", + "severity": finding.severity, + "occurred_at": occurred_at, + "last_seen_at": last_seen_at, + "entity_ref": {"kind": finding.entity_kind, "id": finding.entity_id}, + "title": finding.title, + "detail": finding.detail, + "status": status, + "actions": actions, + } + + +def _gather_exception_items(request: Request) -> tuple[list[dict[str, Any]], str]: + """Run R1/R2, upsert/auto-resolve the overlay, and assemble every current + item across all three sources (§4.1), unfiltered — the list and stats + endpoints both start from this same picture, so counts never drift + between them within one request.""" + store = get_control_plane_store(request.app) + engine = request.app.state.query_engine + tenant_id = _tenant_id(request) + catalog = request.app.state.catalog + order_def = catalog.entities.get("order") + stage_budgets = (getattr(order_def, "stages", None) or []) if order_def else [] + now = datetime.now(UTC) + + # Source 2: webhook dead deliveries — overlay-backed (§4.1 #2). + dead_deliveries = store.list_dead_webhook_deliveries(tenant_id) + webhook_seen_ids = [f"wh:{row['webhook_id']}:{row['event_id']}" for row in dead_deliveries] + for row, item_id in zip(dead_deliveries, webhook_seen_ids, strict=True): + seen_at = coerce_dt(row.get("updated_at")) or now + store.upsert_triage_finding( + item_id=item_id, tenant_id=tenant_id, source="webhook_delivery", seen_at=seen_at + ) + store.auto_resolve_missing_triage_findings( + tenant_id=tenant_id, + source="webhook_delivery", + seen_item_ids=webhook_seen_ids, + resolved_at=now, + ) + + # Source 3: reconciliation findings — overlay-backed (§4.1 #3). + findings = [ + *check_journal_vs_store(engine, tenant_id, stage_budgets), + *check_stuck_replay(store, tenant_id, older_than_seconds=stuck_replay_threshold_seconds()), + ] + reconciliation_seen_ids = [f"rc:{finding.dedupe_key}" for finding in findings] + for finding, item_id in zip(findings, reconciliation_seen_ids, strict=True): + store.upsert_triage_finding( + item_id=item_id, + tenant_id=tenant_id, + source="reconciliation", + seen_at=finding.occurred_at, + ) + store.auto_resolve_missing_triage_findings( + tenant_id=tenant_id, + source="reconciliation", + seen_item_ids=reconciliation_seen_ids, + resolved_at=now, + ) + + overlay_states = { + state.item_id: state for state in store.list_triage_states(tenant_id=tenant_id) + } + + items: list[dict[str, Any]] = [ + _deadletter_row_to_item(row) for row in store.list_dead_letter_events_for_inbox(tenant_id) + ] + items.extend( + _webhook_delivery_row_to_item(row, overlay_states.get(item_id), now) + for row, item_id in zip(dead_deliveries, webhook_seen_ids, strict=True) + ) + items.extend( + _reconciliation_finding_to_item(finding, overlay_states.get(item_id)) + for finding, item_id in zip(findings, reconciliation_seen_ids, strict=True) + ) + return items, tenant_id + + +def _build_exceptions_list_payload( + request: Request, + source: str | None, + status: str | None, + page: int, + page_size: int, +) -> dict[str, Any]: + items, _tenant = _gather_exception_items(request) + + if source is not None: + items = [item for item in items if item["source"] == source] + if status is not None: + items = [item for item in items if item["status"] == status] + else: + # List params default: everything not `resolved` (§4.4). + items = [item for item in items if item["status"] != "resolved"] + + items.sort( + key=lambda item: ( + _SEVERITY_RANK.get(item["severity"], 3), + -item["occurred_at"].timestamp(), + ) + ) + + total = len(items) + start = (page - 1) * page_size + page_items = items[start : start + page_size] + + return { + "items": page_items, + "pagination": { + "page": page, + "page_size": page_size, + "total": total, + "pages": math.ceil(total / page_size) if total else 0, + }, + } + + +def _build_exceptions_stats_payload(request: Request) -> dict[str, Any]: + items, tenant_id = _gather_exception_items(request) + store = get_control_plane_store(request.app) + now = datetime.now(UTC) + + by_source: dict[str, dict[str, int]] = {} + for item in items: + source_counts = by_source.setdefault(item["source"], {}) + source_counts[item["status"]] = source_counts.get(item["status"], 0) + 1 + + last_24h = sum(1 for item in items if item["occurred_at"] >= now - timedelta(hours=24)) + manual_resolutions = store.count_dead_letter_manual_actions( + tenant_id + ) + store.count_triage_manual_actions(tenant_id) + + return { + "by_source": by_source, + "last_24h": last_24h, + "manual_resolutions": manual_resolutions, + } + + +def _require_exceptions_write_access(request: Request) -> None: + tenant_key = getattr(request.state, "tenant_key", None) + if tenant_key is None: + raise HTTPException(status_code=401, detail="Invalid or missing API key.") + if tenant_key.allowed_entity_types is not None: + raise HTTPException( + status_code=403, + detail="This API key has read-only access to exception-inbox operations.", + ) + + +def _set_exception_state( + request: Request, item_id: str, status: str, payload: TriageActionRequest | None +) -> TriageActionResponse: + _require_exceptions_write_access(request) + if item_id.startswith("dl:"): + raise HTTPException( + status_code=409, + detail=( + "Dead-letter items are native — use /v1/deadletter/{event_id}/replay " + "or /dismiss instead of the exception-inbox overlay." + ), + ) + store = get_control_plane_store(request.app) + tenant_id = _tenant_id(request) + note = payload.note if payload is not None else None + updated = store.set_triage_state(item_id=item_id, tenant_id=tenant_id, status=status, note=note) + if not updated: + raise HTTPException(status_code=404, detail=f"Exception item '{item_id}' not found.") + return TriageActionResponse(item_id=item_id, status=status) + + +@router.get("/exceptions", response_model=ExceptionsListResponse) +async def list_exceptions( + request: Request, + source: Literal["deadletter", "webhook_delivery", "reconciliation"] | None = Query( + default=None + ), + status: Literal["open", "in_progress", "acknowledged", "resolved"] | None = Query(default=None), + page: int = Query(default=1, ge=1), + page_size: int = Query(default=50, ge=1, le=100), +) -> ExceptionsListResponse: + payload = await run_in_threadpool( + _build_exceptions_list_payload, request, source, status, page, page_size + ) + return ExceptionsListResponse.model_validate(payload) + + +@router.get("/exceptions/stats", response_model=ExceptionsStatsResponse) +async def exceptions_stats(request: Request) -> ExceptionsStatsResponse: + payload = await run_in_threadpool(_build_exceptions_stats_payload, request) + return ExceptionsStatsResponse.model_validate(payload) + + +@router.post("/exceptions/{item_id}/acknowledge", response_model=TriageActionResponse) +async def acknowledge_exception( + item_id: str, request: Request, payload: TriageActionRequest | None = None +) -> TriageActionResponse: + return await run_in_threadpool(_set_exception_state, request, item_id, "acknowledged", payload) + + +@router.post("/exceptions/{item_id}/resolve", response_model=TriageActionResponse) +async def resolve_exception( + item_id: str, request: Request, payload: TriageActionRequest | None = None +) -> TriageActionResponse: + return await run_in_threadpool(_set_exception_state, request, item_id, "resolved", payload) diff --git a/src/serving/backends/duckdb_backend.py b/src/serving/backends/duckdb_backend.py index f6e1052..182442c 100644 --- a/src/serving/backends/duckdb_backend.py +++ b/src/serving/backends/duckdb_backend.py @@ -10,6 +10,7 @@ from sqlglot import exp from src.serving.backends import BackendExecutionError, BackendMissingTableError, ServingBackend +from src.serving.control_plane import ensure_dead_letter_table from src.serving.db_pool import DuckDBPool from src.serving.duckdb_connection import connect_duckdb @@ -291,6 +292,29 @@ def initialize_demo_data(self) -> None: 'order.served', 4, NOW() - INTERVAL '1 minute') """) + # Dead-letter store counterparts for the two seeded journal rows above + # (ops-surfaces-spec.md §4.6, D4): the exception inbox's native + # source aggregates `dead_letter_events` (control-plane state, always + # on this connection regardless of SERVING_BACKEND — see + # QueryEngine.__init__), not `pipeline_events` — without these rows + # the demo inbox would be empty even though the journal already + # carries two 'events.deadletter' entries (I7). + ensure_dead_letter_table(self._conn) + self._conn.execute(""" + INSERT INTO dead_letter_events + (event_id, tenant_id, event_type, payload, failure_reason, + failure_detail, received_at, retry_count, last_retried_at, status) + VALUES + ('evt-004', 'default', 'order.created', + '{"order_id": "ORD-DRAFT-004", "total_amount": 0}', 'schema_validation', + 'total_amount is below the minimum order threshold', + NOW() - INTERVAL '7 minutes', 0, NULL, 'failed'), + ('evt-009', 'default', 'order.updated', + '{"order_id": "ORD-20260404-1002", "status": "shipped"}', 'duplicate_event', + 'event_id already processed at an earlier journal offset', + NOW() - INTERVAL '2 minutes', 0, NULL, 'failed') + """) + # Stage-entry trails (ops-surfaces-spec.md §1.6): topic='orders.status', # one row per stage transition, back-dated between each order's # created_at and now. ORD-20260404-1004's single pending entry at diff --git a/src/serving/control_plane/__init__.py b/src/serving/control_plane/__init__.py index 3510267..68ee7e6 100644 --- a/src/serving/control_plane/__init__.py +++ b/src/serving/control_plane/__init__.py @@ -16,6 +16,7 @@ ensure_api_usage_table, ensure_dead_letter_table, ensure_outbox_table, + ensure_triage_table, ensure_webhook_deliveries_table, ensure_webhook_delivery_queue_table, ) @@ -24,9 +25,11 @@ CONTROL_PLANE_STORE_ENV, ControlPlaneStore, OutboxEntry, + TriageState, WebhookQueueRow, control_plane_store_kind, get_control_plane_store, + stuck_replay_threshold_seconds, ) __all__ = [ @@ -35,6 +38,7 @@ "ControlPlaneStore", "EmbeddedControlPlaneStore", "OutboxEntry", + "TriageState", "WebhookQueueRow", "control_plane_store_kind", "ensure_alert_history_table", @@ -42,7 +46,9 @@ "ensure_api_usage_table", "ensure_dead_letter_table", "ensure_outbox_table", + "ensure_triage_table", "ensure_webhook_deliveries_table", "ensure_webhook_delivery_queue_table", "get_control_plane_store", + "stuck_replay_threshold_seconds", ] diff --git a/src/serving/control_plane/embedded.py b/src/serving/control_plane/embedded.py index 0c1ecbc..052dc14 100644 --- a/src/serving/control_plane/embedded.py +++ b/src/serving/control_plane/embedded.py @@ -28,7 +28,7 @@ from src.db_concurrency import catalog_ddl_lock from src.serving.duckdb_connection import connect_duckdb -from .store import ControlPlaneStore, OutboxEntry, WebhookQueueRow +from .store import AUTO_RESOLVE_NOTE, ControlPlaneStore, OutboxEntry, TriageState, WebhookQueueRow logger = structlog.get_logger() @@ -180,6 +180,28 @@ def ensure_dead_letter_table(conn: duckdb.DuckDBPyConnection) -> None: ) +def ensure_triage_table(conn: duckdb.DuckDBPyConnection) -> None: + """``ops_exception_triage`` (ops-surfaces-spec.md §4.2) — control-plane + state class 7, extending ADR 0010's inventory. Overlay for + ``webhook_delivery``/``reconciliation`` findings only; dead-letter items + get no overlay row (I6).""" + with catalog_ddl_lock: + conn.execute( + """ + CREATE TABLE IF NOT EXISTS ops_exception_triage ( + item_id TEXT PRIMARY KEY, + tenant_id TEXT, + source TEXT NOT NULL, + status TEXT NOT NULL DEFAULT 'open', + first_seen_at TIMESTAMP, + last_seen_at TIMESTAMP, + resolved_at TIMESTAMP, + note TEXT + ) + """ + ) + + def ensure_api_usage_table(conn: duckdb.DuckDBPyConnection) -> None: # Moved verbatim from auth/usage_table.py in ADR 0010 slice 4. Runs on a # dedicated per-call connection (never the shared query_engine conn — see @@ -1036,6 +1058,277 @@ def get_dead_letter_stats(self, tenant_id: str) -> dict: ], } + def list_dead_letter_events_for_inbox(self, tenant_id: str) -> list[dict]: + cursor = self._conn.cursor() + try: + ensure_dead_letter_table(cursor) + rows = cursor.execute( + """ + SELECT + event_id, + event_type, + failure_reason, + failure_detail, + received_at, + retry_count, + last_retried_at, + status + FROM dead_letter_events + WHERE COALESCE(tenant_id, 'default') = ? + ORDER BY received_at DESC + """, + [tenant_id], + ).fetchall() + finally: + cursor.close() + return [ + { + "event_id": row[0], + "event_type": row[1], + "failure_reason": row[2], + "failure_detail": row[3], + "received_at": row[4], + "retry_count": int(row[5] or 0), + "last_retried_at": row[6], + "status": row[7], + } + for row in rows + ] + + def list_stuck_replay_dead_letter_events( + self, tenant_id: str, *, older_than_seconds: float + ) -> list[dict]: + cursor = self._conn.cursor() + try: + ensure_dead_letter_table(cursor) + cutoff = datetime.now(UTC) - timedelta(seconds=older_than_seconds) + rows = cursor.execute( + """ + SELECT + event_id, + event_type, + failure_reason, + failure_detail, + received_at, + retry_count, + last_retried_at, + status + FROM dead_letter_events + WHERE COALESCE(tenant_id, 'default') = ? + AND status = 'replay_pending' + AND last_retried_at IS NOT NULL + AND last_retried_at < ? + ORDER BY last_retried_at ASC + """, + [tenant_id, cutoff], + ).fetchall() + finally: + cursor.close() + return [ + { + "event_id": row[0], + "event_type": row[1], + "failure_reason": row[2], + "failure_detail": row[3], + "received_at": row[4], + "retry_count": int(row[5] or 0), + "last_retried_at": row[6], + "status": row[7], + } + for row in rows + ] + + def count_dead_letter_manual_actions(self, tenant_id: str) -> int: + cursor = self._conn.cursor() + try: + ensure_dead_letter_table(cursor) + row = cursor.execute( + """ + SELECT COUNT(*) + FROM dead_letter_events + WHERE COALESCE(tenant_id, 'default') = ? + AND status IN ('replayed', 'dismissed') + """, + [tenant_id], + ).fetchone() + finally: + cursor.close() + return int(row[0]) if row and row[0] is not None else 0 + + # --- exception-inbox triage overlay --------------------------------------- + + def ensure_triage_schema(self) -> None: + ensure_triage_table(self._conn) + + def list_triage_states(self, *, tenant_id: str, source: str | None = None) -> list[TriageState]: + cursor = self._conn.cursor() + try: + ensure_triage_table(cursor) + select = ( + "SELECT item_id, tenant_id, source, status, first_seen_at, " + "last_seen_at, resolved_at, note FROM ops_exception_triage " + "WHERE tenant_id = ?" + ) + if source is not None: + rows = cursor.execute(select + " AND source = ?", [tenant_id, source]).fetchall() + else: + rows = cursor.execute(select, [tenant_id]).fetchall() + finally: + cursor.close() + return [ + TriageState( + item_id=row[0], + tenant_id=row[1], + source=row[2], + status=row[3], + first_seen_at=row[4], + last_seen_at=row[5], + resolved_at=row[6], + note=row[7], + ) + for row in rows + ] + + def upsert_triage_finding( + self, *, item_id: str, tenant_id: str, source: str, seen_at: datetime + ) -> None: + conn = self._conn + ensure_triage_table(conn) + existing = conn.execute( + "SELECT status FROM ops_exception_triage WHERE item_id = ?", + [item_id], + ).fetchone() + if existing is None: + conn.execute( + """ + INSERT INTO ops_exception_triage + (item_id, tenant_id, source, status, first_seen_at, last_seen_at, + resolved_at, note) + VALUES (?, ?, ?, 'open', ?, ?, NULL, NULL) + """, + [item_id, tenant_id, source, seen_at, seen_at], + ) + return + (status,) = existing + if status != "resolved": + conn.execute( + "UPDATE ops_exception_triage SET last_seen_at = ? WHERE item_id = ?", + [seen_at, item_id], + ) + return + # Resolved: reopen only if this occurrence is strictly after + # resolved_at — the comparison runs in SQL (not Python) so DuckDB's + # own aware-to-local-naive coercion applies identically to both + # sides, whether the caller passed an aware or naive `seen_at`. + conn.execute( + """ + UPDATE ops_exception_triage + SET status = 'open', last_seen_at = ?, resolved_at = NULL, note = NULL + WHERE item_id = ? AND resolved_at IS NOT NULL AND CAST(? AS TIMESTAMP) > resolved_at + """, + [seen_at, item_id, seen_at], + ) + + def auto_resolve_missing_triage_findings( + self, + *, + tenant_id: str, + source: str, + seen_item_ids: Sequence[str], + resolved_at: datetime, + ) -> None: + conn = self._conn + ensure_triage_table(conn) + seen = set(seen_item_ids) + rows = conn.execute( + """ + SELECT item_id FROM ops_exception_triage + WHERE tenant_id = ? AND source = ? AND status != 'resolved' + """, + [tenant_id, source], + ).fetchall() + for (item_id,) in rows: + if item_id in seen: + continue + conn.execute( + """ + UPDATE ops_exception_triage + SET status = 'resolved', resolved_at = ?, note = ? + WHERE item_id = ? AND tenant_id = ? + """, + [resolved_at, AUTO_RESOLVE_NOTE, item_id, tenant_id], + ) + + def set_triage_state( + self, *, item_id: str, tenant_id: str, status: str, note: str | None = None + ) -> bool: + conn = self._conn + ensure_triage_table(conn) + existing = conn.execute( + "SELECT 1 FROM ops_exception_triage WHERE item_id = ? AND tenant_id = ?", + [item_id, tenant_id], + ).fetchone() + if existing is None: + return False + resolved_at = datetime.now(UTC) if status == "resolved" else None + conn.execute( + """ + UPDATE ops_exception_triage + SET status = ?, resolved_at = ?, note = COALESCE(?, note) + WHERE item_id = ? AND tenant_id = ? + """, + [status, resolved_at, note, item_id, tenant_id], + ) + return True + + def count_triage_manual_actions(self, tenant_id: str) -> int: + # Excludes rows auto-resolved by `auto_resolve_missing_triage_findings` + # (note == AUTO_RESOLVE_NOTE) — the KPI counts human decisions only. + conn = self._conn + ensure_triage_table(conn) + row = conn.execute( + """ + SELECT COUNT(*) FROM ops_exception_triage + WHERE tenant_id = ? + AND (status = 'acknowledged' + OR (status = 'resolved' AND (note IS NULL OR note != ?))) + """, + [tenant_id, AUTO_RESOLVE_NOTE], + ).fetchone() + return int(row[0]) if row and row[0] is not None else 0 + + # --- webhook dead deliveries for the exception inbox ---------------------- + + def list_dead_webhook_deliveries(self, tenant_id: str | None = None) -> list[dict]: + conn = self._conn + ensure_webhook_delivery_queue_table(conn) + select = ( + "SELECT webhook_id, event_id, tenant, event_type, body, attempts, " + "last_status_code, last_error, created_at, updated_at " + "FROM webhook_delivery_queue WHERE status = 'dead'" + ) + if tenant_id is not None: + rows = conn.execute( + select + " AND tenant = ? ORDER BY updated_at DESC", [tenant_id] + ).fetchall() + else: + rows = conn.execute(select + " ORDER BY updated_at DESC").fetchall() + return [ + { + "webhook_id": row[0], + "event_id": row[1], + "tenant": row[2], + "event_type": row[3], + "body": row[4], + "attempts": row[5], + "last_status_code": row[6], + "last_error": row[7], + "created_at": row[8], + "updated_at": row[9], + } + for row in rows + ] + # --- API usage accounting ------------------------------------------------- def ensure_usage_schema(self) -> None: diff --git a/src/serving/control_plane/postgres.py b/src/serving/control_plane/postgres.py index b3f6326..0d1da6f 100644 --- a/src/serving/control_plane/postgres.py +++ b/src/serving/control_plane/postgres.py @@ -55,9 +55,11 @@ import structlog from .store import ( + AUTO_RESOLVE_NOTE, CONTROL_PLANE_PG_DSN_ENV, ControlPlaneStore, OutboxEntry, + TriageState, WebhookQueueRow, ) @@ -192,6 +194,18 @@ ) """, """ + CREATE TABLE IF NOT EXISTS ops_exception_triage ( + item_id TEXT PRIMARY KEY, + tenant_id TEXT, + source TEXT NOT NULL, + status TEXT NOT NULL DEFAULT 'open', + first_seen_at TIMESTAMPTZ, + last_seen_at TIMESTAMPTZ, + resolved_at TIMESTAMPTZ, + note TEXT + ) + """, + """ CREATE TABLE IF NOT EXISTS api_usage ( tenant TEXT, key_name TEXT, @@ -963,6 +977,241 @@ def get_dead_letter_stats(self, tenant_id: str) -> dict: ], } + def list_dead_letter_events_for_inbox(self, tenant_id: str) -> list[dict]: + with self._connect() as conn: + rows = conn.execute( + """ + SELECT event_id, event_type, failure_reason, failure_detail, + received_at, retry_count, last_retried_at, status + FROM dead_letter_events + WHERE COALESCE(tenant_id, 'default') = %s + ORDER BY received_at DESC + """, + (tenant_id,), + ).fetchall() + return [ + { + "event_id": row[0], + "event_type": row[1], + "failure_reason": row[2], + "failure_detail": row[3], + "received_at": row[4], + "retry_count": int(row[5] or 0), + "last_retried_at": row[6], + "status": row[7], + } + for row in rows + ] + + def list_stuck_replay_dead_letter_events( + self, tenant_id: str, *, older_than_seconds: float + ) -> list[dict]: + cutoff = datetime.now(UTC) - timedelta(seconds=older_than_seconds) + with self._connect() as conn: + rows = conn.execute( + """ + SELECT event_id, event_type, failure_reason, failure_detail, + received_at, retry_count, last_retried_at, status + FROM dead_letter_events + WHERE COALESCE(tenant_id, 'default') = %s + AND status = 'replay_pending' + AND last_retried_at IS NOT NULL + AND last_retried_at < %s + ORDER BY last_retried_at ASC + """, + (tenant_id, cutoff), + ).fetchall() + return [ + { + "event_id": row[0], + "event_type": row[1], + "failure_reason": row[2], + "failure_detail": row[3], + "received_at": row[4], + "retry_count": int(row[5] or 0), + "last_retried_at": row[6], + "status": row[7], + } + for row in rows + ] + + def count_dead_letter_manual_actions(self, tenant_id: str) -> int: + with self._connect() as conn: + row = conn.execute( + """ + SELECT COUNT(*) FROM dead_letter_events + WHERE COALESCE(tenant_id, 'default') = %s + AND status IN ('replayed', 'dismissed') + """, + (tenant_id,), + ).fetchone() + return int(row[0]) if row and row[0] is not None else 0 + + # --- exception-inbox triage overlay --------------------------------------- + + def ensure_triage_schema(self) -> None: + self._ensure_schema() + + def list_triage_states(self, *, tenant_id: str, source: str | None = None) -> list[TriageState]: + select = ( + "SELECT item_id, tenant_id, source, status, first_seen_at, " + "last_seen_at, resolved_at, note FROM ops_exception_triage " + "WHERE tenant_id = %s" + ) + with self._connect() as conn: + if source is not None: + rows = conn.execute(select + " AND source = %s", (tenant_id, source)).fetchall() + else: + rows = conn.execute(select, (tenant_id,)).fetchall() + return [ + TriageState( + item_id=row[0], + tenant_id=row[1], + source=row[2], + status=row[3], + first_seen_at=row[4], + last_seen_at=row[5], + resolved_at=row[6], + note=row[7], + ) + for row in rows + ] + + def upsert_triage_finding( + self, *, item_id: str, tenant_id: str, source: str, seen_at: datetime + ) -> None: + with self._connect() as conn: + existing = conn.execute( + "SELECT status FROM ops_exception_triage WHERE item_id = %s", + (item_id,), + ).fetchone() + if existing is None: + conn.execute( + """ + INSERT INTO ops_exception_triage + (item_id, tenant_id, source, status, first_seen_at, last_seen_at, + resolved_at, note) + VALUES (%s, %s, %s, 'open', %s, %s, NULL, NULL) + """, + (item_id, tenant_id, source, seen_at, seen_at), + ) + return + (status,) = existing + if status != "resolved": + conn.execute( + "UPDATE ops_exception_triage SET last_seen_at = %s WHERE item_id = %s", + (seen_at, item_id), + ) + return + # Resolved: reopen only if this occurrence is strictly after + # resolved_at — compared in SQL, same reasoning as the embedded + # adapter (keeps both adapters' comparison semantics identical + # regardless of whether the caller's `seen_at` is naive or aware). + conn.execute( + """ + UPDATE ops_exception_triage + SET status = 'open', last_seen_at = %s, resolved_at = NULL, note = NULL + WHERE item_id = %s AND resolved_at IS NOT NULL AND %s > resolved_at + """, + (seen_at, item_id, seen_at), + ) + + def auto_resolve_missing_triage_findings( + self, + *, + tenant_id: str, + source: str, + seen_item_ids: Sequence[str], + resolved_at: datetime, + ) -> None: + seen = set(seen_item_ids) + with self._connect() as conn: + rows = conn.execute( + """ + SELECT item_id FROM ops_exception_triage + WHERE tenant_id = %s AND source = %s AND status != 'resolved' + """, + (tenant_id, source), + ).fetchall() + for (item_id,) in rows: + if item_id in seen: + continue + conn.execute( + """ + UPDATE ops_exception_triage + SET status = 'resolved', resolved_at = %s, note = %s + WHERE item_id = %s AND tenant_id = %s + """, + (resolved_at, AUTO_RESOLVE_NOTE, item_id, tenant_id), + ) + + def set_triage_state( + self, *, item_id: str, tenant_id: str, status: str, note: str | None = None + ) -> bool: + with self._connect() as conn: + existing = conn.execute( + "SELECT 1 FROM ops_exception_triage WHERE item_id = %s AND tenant_id = %s", + (item_id, tenant_id), + ).fetchone() + if existing is None: + return False + resolved_at = datetime.now(UTC) if status == "resolved" else None + conn.execute( + """ + UPDATE ops_exception_triage + SET status = %s, resolved_at = %s, note = COALESCE(%s, note) + WHERE item_id = %s AND tenant_id = %s + """, + (status, resolved_at, note, item_id, tenant_id), + ) + return True + + def count_triage_manual_actions(self, tenant_id: str) -> int: + # Excludes rows auto-resolved by `auto_resolve_missing_triage_findings` + # (note == AUTO_RESOLVE_NOTE) — the KPI counts human decisions only. + with self._connect() as conn: + row = conn.execute( + """ + SELECT COUNT(*) FROM ops_exception_triage + WHERE tenant_id = %s + AND (status = 'acknowledged' + OR (status = 'resolved' AND (note IS NULL OR note != %s))) + """, + (tenant_id, AUTO_RESOLVE_NOTE), + ).fetchone() + return int(row[0]) if row and row[0] is not None else 0 + + # --- webhook dead deliveries for the exception inbox ---------------------- + + def list_dead_webhook_deliveries(self, tenant_id: str | None = None) -> list[dict]: + select = ( + "SELECT webhook_id, event_id, tenant, event_type, body, attempts, " + "last_status_code, last_error, created_at, updated_at " + "FROM webhook_delivery_queue WHERE status = 'dead'" + ) + with self._connect() as conn: + if tenant_id is not None: + rows = conn.execute( + select + " AND tenant = %s ORDER BY updated_at DESC", (tenant_id,) + ).fetchall() + else: + rows = conn.execute(select + " ORDER BY updated_at DESC").fetchall() + return [ + { + "webhook_id": row[0], + "event_id": row[1], + "tenant": row[2], + "event_type": row[3], + "body": row[4], + "attempts": row[5], + "last_status_code": row[6], + "last_error": row[7], + "created_at": row[8], + "updated_at": row[9], + } + for row in rows + ] + # --- API usage accounting ------------------------------------------------- def ensure_usage_schema(self) -> None: diff --git a/src/serving/control_plane/store.py b/src/serving/control_plane/store.py index f4f4566..9a5c41d 100644 --- a/src/serving/control_plane/store.py +++ b/src/serving/control_plane/store.py @@ -132,6 +132,49 @@ class OutboxEntry: retry_count: int +@dataclass(frozen=True) +class TriageState: + """One ``ops_exception_triage`` overlay row (ops-surfaces-spec.md §4.2) — + control-plane state class 7. Only the ``webhook_delivery`` and + ``reconciliation`` sources get overlay rows; dead-letter items are native + and never tracked here (invariant I6).""" + + item_id: str + tenant_id: str + source: str + status: str + first_seen_at: datetime + last_seen_at: datetime + resolved_at: datetime | None + note: str | None + + +AUTO_RESOLVE_NOTE = "auto-resolved: no longer reproduces" +"""Sentinel ``note`` an auto-resolved overlay row carries (§4.2/§4.3) — +distinguishes a system-cleared finding from an operator's ``resolve`` call so +``count_triage_manual_actions`` counts only genuine human triage decisions +(the ``manual_resolutions`` KPI, §4.5).""" + + +def stuck_replay_threshold_seconds() -> float: + """Staleness threshold for R2 ``stuck_replay`` (ops-surfaces-spec.md + §4.3): "the control-plane lease interval; env-tunable". Reuses + ``AGENTFLOW_CONTROLPLANE_LEASE_SECONDS`` — the same env var + ``postgres.py``'s ``DEFAULT_CLAIM_LEASE_SECONDS`` reads for its claim + leases — rather than a second knob; the embedded profile has no lease + concept of its own but a replay can still get stuck locally (an inline + ``EventReplayer.replay()`` that dies before its outbox entry resolves).""" + lease_env = (os.getenv("AGENTFLOW_CONTROLPLANE_LEASE_SECONDS") or "").strip() + if not lease_env: + return 300.0 + try: + return float(lease_env) + except ValueError: + raise ValueError( + f"AGENTFLOW_CONTROLPLANE_LEASE_SECONDS must be a number of seconds, got {lease_env!r}." + ) from None + + class ControlPlaneStore(ABC): """Port for control-plane state (ADR 0010). See the module docstring for the claim-semantics contract adapters must uphold.""" @@ -378,6 +421,92 @@ def get_dead_letter_stats(self, tenant_id: str) -> dict: """``{"counts": {reason: count}, "last_24h": int, "trend": [...]}`` for one tenant's active (``failed``) dead-letter events.""" + @abstractmethod + def list_dead_letter_events_for_inbox(self, tenant_id: str) -> list[dict]: + """Every dead-letter row for one tenant, any status, newest first — + the exception inbox's native source (§4.1 #1). Unlike + ``list_dead_letter_events`` (the public ``/v1/deadletter`` route: + ``status='failed'`` only, paginated), the inbox aggregates and + paginates across three heterogeneous sources itself.""" + + @abstractmethod + def list_stuck_replay_dead_letter_events( + self, tenant_id: str, *, older_than_seconds: float + ) -> list[dict]: + """``replay_pending`` dead-letter rows whose ``last_retried_at`` is + older than ``older_than_seconds`` — R2 ``stuck_replay`` (§4.3): a + replay was requested but its outbox entry never completed the + invariant-8 flip.""" + + @abstractmethod + def count_dead_letter_manual_actions(self, tenant_id: str) -> int: + """Count of ``replayed``/``dismissed`` dead-letter rows for one + tenant — the native half of the ``manual_resolutions`` KPI (§4.5). + Dismissal carries no per-action timestamp in this schema, so this is + a cumulative count, not time-windowed — an honest limitation, not + faked precision.""" + + # --- exception-inbox triage overlay (control-plane state class 7) -------- + + @abstractmethod + def ensure_triage_schema(self) -> None: + """Idempotently create the ``ops_exception_triage`` table.""" + + @abstractmethod + def list_triage_states(self, *, tenant_id: str, source: str | None = None) -> list[TriageState]: + """Every overlay row for one tenant, optionally filtered by + ``source`` (``'webhook_delivery'`` or ``'reconciliation'``).""" + + @abstractmethod + def upsert_triage_finding( + self, *, item_id: str, tenant_id: str, source: str, seen_at: datetime + ) -> None: + """Insert an ``open`` row for a first-seen finding, or refresh + ``last_seen_at`` for one already ``open``/``acknowledged``. A row an + operator resolved stays ``resolved`` unless ``seen_at`` is after its + ``resolved_at`` — the finding reproducing post-resolution reopens it + as ``open`` with a fresh ``last_seen_at`` (§4.2).""" + + @abstractmethod + def auto_resolve_missing_triage_findings( + self, + *, + tenant_id: str, + source: str, + seen_item_ids: Sequence[str], + resolved_at: datetime, + ) -> None: + """Resolve every non-``resolved`` overlay row of ``source`` + (tenant-scoped) whose ``item_id`` is absent from ``seen_item_ids`` — + the finding no longer reproduces this run (§4.2/§4.3), noted + ``'auto-resolved: no longer reproduces'``.""" + + @abstractmethod + def set_triage_state( + self, *, item_id: str, tenant_id: str, status: str, note: str | None = None + ) -> bool: + """Set one overlay row's status (``'acknowledged'`` or + ``'resolved'``), stamping ``resolved_at`` when resolving — one + transactional ``UPDATE`` (single-row, so DuckDB's autocommit and a + PostgreSQL connection's implicit transaction both satisfy this + without extra ceremony). Returns ``False`` if no row exists for + ``(item_id, tenant_id)`` — the caller 404s; this method never creates + a row (only ``upsert_triage_finding`` does, from a live detection).""" + + @abstractmethod + def count_triage_manual_actions(self, tenant_id: str) -> int: + """Count of overlay rows in ``acknowledged``/``resolved`` status for + one tenant — the overlay half of the ``manual_resolutions`` KPI + (§4.5).""" + + # --- webhook dead deliveries for the exception inbox ---------------------- + + @abstractmethod + def list_dead_webhook_deliveries(self, tenant_id: str | None = None) -> list[dict]: + """Every ``webhook_delivery_queue`` row parked ``dead``, optionally + scoped to one tenant — the exception inbox's overlay source #2 + (§4.1).""" + # --- API usage accounting (per-tenant/per-key request counters) ---------- @abstractmethod diff --git a/src/serving/semantic_layer/reconciliation.py b/src/serving/semantic_layer/reconciliation.py new file mode 100644 index 0000000..ffc21b0 --- /dev/null +++ b/src/serving/semantic_layer/reconciliation.py @@ -0,0 +1,136 @@ +"""Reconciliation checks — R1 ``journal_vs_store``, R2 ``stuck_replay`` +(ops-surfaces-spec.md §4.3), the exception inbox's third source (D4). + +Pure, read-only cross-store consistency probes: both functions read via the +QueryEngine/ControlPlaneStore ports only and never write serving state (I10) +— the caller (``routers/ops.py``) turns findings into overlay upserts and +owns the dedupe-key -> ``item_id`` mapping (``rc:``, §4.4). +""" + +from __future__ import annotations + +from dataclasses import dataclass +from datetime import UTC, datetime +from typing import TYPE_CHECKING, Any + +from src.serving.semantic_layer.stage_clock import coerce_dt, ladder_stage_names + +if TYPE_CHECKING: + from src.serving.control_plane import ControlPlaneStore + from src.serving.semantic_layer.query.engine import QueryEngine + +_STATUS_EVENT_PREFIX = "order.status." + + +@dataclass(frozen=True) +class ReconciliationFinding: + """One live detection from an R1/R2 check — not yet an overlay row.""" + + dedupe_key: str + severity: str + title: str + detail: str + entity_kind: str + entity_id: str + occurred_at: datetime + + +def check_journal_vs_store( + engine: QueryEngine, + tenant_id: str | None, + stage_budgets: list[dict[str, Any]] | None, +) -> list[ReconciliationFinding]: + """R1: for every ``entity_id`` seen in ``orders.status`` journal rows, the + serving store's order status must not be *behind* the latest journal + stage (§4.3). Scoped to the non-terminal ladder — comparing against a + terminal store/journal status is out of scope for v1 (an order already + ``delivered``/``cancelled`` is never "behind" anything); a status outside + the ladder is tolerated and skipped, never a crash (I4's spirit extended + here). + """ + ladder = ladder_stage_names(stage_budgets) + if not ladder: + return [] + rank = {name: index for index, name in enumerate(ladder)} + terminal_names = [ + entry["name"] + for entry in stage_budgets or [] + if isinstance(entry, dict) and entry.get("name") and entry.get("terminal") + ] + + stage_rows = engine.fetch_pipeline_events( + tenant_id=tenant_id, topic="orders.status", newest_first=False + ) + latest_by_entity: dict[str, tuple[str, datetime | None]] = {} + for row in stage_rows: + entity_id = row.get("entity_id") + event_type = row.get("event_type") or "" + if not entity_id or not event_type.startswith(_STATUS_EVENT_PREFIX): + continue + status = event_type[len(_STATUS_EVENT_PREFIX) :] + # Ascending iteration (`newest_first=False`): the last write per + # entity_id wins, matching the stuck-orders worklist's own scan. + latest_by_entity[str(entity_id)] = (status, coerce_dt(row.get("processed_at"))) + if not latest_by_entity: + return [] + + order_rows = engine.fetch_orders_by_status([*ladder, *terminal_names], tenant_id=tenant_id) + store_status_by_id = {str(row.get("order_id")): row.get("status") for row in order_rows} + + findings: list[ReconciliationFinding] = [] + for entity_id, (journal_status, processed_at) in latest_by_entity.items(): + if journal_status not in rank: + continue + store_status = store_status_by_id.get(entity_id) + if store_status is None or store_status not in rank: + # Missing from the store, or already at/behind a terminal status + # (never "behind" by definition) — tolerated, not flagged. + continue + if rank[store_status] < rank[journal_status]: + findings.append( + ReconciliationFinding( + dedupe_key=f"r1:{entity_id}:{journal_status}", + severity="high", + title=f"Order {entity_id} behind its journal stage", + detail=( + f"Journal's latest stage is '{journal_status}' but the serving " + f"store still shows '{store_status}' — an event landed but the " + "serving projection didn't (or forked)." + ), + entity_kind="order", + entity_id=entity_id, + occurred_at=processed_at or datetime.now(UTC), + ) + ) + return findings + + +def check_stuck_replay( + store: ControlPlaneStore, tenant_id: str, *, older_than_seconds: float +) -> list[ReconciliationFinding]: + """R2: dead-letter rows sitting in ``replay_pending`` longer than + ``older_than_seconds`` — a replay was requested but its outbox entry + never completed the invariant-8 flip (§4.3).""" + rows = store.list_stuck_replay_dead_letter_events( + tenant_id, older_than_seconds=older_than_seconds + ) + findings: list[ReconciliationFinding] = [] + for row in rows: + event_id = row["event_id"] + last_retried_at = coerce_dt(row.get("last_retried_at")) or datetime.now(UTC) + findings.append( + ReconciliationFinding( + dedupe_key=f"r2:{event_id}", + severity="medium", + title=f"Replay stuck for {event_id}", + detail=( + f"Dead-letter event '{event_id}' has been 'replay_pending' since " + f"{last_retried_at.isoformat()} — a replay was requested but its " + "outbox entry never completed." + ), + entity_kind="event", + entity_id=event_id, + occurred_at=last_retried_at, + ) + ) + return findings diff --git a/tests/integration/test_exceptions_inbox.py b/tests/integration/test_exceptions_inbox.py new file mode 100644 index 0000000..1c59198 --- /dev/null +++ b/tests/integration/test_exceptions_inbox.py @@ -0,0 +1,350 @@ +"""Integration tests for the exception inbox — +GET/POST /v1/ops/exceptions* (ops-surfaces-spec.md §4, D4). Exercises the +demo story pin (I7: the two seeded dead-letter rows make a non-empty inbox), +native dead-letter lifecycle mirroring with no overlay row (I6), stable item +ids (I5), the R1/R2 reconciliation checks, idempotent concurrent reads (I10), +and the manual_resolutions/last_24h re-pin by arithmetic (I9). Tenant scoping +(I8) and the no-third-path ratchet (I1) live in test_tenant_isolation.py and +test_control_plane_store.py respectively — same pattern as the other ops +surfaces. +""" + +from __future__ import annotations + +from pathlib import Path + +import pytest +from fastapi.testclient import TestClient + +from src.serving.api.auth import AuthManager +from src.serving.api.main import app + +pytestmark = pytest.mark.integration + +_PII_FIELD_NAMES = ("first_name", "last_name", "email", "phone", "birth_date") + + +@pytest.fixture +def client(tmp_path: Path, monkeypatch: pytest.MonkeyPatch): + db_path = tmp_path / "exceptions-inbox.duckdb" + monkeypatch.setenv("DUCKDB_PATH", str(db_path)) + monkeypatch.setenv("SERVING_BACKEND", "duckdb") + monkeypatch.setenv("AGENTFLOW_AUTH_DISABLED", "true") + + with TestClient(app) as c: + yield c + + +@pytest.fixture +def authed_client(tmp_path: Path, monkeypatch: pytest.MonkeyPatch): + db_path = tmp_path / "exceptions-inbox-authed.duckdb" + monkeypatch.setenv("DUCKDB_PATH", str(db_path)) + monkeypatch.setenv("SERVING_BACKEND", "duckdb") + + api_keys_path = tmp_path / "config" / "api_keys.yaml" + api_keys_path.parent.mkdir(parents=True, exist_ok=True) + api_keys_path.write_text( + ( + "keys:\n" + ' - key: "exceptions-readonly-key"\n' + ' name: "Readonly Agent"\n' + ' tenant: "default"\n' + " rate_limit_rpm: 100\n" + ' allowed_entity_types: ["order"]\n' + ' created_at: "2026-07-04"\n' + ' - key: "exceptions-ops-key"\n' + ' name: "Ops Agent"\n' + ' tenant: "default"\n' + " rate_limit_rpm: 100\n" + " allowed_entity_types: null\n" + ' created_at: "2026-07-04"\n' + ), + encoding="utf-8", + ) + + with TestClient(app) as c: + manager = AuthManager( + api_keys_path=api_keys_path, + db_path=tmp_path / "usage.duckdb", + admin_key="admin-secret", + ) + manager.load() + manager.ensure_usage_table() + c.app.state.auth_manager = manager + yield c + + +@pytest.fixture +def auth_headers(): + return { + "readonly": {"X-API-Key": "exceptions-readonly-key"}, + "ops": {"X-API-Key": "exceptions-ops-key"}, + } + + +# --- demo story (I7) / stats re-pin (I9) ------------------------------------ + + +def test_default_view_is_non_empty_with_the_seeded_deadletter_items(client: TestClient): + response = client.get("/v1/ops/exceptions") + + assert response.status_code == 200 + data = response.json() + item_ids = {item["item_id"] for item in data["items"]} + + assert item_ids == {"dl:evt-004", "dl:evt-009"} + for item in data["items"]: + assert item["source"] == "deadletter" + assert item["severity"] == "high" + assert item["status"] == "open" + assert {action["action"] for action in item["actions"]} == {"replay", "dismiss"} + + +def test_stats_reflect_the_demo_seed(client: TestClient): + response = client.get("/v1/ops/exceptions/stats") + + assert response.status_code == 200 + data = response.json() + + assert data["by_source"] == {"deadletter": {"open": 2}} + assert data["last_24h"] == 2 + assert data["manual_resolutions"] == 0 + + +def test_no_pii_field_names_leak_into_the_inbox(client: TestClient): + response = client.get("/v1/ops/exceptions") + + assert response.status_code == 200 + for field_name in _PII_FIELD_NAMES: + assert field_name not in response.text + + +def test_source_filter_narrows_to_deadletter(client: TestClient): + response = client.get("/v1/ops/exceptions", params={"source": "webhook_delivery"}) + + assert response.status_code == 200 + assert response.json()["items"] == [] + + +def test_pagination_shape(client: TestClient): + response = client.get("/v1/ops/exceptions", params={"page": 1, "page_size": 1}) + + assert response.status_code == 200 + data = response.json() + assert len(data["items"]) == 1 + assert data["pagination"] == {"page": 1, "page_size": 1, "total": 2, "pages": 2} + + +# --- I5: stable item ids ----------------------------------------------------- + + +def test_item_ids_are_stable_across_calls(client: TestClient): + first = client.get("/v1/ops/exceptions").json() + second = client.get("/v1/ops/exceptions").json() + + first_ids = sorted(item["item_id"] for item in first["items"]) + second_ids = sorted(item["item_id"] for item in second["items"]) + assert first_ids == second_ids == ["dl:evt-004", "dl:evt-009"] + + +# --- reconciliation: R1 journal_vs_store / R2 stuck_replay ------------------- + + +def test_r1_detects_a_serving_projection_behind_its_journal_stage(client: TestClient): + conn = client.app.state.query_engine._conn + conn.execute( + """ + INSERT INTO pipeline_events + (event_id, topic, tenant_id, entity_id, event_type, latency_ms, processed_at) + VALUES ('evt-r1-test', 'orders.status', 'default', 'ORD-20260404-1003', + 'order.status.shipped', NULL, NOW()) + """ + ) + + response = client.get("/v1/ops/exceptions", params={"source": "reconciliation"}) + + assert response.status_code == 200 + items = response.json()["items"] + assert [item["item_id"] for item in items] == ["rc:r1:ORD-20260404-1003:shipped"] + item = items[0] + assert item["severity"] == "high" + assert item["entity_ref"] == {"kind": "order", "id": "ORD-20260404-1003"} + assert item["status"] == "open" + + +def test_r1_finding_disappears_once_the_serving_projection_catches_up(client: TestClient): + conn = client.app.state.query_engine._conn + conn.execute( + """ + INSERT INTO pipeline_events + (event_id, topic, tenant_id, entity_id, event_type, latency_ms, processed_at) + VALUES ('evt-r1-test', 'orders.status', 'default', 'ORD-20260404-1003', + 'order.status.shipped', NULL, NOW()) + """ + ) + assert ( + len(client.get("/v1/ops/exceptions", params={"source": "reconciliation"}).json()["items"]) + == 1 + ) + + conn.execute("UPDATE orders_v2 SET status = 'shipped' WHERE order_id = 'ORD-20260404-1003'") + + response = client.get("/v1/ops/exceptions", params={"source": "reconciliation"}) + assert response.json()["items"] == [] + + +def test_r2_detects_a_stuck_replay(client: TestClient): + conn = client.app.state.query_engine._conn + conn.execute( + """ + INSERT INTO dead_letter_events + (event_id, tenant_id, event_type, payload, failure_reason, failure_detail, + received_at, retry_count, last_retried_at, status) + VALUES ('evt-stuck-replay', 'default', 'order.created', '{}', 'kafka_error', 'x', + NOW() - INTERVAL '20 minutes', 1, NOW() - INTERVAL '15 minutes', 'replay_pending') + """ + ) + + response = client.get("/v1/ops/exceptions", params={"source": "reconciliation"}) + + assert response.status_code == 200 + items = response.json()["items"] + assert [item["item_id"] for item in items] == ["rc:r2:evt-stuck-replay"] + assert items[0]["severity"] == "medium" + assert items[0]["entity_ref"] == {"kind": "event", "id": "evt-stuck-replay"} + + +def test_reconciliation_reads_are_idempotent_under_repeated_calls(client: TestClient): + # I10: running the checks on every read must never write serving state or + # duplicate overlay rows. + conn = client.app.state.query_engine._conn + conn.execute( + """ + INSERT INTO pipeline_events + (event_id, topic, tenant_id, entity_id, event_type, latency_ms, processed_at) + VALUES ('evt-r1-test', 'orders.status', 'default', 'ORD-20260404-1003', + 'order.status.shipped', NULL, NOW()) + """ + ) + + first = client.get("/v1/ops/exceptions", params={"source": "reconciliation"}).json() + second = client.get("/v1/ops/exceptions", params={"source": "reconciliation"}).json() + third = client.get("/v1/ops/exceptions", params={"source": "reconciliation"}).json() + + assert first == second == third + assert len(first["items"]) == 1 + + +# --- mutations: auth, 409/404, acknowledge/resolve, auto-resolve ------------ + + +def test_acknowledge_requires_an_api_key(authed_client: TestClient): + response = authed_client.post("/v1/ops/exceptions/wh:hook:evt/acknowledge") + assert response.status_code == 401 + + +def test_readonly_key_cannot_mutate(authed_client: TestClient, auth_headers): + response = authed_client.post( + "/v1/ops/exceptions/wh:hook:evt/acknowledge", headers=auth_headers["readonly"] + ) + assert response.status_code == 403 + + +def test_deadletter_item_mutation_is_rejected_with_409(authed_client: TestClient, auth_headers): + response = authed_client.post( + "/v1/ops/exceptions/dl:evt-004/acknowledge", headers=auth_headers["ops"] + ) + assert response.status_code == 409 + + +def test_unknown_item_id_returns_404(authed_client: TestClient, auth_headers): + response = authed_client.post( + "/v1/ops/exceptions/rc:does-not-exist/acknowledge", headers=auth_headers["ops"] + ) + assert response.status_code == 404 + + +def test_acknowledge_then_resolve_lifecycle(authed_client: TestClient, auth_headers): + conn = authed_client.app.state.query_engine._conn + conn.execute( + """ + INSERT INTO pipeline_events + (event_id, topic, tenant_id, entity_id, event_type, latency_ms, processed_at) + VALUES ('evt-r1-test', 'orders.status', 'default', 'ORD-20260404-1003', + 'order.status.shipped', NULL, NOW()) + """ + ) + ops = auth_headers["ops"] + item_id = "rc:r1:ORD-20260404-1003:shipped" + # First GET seeds the overlay row for this finding. + authed_client.get("/v1/ops/exceptions", params={"source": "reconciliation"}, headers=ops) + + ack = authed_client.post(f"/v1/ops/exceptions/{item_id}/acknowledge", headers=ops) + assert ack.status_code == 200 + assert ack.json() == {"item_id": item_id, "status": "acknowledged"} + + still_listed = authed_client.get( + "/v1/ops/exceptions", params={"source": "reconciliation"}, headers=ops + ) + assert [item["status"] for item in still_listed.json()["items"]] == ["acknowledged"] + + resolve = authed_client.post( + f"/v1/ops/exceptions/{item_id}/resolve", headers=ops, json={"note": "fixed manually"} + ) + assert resolve.status_code == 200 + assert resolve.json() == {"item_id": item_id, "status": "resolved"} + + default_view = authed_client.get( + "/v1/ops/exceptions", params={"source": "reconciliation"}, headers=ops + ) + assert default_view.json()["items"] == [] + + resolved_view = authed_client.get( + "/v1/ops/exceptions", + params={"source": "reconciliation", "status": "resolved"}, + headers=ops, + ) + assert [item["item_id"] for item in resolved_view.json()["items"]] == [item_id] + + stats = authed_client.get("/v1/ops/exceptions/stats", headers=ops).json() + assert stats["manual_resolutions"] == 1 + + +def test_auto_resolve_when_the_finding_no_longer_reproduces( + authed_client: TestClient, auth_headers +): + # A reconciliation finding is computed on read, not a persistent row like + # a dead-letter/webhook item — once the underlying mismatch clears there + # is nothing left to render (no live title/detail/entity_ref), so the + # item disappears from the feed entirely rather than reappearing under + # `status=resolved`. The overlay row itself is still auto-resolved + # behind the scenes, verified here through `manual_resolutions` — a raw + # count over the overlay table, independent of what's currently live. + conn = authed_client.app.state.query_engine._conn + conn.execute( + """ + INSERT INTO pipeline_events + (event_id, topic, tenant_id, entity_id, event_type, latency_ms, processed_at) + VALUES ('evt-r1-test', 'orders.status', 'default', 'ORD-20260404-1003', + 'order.status.shipped', NULL, NOW()) + """ + ) + ops = auth_headers["ops"] + authed_client.get("/v1/ops/exceptions", params={"source": "reconciliation"}, headers=ops) + + conn.execute("UPDATE orders_v2 SET status = 'shipped' WHERE order_id = 'ORD-20260404-1003'") + + default_view = authed_client.get( + "/v1/ops/exceptions", params={"source": "reconciliation"}, headers=ops + ) + resolved_view = authed_client.get( + "/v1/ops/exceptions", + params={"source": "reconciliation", "status": "resolved"}, + headers=ops, + ) + assert default_view.json()["items"] == [] + assert resolved_view.json()["items"] == [] + + stats = authed_client.get("/v1/ops/exceptions/stats", headers=ops).json() + # Auto-resolved, not a human decision — excluded from the manual KPI. + assert stats["manual_resolutions"] == 0 diff --git a/tests/integration/test_tenant_isolation.py b/tests/integration/test_tenant_isolation.py index 21a2839..2ffdb3a 100644 --- a/tests/integration/test_tenant_isolation.py +++ b/tests/integration/test_tenant_isolation.py @@ -236,6 +236,39 @@ def test_cross_tenant_stuck_orders_are_scoped_to_tenant_schema(client: TestClien assert demo_shared["user_id"] == "USR-DEMO" +def test_cross_tenant_exceptions_inbox_is_scoped_to_tenant(client: TestClient): + # ops-surfaces-spec.md §1.7 / invariant I8: the exception inbox's native + # dead-letter source scopes by the request tenant, same as every other + # ops surface. `dead_letter_events` is control-plane state on the shared + # connection (a `tenant_id` column, not a per-tenant schema) — the demo + # seed's own `default`-tenant rows (evt-004/evt-009) must not leak here. + conn = client.app.state.query_engine._conn + conn.execute( + """ + INSERT INTO dead_letter_events + (event_id, tenant_id, event_type, payload, failure_reason, failure_detail, + received_at, retry_count, last_retried_at, status) + VALUES + ('evt-acme-1', 'acme', 'order.created', '{}', 'schema_validation', 'x', + NOW(), 0, NULL, 'failed'), + ('evt-demo-1', 'demo', 'order.created', '{}', 'schema_validation', 'x', + NOW(), 0, NULL, 'failed') + """ + ) + + acme_response = client.get("/v1/ops/exceptions", headers={"X-API-Key": "acme-key"}) + demo_response = client.get("/v1/ops/exceptions", headers={"X-API-Key": "demo-key"}) + + assert acme_response.status_code == 200 + assert demo_response.status_code == 200 + + acme_ids = {item["item_id"] for item in acme_response.json()["items"]} + demo_ids = {item["item_id"] for item in demo_response.json()["items"]} + + assert acme_ids == {"dl:evt-acme-1"} + assert demo_ids == {"dl:evt-demo-1"} + + def test_metric_cache_does_not_leak_across_tenants(client: TestClient): class FakeRedis: def __init__(self) -> None: diff --git a/tests/unit/test_control_plane_store.py b/tests/unit/test_control_plane_store.py index dfddadd..626334e 100644 --- a/tests/unit/test_control_plane_store.py +++ b/tests/unit/test_control_plane_store.py @@ -831,6 +831,279 @@ def test_get_store_follows_a_swapped_query_engine( second.close() +# --- exception inbox: dead-letter/webhook reads + triage overlay (D4) ------------ + + +def _seed_webhook_dead( + conn: duckdb.DuckDBPyConnection, + *, + webhook_id: str, + event_id: str, + tenant: str = "acme", + last_error: str = "connection refused", + updated_at: datetime | None = None, +) -> None: + from src.serving.control_plane.embedded import ensure_webhook_delivery_queue_table + + ensure_webhook_delivery_queue_table(conn) + now = updated_at or datetime.now(UTC) + conn.execute( + """ + INSERT INTO webhook_delivery_queue + (webhook_id, event_id, tenant, event_type, body, status, attempts, + last_error, created_at, updated_at) + VALUES (?, ?, ?, 'order.created', '{}', 'dead', 5, ?, ?, ?) + """, + [webhook_id, event_id, tenant, last_error, now, now], + ) + + +def test_list_dead_webhook_deliveries_returns_only_dead_rows_for_the_tenant( + store: EmbeddedControlPlaneStore, conn: duckdb.DuckDBPyConnection +) -> None: + _seed_webhook_dead(conn, webhook_id="wh-1", event_id="evt-a", tenant="acme") + _seed_webhook_dead(conn, webhook_id="wh-2", event_id="evt-b", tenant="beta") + _enqueue(store, "wh-3", "evt-c") # status stays 'pending' — not dead + + acme_rows = store.list_dead_webhook_deliveries("acme") + all_rows = store.list_dead_webhook_deliveries() + + assert [row["webhook_id"] for row in acme_rows] == ["wh-1"] + assert acme_rows[0]["last_error"] == "connection refused" + assert {row["webhook_id"] for row in all_rows} == {"wh-1", "wh-2"} + + +def test_list_dead_letter_events_for_inbox_returns_every_status( + outbox_store: EmbeddedControlPlaneStore, conn: duckdb.DuckDBPyConnection +) -> None: + _seed_dead_letter(conn, event_id="evt-failed", status="failed") + _seed_dead_letter(conn, event_id="evt-dismissed", status="dismissed") + _seed_dead_letter(conn, event_id="evt-other-tenant", tenant_id="beta", status="failed") + + rows = outbox_store.list_dead_letter_events_for_inbox("acme") + + assert {row["event_id"] for row in rows} == {"evt-failed", "evt-dismissed"} + + +def test_list_stuck_replay_dead_letter_events_filters_by_age_and_status( + outbox_store: EmbeddedControlPlaneStore, conn: duckdb.DuckDBPyConnection +) -> None: + _seed_dead_letter(conn, event_id="evt-stuck", status="replay_pending") + conn.execute( + "UPDATE dead_letter_events SET last_retried_at = ? WHERE event_id = ?", + [datetime.now(UTC) - timedelta(seconds=600), "evt-stuck"], + ) + _seed_dead_letter(conn, event_id="evt-fresh", status="replay_pending") + conn.execute( + "UPDATE dead_letter_events SET last_retried_at = ? WHERE event_id = ?", + [datetime.now(UTC) - timedelta(seconds=5), "evt-fresh"], + ) + _seed_dead_letter(conn, event_id="evt-failed", status="failed") + + stuck = outbox_store.list_stuck_replay_dead_letter_events("acme", older_than_seconds=300) + + assert [row["event_id"] for row in stuck] == ["evt-stuck"] + + +def test_count_dead_letter_manual_actions_counts_replayed_and_dismissed_only( + outbox_store: EmbeddedControlPlaneStore, conn: duckdb.DuckDBPyConnection +) -> None: + _seed_dead_letter(conn, event_id="evt-1", status="replayed") + _seed_dead_letter(conn, event_id="evt-2", status="dismissed") + _seed_dead_letter(conn, event_id="evt-3", status="failed") + _seed_dead_letter(conn, event_id="evt-4", tenant_id="beta", status="replayed") + + assert outbox_store.count_dead_letter_manual_actions("acme") == 2 + + +def test_upsert_triage_finding_inserts_open_row_on_first_sight( + store: EmbeddedControlPlaneStore, +) -> None: + # Naive (DuckDB's own TIMESTAMP round-trip, local wall-clock, is naive — + # an aware datetime would come back converted, breaking equality here). + seen_at = datetime.now() + store.upsert_triage_finding( + item_id="rc:r1:ORD-1:shipped", tenant_id="acme", source="reconciliation", seen_at=seen_at + ) + + states = store.list_triage_states(tenant_id="acme") + assert len(states) == 1 + assert states[0].status == "open" + assert states[0].first_seen_at == seen_at + assert states[0].last_seen_at == seen_at + + +def test_upsert_triage_finding_refreshes_last_seen_at_while_open( + store: EmbeddedControlPlaneStore, +) -> None: + first_seen = datetime.now() - timedelta(minutes=10) + later = datetime.now() + store.upsert_triage_finding( + item_id="wh:hook:evt", tenant_id="acme", source="webhook_delivery", seen_at=first_seen + ) + store.upsert_triage_finding( + item_id="wh:hook:evt", tenant_id="acme", source="webhook_delivery", seen_at=later + ) + + state = store.list_triage_states(tenant_id="acme")[0] + assert state.first_seen_at == first_seen + assert state.last_seen_at == later + assert state.status == "open" + + +def test_upsert_triage_finding_stays_resolved_for_the_same_resolved_occurrence( + store: EmbeddedControlPlaneStore, +) -> None: + seen_at = datetime.now(UTC) - timedelta(minutes=10) + store.upsert_triage_finding( + item_id="rc:r2:evt-9", tenant_id="acme", source="reconciliation", seen_at=seen_at + ) + resolved = store.set_triage_state(item_id="rc:r2:evt-9", tenant_id="acme", status="resolved") + assert resolved is True + + # Re-detecting the exact same (older) occurrence must not reopen it — an + # operator's resolve is sticky against a fact that hasn't moved. + store.upsert_triage_finding( + item_id="rc:r2:evt-9", tenant_id="acme", source="reconciliation", seen_at=seen_at + ) + + state = store.list_triage_states(tenant_id="acme")[0] + assert state.status == "resolved" + + +def test_upsert_triage_finding_reopens_when_the_finding_reproduces_after_resolution( + store: EmbeddedControlPlaneStore, +) -> None: + seen_at = datetime.now(UTC) - timedelta(minutes=10) + store.upsert_triage_finding( + item_id="rc:r2:evt-9", tenant_id="acme", source="reconciliation", seen_at=seen_at + ) + store.set_triage_state(item_id="rc:r2:evt-9", tenant_id="acme", status="resolved") + + fresh_occurrence = datetime.now() + store.upsert_triage_finding( + item_id="rc:r2:evt-9", tenant_id="acme", source="reconciliation", seen_at=fresh_occurrence + ) + + state = store.list_triage_states(tenant_id="acme")[0] + assert state.status == "open" + assert state.resolved_at is None + assert state.last_seen_at == fresh_occurrence + + +def test_auto_resolve_missing_triage_findings_resolves_only_the_absent_rows( + store: EmbeddedControlPlaneStore, +) -> None: + now = datetime.now(UTC) + store.upsert_triage_finding( + item_id="rc:r1:ORD-1:shipped", tenant_id="acme", source="reconciliation", seen_at=now + ) + store.upsert_triage_finding( + item_id="rc:r1:ORD-2:shipped", tenant_id="acme", source="reconciliation", seen_at=now + ) + store.upsert_triage_finding( + item_id="wh:hook:evt", tenant_id="acme", source="webhook_delivery", seen_at=now + ) + + resolved_at = datetime.now(UTC) + store.auto_resolve_missing_triage_findings( + tenant_id="acme", + source="reconciliation", + seen_item_ids=["rc:r1:ORD-1:shipped"], + resolved_at=resolved_at, + ) + + states = {state.item_id: state for state in store.list_triage_states(tenant_id="acme")} + assert states["rc:r1:ORD-1:shipped"].status == "open" + assert states["rc:r1:ORD-2:shipped"].status == "resolved" + assert states["rc:r1:ORD-2:shipped"].note == "auto-resolved: no longer reproduces" + # A different source's row is untouched by this call. + assert states["wh:hook:evt"].status == "open" + + +def test_set_triage_state_returns_false_for_an_unknown_item( + store: EmbeddedControlPlaneStore, +) -> None: + result = store.set_triage_state( + item_id="rc:does-not-exist", tenant_id="acme", status="resolved" + ) + assert result is False + + +def test_set_triage_state_acknowledge_then_resolve_stores_note( + store: EmbeddedControlPlaneStore, +) -> None: + store.upsert_triage_finding( + item_id="wh:hook:evt", + tenant_id="acme", + source="webhook_delivery", + seen_at=datetime.now(UTC), + ) + + acknowledged_ok = store.set_triage_state( + item_id="wh:hook:evt", tenant_id="acme", status="acknowledged" + ) + assert acknowledged_ok is True + acknowledged = store.list_triage_states(tenant_id="acme")[0] + assert acknowledged.status == "acknowledged" + assert acknowledged.resolved_at is None + + assert store.set_triage_state( + item_id="wh:hook:evt", tenant_id="acme", status="resolved", note="fixed the webhook" + ) + resolved = store.list_triage_states(tenant_id="acme")[0] + assert resolved.status == "resolved" + assert resolved.resolved_at is not None + assert resolved.note == "fixed the webhook" + + +def test_count_triage_manual_actions_counts_acknowledged_and_resolved( + store: EmbeddedControlPlaneStore, +) -> None: + now = datetime.now(UTC) + for item_id in ("rc:a", "rc:b", "rc:c"): + store.upsert_triage_finding( + item_id=item_id, tenant_id="acme", source="reconciliation", seen_at=now + ) + store.set_triage_state(item_id="rc:a", tenant_id="acme", status="acknowledged") + store.set_triage_state(item_id="rc:b", tenant_id="acme", status="resolved") + # rc:c stays open. + + assert store.count_triage_manual_actions("acme") == 2 + + +def test_count_triage_manual_actions_excludes_auto_resolved_rows( + store: EmbeddedControlPlaneStore, +) -> None: + now = datetime.now(UTC) + store.upsert_triage_finding( + item_id="rc:auto", tenant_id="acme", source="reconciliation", seen_at=now + ) + store.upsert_triage_finding( + item_id="rc:manual", tenant_id="acme", source="reconciliation", seen_at=now + ) + # rc:auto is absent from this run's live findings -> auto-resolved. + store.auto_resolve_missing_triage_findings( + tenant_id="acme", source="reconciliation", seen_item_ids=["rc:manual"], resolved_at=now + ) + store.set_triage_state(item_id="rc:manual", tenant_id="acme", status="resolved") + + assert store.count_triage_manual_actions("acme") == 1 + + +def test_list_triage_states_filters_by_source(store: EmbeddedControlPlaneStore) -> None: + now = datetime.now(UTC) + store.upsert_triage_finding( + item_id="rc:x", tenant_id="acme", source="reconciliation", seen_at=now + ) + store.upsert_triage_finding( + item_id="wh:y", tenant_id="acme", source="webhook_delivery", seen_at=now + ) + + reconciliation_states = store.list_triage_states(tenant_id="acme", source="reconciliation") + assert [state.item_id for state in reconciliation_states] == ["rc:x"] + + # --- structural ratchet ----------------------------------------------------------