diff --git a/contracts/entities/order.yaml b/contracts/entities/order.yaml index 6a7c7e22..f850a6e2 100644 --- a/contracts/entities/order.yaml +++ b/contracts/entities/order.yaml @@ -11,3 +11,20 @@ fields: created_at: Order creation timestamp relationships: user: user_id +stages: + - name: pending + sla_minutes: 30 + description: Confirmation SLA — marketplace orders auto-confirm within + minutes; a pending order older than this needs a payment/CRM decision + - name: confirmed + sla_minutes: 1440 + description: Ship-by SLA — FBS marketplace shipping deadlines drive the + 24h budget for warehouse handover + - name: shipped + sla_minutes: 7200 + description: Delivery SLA — 5 days courier/pickup before a customer + contact is due + - name: delivered + terminal: true + - name: cancelled + terminal: true diff --git a/docs/agent-tools/claude-tools.json b/docs/agent-tools/claude-tools.json index f0208e11..18972278 100644 --- a/docs/agent-tools/claude-tools.json +++ b/docs/agent-tools/claude-tools.json @@ -844,6 +844,44 @@ ] } }, + { + "name": "get_stuck_orders_v1_ops_stuck_orders_get", + "description": "Get Stuck Orders", + "input_schema": { + "type": "object", + "properties": { + "stage": { + "anyOf": [ + { + "type": "string" + }, + { + "type": "null" + } + ], + "title": "Stage" + }, + "include_within_sla": { + "type": "boolean", + "default": false, + "title": "Include Within Sla" + }, + "page": { + "type": "integer", + "minimum": 1, + "default": 1, + "title": "Page" + }, + "page_size": { + "type": "integer", + "maximum": 100, + "minimum": 1, + "default": 50, + "title": "Page Size" + } + } + } + }, { "name": "search_v1_search_get", "description": "Search", diff --git a/docs/agent-tools/openai-tools.json b/docs/agent-tools/openai-tools.json index 5ba5ac6e..3f68f2f2 100644 --- a/docs/agent-tools/openai-tools.json +++ b/docs/agent-tools/openai-tools.json @@ -940,6 +940,47 @@ } } }, + { + "type": "function", + "function": { + "name": "get_stuck_orders_v1_ops_stuck_orders_get", + "description": "Get Stuck Orders", + "parameters": { + "type": "object", + "properties": { + "stage": { + "anyOf": [ + { + "type": "string" + }, + { + "type": "null" + } + ], + "title": "Stage" + }, + "include_within_sla": { + "type": "boolean", + "default": false, + "title": "Include Within Sla" + }, + "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": { diff --git a/docs/openapi.json b/docs/openapi.json index 76bf67cd..f6a91ce4 100644 --- a/docs/openapi.json +++ b/docs/openapi.json @@ -1775,6 +1775,88 @@ } } }, + "/v1/ops/stuck-orders": { + "get": { + "tags": [ + "ops" + ], + "summary": "Get Stuck Orders", + "operationId": "get_stuck_orders_v1_ops_stuck_orders_get", + "parameters": [ + { + "name": "stage", + "in": "query", + "required": false, + "schema": { + "anyOf": [ + { + "type": "string" + }, + { + "type": "null" + } + ], + "title": "Stage" + } + }, + { + "name": "include_within_sla", + "in": "query", + "required": false, + "schema": { + "type": "boolean", + "default": false, + "title": "Include Within Sla" + } + }, + { + "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/StuckOrdersResponse" + } + } + } + }, + "422": { + "description": "Validation Error", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/HTTPValidationError" + } + } + } + } + } + } + }, "/v1/search": { "get": { "tags": [ @@ -3724,6 +3806,159 @@ ], "title": "SearchResult" }, + "StuckOrderItem": { + "properties": { + "order_id": { + "type": "string", + "title": "Order Id" + }, + "user_id": { + "anyOf": [ + { + "type": "string" + }, + { + "type": "null" + } + ], + "title": "User Id" + }, + "status": { + "type": "string", + "title": "Status" + }, + "entered_at": { + "anyOf": [ + { + "type": "string", + "format": "date-time" + }, + { + "type": "null" + } + ], + "title": "Entered At" + }, + "in_stage_seconds": { + "anyOf": [ + { + "type": "number" + }, + { + "type": "null" + } + ], + "title": "In Stage Seconds" + }, + "sla_minutes": { + "anyOf": [ + { + "type": "integer" + }, + { + "type": "null" + } + ], + "title": "Sla Minutes" + }, + "overshoot_ratio": { + "anyOf": [ + { + "type": "number" + }, + { + "type": "null" + } + ], + "title": "Overshoot Ratio" + }, + "clock": { + "type": "string", + "enum": [ + "journal", + "fallback" + ], + "title": "Clock" + }, + "total_amount": { + "anyOf": [ + { + "type": "number" + }, + { + "type": "null" + } + ], + "title": "Total Amount" + }, + "currency": { + "anyOf": [ + { + "type": "string" + }, + { + "type": "null" + } + ], + "title": "Currency" + } + }, + "type": "object", + "required": [ + "order_id", + "status", + "clock" + ], + "title": "StuckOrderItem" + }, + "StuckOrdersResponse": { + "properties": { + "items": { + "items": { + "$ref": "#/components/schemas/StuckOrderItem" + }, + "type": "array", + "title": "Items" + }, + "summary": { + "$ref": "#/components/schemas/StuckOrdersSummary" + }, + "pagination": { + "additionalProperties": { + "type": "integer" + }, + "type": "object", + "title": "Pagination" + } + }, + "type": "object", + "required": [ + "items", + "summary", + "pagination" + ], + "title": "StuckOrdersResponse" + }, + "StuckOrdersSummary": { + "properties": { + "open_by_stage": { + "additionalProperties": { + "type": "integer" + }, + "type": "object", + "title": "Open By Stage" + }, + "breached_by_stage": { + "additionalProperties": { + "type": "integer" + }, + "type": "object", + "title": "Breached By Stage" + } + }, + "type": "object", + "title": "StuckOrdersSummary" + }, "ValidationError": { "properties": { "loc": { diff --git a/src/serving/api/main.py b/src/serving/api/main.py index cbdff2cd..666802dc 100644 --- a/src/serving/api/main.py +++ b/src/serving/api/main.py @@ -41,6 +41,7 @@ from src.serving.api.routers.contracts import router as contracts_router from src.serving.api.routers.deadletter import router as deadletter_router from src.serving.api.routers.lineage import router as lineage_router +from src.serving.api.routers.ops import router as ops_router from src.serving.api.routers.search import router as search_router from src.serving.api.routers.slo import router as slo_router from src.serving.api.routers.stream import router as stream_router @@ -333,6 +334,7 @@ async def demo_mode_guard(request: Request, call_next: RequestResponseEndpoint) app.include_router(contracts_router, prefix="/v1") app.include_router(deadletter_router) app.include_router(lineage_router) +app.include_router(ops_router) app.include_router(search_router, prefix="/v1") app.include_router(slo_router) app.include_router(stream_router) diff --git a/src/serving/api/routers/agent_query.py b/src/serving/api/routers/agent_query.py index 5d2ce341..e2ebd46b 100644 --- a/src/serving/api/routers/agent_query.py +++ b/src/serving/api/routers/agent_query.py @@ -25,6 +25,7 @@ ) from src.serving.cache import ENTITY_TTL_SECONDS, QueryCache, cache_entity_key from src.serving.control_plane import get_control_plane_store +from src.serving.semantic_layer.stage_clock import coerce_dt, resolve_breach, stage_budget logger = structlog.get_logger() tracer = trace.get_tracer("agentflow.api") @@ -114,29 +115,6 @@ def _as_of_iso_text(as_of: datetime | None) -> str | None: return as_of.replace(microsecond=0).isoformat().replace("+00:00", "Z") -def _coerce_dt(value: object) -> datetime | None: - # Naive DuckDB timestamps are local wall-clock (DuckDB's NOW()), not UTC — - # same local_tz convention as EntityQueryMixin.get_entity's _last_updated. - local_tz = datetime.now().astimezone().tzinfo or UTC - if isinstance(value, datetime): - return ( - value.astimezone(UTC) - if value.tzinfo is not None - else value.replace(tzinfo=local_tz).astimezone(UTC) - ) - if isinstance(value, str): - try: - parsed = datetime.fromisoformat(value) - except ValueError: - return None - return ( - parsed.astimezone(UTC) - if parsed.tzinfo is not None - else parsed.replace(tzinfo=local_tz).astimezone(UTC) - ) - return None - - def _transform_payload_for_requested_version( req: Request, payload: dict[str, Any] ) -> dict[str, Any]: @@ -394,36 +372,20 @@ def _build_order_timeline(request: Request, order_id: str) -> dict[str, Any] | N catalog = request.app.state.catalog order_def = catalog.entities.get("order") stage_budgets = (getattr(order_def, "stages", None) or []) if order_def else [] - budget = next( - ( - entry - for entry in stage_budgets - if isinstance(entry, dict) and entry.get("name") == current_status - ), - None, - ) + budget = stage_budget(stage_budgets, current_status) entered_at = None clock = "fallback" target_event_type = f"order.status.{current_status}" for row in reversed(stage_rows): if row.get("event_type") == target_event_type: - entered_at = _coerce_dt(row.get("processed_at")) + entered_at = coerce_dt(row.get("processed_at")) clock = "journal" break if entered_at is None: - entered_at = _coerce_dt(order_row.get("created_at")) + entered_at = coerce_dt(order_row.get("created_at")) - in_stage_seconds = ( - (datetime.now(UTC) - entered_at).total_seconds() if entered_at is not None else None - ) - sla_minutes = budget.get("sla_minutes") if budget else None - is_terminal = bool(budget.get("terminal")) if budget else False - breached = ( - None - if is_terminal or sla_minutes is None or in_stage_seconds is None - else in_stage_seconds > sla_minutes * 60 - ) + in_stage_seconds, sla_minutes, breached = resolve_breach(entered_at=entered_at, budget=budget) customer = None user_id = order_row.get("user_id") diff --git a/src/serving/api/routers/ops.py b/src/serving/api/routers/ops.py new file mode 100644 index 00000000..b87814f8 --- /dev/null +++ b/src/serving/api/routers/ops.py @@ -0,0 +1,190 @@ +"""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). +""" + +from __future__ import annotations + +import math +from datetime import datetime +from typing import Any, Literal, cast + +from fastapi import APIRouter, Query, Request +from pydantic import BaseModel, Field +from starlette.concurrency import run_in_threadpool + +from src.serving.semantic_layer.stage_clock import ( + coerce_dt, + ladder_stage_names, + resolve_breach, + stage_budget, +) + +router = APIRouter(prefix="/v1/ops", tags=["ops"]) + + +class StuckOrderItem(BaseModel): + order_id: str + user_id: str | None = None + status: str + entered_at: datetime | None = None + in_stage_seconds: float | None = None + sla_minutes: int | None = None + overshoot_ratio: float | None = None + clock: Literal["journal", "fallback"] + total_amount: float | None = None + currency: str | None = None + + +class StuckOrdersSummary(BaseModel): + open_by_stage: dict[str, int] = Field(default_factory=dict) + breached_by_stage: dict[str, int] = Field(default_factory=dict) + + +class StuckOrdersResponse(BaseModel): + items: list[StuckOrderItem] + summary: StuckOrdersSummary + pagination: dict[str, int] + + +def _resolve_tenant_id(request: Request) -> str | None: + tenant_key = getattr(request.state, "tenant_key", None) + tenant_id = getattr(request.state, "tenant_id", None) or getattr(tenant_key, "tenant", None) + return cast("str | None", tenant_id) + + +def _build_stuck_order_item( + order_row: dict[str, Any], + latest_stage_row: dict[str, Any] | None, + budget: dict[str, Any] | None, +) -> dict[str, Any]: + """One order's worklist row: stage clock per §1.4, breach per §1.5.""" + entered_at = coerce_dt(latest_stage_row.get("processed_at")) if latest_stage_row else None + clock = "journal" if entered_at is not None else "fallback" + if entered_at is None: + entered_at = coerce_dt(order_row.get("created_at")) + + in_stage_seconds, sla_minutes, breached = resolve_breach(entered_at=entered_at, budget=budget) + overshoot_ratio = ( + in_stage_seconds / (sla_minutes * 60) + if in_stage_seconds is not None and sla_minutes + else None + ) + + return { + "order_id": order_row.get("order_id"), + "user_id": order_row.get("user_id"), + "status": order_row.get("status"), + "entered_at": entered_at, + "in_stage_seconds": in_stage_seconds, + "sla_minutes": sla_minutes, + "overshoot_ratio": overshoot_ratio, + "clock": clock, + "total_amount": order_row.get("total_amount"), + "currency": order_row.get("currency"), + "_breached": breached, + } + + +def _build_stuck_orders_payload( + request: Request, + stage: str | None, + include_within_sla: bool, + page: int, + page_size: int, +) -> dict[str, Any]: + """Sync composition for GET /v1/ops/stuck-orders. + + Runs on a worker thread (matches lineage.py/deadletter.py/the Order 360 + timeline). Ladder + budgets come from the catalog `stages:` block only + (I2) — no stage-name or budget literal here. + """ + engine = request.app.state.query_engine + tenant_id = _resolve_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 [] + ladder = ladder_stage_names(stage_budgets) + + order_rows = engine.fetch_orders_by_status(ladder, tenant_id=tenant_id) + stage_rows = engine.fetch_pipeline_events( + tenant_id=tenant_id, topic="orders.status", newest_first=False + ) + + # Latest journal row per (order, event_type), matching each order's + # *current* status — not merely the most recent row overall (§1.4). + # Ascending iteration means the last write for a key wins. + latest_by_key: dict[tuple[str, str], dict[str, Any]] = {} + for row in stage_rows: + entity_id = row.get("entity_id") + event_type = row.get("event_type") + if not entity_id or not event_type: + continue + latest_by_key[(str(entity_id), str(event_type))] = row + + all_items = [ + _build_stuck_order_item( + order_row, + latest_by_key.get( + (str(order_row.get("order_id")), f"order.status.{order_row.get('status')}") + ), + stage_budget(stage_budgets, order_row.get("status")), + ) + for order_row in order_rows + ] + + open_by_stage: dict[str, int] = {} + breached_by_stage: dict[str, int] = {} + for item in all_items: + status = item["status"] + open_by_stage[status] = open_by_stage.get(status, 0) + 1 + if item["_breached"]: + breached_by_stage[status] = breached_by_stage.get(status, 0) + 1 + + filtered = all_items + if stage is not None: + filtered = [item for item in filtered if item["status"] == stage] + if not include_within_sla: + filtered = [item for item in filtered if item["_breached"]] + + def _sort_key(item: dict[str, Any]) -> float: + ratio = item["overshoot_ratio"] + return ratio if ratio is not None else -1.0 + + filtered.sort(key=_sort_key, reverse=True) + + total = len(filtered) + start = (page - 1) * page_size + page_items = filtered[start : start + page_size] + + return { + "items": [ + {key: value for key, value in item.items() if key != "_breached"} for item in page_items + ], + "summary": {"open_by_stage": open_by_stage, "breached_by_stage": breached_by_stage}, + "pagination": { + "page": page, + "page_size": page_size, + "total": total, + "pages": math.ceil(total / page_size) if total else 0, + }, + } + + +@router.get("/stuck-orders", response_model=StuckOrdersResponse) +async def get_stuck_orders( + request: Request, + stage: str | None = Query(default=None), + include_within_sla: bool = Query(default=False), + page: int = Query(default=1, ge=1), + page_size: int = Query(default=50, ge=1, le=100), +) -> StuckOrdersResponse: + payload = await run_in_threadpool( + _build_stuck_orders_payload, request, stage, include_within_sla, page, page_size + ) + return StuckOrdersResponse.model_validate(payload) diff --git a/src/serving/semantic_layer/catalog.py b/src/serving/semantic_layer/catalog.py index c62ffff9..e6ea1d68 100644 --- a/src/serving/semantic_layer/catalog.py +++ b/src/serving/semantic_layer/catalog.py @@ -19,6 +19,11 @@ class EntityDefinition: fields: dict[str, str] # field_name -> description relationships: dict[str, str] = field(default_factory=dict) contract_version: str | None = None + # Optional SLA-stage ladder (ops-surfaces-spec.md §1.5): list order is + # ladder order, each entry `{name, sla_minutes, description}` or + # `{name, terminal: true}`. None for entities without the block — parsing + # must not require it. + stages: list[dict] | None = None @dataclass diff --git a/src/serving/semantic_layer/entity_type_registry.py b/src/serving/semantic_layer/entity_type_registry.py index 7e817f69..01dcedfa 100644 --- a/src/serving/semantic_layer/entity_type_registry.py +++ b/src/serving/semantic_layer/entity_type_registry.py @@ -103,6 +103,7 @@ def load_entity_contracts( fields=dict(data["fields"]), relationships=dict(data.get("relationships") or {}), contract_version=None, + stages=list(data["stages"]) if data.get("stages") is not None else None, ) ) diff --git a/src/serving/semantic_layer/query/engine.py b/src/serving/semantic_layer/query/engine.py index 2b83ece4..6e0cb36c 100644 --- a/src/serving/semantic_layer/query/engine.py +++ b/src/serving/semantic_layer/query/engine.py @@ -81,6 +81,7 @@ def fetch_pipeline_events( tenant_id: str | None = None, event_type: str | None = None, entity_id: str | None = None, + topic: str | None = None, limit: int | None = None, validated_only: bool = False, newest_first: bool = False, @@ -106,6 +107,11 @@ def fetch_pipeline_events( ``event_type`` accepts the demo event families (``order``, ``payment``, ``clickstream``, ``inventory``) or an exact event type; family semantics mirror what the pipeline produces. + + ``topic`` is an exact-match filter (e.g. ``orders.status`` for the + stage clock, ops-surfaces-spec.md §1.2/§3.2) — orthogonal to + ``event_type``/``validated_only``, usable without an ``entity_id`` + for a bulk scan across many entities in one query. """ # Deliberately uncached (unlike _table_columns): the journal is created # and widened by out-of-process writers, so a scan must see schema @@ -175,6 +181,10 @@ def render(value: str) -> str: where_clauses.append(f"COALESCE(tenant_id, 'default') = {render(str(tenant_id))}") if validated_only and "topic" in columns: where_clauses.append("topic = 'events.validated'") + if topic is not None: + if "topic" not in columns: + return [] + where_clauses.append(f"topic = {render(topic)}") if event_type: if "event_type" not in columns: return [] diff --git a/src/serving/semantic_layer/query/entity_queries.py b/src/serving/semantic_layer/query/entity_queries.py index 8e10a5eb..59327d99 100644 --- a/src/serving/semantic_layer/query/entity_queries.py +++ b/src/serving/semantic_layer/query/entity_queries.py @@ -82,6 +82,56 @@ def get_entity( return entity + def fetch_orders_by_status( + self: QueryExecutionHost, + statuses: list[str], + tenant_id: str | None = None, + ) -> list[dict]: + """Bulk read for the stuck-orders worklist (ops-surfaces-spec.md §3.2). + + Every order whose status is one of ``statuses`` — the caller-supplied + catalog ladder (I2: never a hardcoded stage-name literal here) — in + one query, no per-order round-trips. Journal-side composition (each + order's latest ``orders.status`` row) happens in the caller via + ``fetch_pipeline_events(topic="orders.status")``, the same port + method the Order 360 timeline already uses. + """ + entity_def = self.catalog.entities.get("order") + if entity_def is None or not statuses: + return [] + + table_name = self._qualify_table(entity_def.table, tenant_id) + use_query_params = self._backend_name == self._duckdb_backend.name + params: list[str] = [] + + def render(value: str) -> str: + if use_query_params: + params.append(value) + return "?" + return self._quote_literal(value) + + status_placeholders = ", ".join(render(status) for status in statuses) + sql = ( + # table comes from the catalog allowlist; statuses are the + # caller-supplied catalog ladder, never a literal here + f"SELECT * FROM {table_name} " # nosec B608 + f"WHERE status IN ({status_placeholders}) " + f"ORDER BY {self._quote_identifier(entity_def.primary_key)}" + ) + try: + rows = ( + self._backend.execute(sql, params) + if use_query_params + else self._backend.execute(sql) + ) + except BackendMissingTableError as e: + msg = f"Table '{table_name}' for entity 'order' is not materialized yet" + raise ValueError(msg) from e + except BackendExecutionError as e: + raise ValueError(f"Open-orders lookup failed: {e}") from e + + return [dict(row) for row in rows] + def get_entity_at( self: QueryExecutionHost, entity_type: str, diff --git a/src/serving/semantic_layer/stage_clock.py b/src/serving/semantic_layer/stage_clock.py new file mode 100644 index 00000000..7725931b --- /dev/null +++ b/src/serving/semantic_layer/stage_clock.py @@ -0,0 +1,90 @@ +"""Stage-clock resolution shared by Order 360 timeline (D2) and the +stuck-orders worklist (D3) — ops-surfaces-spec.md §1.4/§1.5. + +Budgets come from exactly one place: the catalog entity's `stages` block +(the contract loaded via ``entity_type_registry.py``). This module is the +single place that reads a budget entry and turns it into breach arithmetic, +so no stage-name or budget literal needs to be duplicated between the +timeline endpoint and the worklist endpoint (invariant I2). +""" + +from __future__ import annotations + +from datetime import UTC, datetime +from typing import Any + + +def coerce_dt(value: object) -> datetime | None: + """Parse a journal/order timestamp (datetime or ISO string) to aware UTC. + + Naive DuckDB timestamps are local wall-clock (DuckDB's ``NOW()``), not + UTC — same ``local_tz`` convention as ``EntityQueryMixin.get_entity``'s + ``_last_updated``. On non-DuckDB backends timestamps arrive as + ISO-format strings (JSON transport), not datetimes. + """ + local_tz = datetime.now().astimezone().tzinfo or UTC + if isinstance(value, datetime): + return ( + value.astimezone(UTC) + if value.tzinfo is not None + else value.replace(tzinfo=local_tz).astimezone(UTC) + ) + if isinstance(value, str): + try: + parsed = datetime.fromisoformat(value) + except ValueError: + return None + return ( + parsed.astimezone(UTC) + if parsed.tzinfo is not None + else parsed.replace(tzinfo=local_tz).astimezone(UTC) + ) + return None + + +def stage_budget( + stage_budgets: list[dict[str, Any]] | None, status: str | None +) -> dict[str, Any] | None: + """Look up the catalog stage-budget entry for `status`, or None if absent.""" + if not stage_budgets: + return None + return next( + ( + entry + for entry in stage_budgets + if isinstance(entry, dict) and entry.get("name") == status + ), + None, + ) + + +def ladder_stage_names(stage_budgets: list[dict[str, Any]] | None) -> list[str]: + """Non-terminal stage names, in catalog list order (§1.5: list order = ladder order).""" + if not stage_budgets: + return [] + return [ + entry["name"] + for entry in stage_budgets + if isinstance(entry, dict) and entry.get("name") and not entry.get("terminal") + ] + + +def resolve_breach( + *, entered_at: datetime | None, budget: dict[str, Any] | None +) -> tuple[float | None, int | None, bool | None]: + """Compute (in_stage_seconds, sla_minutes, breached) for one order's stage entry. + + `breached` is None for terminal/unknown stages or when no budget/clock is + available (I4) — never a crash, never a guess. + """ + in_stage_seconds = ( + (datetime.now(UTC) - entered_at).total_seconds() if entered_at is not None else None + ) + sla_minutes = budget.get("sla_minutes") if budget else None + is_terminal = bool(budget.get("terminal")) if budget else False + breached = ( + None + if is_terminal or sla_minutes is None or in_stage_seconds is None + else in_stage_seconds > sla_minutes * 60 + ) + return in_stage_seconds, sla_minutes, breached diff --git a/tests/integration/test_stuck_orders.py b/tests/integration/test_stuck_orders.py new file mode 100644 index 00000000..8c7df632 --- /dev/null +++ b/tests/integration/test_stuck_orders.py @@ -0,0 +1,153 @@ +"""Integration tests for the stuck-orders worklist — +GET /v1/ops/stuck-orders (ops-surfaces-spec.md §3, D3). Exercises the demo +story pin (I7: default view = exactly ORD-20260404-1004), the SLA-stage +contract block (I2), stage-vocabulary tolerance (I4), and fallback-clock +honesty (I12). 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 test_order_timeline.py. +""" + +from __future__ import annotations + +from pathlib import Path + +import pytest +from fastapi.testclient import TestClient + +from src.serving.api.main import app + +pytestmark = pytest.mark.integration + + +@pytest.fixture +def client(tmp_path: Path, monkeypatch: pytest.MonkeyPatch): + db_path = tmp_path / "stuck-orders.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 + + +def test_default_view_returns_exactly_the_sole_breach(client: TestClient): + # I7: ORD-20260404-1004 is 45 minutes into `pending` against a 30-minute + # budget — the demo's sole breach. + response = client.get("/v1/ops/stuck-orders") + + assert response.status_code == 200 + data = response.json() + + assert [item["order_id"] for item in data["items"]] == ["ORD-20260404-1004"] + item = data["items"][0] + assert item["status"] == "pending" + assert item["clock"] == "journal" + assert item["sla_minutes"] == 30 + assert item["overshoot_ratio"] == pytest.approx(1.5, rel=0.05) + assert item["total_amount"] == 1890.0 + assert item["currency"] == "RUB" + + +def test_summary_reflects_the_full_open_order_set(client: TestClient): + # Summary is always the full open-orders picture, independent of the + # breach-only default filter on `items`. + response = client.get("/v1/ops/stuck-orders") + + assert response.status_code == 200 + summary = response.json()["summary"] + assert summary["open_by_stage"] == {"pending": 2, "confirmed": 2, "shipped": 1} + assert summary["breached_by_stage"] == {"pending": 1} + + +def test_include_within_sla_returns_the_whole_open_worklist(client: TestClient): + response = client.get("/v1/ops/stuck-orders", params={"include_within_sla": "true"}) + + assert response.status_code == 200 + data = response.json() + order_ids = [item["order_id"] for item in data["items"]] + + assert set(order_ids) == { + "ORD-20260404-1002", + "ORD-20260404-1003", + "ORD-20260404-1004", + "ORD-20260404-1007", + "ORD-20260404-1008", + } + # Highest overshoot first — the breach still sorts to the top. + assert order_ids[0] == "ORD-20260404-1004" + ratios = [item["overshoot_ratio"] for item in data["items"]] + assert ratios == sorted(ratios, reverse=True) + + +def test_terminal_orders_never_appear_even_with_include_within_sla(client: TestClient): + response = client.get("/v1/ops/stuck-orders", params={"include_within_sla": "true"}) + + assert response.status_code == 200 + order_ids = {item["order_id"] for item in response.json()["items"]} + # ORD-1001/1005 delivered, ORD-1006 cancelled — terminal, never stuck. + assert order_ids.isdisjoint({"ORD-20260404-1001", "ORD-20260404-1005", "ORD-20260404-1006"}) + + +def test_stage_filter_narrows_to_one_ladder_stage(client: TestClient): + response = client.get( + "/v1/ops/stuck-orders", params={"stage": "confirmed", "include_within_sla": "true"} + ) + + assert response.status_code == 200 + data = response.json() + assert {item["order_id"] for item in data["items"]} == { + "ORD-20260404-1003", + "ORD-20260404-1007", + } + assert all(item["status"] == "confirmed" for item in data["items"]) + + +def test_pagination_shape(client: TestClient): + response = client.get( + "/v1/ops/stuck-orders", params={"include_within_sla": "true", "page": 1, "page_size": 2} + ) + + assert response.status_code == 200 + data = response.json() + assert len(data["items"]) == 2 + assert data["pagination"] == {"page": 1, "page_size": 2, "total": 5, "pages": 3} + + +def test_order_without_stage_rows_reports_fallback_clock(client: TestClient): + # I12: an order written outside the stage-row writer degrades honestly + # to the created_at fallback instead of pretending to have a journal + # clock. + conn = client.app.state.query_engine._conn + conn.execute( + """ + INSERT INTO orders_v2 (order_id, user_id, status, total_amount, currency, created_at) + VALUES ('ORD-BYPASS-1', 'USR-10001', 'confirmed', 500.0, 'RUB', + NOW() - INTERVAL '10 minutes') + """ + ) + + response = client.get("/v1/ops/stuck-orders", params={"include_within_sla": "true"}) + + assert response.status_code == 200 + item = next(i for i in response.json()["items"] if i["order_id"] == "ORD-BYPASS-1") + assert item["clock"] == "fallback" + assert 0 < item["in_stage_seconds"] < 900 + + +def test_order_with_status_outside_the_ladder_never_crashes_or_appears(client: TestClient): + # I4: a status the contract's stages: block doesn't know about is not + # part of the ladder query — it never surfaces and never 500s. + conn = client.app.state.query_engine._conn + conn.execute( + """ + INSERT INTO orders_v2 (order_id, user_id, status, total_amount, currency, created_at) + VALUES ('ORD-WEIRD-STATUS', 'USR-10001', 'on_hold', 500.0, 'RUB', + NOW() - INTERVAL '10 minutes') + """ + ) + + response = client.get("/v1/ops/stuck-orders", params={"include_within_sla": "true"}) + + assert response.status_code == 200 + order_ids = {item["order_id"] for item in response.json()["items"]} + assert "ORD-WEIRD-STATUS" not in order_ids diff --git a/tests/integration/test_tenant_isolation.py b/tests/integration/test_tenant_isolation.py index 78caf12f..21a28396 100644 --- a/tests/integration/test_tenant_isolation.py +++ b/tests/integration/test_tenant_isolation.py @@ -205,6 +205,37 @@ def test_tenant_api_key_reads_own_order_timeline(client: TestClient): assert data["order"]["user_id"] == "USR-ACME-2" +def test_cross_tenant_stuck_orders_are_scoped_to_tenant_schema(client: TestClient): + # ops-surfaces-spec.md §1.7 / invariant I8: the stuck-orders worklist + # scopes the open-orders read by the request tenant, same as the entity + # route and the Order 360 timeline. + acme_response = client.get( + "/v1/ops/stuck-orders", + params={"include_within_sla": "true"}, + headers={"X-API-Key": "acme-key"}, + ) + demo_response = client.get( + "/v1/ops/stuck-orders", + params={"include_within_sla": "true"}, + headers={"X-API-Key": "demo-key"}, + ) + + assert acme_response.status_code == 200 + assert demo_response.status_code == 200 + + acme_items = acme_response.json()["items"] + demo_items = demo_response.json()["items"] + + # ORD-ACME (delivered) is terminal — never in the worklist for anyone. + assert {item["order_id"] for item in acme_items} == {"ORD-SHARED"} + assert {item["order_id"] for item in demo_items} == {"ORD-SHARED", "ORD-DEMO"} + + acme_shared = next(item for item in acme_items if item["order_id"] == "ORD-SHARED") + demo_shared = next(item for item in demo_items if item["order_id"] == "ORD-SHARED") + assert acme_shared["user_id"] == "USR-ACME" + assert demo_shared["user_id"] == "USR-DEMO" + + 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 04e31b32..dfddadd0 100644 --- a/tests/unit/test_control_plane_store.py +++ b/tests/unit/test_control_plane_store.py @@ -854,10 +854,15 @@ def test_webhook_path_does_not_reach_into_the_engine_connection() -> None: def test_ops_timeline_path_does_not_reach_into_the_engine_connection_or_vault() -> None: - """ADR 0011 invariant I1: the ops surfaces (Order 360 timeline first, - D2) compose exactly the QueryEngine/ServingBackend and ControlPlaneStore - ports — no raw ``query_engine._conn`` reach, no vault DSN, ever.""" - source = (PROJECT_ROOT / "src/serving/api/routers/agent_query.py").read_text(encoding="utf-8") - assert "query_engine._conn" not in source - assert "_conn." not in source - assert "VAULT_DSN" not in source + """ADR 0011 invariant I1: the ops surfaces (Order 360 timeline, D2; + stuck-orders worklist, D3) compose exactly the QueryEngine/ServingBackend + and ControlPlaneStore ports — no raw ``query_engine._conn`` reach, no + vault DSN, ever. Covers every module under ``routers/ops*`` per spec §5.""" + for relative in ( + "src/serving/api/routers/agent_query.py", + "src/serving/api/routers/ops.py", + ): + source = (PROJECT_ROOT / relative).read_text(encoding="utf-8") + assert "query_engine._conn" not in source, relative + assert "_conn." not in source, relative + assert "VAULT_DSN" not in source, relative diff --git a/tests/unit/test_entity_type_registry.py b/tests/unit/test_entity_type_registry.py index c8fe88fa..2c99ed3e 100644 --- a/tests/unit/test_entity_type_registry.py +++ b/tests/unit/test_entity_type_registry.py @@ -183,6 +183,43 @@ def test_loader_rejects_missing_directory(tmp_path: Path) -> None: load_entity_contracts(missing) +def test_loader_is_tolerant_of_absent_stages_block(tmp_path: Path) -> None: + # ops-surfaces-spec.md §1.5: entities without a `stages:` block behave as + # today — parsing must not require it. + _write( + tmp_path, + "widget", + { + "name": "widget", + "description": "no SLA ladder", + "table": "widgets", + "primary_key": "widget_id", + "fields": {"widget_id": "pk"}, + }, + ) + + loaded = load_entity_contracts(tmp_path) + + assert loaded[0].stages is None + + +def test_loader_parses_order_stages_ladder() -> None: + # ops-surfaces-spec.md §1.5: list order is ladder order; terminal entries + # carry no sla_minutes. + order = DataCatalog().entities["order"] + + assert [entry["name"] for entry in order.stages] == [ + "pending", + "confirmed", + "shipped", + "delivered", + "cancelled", + ] + assert order.stages[0]["sla_minutes"] == 30 + assert order.stages[3]["terminal"] is True + assert "sla_minutes" not in order.stages[3] + + def test_contracts_dir_shipped_in_repo_is_wellformed() -> None: assert CONTRACT_DIR.is_dir() yaml_files = list(CONTRACT_DIR.glob("*.yaml")) diff --git a/tests/unit/test_security_tooling_policy.py b/tests/unit/test_security_tooling_policy.py index 780cfc7d..df3ff94e 100644 --- a/tests/unit/test_security_tooling_policy.py +++ b/tests/unit/test_security_tooling_policy.py @@ -74,7 +74,12 @@ def test_sql_injection_checks_are_not_globally_suppressed() -> None: "src/serving/backends/duckdb_backend.py": 2, "src/serving/semantic_layer/nl_engine.py": 6, "src/serving/semantic_layer/query/engine.py": 1, - "src/serving/semantic_layer/query/entity_queries.py": 3, + # D3 (reviewed 2026-07-04): fetch_orders_by_status's new stuck-orders + # bulk read follows get_entity's existing pattern in this same file — the + # table name comes from the catalog allowlist (_qualify_table), and every + # status value binds as a query param on DuckDB or is + # _quote_literal-escaped on the non-binding ClickHouse path. + "src/serving/semantic_layer/query/entity_queries.py": 4, "src/serving/semantic_layer/query/nl_queries.py": 3, "src/serving/semantic_layer/search_index.py": 1, }