Repository navigation
Conversation
|
Marking this ready for review — it is the concrete proposal for #1458, and I am keeping it as a draft only until the API shape is settled. Naming / shape questions (why it is still a draft)
What is in the branch (
One thing I found that may want its own decision: on this stack Requesting review from @mdboom and @leofang. If it is easier to review in pieces, I can split device events and system events, or drop the system-event half entirely for a first PR. |
|
Follow-ups pushed after the review request, all CI-visible:
The measurements in the description are from |
e5507ea to
d6bff84
Compare
Re-measurement and clean rebuild are doneFollow-up to the review request above: the tables are final now, and this branch is the tree they were measured on. Branch state. The branch was rebuilt through the Git Data API (git transport to github.com is not available from the test host) and now carries the work-tree commit chain exactly, 9 commits, What changed since the earlier tables. Independent clean rebuild of this same commit (fresh worktree, fresh venv,
The headline row is unchanged from the earlier run — a 1 ms heartbeat in the control loop goes from ~301 ms between beats with the blocking API to 1.11 ms with Two things still need a maintainer's call, both asked above: the |
Second-host verification (8 × RTX 5090)Same commit — this branch head Full test suite. Contract tests. The 27 async-event cases pass unchanged ( Stop latency at scale. The first host only had two free GPUs, so the dispatcher's worker bound was covered by unit tests alone. On eight GPUs, quiet event sets,
Two things worth noting from this:
Caveat: the (N−1) × timeout serialisation is this host's driver behaviour and is not extrapolated to the first host, where only 2 GPUs were available. Direction replication (single-arm, not paired statistics). Re-running the same harness on this host reproduces the headline rows:
The baselines land within a few tenths of a percent of each other across two different GPUs, drivers, Python versions, bindings majors and CUDA header versions, which is what you would expect if these quantities are set by the native wait semantics rather than by host specifics. |
DeviceEvents.wait_async() and RegisteredSystemEvents.wait_async() wait for an NVML event without blocking the event loop. The native wait is issued in slices of at most 100 ms against a monotonic deadline, which keeps cancellation bounded: a cancelled task drains the in-flight slice instead of parking a worker thread and the event-set handle for the caller's whole timeout (indefinitely, when the timeout is 0). The wait state enforces a single consumer per event set, so at most one native wait is in flight, and it keeps an event that a cancelled or expired slice had already consumed in a pending slot for the next consumer; a system-event batch is handed over one buffer_size at a time so no event is dropped. The worker pool used for the blocking native call is private, bounded and created on first use, so importing cuda.core.system still starts no threads and needs no CUDA context. Signed-off-by: 0z5a <Dezhen.lu@student.uni-tuebingen.de>
Deterministic tests drive the wait state with a fake native wait to pin the slice budget, the anti-busy-wait behaviour of an early native timeout, single consumer ownership, cancellation during the in-flight slice, repeated cancellation, and the hand-off of an event that a cancelled slice consumed. Native tests cover a real event set on the installed GPUs: timeout budget, cancel-to-drain latency, two independent event sets, and the sync/async lease conflict. A subprocess test pins that importing cuda.core.system starts no threads. Signed-off-by: 0z5a <Dezhen.lu@student.uni-tuebingen.de>
The work item handed to the worker pool can be cancelled while it is still queued, and its future then raises CancelledError from exception(). The drain path has to read that as "this slice never borrowed the event set" instead of letting a second exception type escape the cancellation path. Signed-off-by: 0z5a <Dezhen.lu@student.uni-tuebingen.de>
A thread pool priced per request dominated the cost of a short wait: every slice paid for a work item, an executor future and a chained asyncio future, and a queued slice could wait behind a long one belonging to another event set. Slices now go through a private queue served by workers that are created on demand - one per in-flight slice, up to the bound - and reused across waits, and the result is delivered straight to the awaiting future. A slice that times out is now a value rather than an exception: only the caller's own deadline raises, once, with the project's NVML timeout type. The result of a slice is kept in a slot on the wait state so that a cancellation can still read it after the future carrying it was cancelled. Measured on 2x L20 with a 1 ms timeout cycle: 926 -> 821 cycles/s against a blocking baseline (was 761 before this change), at 176 us CPU per cycle. Signed-off-by: 0z5a <Dezhen.lu@student.uni-tuebingen.de>
A drain that waits for the borrowed slice through the worker queue can queue up behind another event set's slice, which doubled the measured stop latency. The drain now registers a waiter on the outcome slot and the worker that finishes the slice completes it, so draining costs the rest of the slice and no more. Workers also retire sooner once they go idle, which keeps a process that stops waiting from paying for them at exit. Signed-off-by: 0z5a <Dezhen.lu@student.uni-tuebingen.de>
The extension methods wrapped the wait state machine in one more coroutine, which cost a frame on every wait: on a 1 ms timeout cycle it was ~14 us of the ~28 us the async path still spent over a hand-rolled non-blocking handoff. The methods now hand back the state machine's coroutine with the result converter, so the state machine stays a single pure-Python frame; `await` is unchanged. The conversion of a consumed payload moves into the state, which lets the system-event path turn both a fresh batch and a parked one into SystemEvents through one helper. Signed-off-by: 0z5a <Dezhen.lu@student.uni-tuebingen.de>
stubgen-pyx has no way to annotate a coroutine that a plain def returns, and pre-commit is the gate for review, so DeviceEvents.wait_async and RegisteredSystemEvents.wait_async are real async def methods again: the stubs keep their EventData / SystemEvents return types, and iscoroutinefunction() reports what callers expect. The cost is one more coroutine frame per wait. Cython's annotation typing turns that cdef-class return annotation into a C type the coroutine wrapper cannot return, so the two annotations are marked as hints with @cython.annotation_typing(False), which keeps the generated stubs accurate. Also fixes what pre-commit reported: import order, a bound method that kept the test owner alive through a closure that a later del made look undefined, and a subprocess call that now says why its argv is fixed. Signed-off-by: 0z5a <Dezhen.lu@student.uni-tuebingen.de>
mypy asks for annotations on the queue and the worker list, and the request slot is cleared to None so an idle worker does not pin an event set. Types only: no runtime behaviour changes. Signed-off-by: 0z5a <Dezhen.lu@student.uni-tuebingen.de>
Four cases that the contract implies but the suite did not exercise yet: * a failure after the native event set exists (device_register_events rejects) frees it exactly once instead of leaking it or freeing it twice; * with the worker bound reached, slices are serialised, the queued slice does not reach the driver out of turn, and the queue wait stays inside the caller's budget; * the same event set is reusable across consecutive event loops, while a second loop waiting on it concurrently is rejected by the lease instead of racing the driver for the same event set; * type and range errors come from the annotated signatures (sync and async), including that an async argument error only surfaces when awaited. The suite is 27 cases now. Signed-off-by: 0z5a <Dezhen.lu@student.uni-tuebingen.de>
d6bff84 to
ab5b1a7
Compare
mdboom
left a comment
There was a problem hiding this comment.
Thanks for this PR!
In addition to the code changes below, this will need release notes -- both for the (minor) breaking changes to wait and to mention the new feature of the async APIs.
There was a problem hiding this comment.
This is a pre-existing documentation bug, but it's worth fixing it now.
| The timeout in milliseconds. A default value of 0 means to skip waiting. |
There was a problem hiding this comment.
This is a pre-existing documentation bug, but it's worth fixing it now.
| The timeout in milliseconds. A default value of 0 means to skip waiting. |
| except asyncio.CancelledError: | ||
| result = await _drain(outcome, loop) | ||
| if result is not None: | ||
| self.park(result) |
There was a problem hiding this comment.
From my agent:
cuda_core/cuda/core/system/_system_events.pyx (wait / _take_batch) — HIGH — correctness (event loss)
A cancelled wait_async parks the raw SystemEventData_v1 (park(result) at _async_events.py:304). Sync wait() converts parked data with _take_batch, which assumes a (batch, index) tuple. Only _batch_result handles both shapes, and only the async path uses it. Unpacking a raw SystemEventData_v1 raises ValueError for 1 event and 3 events, and mis-unpacks silently for exactly 2 (confirmed experimentally). take_pending() has already cleared the slot, so the batch is lost.
- Failure scenario: cancel (or time out) a
wait_asyncwhose slice consumed a batch, then callwait(). It raises instead of returning the events, and the events are gone. - Fix: use
_batch_resultin both paths, or park a normalized(batch, 0)tuple. - Related: raw parked batches are returned whole by
_batch_result, ignoring a smallerbuffer_size. That contradicts the docstring ("onebuffer_sizeslice at a time"). Any exception fromconvert, such as thebuffer_size < 1ValueError in_take_batch, also happens after the pending slot was cleared, so it loses the event too. - Test gap:
test_system_batch_remainder_is_preservedonly feeds tuples, so none of this is covered.
| try: | ||
| await asyncio.shield(waiting) | ||
| except asyncio.CancelledError: | ||
| continue |
There was a problem hiding this comment.
From my agent:
If the drained slice failed (for example GpuIsLostError), _drain raises that error from the except CancelledError handler. The task then ends with GpuIsLostError rather than CancelledError (reproduced). This breaks the usual contract that task.cancel() ends in CancelledError, and it confuses asyncio.timeout() and TaskGroup. Consider logging or parking the error and re-raising CancelledError.
| if self._busy: | ||
| raise RuntimeError("an event wait is already in flight for this event set") | ||
| self._busy = True | ||
| try: | ||
| pending, claimed = self._claim(convert) | ||
| return pending if claimed else native_wait(timeout_ms) | ||
| finally: | ||
| self._busy = False |
There was a problem hiding this comment.
From my agent:
The lease check-and-set (if self._busy: ...; self._busy = True) is unlocked. Two threads calling wait(), or one calling wait() and one wait_async() from another loop, can both pass the check. test_event_set_is_reusable_across_event_loops only tests the deterministic ordering
|
|
||
|
|
||
| def test_event_set_is_reusable_across_event_loops(): | ||
| """N22: serial reuse across loops is fine; a concurrent loop is rejected.""" |
There was a problem hiding this comment.
| """N22: serial reuse across loops is fine; a concurrent loop is rejected.""" | |
| """serial reuse across loops is fine; a concurrent loop is rejected.""" |
|
|
||
|
|
||
| def test_parameter_bounds_match_the_signatures(): | ||
| """N23: the type and range errors come from the annotated signatures.""" |
There was a problem hiding this comment.
| """N23: the type and range errors come from the annotated signatures.""" | |
| """the type and range errors come from the annotated signatures.""" |
| with pytest.raises(system.TimeoutError): | ||
| asyncio.run(events.wait_async(timeout_ms=300)) | ||
| elapsed = time.monotonic() - started | ||
| assert 0.3 <= elapsed < 1.0, elapsed |
There was a problem hiding this comment.
| assert 0.3 <= elapsed < 1.0, elapsed | |
| # Timings may be potentially flaky on loaded runners. | |
| # Remove if we see flaky tests in CI. | |
| assert 0.3 <= elapsed < 1.0, elapsed |
| elapsed = time.monotonic() - started | ||
|
|
||
| asyncio.run(main()) | ||
| assert elapsed < 0.5, f"cancel returned only after {elapsed:.3f}s" |
There was a problem hiding this comment.
| assert elapsed < 0.5, f"cancel returned only after {elapsed:.3f}s" | |
| # Timings may be potentially flaky on loaded runners. | |
| # Remove if we see flaky tests in CI. | |
| assert elapsed < 0.5, f"cancel returned only after {elapsed:.3f}s" |
| @@ -0,0 +1,315 @@ | |||
| # SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. | |||
There was a problem hiding this comment.
| # SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. | |
| # SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. |
Serialize waits across threads and event loops, preserve batches and native errors across cancellation, and retain consumed payloads when conversion fails. Cover early timeouts and closed loops, correct synchronous timeout documentation, and document the async APIs and wait changes in the release notes. Signed-off-by: 0z5a <dezhen.lu@student.uni-tuebingen.de>
Signed-off-by: 0z5a <dezhen.lu@student.uni-tuebingen.de>
Signed-off-by: 0z5a <dezhen.lu@student.uni-tuebingen.de>
|
/ok to test 5199f5e |
@mdboom, there was an error processing your request: See the following link for more information: https://docs.gha-runners.nvidia.com/cpr/e/2/ |
|
/ok to test a0faf4a |
@mdboom, there was an error processing your request: See the following link for more information: https://docs.gha-runners.nvidia.com/cpr/e/2/ |
|
No actionable comments were generated in the recent review. 🎉 ℹ️ Recent review info⚙️ Run configuration
📒 Files selected for processing (8)
Included review availability: This review used your included allowance. Your plan provides up to 12 included reviews per hour; 11 remain after this review. 📝 SummarySummary by CodeRabbit
WalkthroughAdds asynchronous NVML event waits for device and system event sets. Waits on each event set are serialized. Bounded native wait slices support cancellation, and consumed events or errors can be handed to a later wait. System-event results are delivered in buffer-sized batches. ChangesNVML event wait coordination
Suggested reviewers: Priority: ⬇️ Low Change: Feature Merge Risk: ⚪ Minimal · up to No merge-blocking issue is established for the async NVML waits. Proceed with normal checks, including native GPU validation when available.
Comment |
|
/ok to test a0faf4a |
|
Incremental review: Files changed against
main.Scope
Related to #1458. Adds
DeviceEvents.wait_async(timeout_ms=0)andRegisteredSystemEvents.wait_async(timeout_ms=0, buffer_size=1)so an asyncio control loop can await NVML events while continuing other work.Async iterators, a public executor setting, background listener management and system hot-plug handling remain outside this proposal.
Behavior
timeout_ms=0means indefinite waiting for the async APIs and skip waiting for synchronouswait().CancelledError, including repeated cancellation. A native error encountered during that drain is saved for the next consumer.buffer_sizepieces, and failed conversion preserves the payload for a retry. Invalid buffer sizes are rejected before consuming an event.cuda.core.systemstarts no threads. Closing a coroutine or its loop keeps the event-set lock held until its outstanding native call finishes and saves its outcome.Review follow-up and current validation
Review fixes in 30feadc address the batch-loss, cancellation-error and unlocked-lease findings, replace concurrent-wait exceptions with serialization, correct the fake early-timeout behavior, and add closed-loop and conversion-failure controls. Release notes cover the async APIs, wait locking overhead and invalid-buffer behavior; generated stubs match the updated docstrings.
The final branch at 5199f5e incorporates
mainat608a37ddc71a37235f2476e57e15a2b3cd7259a9, retaining both sets of release notes. CPU controls, mypy and full Cython code generation were rechecked after this integration. GitHub confirms there are no merge conflicts.asyncio.timeout,TaskGroup, early timeouts and closed-loop cleanup._device.pyxand_system_events.pyxCython code generation, stub generation, SPDX, mempool hygiene, release-note presence and whitespace checks passed locally.Historical GPU measurements
The following measurements were previously reported for the pre-review implementation. They provide context for the proposal; the current review fixes need a fresh native run before these measurements can qualify this revision.
Measured on one host, 2 × NVIDIA L20 (SM89, 46 GB), driver 595.91.07 / NVML 13.595.91.07,
cuda-bindings12.9.8, Python 3.10.12, CUDA toolkit 12.8.61.8b6e9f52db1074a580fc1678c71e7c75c42db2ead6bff844d, before the review follow-up below. These tables were measured afterwait_asyncbecame a realasync def; they have not been re-measured for the current locking and error-handoff changes.evidence/env-patch.json; independent clean rebuild:evidence/env-clean.jsoncuda_core/tests/system/test_system_events_async.pydrives every race with a fake native wait (slice budget, early timeout, single consumer, cancel during the in-flight slice, repeated cancel, hand-off of an event consumed by a cancelled slice, batch remainder, exactly-once free of a failed registration, saturated scheduling, reuse across event loops, argument bounds) and exercises real event sets on both L20s (timeout budget, cancel-to-drain, two independent event sets, sync/async lease conflict)cuda_core/tests/system+test_stream.py: 150 passed, 131 skipped, 8 xfailed, 984 subtests passedcuda-bindings13.4.2, CUDA 13.4 headers), same commit and tree: 27 passed, andcuda_core/testsin full — 4317 passed, 265 skipped, 17 xfailed, 872 subtests passed, 0 failed (4463 collected). 8-GPU fan-in with quiet event sets: stop latency +80.3 % (2 GPUs), +86.8 % (4), +88.6 % (8), with 1+N worker threads and the bound reached exactly at 8; the blocking baseline grows as ~(N-1) x timeout because this host's driver serialises concurrent blocking waits. Single-arm replication of the headline rows: heartbeat p50 300.90 ms -> 1.218 ms, stop p50 1701.7 ms -> 101.0 ms, exit 10015.8 ms -> 164.3 ms (details in the comment below)pip install --no-build-isolation -e ., recorded inevidence/env-clean.jsonwith the SHA256 of all 47 loaded extensions): 27 passed, and the paired results below reproduceEvery row is a paired four-arm quartet (
A→P→P→A/P→A→A→P, a fresh process per arm), 3–5 independent quartets per row, session-level bootstrap intervals, with an A/A placebo measured at −0.0 % / −0.1 %.wait_asyncwait(timeout_ms=1000)baselinewait(timeout_ms=1000)baselinewait(timeout_ms=100)baselinetimeout_ms=0, 10 s watchdog)timeout_ms=1000)timeout_ms=0(the documented "wait indefinitely")timeout_ms=1000Independent clean rebuild at the same revision (4 fresh quartets per configuration): stop latency p50 +93.5 %, p95 +88.2 %, real-event stop +88.0 %, exit +98.8 %, 1 ms churn −11.8 %, idle CPU +4.3 %, real-event count +0.3 %. Repeating the real-event count over runs gives −23.3 % and +0.3 % with intervals that cross zero: it is noise around parity, and no run lost an event.
Notes that keep the numbers honest:
wait()inside an asyncio control loop starves that loop: with the sync API the same 1 ms heartbeat degrades to ~301 ms between beats, because it only gets to run between blocking waits.wait_asynckeeps it at 1.11 ms (p99 1.27 ms) with the same cycle throughput (−0.2 %), and a hand-rolled worker-thread + queue wrapper around the existing sync API reaches the same cadence (−0.1 %, i.e. parity), which is what makes this a fair like-for-like comparison.wait(1 ms)occupies the calling thread for the whole wait; any non-blocking wait pays a wakeup plus an event-loop round trip (a hand-rolled worker-thread + queue wrapper measures −8 % against the same blocking baseline before this API exists, and bareasyncio.sleep(1 ms)on this host costs 1121 µs). This API lands at −3.3 % against that hand-rolled non-blocking wrapper — the lease/deadline/pending/drain machinery plus theasync defframe — and −0.1 % at the 100 ms slice scale the design is built around. Earlier revisions of this branch measured −17.8 % here; the commits contain the fixes for that.timeout_ms=0does not mean "wait indefinitely" on this stack, contrary to the docstring at that time and the NVML C documentation: it returns immediately withTimeoutError(295k calls/s busy loop, 2 s of CPU for 2 s of wall). The async path therefore never passes 0 to the driver — it slices a positive timeout and implements the unbounded budget itself. The synchronous docstrings are corrected in this review follow-up to say that 0 skips waiting.Limits