Skip to content

Reserve pure execution for lifespan loops and native workers - #22

Merged
ancongui merged 4 commits into
mainfrom
fix/background-loop-admission
Oct 8, 2026
Merged

ancongui merged 4 commits into
mainfrom
fix/background-loop-admission

Conversation

@ancongui

@ancongui ancongui commented Oct 8, 2026

Copy link
Copy Markdown
Contributor

Why

A local acceptance run logged "Tenant recovery scan failed; traversal will continue" at the same time Studio requests were getting 429s.

Pure CPU work runs under no-queue leases: WORK_SLOTS = 2 for mutations and CONTROL_SLOTS = 4 for reads and terminal control. Every HTTP request holds one lease for its whole response (BodyBoundary._serve). execute_pure() without an inherited lease takes a fresh one from those same pools and raises CatalogError(429, "WV-OPERATION-CAPACITY") when none is free. The lifespan loops have no lease of their own, so they borrowed request capacity:

Owner Pure work Pool What a refusal did
RecoveryLoop, recovery expired and retry quanta, wait_elapsed and signal_received settlements work "Tenant recovery scan failed", and the tenant's schedule scan was skipped for the cycle
RecoveryLoop, recovery timed_out deadlines, terminal capacity control same
RecoveryLoop, schedules runtime.start work schedule blocked (WV-SCHEDULE-READINESS)
ProviderLoop runtime.start, signals.deliver, settle work, control provider and email receipts blocked
KafkaLoop runtime.start, signals.deliver per record work, control "Broker consumer turn failed", the consumer torn down and the record replayed
NativeDispatcher ConnectorExecutionService.execute held a request work slot for the whole connector call, network I/O included work two concurrent native executions took every mutation slot from Studio; a third was refused at entry and reported as HANDLER_FAILED

The failure logs also gave no cause, so a capacity refusal could not be told apart from any other error.

What changes

  • Logs. runtime/scheduler.py now names the error class and, for a CatalogError, its code in "Tenant recovery scan failed", "Tenant schedule scan failed" and "Recovery catalog scan failed". It never logs the message or a value. Example: Tenant recovery scan failed (CatalogError WV-OPERATION-CAPACITY); traversal will continue.

  • ReservedSlots (operations/execution.py). Inside its execution() context, every execute_pure call takes one slot from the owner's own semaphore instead of the request pools, and any inherited request lease is set aside. execute_pure changes by one line.

    • A call that outlived its cancelled caller keeps its slot until it ends, so the size bounds the owner's pure threads.
    • A spare slot absorbs one such call, as the shared work pool did before.
    • A first version held one lease per cycle. Review showed that one orphaned call would then refuse the rest of the cycle with "Concurrent execution unavailable", including the next provider receipt after a 5 s dispatch_one timeout. That is why slots are now taken per call.
  • Owners:

    Owner Reservation Context
    RecoveryLoop RECOVERY_EXECUTION = ReservedSlots(2) one cycle (recovery and schedules)
    ProviderLoop PROVIDER_EXECUTION = ReservedSlots(2) one provider and email cycle
    KafkaLoop TURN_EXECUTION = ReservedSlots(8) one consumer turn (at most max_clients - 1 turns, max_clients <= 8)
    NativeDispatcher ReservedSlots(sum of executor capacities + 1) one connector execution, passed as execute(..., reservation=)
  • Unchanged:

    • Request admission (WORK_SLOTS, CONTROL_SLOTS) and the public ProcessCapabilities contract.
    • Native claims, heartbeats and settlements stay on request admission with ServiceTransport._replay, in parity with remote workers.
    • Direct callers of ConnectorExecutionService.execute without a reservation are still admitted like a request.
  • Pure thread bound. It rises by at most 2 + 2 + 8 + (native capacity + 1). Native executions used to be capped at 2 by the work pool, with the third failing; the operator's configured capacity is now the ceiling.

Related

  • Start lifespan loops from an empty context #20 creates every lifespan task with a clean contextvars.Context(). Without it, a loop opened by POST .../operations/compatibility/check inherits the request's lease and keeps failing after the response. It touches only the create_task lines of the same files.
  • Keep compatibility rescans from refusing work they confirm #19 adds inventory_execution() for the compatibility rescan. That lease deliberately refuses an overlapping rescan, so it keeps its one-lease semantics and its own slot.
  • Making schedules and provider and email receipts retryable after a capacity 429, instead of blocking them for good, is in a separate session and PR.

Residuals

  • A native execution whose pure call is refused is still reported as HANDLER_FAILED with an unknown outcome, although nothing ran. That only happens when orphaned calls hold every slot, the spare included.
  • WV-OPERATION-CAPACITY is shared with database admission (uow.py, sqlstates WQ001 and 55P03). After this change, one left in a loop's log comes from the database or from orphaned calls that hold a whole reservation.

Verification

  • make check on the branch head: all stages, All requested checks passed.
  • New and changed unit tests:
    • tests/unit/operations/test_execution.py:
      • work and control calls are admitted while requests hold every slot, and requests keep their ceiling;
      • concurrent calls in one context each take a slot, up to the size;
      • an orphaned call keeps its slot, the spare absorbs it, and a second orphan exhausts the reservation;
      • a closed inherited request lease is replaced.
    • tests/unit/runtime/test_recovery_admission.py, through the real RecoveryLoop.cycle → _quanta → _recover → transition_async → execute_pure path:
      • recovery progresses while requests hold every work slot, and a timed_out deadline while they hold every control slot;
      • one orphaned call is tolerated;
      • a full reservation refuses the quantum and is logged once;
      • the three failure logs carry the class and code without values.
    • tests/unit/operations/test_background_admission.py:
      • a provider cycle, two concurrent Kafka turns and four concurrent native executions run pure work while requests hold every slot;
      • the Kafka reservation follows the BrokerPolicy.max_clients cap.
    • tests/unit/connectors/test_local_build.py: the handlers NativeDispatcher.open() registers execute through its reservation.
  • Mutation checks:
    • reverting the one-line execute_pure change fails three reservation tests;
    • dropping reservation=self.reservation from the dispatcher fails the wiring test.
  • Five consecutive runs of the touched areas while other gates loaded the machine: 252 passed each time.

Andres Contreras added 3 commits October 8, 2026 09:51
The recovery and schedule traversal logged only "Tenant recovery scan
failed" and "Tenant schedule scan failed", so a capacity refusal could
not be told apart from any other error. Both lines now name the error
class and, for a CatalogError, its code, never the message or a value.

A unit test drives the real RecoveryLoop cycle through the recovery
quanta into transition_async while two in-flight requests hold every
work slot. The cycle fails with WV-OPERATION-CAPACITY and skips that
tenant's schedule scan; once the requests finish, it reaches the kernel.
The test pins today's admission design, in which lifespan loops borrow
the request work pool.
In-flight requests hold a work or control slot for their whole response,
and lifespan loops borrowed those same no-queue pools. Two Studio
mutations, or four reads, refused recovery quanta, deadline settlements,
schedule starts, provider and email dispatch, and Kafka consumer turns.
Native connector executions held a request work slot for the whole
connector call, so two of them took every mutation slot from Studio and
a third failed at entry.

ReservedSlots gives each owner its own pool. Inside its execution()
context, every pure call takes one slot from that pool instead of the
request pools, and an inherited request lease is set aside. A call that
outlived its cancelled caller keeps its slot until it ends, so the size
bounds the owner's pure threads, and a spare slot absorbs one such call
as the shared work pool did. Recovery and provider cycles get 2 slots,
Kafka turns 8 (at most max_clients - 1 turns, max_clients <= 8), and
the native dispatcher its configured capacity plus one. Native claims,
heartbeats and settlements keep request admission and their replay
policy. Request admission itself is unchanged.

The recovery catalog failure log now names the error class and code as
the tenant logs do.
The ReservedSlots tests now share the request saturation helper in
test_background_admission.py instead of duplicating it at the end of
test_execution.py, which the compatibility rescan change also extends.
@ancongui
ancongui merged commit dd64e2d into main Oct 8, 2026
13 checks passed
@ancongui
ancongui deleted the fix/background-loop-admission branch October 8, 2026 18:40
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