diff --git a/pyproject.toml b/pyproject.toml index a5c7f79282..7523741027 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "hatchling.build" [project] name = "omnimarket" -version = "0.4.73" +version = "0.4.74" description = "OmniMarket - Portable ONEX workflow package registry" authors = [ {name = "OmniNode.ai", email = "contact@omninode.ai"} diff --git a/src/omnimarket/cli/cli_cloud.py b/src/omnimarket/cli/cli_cloud.py index 10ca92e8a9..697514b759 100644 --- a/src/omnimarket/cli/cli_cloud.py +++ b/src/omnimarket/cli/cli_cloud.py @@ -91,6 +91,7 @@ from omnibase_core.errors.model_onex_error import ModelOnexError from pydantic import SecretStr +from omnimarket.cloud.completion_bound import read_declared_completion_bound from omnimarket.cloud.model_cloud_delegation import ( ModelCloudDelegationAck, ModelCloudDelegationReceipt, @@ -558,10 +559,12 @@ def _write_run_files( @click.option( "--timeout", type=click.IntRange(min=1), - default=300, - show_default=True, + default=None, help=( "Total wall-clock budget for the delegation to reach a terminal state. " + "Defaults to the completion bound the delegation contract declares and " + "the runtime enforces, so the client stops asking at the moment the " + "platform stops trying rather than at a number of its own (OMN-18296). " "Waits forced by the gateway's rate limit are spent from this budget; " "they never end the run early." ), @@ -605,7 +608,7 @@ def cloud_delegate( base_url: str | None, api_key_file: Path | None, onex_home: Path | None, - timeout: int, + timeout: int | None, poll_interval: float, max_poll_interval: float, runner_identity: str | None, @@ -637,11 +640,25 @@ def cloud_delegate( transport_factory = _transport_factory_from_context(ctx) + # OMN-18296: an unset --timeout resolves to the contract-declared bound the + # runtime itself enforces. Passing a value still wins, and the error the + # poll raises names which of the two it spent. + if timeout is None: + declared_bound = read_declared_completion_bound() + budget_seconds = declared_bound.max_wall_seconds + budget_source = ( + "the completion bound declared by node_delegation_orchestrator and " + "enforced by the runtime" + ) + else: + budget_seconds = timeout + budget_source = None + try: with transport_factory( base_url=resolved_base_url, api_key=api_key, - timeout_seconds=float(timeout), + timeout_seconds=float(budget_seconds), ) as client: ack = client.submit( prompt=prompt, task_type=task_type, max_tokens=max_tokens @@ -651,7 +668,8 @@ def cloud_delegate( status = client.poll_until_terminal( workflow_id, - deadline_seconds=float(timeout), + deadline_seconds=float(budget_seconds), + deadline_source=budget_source, interval_seconds=poll_interval, max_interval_seconds=max_poll_interval, ) diff --git a/src/omnimarket/cloud/completion_bound.py b/src/omnimarket/cloud/completion_bound.py new file mode 100644 index 0000000000..ca34676017 --- /dev/null +++ b/src/omnimarket/cloud/completion_bound.py @@ -0,0 +1,110 @@ +# SPDX-FileCopyrightText: 2025 OmniNode.ai Inc. +# SPDX-License-Identifier: MIT +"""Reading the delegation completion bound the platform actually enforces. + +The client used to decide on its own how long to wait: a hardcoded 300 second +CLI default, chosen by nobody in particular, with no relationship to how long +the runtime would keep trying. The runtime's own give-up TTL was 900 seconds, +read from an environment variable in a different repository. So a delegation +that the platform was still willing to work on for another ten minutes was +abandoned by its caller at five, and — worse — a delegation the platform had +quietly stopped working on looked, from the client, exactly the same. + +``node_delegation_orchestrator``'s contract now declares that bound once +(``completion_bound.max_wall_seconds``, OMN-18296). The runtime enforces it by +emitting a real terminal event for any workflow that exceeds it. This module is +the client's half: it reads the same declaration, so ``onex cloud delegate`` +stops asking at the moment the platform stops trying, and the typed error it +raises can say where its patience came from instead of quoting a number of its +own invention. + +Read from the contract that ships in this package, not from the gateway: the +contract is the declaration, and a client that asked the server how long to wait +would be trusting the same surface whose silence it is trying to bound. +""" + +from __future__ import annotations + +from pathlib import Path +from typing import Final + +from omnibase_core.enums import EnumCoreErrorCode +from omnibase_core.models.errors import ModelOnexError +from pydantic import BaseModel, ConfigDict, Field, ValidationError + +_CONTRACT_PATH: Final[Path] = ( + Path(__file__).resolve().parent.parent + / "nodes" + / "node_delegation_orchestrator" + / "contract.yaml" +) + + +class ModelDeclaredCompletionBound(BaseModel): + """The client-side view of the contract's ``completion_bound`` block. + + Deliberately a subset. The runtime's own reader + (``omnibase_infra.runtime.state_io.model_completion_bound``) validates the + whole block including the restart policy and the failure attribution, none + of which a caller can act on. What the client needs is the number it should + wait for and the token it should name when it stops; ``extra="ignore"`` lets + the runtime add fields to the block without breaking every installed CLI. + """ + + model_config = ConfigDict(frozen=True, extra="ignore") + + max_wall_seconds: int = Field(..., gt=0) + failure_class: str = Field(..., min_length=1) + + +def read_declared_completion_bound( + contract_path: Path | None = None, +) -> ModelDeclaredCompletionBound: + """Return the delegation completion bound declared in the node contract. + + Raises: + ModelOnexError: the contract is missing, unreadable, or declares no + usable bound. Deliberately fatal rather than falling back to a + built-in number: a silent fallback is how the client came to be + waiting 300 seconds for a 900 second platform in the first place, + and a wrong bound that looks authoritative is worse than a refusal + that names the file it could not read. + """ + path = contract_path if contract_path is not None else _CONTRACT_PATH + try: + import yaml + + raw = yaml.safe_load(path.read_text()) + except FileNotFoundError as exc: + raise ModelOnexError( + f"the delegation contract is missing at {path} — the completion " + "bound this client waits for is declared there and cannot be " + "guessed. Reinstall the omnimarket package.", + error_code=EnumCoreErrorCode.INVALID_INPUT, + ) from exc + except yaml.YAMLError as exc: + raise ModelOnexError( + f"the delegation contract at {path} could not be parsed: {exc}", + error_code=EnumCoreErrorCode.INVALID_INPUT, + ) from exc + block = raw.get("completion_bound") if isinstance(raw, dict) else None + if not isinstance(block, dict): + raise ModelOnexError( + f"the delegation contract at {path} declares no completion_bound " + "block — this client has no declared bound to wait for and will " + "not invent one.", + error_code=EnumCoreErrorCode.INVALID_INPUT, + ) + try: + return ModelDeclaredCompletionBound.model_validate(block) + except ValidationError as exc: + raise ModelOnexError( + f"the completion_bound block in {path} is invalid: {exc}", + error_code=EnumCoreErrorCode.INVALID_INPUT, + ) from exc + + +__all__: list[str] = [ + "ModelDeclaredCompletionBound", + "read_declared_completion_bound", +] diff --git a/src/omnimarket/cloud/transport_cloud_delegation.py b/src/omnimarket/cloud/transport_cloud_delegation.py index 6fe1a73f1f..37fb59be35 100644 --- a/src/omnimarket/cloud/transport_cloud_delegation.py +++ b/src/omnimarket/cloud/transport_cloud_delegation.py @@ -243,6 +243,7 @@ def poll_until_terminal( workflow_id: str, *, deadline_seconds: float, + deadline_source: str | None = None, interval_seconds: float = DEFAULT_POLL_INTERVAL_SECONDS, max_interval_seconds: float = DEFAULT_MAX_POLL_INTERVAL_SECONDS, sleep_fn: Callable[[float], None] = time.sleep, @@ -272,6 +273,13 @@ def poll_until_terminal( workflow_id: The submitted workflow. deadline_seconds: Total wall-clock budget for reaching a terminal state, measured from the first poll. + deadline_source: Where the budget came from, named verbatim in the + timeout error. A caller that took it from the platform's own + declared completion bound says so, so the customer can tell + "my client gave up early" from "the platform was supposed to + have closed this out by now and did not" — two facts with + different owners that a bare number cannot distinguish + (OMN-18296). interval_seconds: The starting cadence. max_interval_seconds: The ceiling the cadence backs off toward. sleep_fn: Injected so tests do not wait. @@ -326,11 +334,22 @@ def poll_until_terminal( if throttled_polls > 0 else "" ) + # A budget the caller chose and the platform's own declared bound are + # different facts about the same elapsed time, and the customer's next + # move differs: the first says wait longer, the second says the runtime + # owed a terminal and did not deliver one. Say which this was. + budget_note = ( + f" The {deadline_seconds:g}s budget is {deadline_source}, so a " + f"workflow still '{observed}' here should already have been closed " + f"out by the runtime; report it if it stays non-terminal." + if deadline_source + else " Retrieve it later with 'onex cloud receipt " + f"{workflow_id}' — it has NOT failed, it has not finished yet." + ) raise ModelOnexError( f"delegation {workflow_id} was still '{observed}' after " - f"{deadline_seconds:g}s — it has NOT failed, it has not finished " - f"yet. Retrieve it later with 'onex cloud receipt " - f"{workflow_id}'.{throttle_note}", + f"{deadline_seconds:g}s and did not reach a terminal state." + f"{budget_note}{throttle_note}", error_code=EnumCoreErrorCode.TIMEOUT_EXCEEDED, ) diff --git a/src/omnimarket/enums/enum_delegation_failure_class.py b/src/omnimarket/enums/enum_delegation_failure_class.py index 71f77e546e..3b58599357 100644 --- a/src/omnimarket/enums/enum_delegation_failure_class.py +++ b/src/omnimarket/enums/enum_delegation_failure_class.py @@ -28,4 +28,14 @@ class EnumDelegationFailureClass(StrEnum): # to a model that is not running (e.g. SGLang echoing an unknown model # string back at HTTP 200). MODEL_ATTRIBUTION_MISMATCH = "model_attribution_mismatch" + # OMN-18296: the process that owned an in-flight leg went away and the leg + # was lost with it. Distinct from TIMEOUT, which is a provider that took too + # long to answer a call we can still see: here no call is outstanding at all + # — the inference command's consumer offset was already committed, so it is + # never redelivered, and no response event, success or failure, will ever + # arrive. Measured on the lab lane 2026-09-13: the effects pod was recreated + # at 09:57:40Z with correlation a2fe0848 in flight since 09:55:04Z; the + # request sat on the bus with no matching response, the FSM row stayed + # ROUTED/in_flight, and the customer's delegation never terminalised. + RUNTIME_RESTART_DURING_DELEGATION = "runtime_restart_during_delegation" UNKNOWN = "unknown" diff --git a/src/omnimarket/nodes/node_delegation_orchestrator/contract.yaml b/src/omnimarket/nodes/node_delegation_orchestrator/contract.yaml index 871d4c03f2..ea6d3d29a1 100644 --- a/src/omnimarket/nodes/node_delegation_orchestrator/contract.yaml +++ b/src/omnimarket/nodes/node_delegation_orchestrator/contract.yaml @@ -43,6 +43,35 @@ state_io: codec: module: "omnimarket.nodes.node_delegation_orchestrator.state_codec" name: "StateIoCodec" +# OMN-18296: how long a delegation may stay non-terminal, and what the runtime +# owes the caller when the process that owned an in-flight leg goes away. +# +# Measured, not assumed. On 2026-09-13 the lab lane's runtime-effects pod was +# recreated at 09:57:40Z while correlation a2fe0848-4b4b-462e-b633-c5f9559afee5 +# (submitted 09:55:04Z) had an inference command in flight. The command's +# consumer offset was already committed, so it was never redelivered; no +# inference response, success or failure, was ever published; the FSM row stayed +# ROUTED with in_flight=TRUE; and because a gateway workflow can only leave +# 'published' when a REAL terminal event is consumed off the bus, the customer's +# delegation never terminalised at all. The runtime's give-up sweep existed but +# was row-only and traffic-gated: it wrote FAILED to the row and told nobody, and +# on an idle lane it never ran a second time. +# +# The bound is declared HERE, once, because two numbers used to govern that run +# and neither was a contract: the runtime's TTL was an env-var default in +# omnibase_infra and the CLI's patience was a hardcoded 300s that did not know +# about it. The runtime enforces this field; `onex cloud delegate` reads the same +# field for its own deadline, so the client stops asking at the moment the +# platform stops trying rather than at an unrelated number of its own. +# +# 900s is the widest real delegation observed on the lab (escalation across +# tiers with compliance repair attempts), rounded up — it bounds abandonment, +# it is not a per-model latency budget. +completion_bound: + max_wall_seconds: 900 + on_runtime_restart: "terminalise_failed" + failure_class: "runtime_restart_during_delegation" + failure_code: "ONEX_MARKET_DELEGATION_RUNTIME_RESTART" # OMN-13629 (WS-F Phase 1): the legacy compat task-delegated.v1 co-writer was # removed. A delegation terminal now emits a single canonical event on # delegation-{completed,failed}.v1. The delegation + savings projections consume diff --git a/src/omnimarket/nodes/node_delegation_orchestrator/state_codec.py b/src/omnimarket/nodes/node_delegation_orchestrator/state_codec.py index 42d7c58287..6679d35179 100644 --- a/src/omnimarket/nodes/node_delegation_orchestrator/state_codec.py +++ b/src/omnimarket/nodes/node_delegation_orchestrator/state_codec.py @@ -43,11 +43,12 @@ from __future__ import annotations import json +import time from collections.abc import Iterator, MutableMapping from typing import cast from uuid import UUID -from pydantic import TypeAdapter +from pydantic import TypeAdapter, ValidationError from omnimarket.nodes.node_delegation_orchestrator.handlers.handler_delegation_workflow import ( DelegationWorkflowState, @@ -67,6 +68,31 @@ # TypeAdapter has no field of that name to accept it). _IN_FLIGHT_KEY = "in_flight" +# OMN-18296: where infra imports the terminal class from, and what it is called. +# Named here rather than passed down from the contract because the CLASS is the +# business shape this codec owns; the contract already owns the TOPIC that class +# resolves to (``published_events``: DelegationFailed -> +# onex.evt.omnibase-infra.delegation-failed.v1). +_TERMINAL_MODULE = "omnibase_core.models.delegation.wire.model_delegation_failed" +_FAILED_TERMINAL_CLASS = "ModelDelegationFailed" +# The word the gateway already renders for a terminal with no serving model. +_NO_MODEL_SERVED = "none" + + +def _terminal_failure_reason(failure_class: str, failure_code: str | None) -> str: + """Render the machine-readable half of the terminal's failure attribution. + + The gateway parses this field with a fixed grammar — + ``(Error|Exception)[: ONEX_CODE]`` — and reports anything else + as carrying no class at all, which is how a typed refusal reads to a customer + as an unexplained failure. Composed from the contract's own vocabulary token + so the two cannot drift: ``runtime_restart_during_delegation`` renders as + ``RuntimeRestartDuringDelegationError``. + """ + camel = "".join(part.capitalize() for part in failure_class.split("_")) + rendered = f"{camel}Error" + return f"{rendered}: {failure_code}" if failure_code else rendered + def encode(state: DelegationWorkflowState) -> bytes: """Serialize workflow state to JSON bytes for durable storage. @@ -321,6 +347,90 @@ def encode(self, state: DelegationWorkflowState) -> bytes: def decode(self, raw: bytes | str) -> DelegationWorkflowState: return decode(raw) + def build_abandoned_terminal( + self, + *, + correlation_id: str, + tenant_id: str, + state: str, + payload_json: str, + failure_class: str, + failure_code: str | None, + max_wall_seconds: int, + ) -> tuple[str, str, dict[str, object]] | None: + """Build the terminal FAILURE event for a row past its completion bound. + + Called by omnibase_infra's state_io wiring (OMN-18296) for a row that is + still ``in_flight``, still non-terminal, carries no re-publishable outbox + batch, and has not advanced within the contract-declared + ``completion_bound.max_wall_seconds``. Infra owns the bus, the envelope + id and the topic; this method owns the only thing infra cannot know — the + business shape of THIS node's terminal. + + Returns ``(module, class_name, payload)`` for ``ModelDelegationFailed``, + or ``None`` when the row's payload cannot be decoded into workflow state + at all (a shape this build does not understand is left alone rather than + closed out on a guess, and the sweep retries it after the next deploy). + + The values are deliberately the honest ones rather than the flattering + ones. There is no content, so ``content`` is empty; there is no verdict, + so ``quality_passed`` is false and ``quality_score`` is 0.0; there is no + answer from any provider, so no tokens are claimed. ``latency_ms`` is the + real elapsed wall time from the request's own start epoch, because "how + long the customer waited" is a true and useful fact even when nothing + came back. ``model_used`` reports ``"none"`` where routing never + selected one — the same word the gateway already renders for a terminal + with no serving model — rather than inventing an attribution for a call + that was never made. + """ + try: + workflow = decode(payload_json) + except (ValueError, ValidationError): + return None + routing = workflow.routing_decision + started_at_ns = workflow.started_at_ns + latency_ms = ( + max(0, (time.time_ns() - started_at_ns) // 1_000_000) + if started_at_ns + else 0 + ) + reason = ( + f"delegation was abandoned in state {state} after the runtime " + f"process holding its in-flight leg went away; no response to that " + f"leg was ever published and its command offset was already " + f"committed, so it is never redelivered. Terminalised by the " + f"contract-declared completion bound of {max_wall_seconds}s. " + f"Resubmit the request." + ) + payload: dict[str, object] = { + "correlation_id": correlation_id, + # ``request`` is Optional on the state dataclass: a row seeded before + # the request was bound carries None. Empty is the honest reading — + # a task type was never recorded, and the wire field is required. + "task_type": ( + workflow.request.task_type if workflow.request is not None else "" + ), + "model_used": ( + routing.selected_model if routing is not None else _NO_MODEL_SERVED + ), + "endpoint_url": routing.endpoint_url if routing is not None else "", + "content": "", + "quality_passed": False, + "quality_score": 0.0, + "latency_ms": latency_ms, + "fallback_to_claude": False, + "failure_reason": reason, + "terminal_failure_reason": _terminal_failure_reason( + failure_class, failure_code + ), + "escalation_count": workflow.escalation_count, + "compliance_attempts": workflow.compliance_attempts, + "context_pack_hash": workflow.context_pack_hash, + "cost_tier_name": workflow.current_tier_name or "", + "tenant_id": tenant_id, + } + return (_TERMINAL_MODULE, _FAILED_TERMINAL_CLASS, payload) + def flush(self, cid: str) -> str | None: """Re-encode the proxy's cached entry for ``cid``, if touched this dispatch. diff --git a/tests/unit/cli/test_cli_cloud.py b/tests/unit/cli/test_cli_cloud.py index e69d3df6b4..699314ebbe 100644 --- a/tests/unit/cli/test_cli_cloud.py +++ b/tests/unit/cli/test_cli_cloud.py @@ -27,6 +27,7 @@ from omnibase_core.errors.model_onex_error import ModelOnexError from omnimarket.cli.cli_cloud import cloud_group +from omnimarket.cloud.completion_bound import read_declared_completion_bound from omnimarket.cloud.model_cloud_delegation import ( ModelCloudDelegationAck, ModelCloudDelegationReceipt, @@ -188,8 +189,12 @@ def poll_until_terminal( deadline_seconds: float, interval_seconds: float, max_interval_seconds: float, + # OMN-18296: names where the deadline came from, so the timeout error can + # distinguish a caller-chosen budget from the platform's declared bound. + deadline_source: str | None = None, ) -> ModelCloudDelegationStatus: self.poll_schedule = (deadline_seconds, interval_seconds, max_interval_seconds) + self.poll_deadline_source = deadline_source return _status(self._terminal_status).model_copy( update={ "workflow_id": self._status_workflow_id or uuid.UUID(workflow_id), @@ -241,6 +246,44 @@ def _logged_in(tmp_path: Path) -> Path: # --------------------------------------------------------------------------- +def test_an_unset_timeout_waits_for_the_contract_declared_bound( + tmp_path: Path, +) -> None: + """OMN-18296: the client's patience comes from the platform, not from itself. + + Before this, ``--timeout`` defaulted to a hardcoded 300 seconds with no + relationship to the 900 second bound the runtime was working to, so a caller + could abandon a delegation the platform was still willing to finish — and + could not tell that from one the platform had silently stopped working on. + """ + home = _logged_in(tmp_path) + factory, made = _factory() + + result = CliRunner().invoke( + cloud_group, + [ + "delegate", + "Summarize what a delegation receipt proves.", + "--task-type", + "summarization", + "--output-dir", + str(tmp_path / "runs"), + "--onex-home", + str(home), + ], + obj={"transport_factory": factory}, + ) + + assert result.exit_code == 0, result.output + declared = read_declared_completion_bound().max_wall_seconds + assert made[0].poll_schedule is not None + assert made[0].poll_schedule[0] == float(declared) + assert made[0].poll_deadline_source is not None, ( + "a budget taken from the platform's declared bound must say so, so the " + "timeout error can name whose bound was spent" + ) + + def test_delegate_prints_the_result_and_saves_it_to_disk(tmp_path: Path) -> None: """AC1 + AC3 + AC4 in one path: output printed, files written, receipt saved.""" home = _logged_in(tmp_path) diff --git a/tests/unit/delegation/test_omn18296_completion_bound.py b/tests/unit/delegation/test_omn18296_completion_bound.py new file mode 100644 index 0000000000..7436d37ab6 --- /dev/null +++ b/tests/unit/delegation/test_omn18296_completion_bound.py @@ -0,0 +1,256 @@ +# SPDX-FileCopyrightText: 2025 OmniNode.ai Inc. +# SPDX-License-Identifier: MIT +"""The declared completion bound, its terminal, and the client that reads it (OMN-18296). + +Cloud delegation ``16eafedc-199c-44c2-a2cb-b9e535839bb9`` was submitted to the +lab lane at 2026-09-13T09:55:04Z. At 09:57:40Z the ``omninode-runtime-effects`` +pod was recreated by a concurrent sanctioned re-apply, with that delegation's +inference command in flight. The command's consumer offset was already committed +(``TOTAL-LAG 0``, current offset 12 of 12), so the record was never redelivered +and no inference response — success or failure — was ever published: the bus +carried 12 requests and 11 responses, and the missing one is this correlation. +The FSM row stayed ``ROUTED`` with ``in_flight = TRUE``, and the gateway row +stayed ``published`` with ``completed_at NULL``, indefinitely. + +Nothing in the system bounded that. The runtime's give-up TTL was an +environment-variable default in another repository and wrote only to the row; +the client's patience was a hardcoded 300 seconds unrelated to it. This suite +pins the three halves of the fix: the bound is declared in the contract, the +terminal it produces is typed and legible to the gateway, and the client waits +for the declared number rather than one of its own. +""" + +from __future__ import annotations + +import json +import uuid +from pathlib import Path +from typing import Any + +import httpx +import pytest +import yaml +from omnibase_core.enums.enum_core_error_code import EnumCoreErrorCode +from omnibase_core.errors.model_onex_error import ModelOnexError +from omnibase_core.models.delegation.wire.model_delegation_failed import ( + ModelDelegationFailed, +) +from pydantic import SecretStr + +from omnimarket.cli.cli_cloud import cloud_delegate +from omnimarket.cloud.completion_bound import read_declared_completion_bound +from omnimarket.cloud.transport_cloud_delegation import ( + CLOUD_DELEGATION_WORKFLOW_TYPE, + TransportCloudDelegation, +) +from omnimarket.enums.enum_delegation_failure_class import EnumDelegationFailureClass +from omnimarket.nodes.node_delegation_orchestrator.state_codec import StateIoCodec + +_CONTRACT = ( + Path(__file__).resolve().parents[3] + / "src" + / "omnimarket" + / "nodes" + / "node_delegation_orchestrator" + / "contract.yaml" +) +_WORKFLOW_ID = "16eafedc-199c-44c2-a2cb-b9e535839bb9" +_CORRELATION = "a2fe0848-4b4b-462e-b633-c5f9559afee5" +_TENANT = "cacacbb1-0e64-4521-9712-ed02ee799907" + +# The gateway's own attribution grammar +# (omninode_infra docker/onex-api/workflow_failure_attribution.py). Reproduced +# here rather than imported because that repository is not a dependency of this +# one; a terminal whose attribution does not match it is reported downstream as +# carrying no failure class at all, which is how a typed refusal reaches a +# customer as an unexplained failure. +_ATTRIBUTION_GRAMMAR = ( + r"^(?P[A-Z][A-Za-z0-9_]*(?:Error|Exception))" + r"(?::[ ](?PONEX_[A-Z0-9_]+))?$" +) + + +def _abandoned_payload() -> str: + """A ``delegation_workflow_state`` payload in the shape the live row had.""" + return json.dumps( + { + "state": "ROUTED", + "in_flight": True, + "tenant_id": _TENANT, + "correlation_id": _CORRELATION, + "started_at_ns": 1789293304198594138, + "escalation_count": 0, + "compliance_attempts": 1, + "context_pack_hash": "", + "current_tier_name": "tenant_overlay", + "request": { + "prompt": "summarise this", + "task_type": "summarization", + "tenant_id": _TENANT, + "correlation_id": _CORRELATION, + "emitted_at": "2026-09-13T09:55:04.089592Z", + }, + "routing_decision": { + "cost_tier": "tenant_byok", + "rationale": "tenant overlay", + "task_type": "summarization", + "tier_name": "tenant_overlay", + "max_tokens": 65536, + "max_context_tokens": 8192, + "api_key_ref": "cred_x", + "endpoint_url": "https://openrouter.ai/api/v1/chat/completions", + "correlation_id": _CORRELATION, + "selected_model": "nvidia/nemotron-3-ultra-550b-a55b:free", + "selected_backend_id": "a83745fa-d9ee-5239-a363-f9349f3bff08", + "selected_backend_ref": "byok-openrouter", + "system_prompt": "You are a summarization assistant.", + }, + } + ) + + +@pytest.mark.unit +def test_the_contract_declares_the_bound_the_runtime_enforces() -> None: + """AC2: the bound is a typed contract field, not a number in one script.""" + block = yaml.safe_load(_CONTRACT.read_text())["completion_bound"] + assert block["max_wall_seconds"] > 0 + assert block["on_runtime_restart"] == "terminalise_failed" + # The declared class must be a real member of this node's failure + # vocabulary, or the terminal names a class nothing else in the system + # recognises. + assert ( + EnumDelegationFailureClass(block["failure_class"]) + is EnumDelegationFailureClass.RUNTIME_RESTART_DURING_DELEGATION + ) + + +@pytest.mark.unit +def test_the_codec_builds_a_typed_restart_terminal_from_an_abandoned_row() -> None: + """AC3: the give-up is a real, typed, gateway-legible terminal event.""" + import re + + built = StateIoCodec().build_abandoned_terminal( + correlation_id=_CORRELATION, + tenant_id=_TENANT, + state="ROUTED", + payload_json=_abandoned_payload(), + failure_class=EnumDelegationFailureClass.RUNTIME_RESTART_DURING_DELEGATION.value, + failure_code="ONEX_MARKET_DELEGATION_RUNTIME_RESTART", + max_wall_seconds=900, + ) + assert built is not None + module, class_name, payload = built + assert class_name == "ModelDelegationFailed" + assert module.endswith("model_delegation_failed") + + terminal = ModelDelegationFailed.model_validate(payload) + assert terminal.quality_passed is False + assert terminal.tenant_id == _TENANT + assert terminal.task_type == "summarization" + # Attribution the gateway can actually parse into a class and a code. + assert terminal.terminal_failure_reason is not None + match = re.match(_ATTRIBUTION_GRAMMAR, terminal.terminal_failure_reason) + assert match is not None, terminal.terminal_failure_reason + assert match.group("failure_class") == "RuntimeRestartDuringDelegationError" + assert match.group("failure_code") == "ONEX_MARKET_DELEGATION_RUNTIME_RESTART" + # Nothing was served, so nothing is claimed. + assert terminal.content == "" + assert terminal.total_tokens == 0 + assert "900s" in terminal.failure_reason + + +@pytest.mark.unit +def test_an_undecodable_row_is_left_alone_rather_than_closed_on_a_guess() -> None: + """A payload shape this build does not understand is not terminalised.""" + assert ( + StateIoCodec().build_abandoned_terminal( + correlation_id=_CORRELATION, + tenant_id=_TENANT, + state="ROUTED", + payload_json='{"not": "a workflow state"}', + failure_class="runtime_restart_during_delegation", + failure_code=None, + max_wall_seconds=900, + ) + is None + ) + + +@pytest.mark.unit +def test_the_client_waits_for_the_declared_bound_not_a_number_of_its_own() -> None: + """AC4: ``--timeout`` defaults to the contract's bound, not a CLI constant.""" + declared = read_declared_completion_bound() + contract_value = yaml.safe_load(_CONTRACT.read_text())["completion_bound"] + assert declared.max_wall_seconds == contract_value["max_wall_seconds"] + + timeout_option = next( + param for param in cloud_delegate.params if param.name == "timeout" + ) + assert timeout_option.default is None, ( + "a hardcoded default here is a second, competing bound — the platform's " + "declared one is the only one a caller should be waiting for" + ) + + +@pytest.mark.unit +def test_the_timeout_error_says_whose_bound_was_spent() -> None: + """AC4: still non-terminal AT the platform's own bound is a different fact. + + A caller-chosen budget running out means "wait longer". The platform's own + declared bound running out means the runtime owed a terminal and did not + deliver one. The typed error must distinguish them, because the customer's + next move differs. + """ + + class _Clock: + def __init__(self) -> None: + self.now = 0.0 + + def monotonic(self) -> float: + return self.now + + def sleep(self, seconds: float) -> None: + self.now += seconds + + def handler(request: httpx.Request) -> httpx.Response: + return httpx.Response( + 200, + json={ + "workflow_id": _WORKFLOW_ID, + "workflow_type": CLOUD_DELEGATION_WORKFLOW_TYPE, + "status": "published", + "envelope_id": str(uuid.uuid4()), + "correlation_id": _CORRELATION, + "command_topic": "onex.cmd.delegation.inference.v1", + "submitted_at": "2026-09-13T09:55:04Z", + "updated_at": "2026-09-13T09:55:04Z", + }, + ) + + def _poll(**extra: Any) -> ModelOnexError: + clock = _Clock() + client = TransportCloudDelegation( + base_url="https://dev.api.omninode.ai", + api_key=SecretStr("onxk_testkey"), + http_client=httpx.Client(transport=httpx.MockTransport(handler)), # type: ignore[arg-type] + ) + with pytest.raises(ModelOnexError) as excinfo: + client.poll_until_terminal( + _WORKFLOW_ID, + deadline_seconds=5.0, + sleep_fn=clock.sleep, + monotonic_fn=clock.monotonic, + **extra, + ) + return excinfo.value + + declared_source = "the completion bound declared by node_delegation_orchestrator" + with_source = _poll(deadline_source=declared_source) + assert with_source.error_code == EnumCoreErrorCode.TIMEOUT_EXCEEDED + assert declared_source in str(with_source) + assert "should already have been closed out by the runtime" in str(with_source) + + # A caller-chosen budget keeps the old, correct reading. + caller_chosen = _poll() + assert "has NOT failed" in str(caller_chosen) + assert declared_source not in str(caller_chosen) diff --git a/tests/unit/models/delegation/llm_cost_routing/test_llm_cost_routing_models.py b/tests/unit/models/delegation/llm_cost_routing/test_llm_cost_routing_models.py index fc963f2c47..505a436678 100644 --- a/tests/unit/models/delegation/llm_cost_routing/test_llm_cost_routing_models.py +++ b/tests/unit/models/delegation/llm_cost_routing/test_llm_cost_routing_models.py @@ -43,10 +43,16 @@ def test_all_values_present(self) -> None: assert "unknown" in values # OMN-16419: the fail-closed model-attribution guard's failure class. assert "model_attribution_mismatch" in values + # OMN-18296: the owning runtime process went away mid-leg and the + # in-flight call was lost with it. Not a TIMEOUT: no call is + # outstanding, because its command offset was already committed and it + # is never redelivered. + assert "runtime_restart_during_delegation" in values - def test_exactly_nine_values(self) -> None: + def test_exactly_eleven_values(self) -> None: # OMN-16419: was 9 — model_attribution_mismatch added. - assert len(EnumDelegationFailureClass) == 10 + # OMN-18296: was 10 — runtime_restart_during_delegation added. + assert len(EnumDelegationFailureClass) == 11 def test_is_str_enum(self) -> None: assert isinstance(EnumDelegationFailureClass.TIMEOUT, str) diff --git a/uv.lock b/uv.lock index 28283a30fc..1e71afec7d 100644 --- a/uv.lock +++ b/uv.lock @@ -2622,7 +2622,7 @@ wheels = [ [[package]] name = "omnimarket" -version = "0.4.73" +version = "0.4.74" source = { editable = "." } dependencies = [ { name = "aiokafka" },