Skip to content
Open
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
13 changes: 12 additions & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,13 @@ the fuller account of each version, including verification notes.
with `", "`), `NaN`/`Infinity` — reports as `None` rather than passing through to break the
`Decimal()` parse the field documents. The field defaults to `None`, so this stays additive
for anything that constructs a `RouterRunResult` by hand.
The vendored Router spec now declares `X-Comfy-Credits-Used` on the run route's `200`, so
`tests/test_router_spec_contract.py` pins this lift against the contract like the other four.
- `QueueBacklogFull` in `comfy_sdk.router_exceptions`, for the Router bucket `queue_backlog_full`:
a queued `submit` refused `429` because the caller already has too many requests waiting. It is
not `ConcurrencyLimitExceeded` — the queue parks a submit at the in-flight limit, and this is the
separate bound on how many may be left waiting. It clears as the caller's own queued requests
finish. Before this, the bucket arrived as a bare `RouterError`.

### Fixed

Expand Down Expand Up @@ -73,8 +80,12 @@ the fuller account of each version, including verification notes.
affected: `raise`, `except` and every attribute a caller reads inside the handler (`.message`,
`.code`, `.http_status`, `.details`, `.request_id`, `.retry_after`) are unchanged.
- `RouterError` is exported from the package root, alongside `CancelRefused` and
`AlreadyCompleted`. The eighteen per-bucket classes still live in
`AlreadyCompleted`. The nineteen per-bucket classes still live in
`comfy_sdk.router_exceptions`.
- `NotEnabled`'s documented meaning widened with the synced Router spec: on a queued `submit` it
can also refuse a *model* whose partner answers a generation directly as bytes (it cannot yet be
queued; nothing is queued or charged; `models.run` serves it). It is still terminal, but on a
submit it no longer proves the caller is not switched on — read `.detail`.
- `ApiError.error_type` records the Router bucket a response named (`X-Comfy-Error-Type`, or the
body's `error_type`), or `None` when it named none — which is also how the SDK tells which
surface answered.
Expand Down
16 changes: 13 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -604,7 +604,16 @@ async with AsyncComfy(api_key="comfyui-...") as client:
This surface is **gated server side**. A caller the queue is not switched on
for is answered `403 not_enabled`, which arrives as
`comfy_sdk.router_exceptions.NotEnabled` — nothing about the request is wrong,
and it is terminal: do not retry it.
and it is terminal: do not retry it. The same `403 not_enabled` also refuses a
*model* whose partner answers a generation directly as bytes: it cannot yet be
queued, nothing is queued or charged, and `client.models.run` serves it instead —
so read `.detail` before concluding the account is not switched on.

A caller with too many queued requests already waiting is refused
`429 queue_backlog_full`, raised as `comfy_sdk.router_exceptions.QueueBacklogFull`.
It is not `ConcurrencyLimitExceeded` (the synchronous route's in-flight bound):
the queue parks a submit at that limit, and this is the separate bound on how many
may be left waiting. It clears as your own queued requests finish.

### Retrying a run without paying for it twice

Expand Down Expand Up @@ -879,12 +888,13 @@ status on `.http_status`; keep an `except ComfyError` outside the clause above
if you need to handle those in the same place.

`RouterError` is exported from the package root because it is the handler most
callers write first. The eighteen per-bucket classes stay in
callers write first. The nineteen per-bucket classes stay in
`comfy_sdk.router_exceptions` — `InvalidInput`, `ContentPolicyViolation`,
`ProviderError`, `ProviderTimeout`, `InsufficientCredits`, `ModelNotFound`,
`Unauthorized`, `Forbidden`, `ConcurrencyLimitExceeded`, `ClientDisconnected`,
`InternalError`, `DeadlineExceeded`, `NotEnabled`, `ServiceUnavailable`,
`RateLimited`, `Cancelled`, `QueueTimeout`, `RequestNotFound` — one import path
`RateLimited`, `Cancelled`, `QueueTimeout`, `RequestNotFound`,
`QueueBacklogFull` — one import path
for the whole set rather than half of it here and half of it there. A bucket added to Router after your installed version
arrives as `RouterError` itself, with the raw value readable on `.error_type`.

Expand Down
98 changes: 83 additions & 15 deletions spec/router-openapi.yaml

Large diffs are not rendered by default.

8 changes: 8 additions & 0 deletions src/comfy_sdk/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -723,6 +723,14 @@ def submit(
for is answered ``403`` ``not_enabled``, which arrives here as
:class:`~comfy_sdk.router_exceptions.NotEnabled`. Nothing about the
request is wrong in that case, and it is terminal — do not retry it.
The same ``not_enabled`` also refuses a *model* whose partner answers a
generation directly as bytes: such a model cannot yet be queued, nothing
is queued or charged, and :meth:`run` serves it instead — so read
``detail`` before concluding the caller is not switched on. A caller
with too many requests already waiting is refused ``429``
``queue_backlog_full``
(:class:`~comfy_sdk.router_exceptions.QueueBacklogFull`), which clears
as its own queued requests finish.
"""
low = cast(ComfyLow, self._low)
key = (
Expand Down
29 changes: 28 additions & 1 deletion src/comfy_sdk/router_exceptions.py
Original file line number Diff line number Diff line change
Expand Up @@ -407,10 +407,18 @@ class NotEnabled(RouterError):
is *not* the same thing, because ``forbidden`` is an entitlement decision
about the caller while this is a state of the rollout. It is **terminal**:
do not retry, and do not treat it as an outage.

The one exception to "about the caller" is the queued submit
(:meth:`~comfy_sdk.models.Models.submit`), which also answers
``not_enabled`` for a *model* whose partner answers a generation directly
as bytes: that model cannot yet be queued, so it is the model and not the
caller that is refused, nothing is queued or charged, and the synchronous
route (:meth:`~comfy_sdk.models.Models.run`) runs it instead. On a submit,
read ``detail`` before concluding the account is not switched on.
"""

error_type = "not_enabled"
_spec_meaning_digest: str = "c4a48688282c"
_spec_meaning_digest: str = "571a30cc6ba0"


class ServiceUnavailable(RouterError):
Expand Down Expand Up @@ -522,6 +530,23 @@ class RequestNotFound(RouterError):
_spec_meaning_digest: str = "385112b3cdcf"


class QueueBacklogFull(RouterError):
"""The caller already has too many queued requests waiting to run, so this
submit was refused.

It shares ``429`` with :class:`ConcurrencyLimitExceeded` and is not the same
thing: that one is the synchronous route's answer for too many calls in
flight at once, whereas the queue accepts a submit at that limit and parks
it, and this bucket is the separate bound on how many a caller may leave
waiting -- so that parking cannot mean enqueuing without end. It clears as
the caller's own queued requests finish, so retry once some of them
complete.
"""

error_type = "queue_backlog_full"
_spec_meaning_digest: str = "50745ff63044"


# -- cancel refusals ---------------------------------------------------------
#
# Deliberately OUTSIDE the closed set below, and carrying no
Expand Down Expand Up @@ -613,6 +638,7 @@ class AlreadyCompleted(CancelRefused):
Cancelled,
QueueTimeout,
RequestNotFound,
QueueBacklogFull,
)

_BY_ERROR_TYPE: dict[str, type[RouterError]] = {cls.error_type: cls for cls in ROUTER_EXCEPTIONS}
Expand Down Expand Up @@ -886,6 +912,7 @@ def _detail_from(entry: Mapping[str, Any]) -> ValidationErrorDetail:
"NotEnabled",
"ProviderError",
"ProviderTimeout",
"QueueBacklogFull",
"QueueTimeout",
"RateLimited",
"RequestNotFound",
Expand Down
1 change: 1 addition & 0 deletions tests/test_exception_modules.py
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,7 @@
(409, "cancelled"),
(504, "queue_timeout"),
(404, "request_not_found"),
(429, "queue_backlog_full"),
(418, "something_invented_later"),
]

Expand Down
21 changes: 20 additions & 1 deletion tests/test_router_exceptions.py
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@
NotEnabled,
ProviderError,
ProviderTimeout,
QueueBacklogFull,
QueueTimeout,
RateLimited,
RequestNotFound,
Expand All @@ -48,7 +49,8 @@
# The statuses are part of the case on purpose: they document the pairing a
# retry policy keys on (502 provider_error vs 504 provider_timeout), and the
# four status collisions the widened set introduced -- 403 forbidden vs
# not_enabled, 429 concurrency_limit_exceeded vs rate_limited, 504
# not_enabled, 429 concurrency_limit_exceeded vs rate_limited (and, since,
# queue_backlog_full), 504
# provider_timeout vs deadline_exceeded, 500 internal_error vs the 503
# service_unavailable it is deliberately NOT merged with.
CASES: list[tuple[str, int, type[RouterError]]] = [
Expand All @@ -70,6 +72,7 @@
("cancelled", 409, Cancelled),
("queue_timeout", 504, QueueTimeout),
("request_not_found", 404, RequestNotFound),
("queue_backlog_full", 429, QueueBacklogFull),
]

# Deliberately not in the set this SDK version knows: a later milestone adds it,
Expand Down Expand Up @@ -405,6 +408,22 @@ def test_a_rate_limited_429_is_not_the_concurrency_429() -> None:
assert type(error_from_response(429, {}, None)) is ConcurrencyLimitExceeded


def test_a_queue_backlog_full_429_is_not_the_concurrency_429() -> None:
# The queue parks a submit at the in-flight limit; this is the separate
# bound on how many a caller may leave waiting, and it clears only as the
# caller's own queued requests finish -- not when one sync call returns.
exc = error_from_response(
429,
{ERROR_TYPE_HEADER: "queue_backlog_full"},
{"detail": "Too many queued requests.", "error_type": "queue_backlog_full"},
)
assert type(exc) is QueueBacklogFull
assert not isinstance(exc, ConcurrencyLimitExceeded)
# A bare 429 from an intermediary still reads as plain throttling: the
# backlog bound is Router's own claim, which a proxy cannot be making.
assert type(error_from_response(429, {}, None)) is ConcurrencyLimitExceeded


# -- never crash the client on a malformed response --------------------------


Expand Down
40 changes: 25 additions & 15 deletions tests/test_router_spec_contract.py
Original file line number Diff line number Diff line change
Expand Up @@ -294,11 +294,22 @@ def test_the_bound_path_has_exactly_the_two_segments_the_binding_fills() -> None
"dropped_params": "X-Comfy-Router-Dropped-Params",
"replayed": "Idempotent-Replayed",
"request_id": "X-Comfy-Request-Id",
"credits_used": "X-Comfy-Credits-Used",
}

#: A value the lift will actually keep, for the lifts that validate what they
#: read. ``credits_used`` reports a non-decimal as ``None`` -- the same as an
#: absent header -- so probing it with an arbitrary string would make a correct
#: lift look like one reading the wrong name. ``12.5`` is the spec's own
#: example for the header. Anything not listed is probed with ``"x"``.
_LIFT_PROBE_VALUES = {"X-Comfy-Credits-Used": "12.5"}

#: Lifted by the SDK but NOT declared on the contract's 200 -- see the tripwire
#: test at the bottom of this file.
_UNDECLARED_HEADER_LIFTS = {"credits_used": "X-Comfy-Credits-Used"}
#: test at the bottom of this file. Empty today: ``credits_used`` sat here until
#: a spec sync declared ``X-Comfy-Credits-Used``, and was moved up into
#: ``_CONTRACT_HEADER_LIFTS`` then. Kept, rather than deleted with its test, so
#: the next lift the SDK reads ahead of the contract has somewhere to go.
_UNDECLARED_HEADER_LIFTS: dict[str, str] = {}


def _declared_run_response_headers() -> set[str]:
Expand Down Expand Up @@ -333,7 +344,7 @@ def test_the_lift_actually_reads_the_declared_name(field: str, header: str) -> N
to fail.
"""
absent = getattr(_run_result({}, {}), field)
present = getattr(_run_result({}, {header: "x"}), field)
present = getattr(_run_result({}, {header: _LIFT_PROBE_VALUES.get(header, "x")}), field)
assert present != absent, (
f"_run_result ignored {header!r}: RouterRunResult.{field} read {absent!r} both with "
f"the header and without it, so the lift is reading some other name."
Expand All @@ -346,18 +357,17 @@ def test_an_undeclared_lift_stays_undeclared_until_someone_reconciles_it(
) -> None:
"""Tripwire, and deliberately asserting the *absence*.

``credits_used`` is lifted from a header the vendored contract does not
declare anywhere -- the 200's only cost headers are the
``X-Committed-Spend-*`` trio, which is a different quantity (USD cents of
in-flight commitment, not the price of this run). Nothing in the suite can
catch a wrong name here, because every test configures its stub to emit the
exact literal the lift reads.

That gap is tracked, not accepted. This test fails the moment a spec sync
declares the header, which is the signal to move the entry up into
``_CONTRACT_HEADER_LIFTS`` and get it pinned like the rest. It also fails
if the header is declared under a *different* name for the same quantity,
because the reconciliation is the same either way.
For a lift the SDK reads from a header the vendored contract does not yet
declare. Nothing in the suite can catch a wrong name for such a lift,
because every test configures its stub to emit the exact literal the lift
reads -- so the gap is tracked, not accepted: this test fails the moment a
spec sync declares the header, which is the signal to move the entry up
into ``_CONTRACT_HEADER_LIFTS`` and get it pinned like the rest.

``credits_used`` (``X-Comfy-Credits-Used``) went through exactly that: it
was lifted ahead of the contract, this test fired on the sync that declared
it, and it now lives in ``_CONTRACT_HEADER_LIFTS``. With nothing left
undeclared the parametrization is empty, which pytest reports as a skip.
"""
declared = _declared_run_response_headers()
assert header not in declared, (
Expand Down
Loading