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):

View source record →

📰 Research Paper
Loading…
⏳ Fetching content…