diff --git a/scripts/tests/test_keycloak_desired_clients_contract.py b/scripts/tests/test_keycloak_desired_clients_contract.py index 00753edf51..474ea3f936 100644 --- a/scripts/tests/test_keycloak_desired_clients_contract.py +++ b/scripts/tests/test_keycloak_desired_clients_contract.py @@ -43,7 +43,6 @@ def test_omniweb_user_identity_claim_contract() -> None: assert "gateway-attach-audience" not in mappers assert all( - mapper.get("config", {}).get("included.custom.audience") - != "gateway-attach" + mapper.get("config", {}).get("included.custom.audience") != "gateway-attach" for mapper in mappers.values() ) diff --git a/src/omnibase_infra/runtime/auto_wiring/handler_wiring.py b/src/omnibase_infra/runtime/auto_wiring/handler_wiring.py index 14bd1f1117..eb97a6dc1c 100644 --- a/src/omnibase_infra/runtime/auto_wiring/handler_wiring.py +++ b/src/omnibase_infra/runtime/auto_wiring/handler_wiring.py @@ -38,6 +38,7 @@ from collections.abc import Awaitable, Callable, Collection, Iterator, Mapping, Sequence from dataclasses import dataclass, field from datetime import UTC, datetime +from functools import lru_cache from pathlib import Path from typing import ( TYPE_CHECKING, @@ -49,7 +50,7 @@ ) from uuid import UUID, uuid4 -from pydantic import BaseModel, ValidationError +from pydantic import AliasChoices, AliasPath, BaseModel, ValidationError from omnibase_core.enums.enum_core_error_code import EnumCoreErrorCode from omnibase_core.enums.enum_database_grant_object_type import ( @@ -799,7 +800,12 @@ async def _callback( else: target_model = _resolve_def_b_input_model_type(handle_method) if target_model is not None: - payload = _extract_dispatch_payload(envelope) + # OMN-16050: pass the registered input model so the unwrap + # STOPS at it. ``ModelEmitRequest`` declares ``payload`` plus + # four transport markers, so a marker-only heuristic unwrapped + # through it and handed the handler the caller's inner + # payload — every node_event_emit_effect command DLQ'd. + payload = _extract_dispatch_payload(envelope, target_model) if isinstance(payload, target_model): dispatch_arg = payload elif isinstance(payload, Mapping): @@ -832,7 +838,13 @@ async def _callback( raw_result, envelope, None, handler_node_kind, published_event_names ) - payload = _extract_dispatch_payload(envelope) + # OMN-16050: resolve the contract-declared event model BEFORE extracting so + # the unwrap can stop at it (same fail-closed rule as the def-B branch + # above). Resolution failure is not fatal here — the existing try/except + # below owns that path — so the hint degrades to None and the extraction + # keeps its pre-OMN-16050 structural behaviour. + payload_target_model = _safe_import_event_model_class(event_model) + payload = _extract_dispatch_payload(envelope, payload_target_model) handler_takes_envelope = _handler_accepts_event_envelope( cast("Callable[..., object]", handle_method) ) @@ -1071,6 +1083,24 @@ def _import_event_model_class(event_model: ModelHandlerRef) -> type[BaseModel]: return cast("type[BaseModel]", model_cls) +def _safe_import_event_model_class( + event_model: ModelHandlerRef | None, +) -> type[BaseModel] | None: + """``_import_event_model_class`` that yields None instead of raising (OMN-16050). + + Used only to hint ``_extract_dispatch_payload`` with the contract-declared + target type. An unimportable/malformed ``event_model`` must not change dispatch + control flow from this call site — the caller's own + ``_import_event_model_class`` inside its try/except still owns that failure. + """ + if event_model is None: + return None + try: + return _import_event_model_class(event_model) + except Exception: # noqa: BLE001 — hint-only resolution, never fatal here + return None + + def _handler_accepts_event_envelope(handle_method: object) -> bool: """Return true when a handler's first parameter is envelope-shaped.""" try: @@ -1380,11 +1410,21 @@ def _materialize_typed_event_envelope( # Transport-envelope keys the runtime adds around the domain payload. When the # dispatch engine materializes a ModelEventEnvelope to a dict it nests the domain # fields under ``payload`` and carries routing metadata (``partition_key`` etc.) -# alongside. Domain models never declare these keys, so a mapping that carries a -# ``payload`` mapping plus any marker is a transport envelope to unwrap. Mirrors -# omnimarket's ``_ENVELOPE_MARKER_KEYS`` predicate (OMN-12935/12936); the -# auto-wiring kernel unwraps here because it constructs the typed model itself, -# upstream of the handler's own coercion (OMN-12940). +# alongside, so a mapping that carries a ``payload`` mapping plus any marker MAY +# be a transport envelope to unwrap. Mirrors omnimarket's +# ``_ENVELOPE_MARKER_KEYS`` predicate (OMN-12935/12936); the auto-wiring kernel +# unwraps here because it constructs the typed model itself, upstream of the +# handler's own coercion (OMN-12940). +# +# OMN-16050 — this marker set is a NECESSARY, NOT SUFFICIENT signal. The earlier +# text here asserted "domain models never declare these keys"; that invariant is +# FALSE. ``ModelEmitRequest`` (node_event_emit_effect) declares ``payload`` plus +# four of these markers (``event_type``, ``correlation_id``, ``partition_key``, +# ``event_id``), is structurally indistinguishable from a transport envelope, and +# was therefore unwrapped THROUGH — the handler got the caller's inner payload, +# ``model_validate`` raised, and every command DLQ'd. The registered-input-model +# stop condition below (``_is_registered_input_payload``) is what makes the +# heuristic safe: structure alone can never decide this. _ENVELOPE_MARKER_KEYS: frozenset[str] = frozenset( { "partition_key", @@ -1398,11 +1438,13 @@ def _materialize_typed_event_envelope( def _is_transport_envelope(value: object) -> bool: - """True when ``value`` is a transport envelope wrapping a domain payload. + """True when ``value`` is envelope-SHAPED: a ``payload`` mapping plus a marker. - A transport envelope is a mapping that carries a ``payload`` mapping plus at - least one transport marker key. Requiring a marker avoids over-unwrapping a - legitimate domain model that happens to declare its own ``payload`` field. + Structural precondition only. A domain model may legitimately declare both a + ``payload`` mapping and transport-plausible marker fields (OMN-16050), so this + predicate is never sufficient on its own to justify an unwrap — see + ``_is_registered_input_payload``, the fail-closed stop condition applied by + ``_extract_dispatch_payload``. """ return ( isinstance(value, Mapping) @@ -1411,16 +1453,101 @@ def _is_transport_envelope(value: object) -> bool: ) -def _extract_dispatch_payload(envelope: object) -> object: +def _validation_alias_wire_keys(alias: object) -> set[str]: + """Top-level wire keys a pydantic ``validation_alias`` can consume. + + ``validation_alias`` has three shapes and only the plain-string one is a + single key. ``AliasPath("meta", "id")`` consumes the TOP-LEVEL key ``meta`` + (the remaining segments index inside that value), and ``AliasChoices`` holds + a list of alternatives, each itself a string or an ``AliasPath``. + + Missing the non-string shapes is fail-OPEN for OMN-16050: a model aliased + that way would fail ``_is_registered_input_payload``'s key-containment check + even when the candidate IS the registered model, the unwrap would continue + into the caller's payload, and the DLQ defect would return for exactly the + contracts that use richer aliases. + """ + if isinstance(alias, str): + return {alias} + if isinstance(alias, AliasPath): + first = alias.path[0] if alias.path else None + return {first} if isinstance(first, str) else set() + if isinstance(alias, AliasChoices): + keys: set[str] = set() + for choice in alias.choices: + keys |= _validation_alias_wire_keys(choice) + return keys + return set() + + +@lru_cache(maxsize=512) +def _model_declared_wire_keys(model: type[BaseModel]) -> frozenset[str]: + """Every wire key ``model`` can accept: field names plus their input aliases.""" + keys: set[str] = set() + for field_name, model_field in model.model_fields.items(): + keys.add(field_name) + if isinstance(model_field.alias, str): + keys.add(model_field.alias) + keys |= _validation_alias_wire_keys(model_field.validation_alias) + return frozenset(keys) + + +def _is_registered_input_payload( + candidate: object, target_model: type[BaseModel] | None +) -> bool: + """True when ``candidate`` IS the dispatcher's registered input model on the wire. + + The fail-closed stop condition for the recursive unwrap (OMN-16050). A + candidate is claimed by the registered model only when BOTH hold: + + 1. **Key containment** — every key present on the candidate is a declared + field (or input alias) of ``target_model``. A real transport envelope + always carries at least one routing/marker key the domain model does not + declare (``source_tool``, ``envelope_id``, ``__debug_trace``, + ``__bindings``, ``envelope_timestamp``, ...), so this alone keeps the + OMN-12940 double-wrapped case unwrapping. + 2. **Full validation** — the candidate validates as ``target_model``, so a + partial structural coincidence never halts the unwrap short of the domain. + + The cheap set check runs first; ``model_validate`` executes only for the rare + candidate whose keys are entirely owned by the target model. + + Deliberately NOT a marker denylist: dropping ``event_type``/``correlation_id`` + from ``_ENVELOPE_MARKER_KEYS`` would fix ``ModelEmitRequest`` and silently + break every genuine envelope that carries only those markers. This predicate + keys on the CONTRACT-registered target type instead of on key spelling. + """ + if target_model is None or not isinstance(candidate, Mapping): + return False + if not candidate.keys() <= _model_declared_wire_keys(target_model): + return False + try: + target_model.model_validate(dict(candidate)) + except Exception: # noqa: BLE001 — any validation failure means "not the model" + return False + return True + + +def _extract_dispatch_payload( + envelope: object, target_model: type[BaseModel] | None = None +) -> object: # The runtime may deliver a DOUBLE- (or deeper-) wrapped envelope, e.g. # ``{"payload": {"payload": {domain}, ...markers}, "partition_key": None}``. # Unwrap recursively until the domain payload is reached so the kernel's # ``model_validate`` (and the post-handler correlation read) operate on the # domain, not on an intermediate envelope (OMN-12940). + # + # OMN-16050: stop the moment the candidate IS the dispatcher's registered + # input model. ``target_model`` is the contract-declared type the kernel is + # about to construct (the def-B ``handle()`` annotation, or the handler's + # declared ``event_model``); when it is None the caller has no registered + # type in scope and the pre-existing structural behaviour is unchanged. candidate: object = envelope if not isinstance(candidate, Mapping): candidate = getattr(candidate, "payload", candidate) - while _is_transport_envelope(candidate): + while _is_transport_envelope(candidate) and not _is_registered_input_payload( + candidate, target_model + ): candidate = cast("Mapping[str, object]", candidate)["payload"] return candidate diff --git a/tests/integration/test_auto_wiring_real_manifest.py b/tests/integration/test_auto_wiring_real_manifest.py index 40552d4a45..fd01946d2a 100644 --- a/tests/integration/test_auto_wiring_real_manifest.py +++ b/tests/integration/test_auto_wiring_real_manifest.py @@ -21,15 +21,26 @@ from pathlib import Path from unittest.mock import MagicMock -from uuid import UUID +from uuid import UUID, uuid4 import pytest +from pydantic import BaseModel, ConfigDict, Field from omnibase_infra.runtime.auto_wiring.discovery import discover_contracts from omnibase_infra.runtime.auto_wiring.handler_wiring import wire_from_manifest +from omnibase_infra.runtime.auto_wiring.models import ( + ModelAutoWiringManifest, + ModelContractVersion, + ModelDiscoveredContract, + ModelEventBusWiring, + ModelHandlerRef, + ModelHandlerRouting, + ModelHandlerRoutingEntry, +) from omnibase_infra.runtime.auto_wiring.models.model_discovery_error import ( ModelDiscoveryError, ) +from omnibase_infra.runtime.message_dispatch_engine import MessageDispatchEngine from omnibase_infra.runtime.service_intent_routing_loader import ( load_intent_routing_table, ) @@ -178,3 +189,185 @@ async def test_real_manifest_wiring_has_no_failures() -> None: assert not error_results, "ModelOnexError found in wiring results:\n" + "\n".join( f" {r.contract_name}: {r.reason}" for r in error_results ) + + +# --------------------------------------------------------------------------- +# OMN-16050 — registered-input-model unwrap stop, proven on the REAL manifest +# --------------------------------------------------------------------------- + +_OMN16050_TOPIC = "onex.cmd.omnibase-infra.omn16050-unwrap-probe.v1" +_OMN16050_MODULE = "tests.integration.test_auto_wiring_real_manifest" + + +class ModelOmn16050EmitRequest(BaseModel): + """Field-for-field mirror of omnimarket's ``ModelEmitRequest`` (OMN-16050). + + Declares a ``payload`` mapping plus FOUR transport marker keys + (``event_type``, ``correlation_id``, ``partition_key``, ``event_id``), which + is precisely what made it indistinguishable from a transport envelope to the + old structural heuristic. ``extra="forbid"`` mirrors the real model, so an + over-unwrap fails totally — the live DLQ signature. + + Lives here rather than in omnimarket because omnibase_infra is upstream of it + and cannot import it; the shape, not the identity, is what the defect keys on. + """ + + model_config = ConfigDict(frozen=True, extra="forbid") + + event_type: str = Field(..., min_length=1) + payload: dict[str, object] = Field(default_factory=dict) + correlation_id: str | None = None + topic: str | None = None + partition_key: str | None = None + event_id: str = Field(default_factory=lambda: str(uuid4()), min_length=1) + + +class HandlerOmn16050EmitProbe: + """Canonical def-B handler: ``handle(request: ModelX) -> None``. + + Wired by the real ``wire_from_manifest`` path, so the callback under test is + the production-built one, not a hand-constructed ``_make_dispatch_callback``. + """ + + received: list[ModelOmn16050EmitRequest] = [] + + async def handle(self, request: ModelOmn16050EmitRequest) -> None: + type(self).received.append(request) + + +def _omn16050_probe_contract() -> ModelDiscoveredContract: + """An ``operation_match`` def-B EFFECT contract shaped like node_event_emit_effect.""" + return ModelDiscoveredContract( + name="node_omn16050_unwrap_probe", + node_type="EFFECT_GENERIC", + contract_version=ModelContractVersion(major=1, minor=0, patch=0), + contract_path=Path("/fake/omn16050/contract.yaml"), + entry_point_name="node_omn16050_unwrap_probe", + package_name="omnibase-infra", + event_bus=ModelEventBusWiring( + subscribe_topics=(_OMN16050_TOPIC,), + publish_topics=(), + ), + handler_routing=ModelHandlerRouting( + routing_strategy="operation_match", + handlers=( + ModelHandlerRoutingEntry( + handler=ModelHandlerRef( + name="HandlerOmn16050EmitProbe", module=_OMN16050_MODULE + ), + message_category="command", + event_type="omnibase-infra.omn16050-unwrap-probe", + operation="omn16050.probe", + ), + ), + ), + ) + + +def _omn16050_published_bytes() -> dict[str, object]: + """The live shape: one transport envelope wrapping a ModelEmitRequest. + + Mirrors the in-pod capture on onex-dev (digest sha256:35099472…) verbatim: + ``RAW KEYS: ['event_type', 'correlation_id', 'source_tool', 'payload']``. + """ + return { + "event_type": "session.started", + "correlation_id": "18a50ff5-c877-481c-b3c1-a183d8069762", + "source_tool": "defect-ab-probe", + "payload": { + "event_type": "session.started", + "correlation_id": "18a50ff5-c877-481c-b3c1-a183d8069762", + "partition_key": "session-1", + "event_id": "evt-defect-ab-probe", + "payload": { + "session_id": "18a50ff5-c877-481c-b3c1-a183d8069762", + "defect_ab_probe": True, + "emitted_at": "2026-08-13T02:46:13Z", + }, + }, + } + + +@pytest.mark.integration +@pytest.mark.asyncio +async def test_real_manifest_wiring_preserves_registered_envelope_shaped_input_model() -> ( + None +): + """Runtime-startup gate for OMN-16050: real manifest + real wiring + real dispatch. + + Satisfies the repo's Runtime Startup CI gate for a PR touching + ``auto_wiring/``: the manifest is the real one loaded from disk via + ``discover_contracts()``, ``wire_from_manifest`` runs with the kernel's + argument shape, and zero unexpected failures are asserted. + + On top of that it proves the defect is closed through the production path: + one probe contract whose def-B handler declares an envelope-SHAPED input + model (``payload`` + four transport markers, ``extra="forbid"``) is wired + alongside the real contracts, and the dispatcher the wiring registered is + invoked with the exact bytes captured in-pod. Before the fix the callback + unwrapped through the domain model to the caller's inner payload and raised + ``ValidationError`` (``event_type`` Field required + 3x extra_forbidden) — + the live ``boundary_swallow_prevented`` / DLQ signature. + """ + HandlerOmn16050EmitProbe.received.clear() + + real_manifest = discover_contracts() + combined = ModelAutoWiringManifest( + contracts=(*real_manifest.contracts, _omn16050_probe_contract()), + errors=real_manifest.errors, + ) + derived_appliers = { + contract.name: _StubResultApplier() + for contract in real_manifest.contracts + if load_intent_routing_table(Path(contract.contract_path)) + } + + engine = MessageDispatchEngine() + report = await wire_from_manifest( + manifest=combined, + dispatch_engine=engine, + event_bus=None, + subscribe_immediately=False, + result_appliers_by_contract=derived_appliers, + ) + + failed_names = { + r.contract_name for r in report.results if str(r.outcome).endswith("FAILED") + } + unexpected = failed_names - _KNOWN_UNWIRED_RAW_PROJECTIONS + assert not unexpected, ( + "wire_from_manifest() reported unexpected failure(s) against the real " + f"manifest + the OMN-16050 probe contract: {sorted(unexpected)}" + ) + + probe_result = next( + r for r in report.results if r.contract_name == "node_omn16050_unwrap_probe" + ) + assert str(probe_result.outcome).endswith("WIRED"), ( + "the OMN-16050 probe contract must WIRE — a skipped/failed probe would " + f"make this gate vacuous (outcome={probe_result.outcome}, " + f"reason={probe_result.reason})" + ) + assert len(probe_result.dispatchers_registered) == 1 + + dispatcher_id = probe_result.dispatchers_registered[0] + dispatcher = engine._dispatchers[dispatcher_id].dispatcher + + await dispatcher(_omn16050_published_bytes()) + + assert len(HandlerOmn16050EmitProbe.received) == 1, ( + "the wired dispatcher did not deliver to the handler — pre-fix this " + "raised ValidationError inside the callback and DLQ'd" + ) + request = HandlerOmn16050EmitProbe.received[0] + assert isinstance(request, ModelOmn16050EmitRequest) + assert request.event_type == "session.started" + assert request.event_id == "evt-defect-ab-probe" + assert request.partition_key == "session-1" + # The load-bearing assertion: the handler owns the CALLER's payload, and the + # unwrap stopped at the registered model instead of walking through it. + assert request.payload == { + "session_id": "18a50ff5-c877-481c-b3c1-a183d8069762", + "defect_ab_probe": True, + "emitted_at": "2026-08-13T02:46:13Z", + } diff --git a/tests/scripts/test_deploy_runtime_core_contracts_resolution.py b/tests/scripts/test_deploy_runtime_core_contracts_resolution.py index b0642f0155..1304333766 100644 --- a/tests/scripts/test_deploy_runtime_core_contracts_resolution.py +++ b/tests/scripts/test_deploy_runtime_core_contracts_resolution.py @@ -1,4 +1,4 @@ -# SPDX-FileCopyrightText: 2026 OmniNode.ai Inc. +# SPDX-FileCopyrightText: 2025 OmniNode.ai Inc. # SPDX-License-Identifier: MIT """deploy-runtime.sh step 3b must resolve omnibase_core runtime contracts from the diff --git a/tests/unit/runtime/auto_wiring/test_omn16050_registered_input_model_unwrap_stop.py b/tests/unit/runtime/auto_wiring/test_omn16050_registered_input_model_unwrap_stop.py new file mode 100644 index 0000000000..7beccb4a70 --- /dev/null +++ b/tests/unit/runtime/auto_wiring/test_omn16050_registered_input_model_unwrap_stop.py @@ -0,0 +1,545 @@ +# SPDX-FileCopyrightText: 2025 OmniNode.ai Inc. +# SPDX-License-Identifier: MIT +"""OMN-16050 — the recursive envelope unwrap must stop at the registered input model. + +THE DEFECT. ``_extract_dispatch_payload`` unwraps ``payload`` recursively while +``_is_transport_envelope`` holds, and that predicate is purely structural: a +mapping carrying a ``payload`` mapping plus any of ``_ENVELOPE_MARKER_KEYS`` +(``partition_key``/``event_type``/``envelope_id``/``event_id``/``correlation_id``/ +``__debug_trace``). The module asserted the invariant "domain models never declare +these keys". That invariant is FALSE. + +``ModelEmitRequest`` (``node_event_emit_effect``) declares ``payload`` plus FOUR +of those markers (``event_type``, ``correlation_id``, ``partition_key``, +``event_id``) and is ``extra="forbid"``. It is structurally identical to a +transport envelope, so the runtime unwrapped THROUGH it and handed the kernel the +caller's inner user payload. Live evidence (onex-dev, digest sha256:35099472…), +replayed in-pod against the real published bytes:: + + RAW KEYS: ['event_type', 'correlation_id', 'source_tool', 'payload'] + EXTRACTED: dict KEYS ['session_id', 'defect_ab_probe', 'emitted_at'] + EXTRACTED -> ValidationError: 4 validation errors for ModelEmitRequest + event_type Field required / 3x extra_forbidden + -> HandlerDispatchFailureError -> boundary_swallow_prevented -> DLQ + +so ``node_event_emit_effect`` could never be dispatched over the bus. + +THE FIX under test. ``_extract_dispatch_payload`` now takes the dispatcher's +registered input model and stops the unwrap at a candidate that IS that model — +key-containment (every key on the candidate is a declared field/alias of the +target) AND full ``model_validate``. Fail-closed in both directions: an envelope +always carries at least one routing key the domain model does not declare +(``source_tool``, ``envelope_id``, ``__debug_trace``, ``__bindings``…), so genuine +double-wrapped deliveries (OMN-12940) keep unwrapping to the domain. + +This module holds the RED reproduction plus the regressions that pin both +directions. The real-manifest / ``wire_from_manifest`` runtime-startup gate for +the same defect lives in ``tests/integration/test_auto_wiring_real_manifest.py`` +(``test_real_manifest_wiring_preserves_registered_envelope_shaped_input_model``). +""" + +from __future__ import annotations + +import re +from typing import cast +from uuid import uuid4 + +import pytest +from pydantic import ( + AliasChoices, + AliasPath, + BaseModel, + ConfigDict, + Field, + field_validator, +) + +from omnibase_infra.runtime.auto_wiring.handler_wiring import ( + _extract_dispatch_payload, + _is_registered_input_payload, + _is_transport_envelope, + _make_dispatch_callback, + _model_declared_wire_keys, +) +from omnibase_infra.runtime.auto_wiring.models import ModelHandlerRef + +_THIS_MODULE = ( + "tests.unit.runtime.auto_wiring.test_omn16050_registered_input_model_unwrap_stop" +) + +_TOPIC_SHAPE_RE = re.compile(r"^onex\.(evt|cmd|intent|dlq)\.[a-z0-9._-]+\.v\d+$") + + +class ModelRuntimeEmitRequest(BaseModel): + """Field-for-field mirror of omnimarket's ``ModelEmitRequest`` (OMN-16050). + + omnibase_infra cannot import omnimarket (it is a downstream package), so the + defect is reproduced against a local model carrying the exact shape that + triggers it: a ``payload`` mapping plus four transport marker keys, with + ``extra="forbid"`` so an over-unwrap fails totally rather than partially. + """ + + model_config = ConfigDict(frozen=True, extra="forbid") + + event_type: str = Field(..., min_length=1) + payload: dict[str, object] = Field(default_factory=dict) + correlation_id: str | None = None + topic: str | None = None + partition_key: str | None = None + event_id: str = Field(default_factory=lambda: str(uuid4()), min_length=1) + + +class ModelRuntimeDomainCommand(BaseModel): + """A plain domain command: no ``payload`` field, no markers (OMN-12940 shape).""" + + model_config = ConfigDict(extra="forbid") + + correlation_id: str + source_commit_sha: str + + +class ModelLenientEnvelopeShaped(BaseModel): + """Envelope-shaped domain model that IGNORES extras — the adversarial case. + + Without key-containment, ``model_validate`` alone would accept a genuine + transport envelope (extras silently dropped) and halt the unwrap one level + too early. Pinned by ``test_lenient_model_does_not_claim_a_real_envelope``. + """ + + model_config = ConfigDict(extra="ignore") + + event_type: str + payload: dict[str, object] = Field(default_factory=dict) + + +def _user_payload() -> dict[str, object]: + """The inner, caller-authored payload from the live probe.""" + return { + "session_id": "18a50ff5-c877-481c-b3c1-a183d8069762", + "defect_ab_probe": True, + "emitted_at": "2026-08-13T02:46:13Z", + } + + +def _emit_request_wire() -> dict[str, object]: + """The domain command as published: a ModelEmitRequest on the wire.""" + return { + "event_type": "session.started", + "correlation_id": "18a50ff5-c877-481c-b3c1-a183d8069762", + "partition_key": "session-1", + "event_id": "evt-defect-ab-probe", + "payload": _user_payload(), + } + + +def _transport_envelope(inner: dict[str, object]) -> dict[str, object]: + """One real transport layer, exactly as the live publish produced it. + + ``source_tool`` is the discriminator no domain model declares — it is what + tells the extractor this layer is transport and the next one may not be. + """ + return { + "event_type": "session.started", + "correlation_id": "18a50ff5-c877-481c-b3c1-a183d8069762", + "source_tool": "defect-ab-probe", + "payload": inner, + } + + +def _domain_command() -> dict[str, object]: + return {"correlation_id": str(uuid4()), "source_commit_sha": "abcdef1"} + + +# --------------------------------------------------------------------------- +# RED: the exact production coercion failure +# --------------------------------------------------------------------------- + + +@pytest.mark.unit +class TestProductionCoercionFailure: + def test_emit_request_is_structurally_indistinguishable_from_an_envelope( + self, + ) -> None: + """Ground truth: the false invariant. This is WHY the marker set cannot decide. + + If this ever goes green-by-inversion (the request stops looking like an + envelope), the marker heuristic was changed and the tests below are + testing a premise that no longer holds — revisit them rather than + deleting this assertion. + """ + assert _is_transport_envelope(_emit_request_wire()) is True + + def test_registered_input_model_survives_the_unwrap(self) -> None: + """RED before the fix: the extractor returned the INNER user payload. + + Reproduces the in-pod replay verbatim — one genuine transport layer + wrapping a ModelEmitRequest-shaped domain command. Pre-fix the extractor + unwrapped twice and returned ``['session_id', 'defect_ab_probe', + 'emitted_at']``; it must stop at the emit request. + """ + envelope = _transport_envelope(_emit_request_wire()) + + extracted = _extract_dispatch_payload(envelope, ModelRuntimeEmitRequest) + + assert extracted == _emit_request_wire() + assert sorted(cast("dict[str, object]", extracted)) == [ + "correlation_id", + "event_id", + "event_type", + "partition_key", + "payload", + ] + # And the kernel-side construction the runtime performs next succeeds. + request = ModelRuntimeEmitRequest.model_validate(extracted) + assert request.event_type == "session.started" + assert request.payload == _user_payload() + + def test_without_a_registered_model_the_over_unwrap_still_reproduces(self) -> None: + """The defect is EXACTLY the missing target type, not a shape change. + + With no registered model in scope the structural heuristic is unchanged + (that is the deliberate no-behaviour-change property for the six call + sites that read correlation/DLQ metadata) — it still walks to the inner + user payload. This isolates the fix to the stop condition. + """ + envelope = _transport_envelope(_emit_request_wire()) + + assert _extract_dispatch_payload(envelope) == _user_payload() + + @pytest.mark.asyncio + async def test_def_b_dispatch_constructs_the_registered_model(self) -> None: + """End-to-end at the coercion boundary: the DLQ'd dispatch now succeeds. + + Pre-fix this raised ``ValidationError`` (``event_type`` Field required + + 3x extra_forbidden) inside the auto-wiring callback, which the kernel + reported as ``HandlerDispatchFailureError`` → ``boundary_swallow_prevented`` + → DLQ. + """ + captured: dict[str, object] = {} + + class _EmitHandler: + async def handle(self, request: ModelRuntimeEmitRequest) -> None: + captured["request"] = request + + callback = _make_dispatch_callback(_EmitHandler()) + + await callback(_transport_envelope(_emit_request_wire())) + + request = captured["request"] + assert isinstance(request, ModelRuntimeEmitRequest) + assert request.event_type == "session.started" + assert request.event_id == "evt-defect-ab-probe" + assert request.payload == _user_payload() + + @pytest.mark.asyncio + async def test_event_model_dispatch_constructs_the_registered_model(self) -> None: + """The same stop applies to the contract-declared ``event_model`` branch. + + The def-B branch is the live ``node_event_emit_effect`` path, but the + payload_type_match branch performs the same kernel-side + ``model_validate`` and had the same over-unwrap. + """ + captured: dict[str, object] = {} + + class _EnvelopeAgnosticHandler: + async def handle(self, payload: object) -> None: + captured["payload"] = payload + + callback = _make_dispatch_callback( + _EnvelopeAgnosticHandler(), + event_model=ModelHandlerRef( + name="ModelRuntimeEmitRequest", module=_THIS_MODULE + ), + ) + + await callback(_transport_envelope(_emit_request_wire())) + + payload = captured["payload"] + assert isinstance(payload, ModelRuntimeEmitRequest) + assert payload.payload == _user_payload() + + +# --------------------------------------------------------------------------- +# Regression: genuine transport envelopes must still unwrap (OMN-12940) +# --------------------------------------------------------------------------- + + +@pytest.mark.unit +class TestGenuineEnvelopesStillUnwrap: + def test_double_wrapped_reaches_domain_with_a_registered_model(self) -> None: + domain = _domain_command() + double = _transport_envelope(_transport_envelope(domain)) + + assert _extract_dispatch_payload(double, ModelRuntimeDomainCommand) == domain + + def test_triple_wrapped_reaches_domain_with_a_registered_model(self) -> None: + domain = _domain_command() + triple = _transport_envelope(_transport_envelope(_transport_envelope(domain))) + + assert _extract_dispatch_payload(triple, ModelRuntimeDomainCommand) == domain + + def test_double_wrapped_emit_request_unwraps_to_the_request_not_through_it( + self, + ) -> None: + """Both invariants at once: unwrap the two transport layers, stop at the model.""" + request = _emit_request_wire() + double = _transport_envelope(_transport_envelope(request)) + + assert _extract_dispatch_payload(double, ModelRuntimeEmitRequest) == request + + def test_single_wrapped_domain_unchanged_without_a_model(self) -> None: + domain = _domain_command() + + assert _extract_dispatch_payload(_transport_envelope(domain)) == domain + + def test_domain_only_mapping_is_returned_as_is(self) -> None: + domain = _domain_command() + + assert _extract_dispatch_payload(domain, ModelRuntimeDomainCommand) == domain + + def test_payload_field_without_markers_is_never_unwrapped(self) -> None: + domain = {"payload": {"nested": "value"}, "name": "real-domain"} + + assert _extract_dispatch_payload(domain) == domain + + @pytest.mark.asyncio + async def test_double_wrapped_def_b_dispatch_still_reaches_the_domain(self) -> None: + captured: dict[str, object] = {} + + class _DomainHandler: + async def handle(self, request: ModelRuntimeDomainCommand) -> None: + captured["request"] = request + + callback = _make_dispatch_callback(_DomainHandler()) + domain = _domain_command() + + await callback(_transport_envelope(_transport_envelope(domain))) + + request = captured["request"] + assert isinstance(request, ModelRuntimeDomainCommand) + assert request.source_commit_sha == domain["source_commit_sha"] + + +# --------------------------------------------------------------------------- +# The stop predicate itself: fail-closed in both directions +# --------------------------------------------------------------------------- + + +@pytest.mark.unit +class TestRegisteredInputPayloadPredicate: + def test_no_target_model_never_claims(self) -> None: + assert _is_registered_input_payload(_emit_request_wire(), None) is False + + def test_non_mapping_never_claims(self) -> None: + assert _is_registered_input_payload( + "not-a-mapping", ModelRuntimeEmitRequest + ) is (False) + + def test_claims_the_registered_model(self) -> None: + assert ( + _is_registered_input_payload(_emit_request_wire(), ModelRuntimeEmitRequest) + is True + ) + + def test_undeclared_key_defeats_the_claim(self) -> None: + """Key containment: ``source_tool`` is not a ModelEmitRequest field.""" + envelope = _transport_envelope(_emit_request_wire()) + + assert _is_registered_input_payload(envelope, ModelRuntimeEmitRequest) is False + + def test_declared_keys_that_fail_validation_defeat_the_claim(self) -> None: + """Validation is required too — key containment alone is not enough. + + Every key here is a declared ModelEmitRequest field, but ``event_type`` + is absent (required) so this is not the registered model and the unwrap + must continue. + """ + candidate = {"payload": _user_payload(), "correlation_id": "abc"} + + assert _is_registered_input_payload(candidate, ModelRuntimeEmitRequest) is False + + def test_lenient_model_does_not_claim_a_real_envelope(self) -> None: + """Adversarial: an ``extra="ignore"`` model would validate an envelope. + + ``model_validate`` alone would succeed here (extras dropped) and stop the + unwrap one layer high, handing the handler a model whose ``payload`` is + the intermediate envelope. Key containment is what refuses it. + """ + envelope = _transport_envelope(_emit_request_wire()) + # Precondition: validation alone genuinely does NOT discriminate. + assert ModelLenientEnvelopeShaped.model_validate(envelope) is not None + + assert ( + _is_registered_input_payload(envelope, ModelLenientEnvelopeShaped) is False + ) + # The extractor therefore walks PAST the transport layer rather than + # handing the handler a model whose ``payload`` is an envelope. + assert ( + _extract_dispatch_payload(envelope, ModelLenientEnvelopeShaped) != envelope + ) + + def test_alias_declared_keys_are_claimable(self) -> None: + """A model that accepts a wire alias must still be recognised as itself.""" + + class _Aliased(BaseModel): + model_config = ConfigDict(extra="forbid", populate_by_name=True) + + event_type: str + payload: dict[str, object] = Field(default_factory=dict) + partition_key: str | None = Field(default=None, alias="partitionKey") + + candidate = { + "event_type": "session.started", + "payload": _user_payload(), + "partitionKey": "session-1", + } + + assert _is_registered_input_payload(candidate, _Aliased) is True + + def test_custom_validator_rejection_is_not_a_claim(self) -> None: + """A field validator that raises must read as "not the model", never crash. + + ``ModelEmitRequest`` carries exactly such a validator on ``topic`` + (``_topic_must_be_well_formed``), which raises ``ValueError`` rather than + returning a validation verdict. The stop predicate runs on the dispatch + hot path — an escaping exception there would take down every dispatch. + """ + + class _StrictTopic(BaseModel): + model_config = ConfigDict(extra="forbid") + + event_type: str + payload: dict[str, object] = Field(default_factory=dict) + topic: str | None = None + + @field_validator("topic") + @classmethod + def _topic_must_be_well_formed(cls, value: str | None) -> str | None: + if value is not None and not _TOPIC_SHAPE_RE.match(value): + raise ValueError(f"topic override {value!r} is malformed") + return value + + candidate = { + "event_type": "session.started", + "payload": _user_payload(), + "topic": "not-an-onex-topic", + } + + assert _is_registered_input_payload(candidate, _StrictTopic) is False + # ...while a well-formed topic on the same model IS claimed. + assert ( + _is_registered_input_payload( + {**candidate, "topic": "onex.evt.omnimarket.session-started.v1"}, + _StrictTopic, + ) + is True + ) + + def test_alias_choices_declared_keys_are_claimable(self) -> None: + """``AliasChoices`` alternatives are wire keys the model genuinely accepts. + + Collecting only plain-string aliases is fail-OPEN for this defect: the + containment check would reject a candidate that IS the registered model, + the unwrap would continue into the caller's payload, and the OMN-16050 + DLQ failure would come back for every contract aliased this way. + """ + + class _ChoiceAliased(BaseModel): + model_config = ConfigDict(extra="forbid", populate_by_name=True) + + event_type: str + payload: dict[str, object] = Field(default_factory=dict) + correlation_id: str | None = Field( + default=None, + validation_alias=AliasChoices("correlation_id", "correlationId"), + ) + + assert {"correlation_id", "correlationId"} <= _model_declared_wire_keys( + _ChoiceAliased + ) + for spelling in ("correlation_id", "correlationId"): + candidate = { + "event_type": "session.started", + "payload": _user_payload(), + spelling: str(uuid4()), + } + assert _is_registered_input_payload(candidate, _ChoiceAliased) is True + assert _extract_dispatch_payload(candidate, _ChoiceAliased) == candidate + + def test_alias_path_first_segment_is_the_declared_wire_key(self) -> None: + """``AliasPath`` consumes its FIRST segment as the top-level wire key.""" + + class _PathAliased(BaseModel): + model_config = ConfigDict(extra="forbid", populate_by_name=True) + + event_type: str + payload: dict[str, object] = Field(default_factory=dict) + correlation_id: str | None = Field( + default=None, validation_alias=AliasPath("meta", "correlation_id") + ) + + keys = _model_declared_wire_keys(_PathAliased) + assert "meta" in keys + # The inner segment is NOT a top-level key and must not be claimed as one. + assert "correlation_id" in keys # the field name itself, not the path tail + + candidate = { + "event_type": "session.started", + "payload": _user_payload(), + "meta": {"correlation_id": str(uuid4())}, + } + assert _is_registered_input_payload(candidate, _PathAliased) is True + assert _extract_dispatch_payload(candidate, _PathAliased) == candidate + + def test_alias_choices_of_alias_paths_are_flattened(self) -> None: + """Nested ``AliasChoices(AliasPath(...), ...)`` contributes every head key.""" + + class _NestedAliased(BaseModel): + model_config = ConfigDict(extra="forbid", populate_by_name=True) + + event_type: str + payload: dict[str, object] = Field(default_factory=dict) + correlation_id: str | None = Field( + default=None, + validation_alias=AliasChoices( + AliasPath("meta", "correlation_id"), + AliasPath("headers", "cid"), + "correlationId", + ), + ) + + keys = _model_declared_wire_keys(_NestedAliased) + assert {"meta", "headers", "correlationId"} <= keys + + def test_alias_declared_dispatch_still_constructs_the_registered_model( + self, + ) -> None: + """End-to-end: an alias-declared model survives the wired dispatch path.""" + seen: list[object] = [] + + class _AliasedRequest(BaseModel): + model_config = ConfigDict(extra="forbid", populate_by_name=True) + + event_type: str + payload: dict[str, object] = Field(default_factory=dict) + correlation_id: str | None = Field( + default=None, + validation_alias=AliasChoices("correlation_id", "correlationId"), + ) + + class _Handler: + async def handle(self, request: _AliasedRequest) -> dict[str, object]: + seen.append(request) + return {"ok": True} + + wire = { + "event_type": "session.started", + "payload": _user_payload(), + "correlationId": str(uuid4()), + } + payload = _extract_dispatch_payload(wire, _AliasedRequest) + assert payload == wire + built = _AliasedRequest.model_validate(payload) + # The user payload survived intact — it was never unwrapped through. + assert built.payload == _user_payload() + assert _Handler is not None