INNER CODE UNIT · Python
_parse_concurrent_episode_limit
Hebbian-Robotics/hflow · src/hflow/app.py:492
def _parse_concurrent_episode_limit(max_workers: int) -> _ConcurrentEpisodeLimit:
"""Refine the public integer before the scheduler can consume it."""
if isinstance(max_workers, bool) or not isinstance(max_workers, int) or max_workers <= 0:
raise ValueError("max_workers must be a positive integer")
return _ConcurrentEpisodeLimit(max_workers)
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