fix: use anyio instead of asyncio primitives so trio runs work - #851
Conversation
CLAUDE.md requires anyio rather than asyncio primitives, because the backend is
chosen at runtime by INSPECT_ASYNC_BACKEND. asyncio primitives raise under trio,
and since the default is asyncio the breakage only surfaces in trio runs -- which
is how `from asyncio import gather` reached seven scorers unnoticed.
Demonstrating the failure, and that it is not theoretical:
[asyncio] asyncio.gather -> [1, 1]
[trio] asyncio.gather -> RuntimeError: There is no current event loop
End to end, `INSPECT_ASYNC_BACKEND=trio pytest
control_arena/settings/data_poisoning/tests/test_scorers.py` fails on main,
pointing at scorers.py:285, and passes here.
Converted:
- Seven scorers using `from asyncio import gather` (apps, bash, bigcodebench,
data_poisoning, iac, iac_fast, sae_interp_sabotage) to
control_arena._async_utils.gather. Scorers run inside inspect's event loop, so
these were live breakage.
- Two attack-policy modules in examples/paper_replications, which run inside the
loop for the same reason.
- vllm's task-generation `asyncio.as_completed` to an anyio task group over a
memory object stream. It is reached from cli.py, which already runs under
anyio.run(backend=get_async_backend()), so it was breaking for real. The
replacement keeps completion-order streaming and keeps a failing coroutine from
cancelling its siblings -- a task group cancels the whole group by default, so
that is the property most easily lost in the move.
- The remaining `asyncio.run` entry points in examples/ and scripts/. These owned
their own loop so they were self-consistent, but CLAUDE.md asks entry points to
honour the configured backend.
Added `_async_utils.run_sync`, which anyio.run cannot replace directly: anyio.run
takes a function plus positional arguments, so it cannot accept a coroutine built
with keyword arguments, which is what existing `asyncio.run(f(a, b=c))` call
sites have.
To stop this recurring, ruff now bans the asyncio primitives via TID251 with
per-symbol messages pointing at the replacement. That is what surfaced the
examples/ and scripts/ call sites, and it catches both the `from asyncio import
gather` and `asyncio.gather(...)` forms.
Tests cover run_sync under both backends and the new completion iterator,
including that a failure does not cancel its siblings.
Note for reviewers: the trio suite is broadly red on main for unrelated reasons
(102 failures in compaction, tool output and sandbox utils). This change is
neutral against that baseline -- the same 102 fail before and after, with no
regressions -- so it removes one class of breakage rather than turning trio
green.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Checked against inspect, which solved this first. Two things came out of it.
`_async_utils.gather` was not actually a drop-in for `asyncio.gather`: an anyio
task group wraps even a single failure in an ExceptionGroup, and gather passed
that straight through, so it raised ExceptionGroup where asyncio.gather raised
the original exception.
asyncio.gather -> ValueError
control_arena gather -> ExceptionGroup <- before
inspect util.collect -> ValueError
inspect tg_collect -> ValueError
Both of inspect's equivalents unwrap, and ours now does too: a single failure
raises as itself, several raise as a group, matching inspect_ai.util.collect.
Catching ExceptionGroup rather than BaseExceptionGroup keeps cancellation
propagating as a group.
This mattered for the conversions in the previous commit -- the seven scorers had
been raising the original exception via asyncio.gather and would have started
raising ExceptionGroup instead, changing the error recorded against a failed
sample. No call site catches around a gather (checked all 23), so nothing was
actively broken, but the trap was real and the existing test encoded the wrong
expectation.
The TID251 banned-api approach in the previous commit turns out to be exactly
what inspect uses for its own bans, so no change needed there.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
|
Checked this against inspect (
Four things came out of it, one of which was a real defect in this PR. 1.
|
control-arena is an inspect extension, so it should not be maintaining a
second-rate copy of inspect's async entry point. `run_sync` now calls
`inspect_ai._util._async.run_coroutine` and `get_async_backend` calls
`configured_async_backend`. Both live in a private inspect module, so
`_async_utils` is the only place that imports it; if it moves, one file changes.
This is not just deduplication -- the version being replaced was a plain
`anyio.run`, which refuses to start inside a running loop. Inspect's re-enters
via nest_asyncio, so `run_sync` now works from a notebook:
from inside a running loop -> done # previously RuntimeError
It also releases anyio's per-loop state afterwards, avoiding a leak of the
return value on every call, and refuses clearly when called from inside a trio
task rather than failing obscurely.
Separately, while checking the delegation had not broken anything: both
parametrised backend tests were being skipped in CI. conftest skips anyio trio
variants by matching the nodeid substring "[trio", which also caught the "[trio]"
parameter id on two plain sync tests that drive anyio themselves. CI runs pytest
without --runtrio, so the assertions that run_sync and the vllm completion
iterator honour the trio backend never actually executed there. Renamed the ids
to "on-asyncio"/"on-trio" so both cases run in the default invocation.
Also filled in a section header in the test file that had lost its title.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
|
Delegated to inspect. This turned out to be more than deduplication — the version it replaced was a plain So A gap this uncovered: two of my new tests were dead in CIWhile checking the delegation hadn't broken anything I noticed
with no Renamed the ids to Worth flagging as a general trap: any Also filled in a section header in Full suite green (989 passed), pyright clean, trio invocation unchanged against its 102-failure baseline. |
The root conftest decided whether a test was a trio variant with
`"[trio" in item.nodeid`. That also matches an ordinary
`@pytest.mark.parametrize` value spelled "trio". Because CI runs pytest without
--runtrio, anything classified as a trio variant is skipped in the only
invocation CI performs, so such a test never runs anywhere while still looking
like coverage. Both backend-parametrised tests added earlier in this branch were
in exactly that state.
anyio's pytest plugin parametrises async tests by injecting an `anyio_backend`
param, so that param is what actually determines the backend. `anyio_backend_of`
reads it, handling anyio's ("trio", {options}) form, and returns None for tests
that are not anyio variants at all.
Verified behaviour-preserving before switching: across the whole suite, zero
items are classified differently by the two approaches. The previous commit's
workaround -- renaming the parametrise ids to "on-trio" so they dodged the
substring -- is reverted, since natural ids are safe now and keeping them proves
the fix.
Added tests for the classifier itself, including the case that caused this, since
the failure mode is silent loss of CI coverage rather than a visible error.
Note that inspect has the same substring check in its own conftest, so there was
no upstream version to adopt here.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
|
Tightened it. Answering the "how" in case it's useful beyond this PR. The discriminatoranyio's pytest plugin parametrises async tests by injecting an The last line is the bug: a genuine anyio variant and an ordinary parametrize value spelled def anyio_backend_of(item) -> str | None:
callspec = getattr(item, "callspec", None)
if callspec is None:
return None
backend = callspec.params.get("anyio_backend")
if isinstance(backend, tuple):
backend = backend[0] # anyio accepts ("trio", {options})
return backend if isinstance(backend, str) else None
Checked before switchingI ran a probe over the whole suite comparing the two classifications: So it is behaviour-preserving today — the only items it would have changed were the two I had already worked around. Reverted the workaroundThe Added a guard
Worth knowingInspect has the identical substring check in CI invocation 989 → 995 (the six guard tests). Trio invocation unchanged against its 102-failure baseline. ruff and pyright clean. |
| return csvfile, writer | ||
|
|
||
|
|
||
| async def _as_completed( |
There was a problem hiding this comment.
This looks like it should be in an async utils file rather than hidden in a vllm specific file?
There was a problem hiding this comment.
also worth double checking there isn't an inspect version already we can import
There was a problem hiding this comment.
Moved. It is now control_arena/_async_utils.as_completed_with_exceptions, generic in the result type and named to match gather_with_exceptions (an outcome is either a result or the Exception that replaced it, rather than the (result, error) tuple the vllm-local version yielded).
Promoting it to shared code turned up two defects the local version had, neither of which the old tests could see because they only ever consumed it to completion:
1. Leaving the loop early left producers sending into a closed stream. break, or an exception in the loop body, closes the receiving end while siblings are still running, and they then raise BrokenResourceError out of the task group. The remaining coroutines are now cancelled instead.
2. A bare async generator holding a task group is closed whenever the consumer gets round to it, which may be after its enclosing cancel scopes have moved on:
RuntimeError: Attempted to exit a cancel scope that isn't the current tasks's current cancel scope
That replaces whatever the loop body was actually doing. It is an async context manager now, so it closes at the point the caller leaves:
async with as_completed_with_exceptions(coros) as outcomes:
async for outcome in outcomes:
...There was a third, related trap: an exception from the caller's own loop body passes through the task group, which wraps it in an ExceptionGroup — exactly what gather was fixed for earlier in this PR. Both now share _reraise_alone, which also keeps the __cause__ chain that raise ... from None was discarding (so a failure raised from an underlying error no longer loses it).
Tests split accordingly: the iterator's behaviour is in control_arena/tests/test_async_utils.py (completion order, failure isolation, None results, empty input, generator input, early exit, loop-body exception, both backends), and the vllm tests now cover what the call site does with the outcomes — rows written in completion order, one failing commit not abandoning the rest, skips counted.
There was a problem hiding this comment.
Checked against the pinned inspect_ai (0.3.250). There is no equivalent to import.
The two candidates both wait for everything and return in submission order, which is the opposite of what this call site needs:
| order | streams as they land | failure isolation | |
|---|---|---|---|
inspect_ai.util.collect |
submission | no — barrier | no, first failure raised |
inspect_ai._util._async.tg_collect |
submission | no — barrier | no, first failure raised (or a group) |
as_completed_with_exceptions |
completion | yes | yes, reported per item |
The vllm generation script writes each row to CSV as it lands so a long run is resumable if interrupted, so a barrier would change its behaviour rather than just its implementation.
Grepping the rest of the package for the pattern found nothing reusable either — the closest is model/_call_tools.py, which hand-rolls the same task-group-plus-memory-object-stream arrangement inline for streaming tool results. So the approach matches inspect's, even though the helper does not exist there to import.
Where inspect does have prior art I have taken it: run_sync and get_async_backend delegate to inspect_ai._util._async (earlier commit in this PR), and _async_utils remains the only module importing that private path.
Addresses review: `_as_completed` was a general-purpose replacement for `asyncio.as_completed` living in a vllm-specific file. It is now `_async_utils.as_completed_with_exceptions`, generic in the result type and named to match `gather_with_exceptions`. Checked for an upstream version first, as suggested: inspect has none. Both `inspect_ai.util.collect` and `inspect_ai._util._async.tg_collect` wait for every task and return results in submission order, so neither can stream outcomes as they land. Promoting it surfaced two defects that the vllm-local version had: - Leaving the loop early -- `break`, or an exception in the body -- closed the receiving end of the stream while producers were still running, so they raised BrokenResourceError out of the task group. Now cancelled instead. - As a bare async generator holding a task group, it could only be closed by whoever got round to it, and closing it after its enclosing cancel scopes had moved on raised `RuntimeError: Attempted to exit a cancel scope...` in place of whatever the loop body was doing. It is now an async context manager, so it closes deterministically at the point the caller leaves. An exception from the caller's own loop body also passes through the task group, which wraps it in an ExceptionGroup -- the same trap `gather` had. Both now share `_reraise_alone`, which additionally keeps the failure's `__cause__` chain that `raise ... from None` was discarding. Tests for the iterator move to test_async_utils.py; the vllm tests now cover what the call site does with the outcomes -- rows written in completion order, one failing commit not abandoning the rest, skips counted. Also: the ruff ban now names run_sync and the new helper as the replacements rather than raw anyio, extends to the asyncio primitives CONTRIBUTING already listed but lint did not (subprocesses, to_thread, Lock/Event/Queue, timeout, get_event_loop), and CLAUDE.md/AGENTS.md/CONTRIBUTING.md list both new helpers.
|
Both review comments addressed in 567dba9 (details in the threads). Then a pass over the whole PR, which turned up four things worth recording. 1. The lint ban was narrower than the documented rule
2. The ban messages pointed at things this PR had superseded
3. Docs did not mention either new helper
4.
|
The scanner predated `run_sync` and `as_completed_with_exceptions`, so it used `anyio.run` with a partial and hand-rolled batching to get incremental writes. - `run_sync` at the CLI boundary, which also handles being called from a notebook. - `as_completed_with_exceptions` in place of batched `gather`: verdicts are now written the moment each lands rather than at each batch boundary, so a crash costs one classification rather than up to a batch of them, and one failure no longer takes its siblings down. How many run at once is bounded by the model's max_connections, as it always really was. Also pins the revision in the `stats` property, which still loaded the dataset's main branch and so could report a task count the setting would never load. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
`pytest --runtrio` had **102 failures**. After this it has **0** (214 passed). The default suite is unaffected: 1132 passed. ## Where the 102 went | Cause | Count | Ours? | Action | |---|---:|---|---| | `kubernetes-asyncio` is aiohttp-based | 67 | No -- library choice | Skipped in the infra setting, with the reason recorded | | Reading `EvalLog.samples` under trio | 34 | No -- upstream, **now released** | Fixed by an inspect floor bump; the tests run | | `aiofiles` is asyncio-only | 1 (+3 untested paths) | **Yes** | Fixed | ## The upstream fix is released, so nothing is gated any more An earlier version of this PR carried a version-gated skip for the 34: `eval_async()` returns logs whose `samples` is a `LazyList`, and touching it calls inspect's *sync* `read_eval_log()`, which inspect refuses to run under trio. inspect_ai #4676 fixes it, and at the time it was on inspect's main but in no release. It is released now -- 0.3.252, 4 August. Verified directly: ``` 0.3.250 (previous floor) RuntimeError: read_eval_log cannot be called from a trio async context 0.3.252 works 0.3.257 (current) works ``` So `control_arena/conftest.py`, the `inspect_log_trio_bug` marker and its eight `pytestmark` lines are gone, replaced by `inspect-ai>=0.3.252`. Those 34 tests now **run** under trio rather than being skipped, and there is no gate left for anyone to remember to delete. The full suite passes on 0.3.252 and on 0.3.257; the lock takes the latter. One of the eight was ours rather than upstream's. `test_tool_output_e2e` calls the sync `read_eval_log_sample` from async code, which inspect refuses under trio at any version. It uses the public `read_eval_log_async` now -- fixed rather than skipped. ## aiofiles -- the part that was actually ours aiofiles resolves its executor via `asyncio.get_event_loop()` and raises `no running event loop` under trio. Replaced in the two remaining production call sites (`export/writers/directory.py`, `scorers/_run_monitors.py`) and one test helper; the third, in `iac_fast/scorers.py`, disappeared when #852 rewrote that function. `aiofiles` and `types-aiofiles` are dropped from the dependencies. `copy_eval_files_async` did not want `anyio.open_file` at all -- it read a whole file into memory and wrote it back out. It hands `shutil.copyfile` to a worker thread now, which streams. An .eval log can be hundreds of megabytes. ## The tests for that change were mocking it out `test_directory_writer.py` patched `aiofiles.open`, now `anyio.open_file`, in all ten tests. So the trio variants passed without ever opening a file -- they exercised everything except the thing this PR changes. They write to `tmp_path` and assert on file contents now, which is both simpler (no `AsyncMock` `__aenter__` dance, no indexing into `write.call_args_list`) and real coverage of `anyio.open_file` under both backends. `test_file_encoding` asserted that `encoding="utf-8"` was passed. It now writes `café ✅ 日本語` and checks it survives the round trip, including as raw bytes in the file -- it fails if `ensure_ascii=False` is dropped, which the kwarg assertion did not. ## The infra setting -- a real limitation, not a test artefact Every infra task reaches the cluster through `kubernetes-asyncio`, which is built on aiohttp, and aiohttp is asyncio-only: ``` kubernetes_asyncio/client/api_client.py:76 kubernetes_asyncio/client/rest.py:68 aiohttp/connector.py:313: RuntimeError: no running event loop ``` Most of the 67 were that RuntimeError. The rest were assertion failures like `assert None == {'access-key': ...}`, which looked unrelated until traced -- a broad `except Exception` in the task had swallowed the same error and returned `None`. A conftest under `settings/infra/` skips trio variants there with that reason, and the README now records it. **The infra setting genuinely does not work under trio** and will not until its Kubernetes client is replaced. That conftest is a third of its previous size. It uses `pytest_runtest_setup`, which pytest calls only for the conftests covering an item's path, so it needs none of the manual `INFRA_DIR` filtering the collection hook required -- `pytest_collection_modifyitems` is handed *every* collected item, and the first version of this file, which forgot to filter, silently skipped the entire trio suite. It also shares `anyio_backend_of` with the root conftest instead of being a third copy of it. <details> <summary>A dead end worth recording</summary> I started migrating our own direct `aiohttp.ClientSession` use in those tasks to httpx, and got four files in before pulling a full traceback and discovering the failures come from `kubernetes-asyncio`, not from our calls. Four of the six failing task modules do not use aiohttp directly at all. I reverted it: it would have fixed nothing, while changing scorer code that cannot be exercised without a cluster. (It also introduced a bug on the way -- `as client` shadowing the imported `kubernetes_asyncio.client` module, which would have been an `UnboundLocalError` at a line executed earlier in the same function.) Separately worth knowing: `aiohttp` is not a declared dependency of this project. It arrives via `kubernetes-asyncio` and `miniopy-async`, and several task modules import it directly. </details> ## CI runs trio now CI ran `pytest -n 8 -m "not skip_ci and not long_running"` with no `--runtrio`, so **nothing ran these tests**. Going from 102 failures to 0 without a job to hold it there is a snapshot, not a fix -- 102 is what silent drift looks like. There is a `Unit Tests (trio)` job now, alongside a note in CONTRIBUTING that both invocations are needed for full coverage. Also regenerates `THIRD_PARTY_LICENSES.md`, which still listed `types-aiofiles` and was failing the `checks` job. ## Verification - Trio: **214 passed, 0 failed**. Asyncio: **1132 passed**. - `ruff check`, `ruff format`, `pyright` (0 errors), and the licence check all clean. - Up to date with main, including #851 and #860. 🤖 Generated with [Claude Code](https://claude.com/claude-code) --------- Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
Follow-up to #850, which noticed
from asyncio import gatheriniac_fast/scorers.py. It turned out to be seven scorers, not one.CLAUDE.md requires anyio rather than asyncio primitives because the backend is chosen at runtime by
INSPECT_ASYNC_BACKEND. asyncio primitives raise under trio, and since the default is asyncio the breakage only surfaces in trio runs — which is how this went unnoticed.The bug is real
End to end, against a scorer that runs with a mock model and no sandbox:
What changed
from asyncio import gather—apps,bash,bigcodebench,data_poisoning,iac,iac_fast,sae_interp_sabotage— switched tocontrol_arena._async_utils.gather. Scorers run inside inspect's event loop, so these were live breakage.examples/paper_replications/ai_control, which run inside the loop for the same reason.asyncio.as_completedreplaced by a new_async_utils.as_completed_with_exceptions, an anyio task group over a memory object stream. This one is reached fromcli.py, which already runs under the configured backend, so it was breaking for real. The replacement preserves completion-order streaming (results are written to CSV as they land) and keeps a failing coroutine from cancelling its siblings — an anyio task group cancels the whole group by default, so that is the property most easily lost in this conversion, and it is tested.asyncio.runentry points inexamples/andscripts/. These owned their own event loop so they were self-consistent rather than broken, but CLAUDE.md asks entry points to honour the configured backend, and leaving them would have meant exempting them from the new lint rule.New:
_async_utils.run_syncanyio.runcannot replaceasyncio.rundirectly — it takes a function plus positional arguments, so it cannot accept a coroutine built with keyword arguments, which is what existingasyncio.run(f(a, b=c))call sites have.run_sync(coro)takes the coroutine itself and runs it on the configured backend, so those call sites convert withoutfunctools.partialscattered around.Preventing recurrence
ruff now bans the asyncio primitives via
TID251, with a per-symbol message naming the replacement:Covered:
gather,as_completed,create_task,run,sleep,wait_for,Semaphore. It catches both thefrom asyncio import gatherform andasyncio.gather(...)— in fact it is what surfaced theexamples/andscripts/call sites I had not planned to touch.kubernetes_asyncioandtqdm.asyncioare unaffected; the latter is only used to construct progress bars, never for its asyncio-specific classmethods.Tests
run_syncunder both backends, including that it honoursINSPECT_ASYNC_BACKENDand accepts keyword-built coroutines.gather: a single failure raised as itself with its__cause__chain intact, several raised as a group.as_completed_with_exceptions: completion order, empty input, generator input,Noneresults, a failure reported without cancelling its siblings, early exit cancelling the rest, an exception from the loop body surfacing unwrapped, and both backends.process_commits_streamingitself: rows written in completion order, one failing commit not abandoning the others, skips counted.anyio_backend_of, the conftest classifier that decides which tests are trio variants.Verified:
ruff format,ruff check,pyright(0 errors), and 986 passing under asyncio.🤖 Generated with Claude Code