INNER CODE UNIT · Python
_run_with_bounded_concurrency
Hebbian-Robotics/hflow · src/hflow/app.py:500
async def _run_with_bounded_concurrency(
inputs: Sequence[_ConcurrentInputT],
operation: Callable[[_ConcurrentInputT], Awaitable[_ConcurrentResultT]],
*,
concurrency_limit: _ConcurrentEpisodeLimit,
on_completion: Callable[[int, _ConcurrentResultT, int, int], None] | None = None,
) -> list[_ConcurrentResultT]:
"""Bound admission, preserve input order, and drain cancelled siblings on failure."""
results_by_input_index: dict[int, _ConcurrentResultT] = {}
next_input_index = 0
input_index_by_task: dict[asyncio.Task[_ConcurrentResultT], int] = {}
async def execute_input(value: _ConcurrentInputT) -> _ConcurrentResultT:
return await operation(value)
def submit_next_input() -> None:
nonlocal next_input_index
if next_input_index < len(inputs):