Skip to content

Keep outbox deliveries and native tasks retryable after a capacity refusal - #25

Merged
ancongui merged 9 commits into
mainfrom
fix/capacity-consumers
Oct 8, 2026
Merged

ancongui merged 9 commits into
mainfrom
fix/capacity-consumers

Conversation

@ancongui

@ancongui ancongui commented Oct 8, 2026 •

Copy link
Copy Markdown
Contributor

Summary

A capacity rejection (CatalogError(429, "WV-OPERATION-CAPACITY" | "WV-REQUEST-CAPACITY")) means the platform committed nothing and the call may be sent again. Two consumers still turned it into durable failures:

  • Outbox delivery (operations/event_delivery.py, _attempt). A refusal of the transaction that records the provider version (before the POST) or of _settle (after it) fell into the except Exception fallback. That fallback re-opened, passed _checked, and settled DELIVERY_FAILED, which became an incident on the last attempt. If the re-open was refused too, the delivery became a terminal AUTHORITY_REVOKED incident on its first attempt. It also recorded DELIVERY_FAILED after a receiver had acknowledged the request.
  • Native connector worker (connectors/execution.py). A refusal of the invocation check, or of the authorize() and credentials callbacks, reached sdk/worker.py's except Exception. That sent fail(HANDLER_FAILED). The kernel then opened WV-TASK-FAILED and suspended the run, and recovery treats HANDLER_FAILED as non-transient, so the result was always WV-TASK-AMBIGUOUS or WV-TASK-RETRIES-EXHAUSTED, even though nothing had run.

Changes

  1. Shared classifier. capacity_rejected() and CAPACITY_CODES live in definitions/models.py. The commit is byte-identical to 61224ad from Keep schedules and receipts retryable after a capacity refusal #21, which is now on main, so it adds nothing after the merge.
  2. Outbox. A capacity refusal anywhere in an attempt, including the recheck re-open, returns retry without settling. The fenced lease expires, and the next dispatch records ACK_UNKNOWN and redelivers the same event ID. Every other error path is unchanged.
    • The attempt still counts. attempts increments at claim time, and delivery_attempts is append-only (no DELETE grant, (delivery_id, generation) primary key). A delivery refused on every attempt ends in a DELIVERY_EXHAUSTED incident. Refunding an attempt would need a migration or a reclaim rule, which is an owner decision and is not part of this PR.
  3. Native handler. _admitted() replays capacity refusals of platform calls that committed nothing.
    • Before the connector starts. The first invocation check. The owning handler retries while its lease watchdog stays authoritative. This is the policy sdk/transport.py applies to a remote worker's context and credentials.
    • While the connector runs. authorize() (the invocation check plus credential_authority) and credentials. Each call makes at most three attempts within one second, so a refusal cannot keep a connector's own resources open. That covers an uncommitted SQL transaction, a machine-token slot, and a quota-full audit insert that would otherwise repeat a provider read. Every attempt rechecks that the handler is still active, so a retry cannot outlive it.
  4. Merge with main (Keep schedules and receipts retryable after a capacity refusal #21, Reserve pure execution for lifespan loops and native workers #22). execute() keeps Reserve pure execution for lifespan loops and native workers #22's request_execution() if reservation is None else reservation.execution(). Native workers bring the dispatcher's reservation, which never refuses at entry, so this PR's retry around taking a request work slot was unreachable and is removed. The invocation, authority, and credential retries stay.

Docs: reference/integration-events.md and reference/http-and-webhooks.md, each in a "Busy platform" paragraph. The latter now says in-process executors run up to their configured capacity, apart from request execution capacity.

CHANGELOG (Unreleased):

Tests

The unit tests and the outbox integration test failed on the base code before the fix:

  • tests/unit/operations/test_outbox_capacity.py. Refusal before send, refused recheck, and refused settlement after send. It also guards that other failures still settle DELIVERY_FAILED or AUTHORITY_REVOKED.
  • tests/unit/connectors/test_native_handler_capacity.py. Uses the real ConnectorExecutionService, ServiceTransport, Worker, and the NativeDispatcher reservation, called the way the dispatcher calls it. It covers a task that completes, including its pure work, while requests hold every work slot and the invocation check is refused once; a refused invocation check that outlasts the one-second direct window; refused authorization and credentials; the short window while a connector runs; a retry that cannot outlive the handler; and non-capacity errors that still fail. The earlier work-slot wait test is replaced, since production never takes that slot.
  • tests/integration/test_outbox.py::test_capacity_refusal_before_send_leaves_the_lease_for_recovery. Real contention: another writer holds the project operations advisory lock past the 250 ms timeout. On the old code the delivery ended up retry/DELIVERY_FAILED. With the fix it stays leased, and on expiry it is redelivered once, with outcomes {ACK_UNKNOWN, ACK}.

Verification

On the merged tree (6af9dfc, same tree as ff0569c):

  • make check: all requested checks passed.
  • Unit: 198 passed across the native handler, outbox, background admission, invocation fence, execution, email, provider, scheduler, recovery, worker capacity, and native worker suites.
  • Integration against disposable PostgreSQL 17 containers (colima-weave-tests): 216 of 217 passed. Only the built-image case of connectors/test_postgresql_native did not run, because it needs WEAVE_D1_IMAGE_ID. The suites were test_outbox, test_connector_execution, test_http_connector, test_leases, test_restart, test_terminal_capacity, test_incidents, test_schedules, test_email, connectors/test_postgresql, connectors/test_postgresql_native, providers/test_inbox, and providers/test_teams.
  • The outbox capacity integration test passed 6 of 6 runs.
  • Not run: the Kafka and TLS connector integration lanes, which need the owned release backends.

Follow-ups (not in this PR)

  • Connector wrappers (http_profiles, machine_tokens, http, postgresql, kafka and others) still turn an exhausted capacity refusal of a callback into their own permanent failure codes.
  • ServiceTransport.complete/fail use fixed 48-attempt/10-second windows and ignore the lease settlement scope. The remote transport retries until the task deadline instead.
  • email_submit (email.queue and email.execute) does not replay capacity refusals.
  • UnitOfWork maps every WQ001, including allocation quota limits, to WV-OPERATION-CAPACITY. A quota limit is not a burst.

Andres Contreras added 9 commits October 8, 2026 10:42
 into fix/capacity-consumers

# Conflicts:
#	src/firefly_weave/connectors/execution.py
Keep the reservation that native workers bring to a connector call and
drop the retry around the request work slot, which those workers no
longer take. The invocation, authority, and credential retries stay.
Run the native handler tests through the dispatcher's reservation, and
check that a task completes while requests hold every work slot.
Keep this branch's resolution of the main merge. Native workers run
from their reservation, which never refuses at entry, so the retry
around taking a work slot is removed instead of wrapped around the
reservation. The handler tests run through the dispatcher's
reservation, and the busy-platform docs keep both capacity codes,
since both are sent again.
@ancongui
ancongui merged commit 6408ed7 into main Oct 8, 2026
13 checks passed
@ancongui
ancongui deleted the fix/capacity-consumers branch October 8, 2026 21:24
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant