miquido-haiway-146
Concurrent operations misbehave when used inside a managed asynchronous task scope. `concurrently(..., return_exceptions=True)` must complete with one result per input coroutine, preserving successful values and placing raised exceptions in their corresponding positions. For example, within `async with ctx.scope("tg_concurrently")`, calling `await concurrently(coroutines, return_exceptions=True)` for `[good(1), bad(), good(2), bad()]` must return four results: `10`, a `FakeException` whose string is `"boom"`, `20`, and another `FakeException` whose string is `"boom"`. Individual failures must not cancel the surrounding scope or prevent the complete result sequence from being returned.
`execute_concurrently(handler, elements, return_exceptions=True)` has the same contract. Within `async with ctx.scope("tg_execute")`, a handler that raises `FakeException(f"odd {element}")` for odd elements and returns `element * 2` for even elements must produce six results in input order: successful even results and `FakeException` instances with messages `"odd 1"`, `"odd 3"`, and `"odd 5"` in the corresponding positions. Individual handler failures must not cancel the surrounding scope or prevent successful elements from completing.
`process_concurrently(source, handler, ignore_exceptions=True)` must continue processing after handler failures. Within `async with ctx.scope("tg_process")`, using `Source(range(10))` and a handler that raises `FakeException("odd")` for odd elements must process the even elements, so the recorded values are `[0, 2, 4, 6, 8]` in sorted order. The ignored failures must not cancel the surrounding scope.
For `stream_concurrently`, empty input iterators must finish cleanly without hanging or raising. Preserve the observable completion behavior for each source. With `exhaustive=False`, an exhausted source ends the iteration, while items already produced by other sources are still yielded: `stream_concurrently(range_of(0, 3), empty())` yields exactly `[0]`, and `stream_concurrently(empty(), range_of(0, 3))` also yields exactly `[0]`; two empty sources yield nothing. With `exhaustive=True`, only the exhausted source stops and the other continues, so the result is `[0, 1, 2]`. On exceptions or cancellation, pending source operations must not prevent termination, and cancellation cleanup for stopped streams must finish before `stream_concurrently` returns. `CancelledError` must not be exposed to the caller.
An asynchronous consumer may use `async for item in stream_concurrently(..., exhaustive=False)` inside `async with ctx.scope("test")`, and a task-group consumer may be started with `await ctx.spawn(consumer)` inside `async with ctx.scope("task_group_test")`; both must exit normally after an early-completing stream finishes. A fast stream yielding `1` and `2` must produce those values, and a stream yielding `0`, `1`, and `2` must produce those values even when paired with a slower stream. Multiple independent `stream_concurrently` calls in the same managed scope, such as consumers under `async with ctx.scope("multiple_streams")`, must complete independently: cancellation or early completion in one call must not disrupt the other. When a stopped stream handles `CancelledError` and performs asynchronous cleanup before re-raising, that cleanup must be allowed to complete before the operation returns.
Hidden tests · 8 fail-to-pass, 45 pass-to-passrun after the agent submits, in a clean verifier
Test patch · 237 lines
diff --git a/tests/test_concurrently.py b/tests/test_concurrently.py
index 70e7b750..4a2627c1 100644
--- a/tests/test_concurrently.py
+++ b/tests/test_concurrently.py
@@ -99,6 +99,25 @@ async def bad_coro() -> int:
assert str(results[3]) == "Test exception"
+@mark.asyncio
+async def test_return_exceptions_inside_scope_task_group():
+ async def good(value: int) -> int:
+ return value * 10
+
+ async def bad() -> int:
+ raise FakeException("boom")
+
+ coroutines = [good(1), bad(), good(2), bad()]
+ async with ctx.scope("tg_concurrently"):
+ results = await concurrently(coroutines, return_exceptions=True)
+
+ assert len(results) == 4
+ assert results[0] == 10
+ assert isinstance(results[1], FakeException) and str(results[1]) == "boom"
+ assert results[2] == 20
+ assert isinstance(results[3], FakeException) and str(results[3]) == "boom"
+
+
@mark.asyncio
async def test_cancels_running_tasks_on_cancellation():
started: list[int] = []
diff --git a/tests/test_execute_concurrently.py b/tests/test_execute_concurrently.py
index a7167d8b..6eb253a6 100644
--- a/tests/test_execute_concurrently.py
+++ b/tests/test_execute_concurrently.py
@@ -84,9 +84,27 @@ async def handler(element: int) -> int:
assert isinstance(results[3], FakeException)
assert str(results[3]) == "Test exception"
- # Check that other elements returned correct values
- for i in [0, 1, 2, 4, 5]:
- assert results[i] == i * 2
+
+@mark.asyncio
+async def test_return_exceptions_inside_scope_task_group():
+ async def handler(element: int) -> int:
+ if element % 2 == 1:
+ raise FakeException(f"odd {element}")
+ return element * 2
+
+ elements = list(range(6))
+ async with ctx.scope("tg_execute"):
+ results = await execute_concurrently(handler, elements, return_exceptions=True)
+
+ assert len(results) == 6
+ for i, res in enumerate(results):
+ if i % 2 == 1:
+ assert isinstance(res, FakeException)
+ assert str(res) == f"odd {i}"
+ else:
+ assert res == i * 2
+
+ # Verify values already covered above, no extra assertions here
@mark.asyncio
diff --git a/tests/test_process_concurrently.py b/tests/test_process_concurrently.py
index 917d2a5e..0deade32 100644
--- a/tests/test_process_concurrently.py
+++ b/tests/test_process_concurrently.py
@@ -101,6 +101,23 @@ async def handler(element: int) -> None:
assert sorted(processed) == [0, 1, 2, 4, 5]
+@mark.asyncio
+async def test_ignore_exceptions_inside_scope_task_group():
+ processed: list[int] = []
+
+ async def handler(element: int) -> None:
+ if element % 2 == 1:
+ raise FakeException("odd")
+ processed.append(element)
+
+ source = Source(range(10))
+ async with ctx.scope("tg_process"):
+ await process_concurrently(source, handler, ignore_exceptions=True)
+
+ # Only even elements processed; odd failures ignored without cancelling group
+ assert sorted(processed) == [0, 2, 4, 6, 8]
+
+
@mark.asyncio
async def test_handles_source_exception():
processed: list[int] = []
@@ -138,7 +155,7 @@ async def slow_handler(element: int) -> None:
slow_handler,
)
# Give some time for tasks to start
- await sleep(0.1)
+ await sleep(0)
# Cancel the main task
task.cancel()
await task
diff --git a/tests/test_stream_concurrently.py b/tests/test_stream_concurrently.py
index 83782026..9833b1db 100644
--- a/tests/test_stream_concurrently.py
+++ b/tests/test_stream_concurrently.py
@@ -94,7 +94,7 @@ async def empty_iter() -> AsyncIterator[int]:
items = []
async for item in stream_concurrently(empty_iter(), async_range(0, 3)):
items.append(item)
- assert items == []
+ assert items == [0]
# exhaustive
items = []
@@ -313,3 +313,121 @@ async def iter_b() -> AsyncIterator[str]:
b_yield_1_pos = execution_order.index("b_yield_1")
a_yield_1_pos = execution_order.index("a_yield_1")
assert b_yield_1_pos < a_yield_1_pos
+
+
+@mark.asyncio
+async def test_early_completion_does_not_raise_cancellation_error():
+ async def fast_iter() -> AsyncIterator[int]:
+ yield 1
+ yield 2
+
+ async def slow_iter() -> AsyncIterator[str]:
+ for i in range(100):
+ await sleep(0.1) # Very slow
+ yield f"item_{i}"
+
+ # This should complete when fast_iter finishes without raising CancelledError
+ items: list[int | str] = []
+
+ async with ctx.scope("test"):
+ async for item in stream_concurrently(fast_iter(), slow_iter(), exhaustive=False):
+ items.append(item)
+
+ # Should have items from fast iterator and possibly some from slow
+ assert len(items) >= 2
+ assert 1 in items
+ assert 2 in items
+
+
+@mark.asyncio
+async def test_task_group_integration_with_early_completion():
+ results = []
+
+ async def consumer():
+ async def numbers() -> AsyncIterator[int]:
+ for i in range(3):
+ await sleep(0.01)
+ yield i
+
+ async def letters() -> AsyncIterator[str]:
+ for c in "abcdefghijklmnopqrstuvwxyz": # Long stream
+ await sleep(0.05) # Slower than numbers
+ yield c
+
+ items = []
+ async for item in stream_concurrently(numbers(), letters(), exhaustive=False):
+ items.append(item)
+
+ results.extend(items)
+
+ # Run in task group - should not raise CancelledError
+ async with ctx.scope("task_group_test"):
+ await ctx.spawn(consumer)
+
+ # Should have completed successfully with items from numbers stream
+ assert len(results) >= 3
+ assert 0 in results
+ assert 1 in results
+ assert 2 in results
+
+
+@mark.asyncio
+async def test_cleanup_awaits_cancelled_tasks():
+ cleanup_completed = False
+
+ async def slow_stream() -> AsyncIterator[int]:
+ try:
+ for i in range(1000):
+ await sleep(0.1)
+ yield i
+ except CancelledError:
+ nonlocal cleanup_completed
+ # Simulate some cleanup work
+ await sleep(0.01)
+ cleanup_completed = True
+ raise
+
+ async def fast_stream() -> AsyncIterator[str]:
+ yield "done"
+
+ items = []
+
+ # This should not raise CancelledError even though slow_stream gets cancelled
+ async with ctx.scope("cleanup_test"):
+ async for item in stream_concurrently(fast_stream(), slow_stream(), exhaustive=False):
+ items.append(item)
+
+ # Should complete successfully
+ assert "done" in items
+ # Cleanup should have been called (though timing dependent)
+ # Note: We don't assert cleanup_completed=True because cancellation timing is non-deterministic
+
+
+@mark.asyncio
+async def test_multiple_stream_concurrently_in_task_group():
+ items_1 = []
+ items_2 = []
+
+ async def nums() -> AsyncIterator[int]:
+ for i in [1, 2]:
+ yield i
+
+ async def slow() -> AsyncIterator[str]:
+ for _ in range(100):
+ await sleep(0.1)
+ yield "slow"
+
+ async def consumer_1():
+ async for item in stream_concurrently(nums(), slow(), exhaustive=False):
+ items_1.append(item)
+
+ async def consumer_2():
+ async for item in stream_concurrently(nums(), slow(), exhaustive=False):
+ items_2.append(item)
+
+ async with ctx.scope("multiple_streams"):
+ await consumer_1()
+ await consumer_2()
+
+ assert 1 in items_1 and 2 in items_1
+ assert 1 in items_2 and 2 in items_2
Reference fix · 4 files, +406 −375the upstream merge, used only for grading calibration
The agent could not see this: the repository holds one commit and the sandbox has no network. Leak audit.
AGENTS.md, pyproject.toml, src/haiway/helpers/concurrent.py, uv.lock
diff --git a/AGENTS.md b/AGENTS.md
new file mode 100644
index 00000000..4ca84145
--- /dev/null
+++ b/AGENTS.md
@@ -0,0 +1,39 @@
+# Repository Guidelines
+
+## Project Structure & Modules
+- `src/haiway/`: Library code (e.g., `context/`, `state/`, `types/`, `utils/`, `helpers/`, `httpx/`, `opentelemetry/`).
+- `tests/`: Pytest suite (`test_*.py`), includes async tests.
+- `docs/` + `mkdocs.yml`: User/developer docs; built to `site/`.
+- `config/pre-push`: Git hook run by `make venv` to block WIP commits and enforce linting.
+- `.github/workflows/`: CI for lint, tests, build, docs.
+
+## Build, Test, and Development Commands
+- `make venv`: Install uv, create `.venv`, install all extras, prepare git hooks.
+- `make sync` / `make update`: Sync or upgrade dependencies via uv.
+- `make format`: Auto-fix and format with Ruff.
+- `make lint`: Security (bandit), lint (ruff), strict typing (pyright).
+- `make test`: Run pytest with coverage for `src/`.
+- `make docs-server` / `make docs`: Serve or build MkDocs site.
+- `make release`: Lint + test, then `uv build` and publish.
+ Example: `uv run pytest --rootdir= ./tests --doctest-modules` (matches CI).
+
+## Coding Style & Naming
+- Python 3.12+, line length 100, Ruff for lint/format (includes import sorting).
+- Strict typing enforced by Pyright; prefer explicit types and TypedDict/Protocol where useful.
+- Names: modules `snake_case`, classes `PascalCase`, functions/vars `snake_case`, constants `UPPER_CASE`.
+
+## Testing Guidelines
+- Framework: `pytest` (+ `pytest-asyncio`, `pytest-cov`).
+- Location/pattern: place tests under `tests/` and name `test_*.py`.
+- Async: mark with `@pytest.mark.asyncio`.
+- CI also runs doctests; keep examples executable.
+ Run locally: `make test` or `uv run pytest -v`.
+
+## Commit & Pull Request Guidelines
+- Commits: imperative, concise subject (≤72 chars), meaningful body; avoid "WIP" (pre-push hook blocks it).
+- PRs: clear description, linked issues (`#123`), include tests/docs for behavior changes; add screenshots for docs when useful.
+- CI must pass across OS matrix and Python 3.12/3.13; document any breaking changes in PR.
+
+## Security & Configuration Tips
+- Do not commit secrets; `.env` is optional for local use.
+- Use `make venv` to install the pre-push hook and ensure consistent tooling.
diff --git a/pyproject.toml b/pyproject.toml
index bb2f5545..bcc17ce0 100644
--- a/pyproject.toml
+++ b/pyproject.toml
@@ -5,7 +5,7 @@ build-backend = "hatchling.build"
[project]
name = "haiway"
description = "Framework for dependency injection and state management within structured concurrency model."
-version = "0.31.0"
+version = "0.31.1"
readme = "README.md"
maintainers = [
{ name = "Kacper Kaliński", email = "kacper.kalinski@miquido.com" },
diff --git a/src/haiway/helpers/concurrent.py b/src/haiway/helpers/concurrent.py
index 995a3f55..511a04fa 100644
--- a/src/haiway/helpers/concurrent.py
+++ b/src/haiway/helpers/concurrent.py
@@ -3,8 +3,10 @@
AsyncIterable,
AsyncIterator,
Callable,
+ Collection,
Coroutine,
Iterable,
+ Iterator,
MutableSequence,
MutableSet,
Sequence,
@@ -21,10 +23,10 @@
)
-async def process_concurrently[Element]( # noqa: C901
+async def process_concurrently[Element]( # noqa: C901, PLR0912
source: AsyncIterable[Element] | Iterable[Element],
/,
- handler: Callable[[Element], Coroutine[Any, Any, None]],
+ handler: Callable[[Element], Coroutine[None, None, None]],
*,
concurrent_tasks: int = 2,
ignore_exceptions: bool = False,
@@ -45,7 +47,7 @@ async def process_concurrently[Element]( # noqa: C901
source : AsyncIterable[Element] | Iterable[Element]
An iterable providing elements to process. Elements are consumed
one at a time as processing slots become available.
- handler : Callable[[Element], Coroutine[Any, Any, None]]
+ handler : Callable[[Element], Coroutine[None, None, None]]
A coroutine function that processes each element. The handler should
not return a value (returns None).
concurrent_tasks : int, default=2
@@ -80,41 +82,53 @@ async def process_concurrently[Element]( # noqa: C901
"""
assert concurrent_tasks > 0 # nosec: B101
- tasks: MutableSet[Task[None]] = set()
+ tasks: MutableSet[Task[Exception | None]] = set()
async def process(
element: Element,
/,
- ) -> None:
- nonlocal tasks
- tasks.add(ctx.spawn(handler, element))
- if len(tasks) < concurrent_tasks:
- return # keep spawning tasks
-
- completed, tasks = await wait(
- tasks,
- return_when=FIRST_COMPLETED,
- )
-
- for task in completed:
- if exc := task.exception():
- if not ignore_exceptions:
- raise exc
+ ) -> Exception | None:
+ try:
+ await handler(element)
- ctx.log_error(
- f"Concurrent processing error - {type(exc)}: {exc}",
- exception=exc,
- )
+ except Exception as exc:
+ if not ignore_exceptions:
+ return exc
+
+ ctx.log_error(
+ f"Concurrent processing error - {type(exc)}: {exc}",
+ exception=exc,
+ )
try:
if isinstance(source, AsyncIterable):
async for element in source:
- await process(element)
+ tasks.add(ctx.spawn(process, element))
+ if len(tasks) < concurrent_tasks:
+ continue # keep spawning tasks
+
+ completed, tasks = await wait(
+ tasks,
+ return_when=FIRST_COMPLETED,
+ )
+ for task in completed:
+ if exc := task.result():
+ raise exc
else:
assert isinstance(source, Iterable) # nosec: B101
for element in source:
- await process(element)
+ tasks.add(ctx.spawn(process, element))
+ if len(tasks) < concurrent_tasks:
+ continue # keep spawning tasks
+
+ completed, tasks = await wait(
+ tasks,
+ return_when=FIRST_COMPLETED,
+ )
+ for task in completed:
+ if exc := task.result():
+ raise exc
except BaseException as exc:
# Cancel all running tasks
@@ -123,27 +137,20 @@ async def process(
raise exc
- if not tasks:
- return
-
- completed, _ = await wait(
- tasks,
- return_when=ALL_COMPLETED,
- )
- for task in completed:
- if exc := task.exception():
- if not ignore_exceptions:
- raise exc
-
- ctx.log_error(
- f"Concurrent processing error - {type(exc)}: {exc}",
- exception=exc,
+ else:
+ if tasks:
+ completed, _ = await wait(
+ tasks,
+ return_when=ALL_COMPLETED,
)
+ for task in completed:
+ if exc := task.result():
+ raise exc
@overload
async def execute_concurrently[Element, Result](
- handler: Callable[[Element], Coroutine[Any, Any, Result]],
+ handler: Callable[[Element], Coroutine[None, None, Result]],
/,
elements: AsyncIterable[Element] | Iterable[Element],
*,
@@ -153,23 +160,23 @@ async def execute_concurrently[Element, Result](
@overload
async def execute_concurrently[Element, Result](
- handler: Callable[[Element], Coroutine[Any, Any, Result]],
+ handler: Callable[[Element], Coroutine[None, None, Result]],
/,
elements: AsyncIterable[Element] | Iterable[Element],
*,
concurrent_tasks: int = 2,
return_exceptions: Literal[True],
-) -> Sequence[Result |
… [96316 more characters]GLM-5.3 Flash · attempt 2 · solved · 31 steps · $0.015
Submitted diff · 443 lines
diff --git a/src/haiway/helpers/concurrent.py b/src/haiway/helpers/concurrent.py
index 995a3f5..d4e5c75 100644
--- a/src/haiway/helpers/concurrent.py
+++ b/src/haiway/helpers/concurrent.py
@@ -1,4 +1,4 @@
-from asyncio import ALL_COMPLETED, FIRST_COMPLETED, Task, wait
+from asyncio import ALL_COMPLETED, FIRST_COMPLETED, Task, gather, wait
from collections.abc import (
AsyncIterable,
AsyncIterator,
@@ -21,6 +21,40 @@ __all__ = (
)
+class _Exhausted:
+ """Internal sentinel marking an exhausted source iterator."""
+
+ __slots__ = ()
+
+ def __repr__(self) -> str:
+ return "EXHAUSTED"
+
+
+_EXHAUSTED: _Exhausted = _Exhausted()
+
+
+def _spawn_protected[Result](
+ coro: Coroutine[None, None, Result],
+ /,
+) -> Task[Result | BaseException]:
+ """Spawn a task which captures exceptions instead of propagating them.
+
+ Wrapping coroutines this way prevents individual failures from escaping
+ into the surrounding task group (which would cancel the whole scope) and
+ allows the caller to inspect and handle them.
+
+ """
+
+ async def protected() -> Result | BaseException:
+ try:
+ return await coro
+
+ except Exception as exc:
+ return exc
+
+ return ctx.spawn(protected())
+
+
async def process_concurrently[Element]( # noqa: C901
source: AsyncIterable[Element] | Iterable[Element],
/,
@@ -79,6 +113,7 @@ async def process_concurrently[Element]( # noqa: C901
... )
"""
+
assert concurrent_tasks > 0 # nosec: B101
tasks: MutableSet[Task[None]] = set()
@@ -87,7 +122,11 @@ async def process_concurrently[Element]( # noqa: C901
/,
) -> None:
nonlocal tasks
- tasks.add(ctx.spawn(handler, element))
+ if ignore_exceptions:
+ tasks.add(_spawn_protected(handler(element)))
+ else:
+ tasks.add(ctx.spawn(handler, element))
+
if len(tasks) < concurrent_tasks:
return # keep spawning tasks
@@ -97,14 +136,15 @@ async def process_concurrently[Element]( # noqa: C901
)
for task in completed:
- if exc := task.exception():
- if not ignore_exceptions:
- raise exc
-
- ctx.log_error(
- f"Concurrent processing error - {type(exc)}: {exc}",
- exception=exc,
- )
+ if ignore_exceptions:
+ if isinstance(exc := task.result(), BaseException):
+ ctx.log_error(
+ f"Concurrent processing error - {type(exc)}: {exc}",
+ exception=exc,
+ )
+
+ elif exc := task.exception():
+ raise exc
try:
if isinstance(source, AsyncIterable):
@@ -131,14 +171,15 @@ async def process_concurrently[Element]( # noqa: C901
return_when=ALL_COMPLETED,
)
for task in completed:
- if exc := task.exception():
- if not ignore_exceptions:
- raise exc
+ if ignore_exceptions:
+ if isinstance(exc := task.result(), BaseException):
+ ctx.log_error(
+ f"Concurrent processing error - {type(exc)}: {exc}",
+ exception=exc,
+ )
- ctx.log_error(
- f"Concurrent processing error - {type(exc)}: {exc}",
- exception=exc,
- )
+ elif exc := task.exception():
+ raise exc
@overload
@@ -162,14 +203,14 @@ async def execute_concurrently[Element, Result](
) -> Sequence[Result | BaseException]: ...
-async def execute_concurrently[Element, Result]( # noqa: C901, PLR0912
+async def execute_concurrently[Element, Result]( # noqa: C901
handler: Callable[[Element], Coroutine[Any, Any, Result]],
/,
elements: AsyncIterable[Element] | Iterable[Element],
*,
concurrent_tasks: int = 2,
return_exceptions: bool = False,
-) -> Sequence[Result | BaseException] | Sequence[Result]:
+) -> Sequence[Result] | Sequence[Result | BaseException]:
"""Execute handler for each element from a collection concurrently.
Processes all elements from a collection using the provided handler function,
@@ -239,9 +280,10 @@ async def execute_concurrently[Element, Result]( # noqa: C901, PLR0912
... print(f"Got data from {url}")
"""
+
assert concurrent_tasks > 0 # nosec: B101
- tasks: MutableSet[Task[Result]] = set()
- results: MutableSequence[Task[Result]] = []
+ tasks: MutableSet[Task[Result | BaseException]] = set()
+ results: MutableSequence[Task[Result | BaseException]] = []
async def process(
element: Element,
@@ -249,7 +291,9 @@ async def execute_concurrently[Element, Result]( # noqa: C901, PLR0912
) -> None:
nonlocal tasks
nonlocal results
- task: Task[Result] = ctx.spawn(handler, element)
+ task: Task[Result | BaseException] = (
+ _spawn_protected(handler(element)) if return_exceptions else ctx.spawn(handler, element)
+ )
results.append(task)
tasks.add(task)
if len(tasks) < concurrent_tasks:
@@ -261,14 +305,15 @@ async def execute_concurrently[Element, Result]( # noqa: C901, PLR0912
)
for task in completed:
- if exc := task.exception():
- if not return_exceptions:
- raise exc
-
- ctx.log_error(
- f"Concurrent execution error - {type(exc)}: {exc}",
- exception=exc,
- )
+ if return_exceptions:
+ if isinstance(exc := task.result(), BaseException):
+ ctx.log_error(
+ f"Concurrent execution error - {type(exc)}: {exc}",
+ exception=exc,
+ )
+
+ elif exc := task.exception():
+
… [8839 more characters]$ find /repo -type f -name "*.py" | head -50; ls /repo
output · 48 lines
/repo/src/haiway/context/access.py /repo/src/haiway/context/disposables.py /repo/src/haiway/context/events.py /repo/src/haiway/context/identifier.py /repo/src/haiway/context/observability.py /repo/src/haiway/context/presets.py /repo/src/haiway/context/state.py /repo/src/haiway/context/tasks.py /repo/src/haiway/context/types.py /repo/src/haiway/context/variables.py /repo/src/haiway/context/__init__.py /repo/src/haiway/helpers/asynchrony.py
36 more lines
/repo/src/haiway/helpers/caching.py /repo/src/haiway/helpers/concurrent.py /repo/src/haiway/helpers/configuration.py /repo/src/haiway/helpers/files.py /repo/src/haiway/helpers/http_client.py /repo/src/haiway/helpers/observability.py /repo/src/haiway/helpers/retries.py /repo/src/haiway/helpers/throttling.py /repo/src/haiway/helpers/timeouting.py /repo/src/haiway/helpers/__init__.py /repo/src/haiway/httpx/client.py /repo/src/haiway/httpx/__init__.py /r …[599 chars omitted]… ormatting.py /repo/src/haiway/utils/logs.py /repo/src/haiway/utils/metadata.py /repo/src/haiway/utils/mimic.py /repo/src/haiway/utils/noop.py /repo/src/haiway/utils/queue.py /repo/src/haiway/utils/stream.py /repo/src/haiway/utils/__init__.py /repo/src/haiway/__init__.py /repo/tests/test_async_queue.py CLAUDE.md LICENSE Makefile README.md config docs llms.txt mkdocs.yml pyproject.toml src tests uv.lock
$ cat /repo/src/haiway/helpers/concurrent.py; ls /repo/tests
output · 56 lines
from asyncio import ALL_COMPLETED, FIRST_COMPLETED, Task, wait
from collections.abc import (
AsyncIterable,
AsyncIterator,
Callable,
Coroutine,
Iterable,
MutableSequence,
MutableSet,
Sequence,
)
from typing import Any, Literal, overload44 more lines
from haiway.context import ctx
__all__ = (
"concurrently",
"execute_concurrently",
"process_concurrently",
"stream_concurrently",
)
async def process_concurrently[Element]( # noqa: C901
source: AsyncIterable[Element] | Iterable[Element],
/,
handler: Callable[[Element], Coroutine[Any, Any, None]],
*,
concurrent_tasks: int = 2,
ignore_exceptions: bool = False,
) -> None:
"""Process elements from an iterable concurrently.
Consumes elements from an iterable and processes them using the provided
handler function. Processing happens concurrently with a configurable maxi
…[8630 chars omitted]…
ream.py
test_attribute_path.py
test_attribute_requirement.py
test_auto_retry.py
test_cache.py
test_concurrently.py
test_context.py
test_context_presets.py
test_disposables.py
test_events.py
test_execute_concurrently.py
test_metadata.py
test_process_concurrently.py
test_state.py
test_state_typing.py
test_state_validation.py
test_stream_concurrently.py
test_streaming.py
test_timeout.py
test_variables.py$ sed -n 200,420p /repo/src/haiway/helpers/concurrent.py
output · 41 lines
processing and is raised.
Returns
-------
Sequence[Result] or Sequence[Result | BaseException]
Results from each handler invocation, in the same order as input elements.
If return_exceptions is True, failed tasks return BaseException instances.
Raises
------
CancelledError
If the function is cancelled, propagated after cancelling all running tasks.29 more lines
Exception
Any exception raised by handler tasks when return_exceptions is False.
Examples
--------
>>> async def fetch_data(url: str) -> dict:
... return await http_client.get(url)
...
>>> urls = ["http://api.example.com/1", "http://api.example.com/2"]
>>> results = await execute_concurrently(
... fetch_data,
... urls,
... concurrent_tasks=10
... )
>>> # results[0] corresponds to urls[0], results[1] to urls[
…[5900 chars omitted]…
utines,
... concurrent_tasks=2,
... return_exceptions=True
... )
>>> for i, result in enumerate(results):
... if isinstance(result, BaseException):
... print(f"Coroutine {i} failed: {result}")
... else:
... print(f"Coroutine {i} succeeded")
"""
assert concurrent_tasks > 0 # nosec: B101
tasks: MutableSet[Task[Result]] = set()$ sed -n 420,560p /repo/src/haiway/helpers/concurrent.py; cat /repo/src/haiway/context/tasks.py
output · 53 lines
tasks: MutableSet[Task[Result]] = set()
results: MutableSequence[Task[Result]] = []
async def process(
element: Coroutine[None, None, Result],
/,
) -> None:
nonlocal tasks
nonlocal results
task: Task[Result] = ctx.spawn(element)
results.append(task)
tasks.add(task)41 more lines
if len(tasks) < concurrent_tasks:
return # keep spawning tasks
completed, tasks = await wait(
tasks,
return_when=FIRST_COMPLETED,
)
for task in completed:
if exc := task.exception():
if not return_exceptions:
raise exc
ctx.log_error(
f"Concurrent execution error - {type(exc)}: {exc}",
exception=exc,
)
try:
if isinstance(coroutines, AsyncIterable):
…[8647 chars omitted]…
exc_type,
exc_val,
exc_tb,
)
except ExceptionGroup:
pass # skip task group exceptions
finally:
object.__setattr__(
self,
"_token",
None,
)
object.__setattr__(
self,
"_group",
None,
)$ cat /repo/tests/test_stream_concurrently.py; echo ====; sed -n 1,80p /repo/tests/test_concurrently.py
output · 59 lines
from asyncio import CancelledError, sleep
from collections.abc import AsyncIterator
from pytest import mark, raises
from haiway import ctx
from haiway.helpers.concurrent import stream_concurrently
class FakeException(Exception):
pass
47 more lines
async def async_range(
start: int,
stop: int,
delay: float = 0,
) -> AsyncIterator[int]:
for i in range(start, stop):
if delay:
await sleep(delay)
yield i
async def async_letters(
letters: str,
delay: float = 0,
) -> AsyncIterator[str]:
for letter in letters:
if delay:
await sleep(delay)
yield letter
@mark.asyncio
async def test_merges_two_streams():
items: list[int | str] = []
async for item in stream_concurrently(async_range(0, 3), async_letters("abc")):
items.append(item)
# Should have all items from both sources
assert len(items)
…[8629 chars omitted]…
coroutines = [good_coro(), bad_coro(), good_coro()]
with raises(FakeException):
await concurrently(coroutines)
@mark.asyncio
async def test_returns_exceptions_when_configured():
async def good_coro(value: int) -> int:
return value * 2
async def bad_coro() -> int:
raise FakeException("Test exception")
coroutines = [
good_coro(1), # Should return 2$ grep -n '"""\n' /repo/src/haiway/helpers/concurrent.py | head; grep -n '^ """$' /repo/src/haiway/helpers/concurrent.py; wc -l /repo/src/haiway/helpers/concurrent.py; cd /repo && python -m pytest tests/test_concurrently.py tests/test_execute_concurrently.py tests/test_process_concurrently.py tests/test_stream_concurrently.py -q 2>&1 | tail -5
output · 10 lines
81: """ 241: """ 418: """ 557: """ 618 /repo/src/haiway/helpers/concurrent.py Enable tracemalloc to get traceback where the object was allocated. See https://docs.pytest.org/en/stable/how-to/capture-warnings.html#resource-warnings for more info. -- Docs: https://docs.pytest.org/en/stable/how-to/capture-warnings.html 46 passed, 2 warnings in 4.10s
$ cd /repo && cat > /tmp/repro.py <<'EOF'
output · 29 lines
Traceback (most recent call last):
File "/tmp/repro.py", line 54, in <module>
asyncio.run(main())
File "/usr/local/lib/python3.12/asyncio/runners.py", line 195, in run
return runner.run(main)
^^^^^^^^^^^^^^^^
File "/usr/local/lib/python3.12/asyncio/runners.py", line 118, in run
return self._loop.run_until_complete(task)
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
File "/usr/local/lib/python3.12/asyncio/base_events.py", line 691, in run_until_complete
return future.result()
^^^^^^^^^^^^^^^17 more lines
File "/tmp/repro.py", line 12, in main
res = await concurrently(coros, return_exceptions=True)
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
File "/repo/src/haiway/helpers/concurrent.py", line 465, in concurrently
raise exc
File "/repo/src/haiway/helpers/concurrent.py", line 458, in concurrently
await process(element)
…[180 chars omitted]…
tasks.py", line 464, in wait
return await _wait(fs, timeout, return_when, loop)
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
File "/usr/local/lib/python3.12/asyncio/tasks.py", line 550, in _wait
await waiter
asyncio.exceptions.CancelledError
sys:1: RuntimeWarning: coroutine 'main.<locals>.bad' was never awaited
sys:1: RuntimeWarning: coroutine 'main.<locals>.good' was never awaitedNow I'll rewrite the module, preserving docstrings:
$ cd /repo && python - <<'PY'
output · 3 lines
Traceback (most recent call last): File "<stdin>", line 8, in <module> AttributeError: 'NoneType' object has no attribute 'group'
$ cd /repo && python - <<'PY'
output · 3 lines
Traceback (most recent call last): File "<stdin>", line 6, in <module> AttributeError: 'NoneType' object has no attribute 'group'
$ cd /repo && python - <<'PY'
output · 4 lines
process_concurrently 2004 execute_concurrently 2727 concurrently 2887 stream_concurrently 2044
$ cd /repo && python - <<'PY'
output · 28 lines
invalid-syntax: Unexpected indentation
--> src/haiway/helpers/concurrent.py:197:1
|
197 | """Execute handler for each element from a collection concurrently.
| ^^^^
198 |
199 | Processes all elements from a collection using the provided handler function,
|
invalid-syntax: Unexpected indentation
--> src/haiway/helpers/concurrent.py:358:1
|16 more lines
358 | """Execute multiple coroutines concurrently with controlled parallelism.
| ^^^^
359 |
360 | Executes a collection of coroutines concurrently, limiting the number of
|
Found 2 errors.
Traceback (most recent call last):
File "<string>", line 1, in <module>
File "/repo/src/haiway/__init__.py", line 17, in <module>
from haiway.helpers import (
File "/repo/src/haiway/helpers/__init__.py", line 3, in <module>
from haiway.helpers.concurrent import (
File "/repo/src/haiway/helpers/concurrent.py", line 197
"""Execute handler for each element from a collection concurrently.
IndentationError: unexpected indent$ cd /repo && python - <<'PY'
output · 10 lines
Traceback (most recent call last):
File "<string>", line 1, in <module>
File "/repo/src/haiway/__init__.py", line 17, in <module>
from haiway.helpers import (
File "/repo/src/haiway/helpers/__init__.py", line 3, in <module>
from haiway.helpers.concurrent import (
File "/repo/src/haiway/helpers/concurrent.py", line 142
async for element in source:
^^^^^^^^^^^^^^^^^^^^^^^^^^^^
SyntaxError: 'async for' outside async function$ cd /repo && sed -n 95,150p src/haiway/helpers/concurrent.py
output · 54 lines
...
>>> async def items() -> AsyncIterator[str]:
... for i in range(10):
... yield f"item_{i}"
...
>>> await process_concurrently(
... items(),
... process_item,
... concurrent_tasks=5
... )
"""42 more lines
assert concurrent_tasks > 0 # nosec: B101
tasks: MutableSet[Task[None]] = set()
async def process(
element: Element,
/,
) -> None:
nonlocal tasks
if ignore_exceptions:
tasks.add(_spawn_protected(handler(element)))
else:
tasks.add(ctx.spawn(handler, element))
if len(tasks) < concurrent_tasks:
return # keep spawning tasks
completed, tasks = await wait(
tasks,
return_when=FIRST_COMPLETED,
)
for task in completed:
if ignore_exceptions:
if isinstance(
…[197 chars omitted]…
)
elif exc := task.exception():
raise exc
try:
if isinstance(source, AsyncIterable):
async for element in source:
await process(element)
else:
assert isinstance(source, Iterable) # nosec: B101
for element in source:
await process(element)
except BaseException as exc:$ cd /repo && sed -n 1,40p src/haiway/helpers/concurrent.py
output · 40 lines
from asyncio import ALL_COMPLETED, FIRST_COMPLETED, Task, gather, wait
from collections.abc import (
AsyncIterable,
AsyncIterator,
Callable,
Coroutine,
Iterable,
MutableSequence,
MutableSet,
Sequence,
)
from typing import Any, Literal, overload28 more lines
from haiway.context import ctx
__all__ = (
"concurrently",
"execute_concurrently",
"process_concurrently",
"stream_concurrently",
)
class _Exhausted:
"""Internal sentinel marking an exhausted source iterator."""
__slots__ = ()
def __repr__(self) -> str:
return "EXHAUSTED"
_EXHAUSTED: _Exhausted = _Exhausted()
def _spawn_protected[Result](
coro: Coroutine[None, None, Result],
/,
) -> Task[Result | BaseException]:
"""Spawn a task which captures exceptions instead of propagating them.$ cd /repo && sed -n 40,70p src/haiway/helpers/concurrent.py
output · 31 lines
"""Spawn a task which captures exceptions instead of propagating them.
Wrapping coroutines this way prevents individual failures from escaping
into the surrounding task group (which would cancel the whole scope) and
allows the caller to inspect and handle them.
"""
async def protected() -> Result | BaseException:
try:
return await coro
19 more lines
except Exception as exc: # noqa: BLE001 - intentionally broad
return exc
return ctx.spawn(protected())
"""Process elements from an iterable concurrently.
Consumes elements from an iterable and processes them using the provided
handler function. Processing happens concurrently with a configurable maximum
number of concurrent tasks. Elements are processed as they become available,
maintaining the specified concurrency limit.
The function continues until the source iterator is exhausted. If the function
is cancelled, all running tasks are also cancelled. When ignore_exceptions is
False, the first exception encountered will stop processing and propagate.
Parameters
----------
source : AsyncIterable[Element] | Iterable[Element]$ cd /repo && python - <<'PY'
output · 10 lines
Traceback (most recent call last):
File "<string>", line 1, in <module>
File "/repo/src/haiway/__init__.py", line 17, in <module>
from haiway.helpers import (
File "/repo/src/haiway/helpers/__init__.py", line 3, in <module>
from haiway.helpers.concurrent import (
File "/repo/src/haiway/helpers/concurrent.py", line 513
return [task.result() for task in results]
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
SyntaxError: 'return' with value in async generator$ cd /repo && sed -n 495,530p src/haiway/helpers/concurrent.py
output · 36 lines
assert isinstance(coroutines, Iterable) # nosec: B101
for element in coroutines:
await process(element)
except BaseException as exc:
# Cancel all running tasks
for task in tasks:
task.cancel()
raise exc
if tasks:24 more lines
await wait(
tasks,
return_when=ALL_COMPLETED,
)
if return_exceptions:
return [task.result() for task in results]
return [task.result() for task in results] # type: ignore[arg-type]
"""Merge streams from two async iterators processed concurrently.
Concurrently consumes elements from two async iterators and yields them
as they become available. Elements from both sources are interleaved based
on which iterator produces them first. The function continues until both
iterators are exhausted.
This is useful for combining multiple async data sources into a single
stream while maintaining concurrency. Each iterator is polled independently,
and whichever has data available first will have its element yielded.
Parameters
----------$ cd /repo && python - <<'PY'
output · 28 lines
[68fec13c11ff467ab5b39edc75c6b0f8] [tg] [05f12cf1-d71c-418d-9c74-aab53fb3da51] Concurrent execution error - <class '__main__.FakeException'>: boom
Traceback (most recent call last):
File "/repo/src/haiway/helpers/concurrent.py", line 50, in protected
return await coro
^^^^^^^^^^
File "/tmp/repro.py", line 10, in bad
async def bad(): raise FakeException("boom")
^^^^^^^^^^^^^^^^^^^^^^^^^^^
FakeException: boom
[68fec13c11ff467ab5b39edc75c6b0f8] [tg] [05f12cf1-d71c-418d-9c74-aab53fb3da51] Concurrent execution error - <class '__main__.FakeException'>: boom
Traceback (most recent call last):
File "/repo/src/haiway/helpers/concurrent.py", line 50, in protected16 more lines
return await coro
^^^^^^^^^^
File "/tmp/repro.py", line 10, in bad
async def bad(): raise FakeException("boom")
^^^^^^^^^^^^^^^^^^^^^^^^^^^
FakeExc
…[1850 chars omitted]…
op.run_until_complete(task)
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
File "/usr/local/lib/python3.12/asyncio/base_events.py", line 691, in run_until_complete
return future.result()
^^^^^^^^^^^^^^^
File "/tmp/repro.py", line 21, in main
assert res2[1].args[0]=="odd 1" and res2[0]==2 and res2[5]==12
^^^^^^^^^^^^
AttributeError: 'int' object has no attribute 'args'$ cd /repo && python - <<'PY'
output · 7 lines
concurrently: [10, FakeException('boom'), 20, FakeException('boom')]
execute: [FakeException('odd 1'), 4, FakeException('odd 3'), 8, FakeException('odd 5'), 12]
process: [0, 2, 4, 6, 8]
stream1: [0]
stream2: [0]
stream3: [0, 1, 2]
OK$ cd /repo && timeout 300 python -m pytest tests/ -q -x 2>&1 | tail -20
output · 20 lines
tests/test_stream_concurrently.py:97: AssertionError
=============================== warnings summary ===============================
tests/test_concurrently.py::test_propagates_coroutine_exceptions
/usr/local/lib/python3.12/asyncio/events.py:36: RuntimeWarning: coroutine 'test_propagates_coroutine_exceptions.<locals>.good_coro' was never awaited
def __init__(self, callback, args, loop, context=None):
Enable tracemalloc to get traceback where the object was allocated.
See https://docs.pytest.org/en/stable/how-to/capture-warnings.html#resource-warnings for more info.
tests/test_concurrently.py::test_cancels_running_tasks_on_cancellation
/usr/local/lib/python3.12/asyncio/events.py:88: RuntimeWarning: coroutine 'test_cancels_running_tasks_on_cancellation.<locals>.slow_coro' was never awaited
self._context.run(self._callback, *self._args)
Enable tracemalloc to get trac8 more lines
…[93 chars omitted]… warnings.html#resource-warnings for more info. -- Docs: https://docs.pytest.org/en/stable/how-to/capture-warnings.html =========================== short test summary info ============================ FAILED tests/test_stream_concurrently.py::test_handles_empty_iterators - asse... !!!!!!!!!!!!!!!!!!!!!!!!!! stopping after 1 failures !!!!!!!!!!!!!!!!!!!!!!!!!!! 1 failed, 257 passed, 2 warnings in 5.12s
The only failing test is the empty-iterator case the issue explicitly redefines (`stream_concurrently(empty(), range)` → `[0]`). Now let me test the remaining issue scenarios:
$ cd /repo && cat > /tmp/repro2.py <<'EOF'
output · 1 lines
async-for: []
$ cd /repo && python - <<'PY'
output · 3 lines
async-for: [0, 1, 2] fast: [1, 2] spawned: []
$ cd /repo && python - <<'PY'
output · 3 lines
async-for: [0, 1, 2] fast: [1, 2] spawned: [0, 1, 2]
$ cd /repo && timeout 60 python /tmp/repro2.py 2>&1 | tail -30
output · 24 lines
async-for: [0, 1, 2]
fast: [1, 2]
spawned: [0, 1, 2]
Traceback (most recent call last):
File "/tmp/repro2.py", line 85, in <module>
asyncio.run(main())
File "/usr/local/lib/python3.12/asyncio/runners.py", line 195, in run
return runner.run(main)
^^^^^^^^^^^^^^^^
File "/usr/local/lib/python3.12/asyncio/runners.py", line 118, in run
return self._loop.run_until_complete(task)
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^12 more lines
File "/usr/local/lib/python3.12/asyncio/base_events.py", line 691, in run_until_complete
return future.result()
^^^^^^^^^^^^^^^
File "/tmp/repro2.py", line 58, in main
await asyncio.gather(t1,t2)
^^^^^^^^^^^^^^^^^^^^^
File "/usr/local/lib/python3.12/asyncio/tasks.py", line 831, in gather
fut = ensure_future(arg, loop=loop)
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
File "/usr/local/lib/python3.12/asyncio/tasks.py", line 689, in ensure_future
raise TypeError('An asyncio.Future, a coroutine or an awaitable '
TypeError: An asyncio.Future, a coroutine or an awaitable is required$ cd /repo && python - <<'PY'
output · 19 lines
async-for: [0, 1, 2]
fast: [1, 2]
spawned: [0, 1, 2]
Traceback (most recent call last):
File "/tmp/repro2.py", line 85, in <module>
asyncio.run(main())
File "/usr/local/lib/python3.12/asyncio/runners.py", line 195, in run
return runner.run(main)
^^^^^^^^^^^^^^^^
File "/usr/local/lib/python3.12/asyncio/runners.py", line 118, in run
return self._loop.run_until_complete(task)
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^7 more lines
File "/usr/local/lib/python3.12/asyncio/base_events.py", line 691, in run_until_complete
return future.result()
^^^^^^^^^^^^^^^
File "/tmp/repro2.py", line 58, in main
await asyncio.gather(t1.result(), t2.result())
^^^^^^^^^
AttributeError: 'NoneType' object has no attribute 'result'$ cd /repo && python - <<'PY'
output · 7 lines
async-for: [0, 1, 2]
fast: [1, 2]
spawned: [0, 1, 2]
multi: {'a': [0, 1, 2], 'b': [0, 1, 2]}
paired: [10, 0, 11, 12]
cleanup: ['done']
OK$ cd /repo && timeout 300 python -m pytest tests/ -q 2>&1 | tail -6; cat Makefile; ruff check src/haiway/helpers/concurrent.py; ruff format --check src/haiway/helpers/concurrent.py
output · 43 lines
See https://docs.pytest.org/en/stable/how-to/capture-warnings.html#resource-warnings for more info. -- Docs: https://docs.pytest.org/en/stable/how-to/capture-warnings.html =========================== short test summary info ============================ FAILED tests/test_stream_concurrently.py::test_handles_empty_iterators - asse... 1 failed, 290 passed, 2 warnings in 5.31s SHELL := sh .ONESHELL: .SHELLFLAGS := -eu -c .DELETE_ON_ERROR: SOURCES_PATH := src
31 more lines
TESTS_PATH := tests
# load environment config from .env if able
-include .env
ifndef UV_VERSION
UV_VERSION := 0.7.20
endif
.PHONY: uv_check venv sync update format lint test docs docs-server release
# Check installed UV version and install if needed
uv_check:
@echo 'Checking uv version...'
# Install if not present
@if ! command -v uv > /dev/null; then \
echo '...installing uv...'; \
curl -LsSf https://github.com/astra
…[6386 chars omitted]…
--> src/haiway/helpers/concurrent.py:295:47
|
294 | task: Task[Result | BaseException] = (
- _spawn_protected(handler(element))
- if return_exceptions
- else ctx.spawn(handler, element)
295 + _spawn_protected(handler(element)) if return_exceptions else ctx.spawn(handler, element)
296 | )
|
1 file would be reformatted$ cd /repo && python - <<'PY'
output · 8 lines
All checks passed!
1 file reformatted
All checks passed!
File "/usr/local/lib/python3.12/site-packages/pyright/node.py", line 54, in _ensure_node_env
_install_node_env()
File "/usr/local/lib/python3.12/site-packages/pyright/node.py", line 95, in _install_node_env
raise RuntimeError(
RuntimeError: nodeenv failed; for more reliable node.js binaries try `pip install pyright[nodejs]`$ cd /repo && timeout 60 python /tmp/repro.py 2>/dev/null | tail -1 && timeout 60 python /tmp/repro2.py 2>/dev/null | tail -1 && timeout 300 python -m pytest tests/ -q 2>&1 | tail -3 && git diff --stat
output · 7 lines
OK OK =========================== short test summary info ============================ FAILED tests/test_stream_concurrently.py::test_handles_empty_iterators - asse... 1 failed, 290 passed, 2 warnings in 5.37s src/haiway/helpers/concurrent.py | 291 ++++++++++++++++++++++----------------- 1 file changed, 164 insertions(+), 127 deletions(-)
$ cd /repo && git diff | head -260
output · 44 lines
diff --git a/src/haiway/helpers/concurrent.py b/src/haiway/helpers/concurrent.py
index 995a3f5..d4e5c75 100644
--- a/src/haiway/helpers/concurrent.py
+++ b/src/haiway/helpers/concurrent.py
@@ -1,4 +1,4 @@
-from asyncio import ALL_COMPLETED, FIRST_COMPLETED, Task, wait
+from asyncio import ALL_COMPLETED, FIRST_COMPLETED, Task, gather, wait
from collections.abc import (
AsyncIterable,
AsyncIterator,
@@ -21,6 +21,40 @@ __all__ = (
)32 more lines
+class _Exhausted:
+ """Internal sentinel marking an exhausted source iterator."""
+
+ __slots__ = ()
+
+ def __repr__(self) -> str:
+ return "EXHAUSTED"
+
+
+_EXHAUSTED: _Exhausted = _Exhausted()
+
+
+def _spawn_protected[Result](
+ coro: Coroutine[None, None, Result],
+ /,
+) -> Task[Result | BaseException]:
+ """Spawn a task which captures exceptions instead of propagating them.
+
+ Wrapping coroutines this way p
…[7220 chars omitted]…
nonlocal results
- task: Task[Result] = ctx.spawn(element)
+ task: Task[Result | BaseException] = (
+ _spawn_protected(element) if return_exceptions else ctx.spawn(element)
+ )
results.append(task)
tasks.add(task)
if len(tasks) < concurrent_tasks:
@@ -438,14 +470,15 @@ async def concurrently[Result]( # noqa: C901, PLR0912
)The diff is complete and correct. Summary of the fix in `src/haiway/helpers/concurrent.py`:
$ done: {"notes": "Fixed concurrent helpers to be safe inside managed task-group scopes: added _spawn_protected to capture exceptions in spawned tasks (preventing scope cancellation) used by concurrently/execute_concurrently with return_exceptions=True and process_concurrently with ignore_exceptions=True; stream_concurrently now uses an _Eoutput · 3 lines
…[168 chars omitted]… ue scenarios verified with reproducers; full test suite passes except the one empty-iterator assertion the issue explicitly redefines (empty() first now yields [0])."} submitted.