Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"}
Expand Down
28 changes: 23 additions & 5 deletions src/omnimarket/cli/cli_cloud.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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."
),
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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
Expand All @@ -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,
)
Expand Down
110 changes: 110 additions & 0 deletions src/omnimarket/cloud/completion_bound.py
Original file line number Diff line number Diff line change
@@ -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",
]
25 changes: 22 additions & 3 deletions src/omnimarket/cloud/transport_cloud_delegation.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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,
)

Expand Down
10 changes: 10 additions & 0 deletions src/omnimarket/enums/enum_delegation_failure_class.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"
29 changes: 29 additions & 0 deletions src/omnimarket/nodes/node_delegation_orchestrator/contract.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading
Loading