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

View source record →

📰 Research Paper
Loading…
⏳ Fetching content…