Skip to content

Await coroutine functions in concatmap/flatmap/switchmap (#129) - #140

Closed
apoorvdarshan wants to merge 1 commit into
vxgmichel:mainfrom
apoorvdarshan:fix/issue-129-flatmap-async-func
Closed

apoorvdarshan wants to merge 1 commit into
vxgmichel:mainfrom
apoorvdarshan:fix/issue-129-flatmap-async-func

Conversation

@apoorvdarshan

Copy link
Copy Markdown

Fixes #129

Root cause

concatmap, flatmap and switchmap build their intermediate stream via
combine.smap.raw — the synchronous map that yields func(*item) without
ever awaiting it:

# aiostream/stream/combine.py
async def smap(source, func, *more_sources):
    ...
    async for item in streamer:
        yield func(*item)          # not awaited

When func is an async def, func(*item) returns a coroutine object rather
than an AsyncIterable. That coroutine is passed downstream to base_combine,
where it trips assert isinstance(result, AsyncIterable) in
aiostream/stream/advanced.py, and a RuntimeWarning: coroutine '...' was never awaited is emitted. stream.map already handles this case by detecting
a coroutine function with asyncio.iscoroutinefunction and routing to the
async path; the advanced *map operators did not.

Reproduction

import asyncio
import warnings
from aiostream import stream, pipe


async def async_processor(x):
    await asyncio.sleep(0)
    return stream.iterate([x, x * 10])


async def main():
    warnings.simplefilter("error", RuntimeWarning)
    source = stream.range(0, 3)
    pipeline = source | pipe.flatmap(async_processor)
    async with pipeline.stream() as streamer:
        result = [item async for item in streamer]
    print("RESULT:", sorted(result))


asyncio.run(main())

Actual output before the fix:

FAILED: AssertionError:
Exception ignored while finalizing coroutine <coroutine object async_processor ...>:
RuntimeWarning: coroutine 'async_processor' was never awaited

Output after the fix (matching the equivalent synchronous function, which was
already supported):

RESULT: [0, 0, 1, 2, 10, 20]

Fix

Added a small _map_sources helper in aiostream/stream/advanced.py that
routes through combine.map.raw when the function is a coroutine function
(detected with asyncio.iscoroutinefunction, the same idiom combine.map
uses) and keeps using combine.smap.raw otherwise. All three advanced *map
operators now go through this helper, and their func annotation is broadened
from SmapCallable to MapCallable to reflect that a coroutine function is
now accepted (again mirroring stream.map).

The change is localized: the synchronous path is byte-for-byte equivalent to
the previous behaviour, and the async path reuses the already-tested map/
amap machinery, so the coroutine is awaited before its AsyncIterable
result is iterated. The task_limit argument is intentionally left to govern
the outer flatten/concat/switch stage, exactly as before.

Tests

Added test_concatmap_with_coroutine_function,
test_flatmap_with_coroutine_function and
test_switchmap_with_coroutine_function in tests/test_advanced.py. Each
passes an async def function to the operator and asserts the flattened
output. All three fail on the unmodified code (with the "coroutine was never
awaited" warning) and pass with the fix. The full existing suite continues to
pass (112 passed), and black, ruff, flake8, mypy --strict and
pyright are clean on the changed files.

Disclosure: prepared with AI assistance; reviewed and verified locally.

The advanced *map operators built the intermediate stream via
`combine.smap.raw`, the synchronous map that never awaits its function.
When passed an `async def` function, `smap` yielded the un-awaited
coroutine object downstream, which then failed the
`assert isinstance(result, AsyncIterable)` check in `base_combine`
(and emitted a "coroutine was never awaited" RuntimeWarning).

Route through `combine.map.raw` when the function is a coroutine function
(detected with `asyncio.iscoroutinefunction`, as `combine.map` already
does), and keep using `combine.smap.raw` otherwise. This mirrors the
existing behaviour of `stream.map` and awaits the coroutine before its
result is iterated.

Fixes vxgmichel#129
@vxgmichel

Copy link
Copy Markdown
Owner

Hi @apoorvdarshan, thank you for taking the time to make this PR :)

It turns out that the fix you presented does not respect the task_limit argument, as calling combine.map would add new tasks not handled by the base_combine operator in flatten.raw(mapped, task_limit=task_limit) and concat.raw(mapped, task_limit=task_limit).

However a proper fix exists, as I discovered while replying in #129.

Would you like to integrate it to your PR? Otherwise, I'll create a new one myself, you let me know :)

@vxgmichel

Copy link
Copy Markdown
Owner

I went ahead and merged #143, thanks for the PR :)

@vxgmichel vxgmichel closed this Sep 8, 2026
@vxgmichel

Copy link
Copy Markdown
Owner

Also, #143 is released in v0.7.2 🎉

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

pipe.flatmap fails AsyncIterable assertion when used with async def function (coroutine not awaited)

2 participants