From c38c0081fdfe0555d3ae8caa8799f825fabf42fc Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Wed, 30 Sep 2026 11:31:46 +0200 Subject: [PATCH 1/7] feat: add a concurrency limit for named child runs --- src/apify/_actor.py | 17 ++ src/apify/_child_runs.py | 136 ++++++++++++-- tests/unit/actor/test_actor_child_runs.py | 215 ++++++++++++++++++++++ 3 files changed, 351 insertions(+), 17 deletions(-) diff --git a/src/apify/_actor.py b/src/apify/_actor.py index cbbf4bd2..d200cd53 100644 --- a/src/apify/_actor.py +++ b/src/apify/_actor.py @@ -1193,6 +1193,8 @@ async def call( run = await self._wait_for_child_run( client.run(started_run.id), started_run, wait=wait, logger=logger, from_start=is_new ) + if run is not None: + await self._child_run_registry.run_finished(name, run) if run is None: raise RuntimeError(f'Failed to call Actor with ID "{actor_id}".') @@ -1270,6 +1272,21 @@ async def child_runs(self) -> dict[str, ChildRunInfo]: """ return await self._child_run_registry.list_runs(self.apify_client) + def set_child_run_limits(self, *, max_concurrent_runs: int | None) -> None: + """Limit the named child runs of this Actor run. + + While `max_concurrent_runs` named child runs are `READY`, `RUNNING`, `ABORTING` or `TIMING-OUT`, a named + `Actor.start` or `Actor.call` that would start or resurrect a run waits until one of them finishes. Reattaching + to a recorded run never waits. Runs started without a `name` are not counted and never wait. A child run not + awaited by `Actor.call` is fetched again before it is counted, if its status is more than 10 seconds old. + + The limit is kept in memory, so call this method again after a migration or resurrection of this Actor run. + + Args: + max_concurrent_runs: How many named child runs may be active at once, or `None` for no limit. + """ + self._child_run_registry.set_max_concurrent_runs(max_concurrent_runs) + @_ensure_context async def call_task( self, diff --git a/src/apify/_child_runs.py b/src/apify/_child_runs.py index 76f53bed..7ec951da 100644 --- a/src/apify/_child_runs.py +++ b/src/apify/_child_runs.py @@ -2,7 +2,9 @@ import asyncio from collections import defaultdict +from contextlib import asynccontextmanager, suppress from dataclasses import dataclass +from datetime import timedelta from logging import getLogger from typing import TYPE_CHECKING @@ -12,7 +14,7 @@ from apify._utils import docs_group if TYPE_CHECKING: - from collections.abc import Awaitable, Callable + from collections.abc import AsyncIterator, Awaitable, Callable from apify_client import ApifyClientAsync from apify_client._models import Run @@ -32,6 +34,11 @@ _ABORTABLE_STATUSES = frozenset({'READY', 'RUNNING'}) +_ACTIVE_STATUSES = frozenset({'READY', 'RUNNING', 'ABORTING', 'TIMING-OUT'}) + +_STATUS_MAX_AGE = timedelta(seconds=10) +"""How long an observed active status counts toward the concurrency limit before the run is fetched again.""" + class ChildRunRecord(BaseModel): """A child run tracked under a name in the child run registry.""" @@ -90,6 +97,19 @@ def __init__(self, open_key_value_store: Callable[[], Awaitable[KeyValueStore]]) self._name_locks: defaultdict[str, asyncio.Lock] = defaultdict(asyncio.Lock) self._clients: dict[str, ApifyClientAsync] = {} """Client each name was last started or reattached with in this process, used to abort its run.""" + self._max_concurrent_runs: int | None = None + self._slots = asyncio.Condition() + self._starting: set[str] = set() + """Names holding a slot for a start or resurrection in flight, not counted from the records yet.""" + self._observed: dict[str, tuple[str, float]] = {} + """Last observed status of the run recorded under each name, with the event loop time it was observed at.""" + self._parent_aborting = False + + def set_max_concurrent_runs(self, max_concurrent_runs: int | None) -> None: + """Set how many recorded runs may be active at once, or remove the limit with `None`.""" + if max_concurrent_runs is not None and max_concurrent_runs < 1: + raise ValueError(f'`max_concurrent_runs` must be at least 1, got {max_concurrent_runs}.') + self._max_concurrent_runs = max_concurrent_runs async def find_or_start( self, @@ -107,6 +127,8 @@ async def find_or_start( An `ABORTED` or `TIMED-OUT` run is resurrected, since Actors are expected to resume from their state. A `FAILED` run, or one the API no longer knows, is replaced by a new run under the same name. + Starting or resurrecting a run waits while the concurrency limit is reached. Reattaching never waits. + Args: name: Name of the child run, unique within the parent run. actor_id: The Actor to start. It must match the Actor already recorded under `name`. @@ -132,13 +154,14 @@ async def find_or_start( self._clients[name] = client if record is None: - run = await self._start( - name, - actor_id=actor_id, - start_run=start_run, - previous_run_ids=[], - abort_with_parent=abort_with_parent, - ) + async with self._slot(name, client): + run = await self._start( + name, + actor_id=actor_id, + start_run=start_run, + previous_run_ids=[], + abort_with_parent=abort_with_parent, + ) return run, True run_client = client.run(record.run_id) @@ -148,25 +171,41 @@ async def find_or_start( run = await run_client.wait_for_finish() if run is None or run.status == 'FAILED': - run = await self._start( - name, - actor_id=actor_id, - start_run=start_run, - previous_run_ids=[*record.previous_run_ids, record.run_id], - abort_with_parent=abort_with_parent, - ) + async with self._slot(name, client): + run = await self._start( + name, + actor_id=actor_id, + start_run=start_run, + previous_run_ids=[*record.previous_run_ids, record.run_id], + abort_with_parent=abort_with_parent, + ) return run, True + self._observe(name, run) + if record.abort_with_parent != abort_with_parent: await self._save(name, record.model_copy(update={'abort_with_parent': abort_with_parent})) if run.status in _RESURRECTABLE_STATUSES: - logger.info(f'Resurrecting child run "{name}"', extra={'run_id': run.id, 'status': run.status}) - return await resurrect_run(run_client), False + async with self._slot(name, client): + logger.info(f'Resurrecting child run "{name}"', extra={'run_id': run.id, 'status': run.status}) + run = await resurrect_run(run_client) + self._observe(name, run) + return run, False logger.info(f'Reattaching to child run "{name}"', extra={'run_id': run.id, 'status': run.status}) return run, False + async def run_finished(self, name: str, run: Run) -> None: + """Record the status of a run under `name` that was awaited, releasing its slot when it is no longer active.""" + records = await self._load() + record = records.get(name) + if record is None or record.run_id != run.id: + return + self._observe(name, run) + async with self._slots: + self._slots.notify_all() + async def list_runs(self, client: ApifyClientAsync) -> dict[str, ChildRunInfo]: """Return every recorded child run by name, with its current state fetched from the API. @@ -176,6 +215,9 @@ async def list_runs(self, client: ApifyClientAsync) -> dict[str, ChildRunInfo]: # Copy the records, since a named start can add one while the runs are fetched. records = dict(await self._load()) runs = await asyncio.gather(*(client.run(record.run_id).get() for record in records.values())) + for (name, record), run in zip(records.items(), runs, strict=True): + if run is not None and run.id == record.run_id: + self._observe(name, run) return { name: ChildRunInfo( actor_id=record.actor_id, @@ -195,6 +237,10 @@ async def abort_runs_with_parent(self, client: ApifyClientAsync) -> None: Args: client: Client used for a name not started or reattached in this process, e.g. after a migration. """ + self._parent_aborting = True + async with self._slots: + self._slots.notify_all() + records = await self._load() # Names with a start in flight are not recorded yet, so their locks are awaited too. await asyncio.gather(*(self._abort(name, client) for name in {*records, *self._name_locks})) @@ -232,8 +278,64 @@ async def _start( abort_with_parent=abort_with_parent, ) await self._save(name, record) + self._observe(name, run) return run + @asynccontextmanager + async def _slot(self, name: str, client: ApifyClientAsync) -> AsyncIterator[None]: + """Hold a slot for starting or resurrecting the run under `name`, waiting while the limit is reached.""" + if self._max_concurrent_runs is None: + yield + return + + async with self._slots: + while await self._count_active(client, exclude=name) >= self._max_concurrent_runs: + if self._parent_aborting: + raise RuntimeError( + f'Child run "{name}" was not started, since this Actor run is being aborted and the limit ' + f'of {self._max_concurrent_runs} concurrent child runs is reached.' + ) + logger.debug(f'Child run "{name}" is waiting for a free slot') + with suppress(TimeoutError): + await asyncio.wait_for(self._slots.wait(), timeout=_STATUS_MAX_AGE.total_seconds()) + self._starting.add(name) + + try: + yield + finally: + self._starting.discard(name) + async with self._slots: + self._slots.notify_all() + + async def _count_active(self, client: ApifyClientAsync, *, exclude: str) -> int: + """Count active recorded runs and slots held by others, fetching runs whose active status is not fresh.""" + records = await self._load() + names = [name for name in records if name != exclude and name not in self._starting] + now = asyncio.get_running_loop().time() + stale = [ + name + for name in names + if name not in self._observed + or ( + self._observed[name][0] in _ACTIVE_STATUSES + and now - self._observed[name][1] >= _STATUS_MAX_AGE.total_seconds() + ) + ] + runs = await asyncio.gather( + *(self._clients.get(name, client).run(records[name].run_id).get() for name in stale) + ) + for name, run in zip(stale, runs, strict=True): + if run is None: + self._observed[name] = ('MISSING', now) + else: + self._observe(name, run) + + active = sum(1 for name in names if self._observed[name][0] in _ACTIVE_STATUSES) + return active + len(self._starting) + + def _observe(self, name: str, run: Run) -> None: + self._observed[name] = (run.status, asyncio.get_running_loop().time()) + async def _load(self) -> dict[str, ChildRunRecord]: async with self._load_lock: if self._records is None: diff --git a/tests/unit/actor/test_actor_child_runs.py b/tests/unit/actor/test_actor_child_runs.py index 95411649..1f2b71ff 100644 --- a/tests/unit/actor/test_actor_child_runs.py +++ b/tests/unit/actor/test_actor_child_runs.py @@ -652,3 +652,218 @@ async def test_removing_all_aborting_listeners_keeps_aborting_child_runs( await apify_event_manager.wait_for_all_listeners_to_complete() assert len(apify_client_async_patcher.calls['run']['abort']) == 1 + + +def make_client(statuses: dict[str, str]) -> Mock: + """A client whose `run(run_id).get()` returns the run with its current status in `statuses`.""" + client = Mock() + + def run(run_id: str) -> Mock: + run_client = Mock() + run_client.get = AsyncMock(side_effect=lambda: make_run(run_id, statuses[run_id])) + run_client.resurrect = AsyncMock(side_effect=lambda: make_run(run_id, 'RUNNING')) + run_client.abort = AsyncMock() + return run_client + + client.run.side_effect = run + return client + + +async def start_child( + registry: ChildRunRegistry, client: Mock, name: str, statuses: dict[str, str], *, run_id: str | None = None +) -> Run: + """Start a named child run with the registry, adding its run to `statuses` as `RUNNING`.""" + + async def start_run() -> Run: + new_run_id = run_id or f'{name}-run' + statuses[new_run_id] = 'RUNNING' + return make_run(new_run_id, 'READY') + + run, _ = await registry.find_or_start( + name, + actor_id='some-actor', + client=client, + start_run=start_run, + resurrect_run=lambda run_client: run_client.resurrect(), + ) + return run + + +async def assert_waiting(task: asyncio.Task) -> None: + await asyncio.sleep(0.05) + assert not task.done() + + +@pytest.mark.parametrize( + 'max_concurrent_runs', + [ + pytest.param(0, id='zero'), + pytest.param(-1, id='negative'), + ], +) +async def test_set_child_run_limits_rejects_non_positive_limit(max_concurrent_runs: int) -> None: + """A concurrency limit below 1 is rejected.""" + async with Actor: + with pytest.raises(ValueError, match='must be at least 1'): + Actor.set_child_run_limits(max_concurrent_runs=max_concurrent_runs) + + +async def test_named_start_waits_while_the_limit_is_reached() -> None: + """A named start waits while the limit is reached and proceeds once an active child run finishes.""" + statuses: dict[str, str] = {} + client = make_client(statuses) + + async with Actor: + registry = ChildRunRegistry(Actor.open_key_value_store) + registry.set_max_concurrent_runs(1) + first = await start_child(registry, client, 'first', statuses) + second_task = asyncio.create_task(start_child(registry, client, 'second', statuses)) + await assert_waiting(second_task) + + statuses[first.id] = 'SUCCEEDED' + await registry.run_finished('first', make_run(first.id, 'SUCCEEDED')) + second = await second_task + + assert second.id == 'second-run' + + +async def test_stale_active_status_is_refreshed_before_counting(monkeypatch: pytest.MonkeyPatch) -> None: + """A child run nobody awaited is fetched again once its status is stale, so a finished one frees its slot.""" + monkeypatch.setattr('apify._child_runs._STATUS_MAX_AGE', timedelta(0)) + statuses: dict[str, str] = {} + client = make_client(statuses) + + async with Actor: + registry = ChildRunRegistry(Actor.open_key_value_store) + registry.set_max_concurrent_runs(1) + first = await start_child(registry, client, 'first', statuses) + second_task = asyncio.create_task(start_child(registry, client, 'second', statuses)) + await assert_waiting(second_task) + + statuses[first.id] = 'SUCCEEDED' + second = await asyncio.wait_for(second_task, timeout=1) + + assert second.id == 'second-run' + + +async def test_child_run_recorded_by_an_earlier_attempt_counts_toward_the_limit() -> None: + """An active child run recorded before a migration holds a slot, since it is fetched before it is counted.""" + statuses = {'old-run': 'RUNNING'} + client = make_client(statuses) + + async with Actor: + await record_child_run('first', 'old-run') + registry = ChildRunRegistry(Actor.open_key_value_store) + registry.set_max_concurrent_runs(1) + second_task = asyncio.create_task(start_child(registry, client, 'second', statuses)) + await assert_waiting(second_task) + second_task.cancel() + + +async def test_reattach_does_not_wait_for_a_slot() -> None: + """Reattaching to an active recorded child run returns it even while the limit is reached.""" + statuses = {'old-run': 'RUNNING'} + client = make_client(statuses) + + async with Actor: + await record_child_run('first', 'old-run') + registry = ChildRunRegistry(Actor.open_key_value_store) + registry.set_max_concurrent_runs(1) + run = await asyncio.wait_for(start_child(registry, client, 'first', statuses), timeout=1) + + assert run.id == 'old-run' + + +async def test_resurrection_waits_for_a_slot() -> None: + """Resurrecting an aborted child run waits while the limit is reached.""" + statuses = {'old-run': 'ABORTED'} + client = make_client(statuses) + + async with Actor: + await record_child_run('first', 'old-run') + registry = ChildRunRegistry(Actor.open_key_value_store) + registry.set_max_concurrent_runs(1) + second = await start_child(registry, client, 'second', statuses) + first_task = asyncio.create_task(start_child(registry, client, 'first', statuses)) + await assert_waiting(first_task) + + statuses[second.id] = 'SUCCEEDED' + await registry.run_finished('second', make_run(second.id, 'SUCCEEDED')) + first = await first_task + + assert first.id == 'old-run' + assert first.status == 'RUNNING' + + +async def test_concurrent_named_starts_respect_the_limit() -> None: + """Concurrent named starts under different names start no more runs than the limit.""" + statuses: dict[str, str] = {} + client = make_client(statuses) + + async with Actor: + registry = ChildRunRegistry(Actor.open_key_value_store) + registry.set_max_concurrent_runs(2) + tasks = [asyncio.create_task(start_child(registry, client, f'child-{i}', statuses)) for i in range(3)] + await asyncio.sleep(0.05) + for task in tasks: + task.cancel() + + assert sum(task.cancelled() for task in tasks) == 1 + assert len(statuses) == 2 + + +async def test_failed_start_releases_its_slot() -> None: + """A named start that raises frees its slot for the next one.""" + statuses: dict[str, str] = {} + client = make_client(statuses) + + async with Actor: + registry = ChildRunRegistry(Actor.open_key_value_store) + registry.set_max_concurrent_runs(1) + with pytest.raises(RuntimeError, match='start failed'): + await registry.find_or_start( + 'first', + actor_id='some-actor', + client=client, + start_run=AsyncMock(side_effect=RuntimeError('start failed')), + resurrect_run=AsyncMock(), + ) + second = await asyncio.wait_for(start_child(registry, client, 'second', statuses), timeout=1) + + assert second.id == 'second-run' + + +async def test_parent_abort_stops_starts_waiting_for_a_slot() -> None: + """A named start waiting for a slot raises once the parent is aborted, without starting a run.""" + statuses: dict[str, str] = {} + client = make_client(statuses) + + async with Actor: + registry = ChildRunRegistry(Actor.open_key_value_store) + registry.set_max_concurrent_runs(1) + await start_child(registry, client, 'first', statuses) + second_task = asyncio.create_task(start_child(registry, client, 'second', statuses)) + await assert_waiting(second_task) + + await asyncio.wait_for(registry.abort_runs_with_parent(client), timeout=1) + with pytest.raises(RuntimeError, match='being aborted'): + await second_task + + assert 'second-run' not in statuses + + +async def test_named_call_frees_its_slot_when_the_run_finishes( + apify_client_async_patcher: ApifyClientAsyncPatcher, +) -> None: + """A child run awaited by a named call to its end frees its slot without being fetched again.""" + apify_client_async_patcher.patch('actor', 'start', return_value=make_run('new-run', 'READY')) + # A fetch would report the run as still running, so only the awaited status can free the slot. + apify_client_async_patcher.patch('run', 'get', return_value=make_run('new-run', 'RUNNING')) + apify_client_async_patcher.patch('run', 'wait_for_finish', return_value=make_run('new-run', 'SUCCEEDED')) + + async with Actor: + Actor.set_child_run_limits(max_concurrent_runs=1) + await Actor.call('some-actor', name='first', logger=None) + await asyncio.wait_for(Actor.start('some-actor', name='second'), timeout=1) + + assert len(apify_client_async_patcher.calls['actor']['start']) == 2 From ab4bb4176875a68555ac25092e4b03667c677a16 Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Wed, 30 Sep 2026 11:31:47 +0200 Subject: [PATCH 2/7] test: add an E2E test for the child run concurrency limit --- tests/e2e/test_actor_child_runs.py | 33 ++++++++++++++++++++++++++++++ 1 file changed, 33 insertions(+) diff --git a/tests/e2e/test_actor_child_runs.py b/tests/e2e/test_actor_child_runs.py index 0881ae7b..e0a825f5 100644 --- a/tests/e2e/test_actor_child_runs.py +++ b/tests/e2e/test_actor_child_runs.py @@ -134,3 +134,36 @@ async def main() -> None: child_run = await apify_client_async.run(child_run_id).wait_for_finish(wait_duration=timedelta(seconds=120)) assert child_run is not None assert child_run.status == 'ABORTED' + + +async def test_named_child_runs_respect_the_concurrency_limit( + make_actor: MakeActorFunction, + run_actor: RunActorFunction, +) -> None: + """Named child runs started concurrently under a limit of one run one after another.""" + + async def main() -> None: + async with Actor: + actor_input = (await Actor.get_input()) or {} + if actor_input.get('is_child') is True: + await asyncio.sleep(10) + return + + actor_id = Actor.configuration.actor_id or '' + Actor.set_child_run_limits(max_concurrent_runs=1) + runs = await asyncio.gather( + *( + Actor.call(actor_id=actor_id, run_input={'is_child': True}, name=f'child-{index}') + for index in range(2) + ) + ) + first, second = sorted(runs, key=lambda run: run.started_at) + assert first.finished_at is not None, 'first.finished_at is None' + assert second.started_at >= first.finished_at, f'first={first}, second={second}' + + actor = await make_actor(label='child-run-limit', main_func=main) + run_result = await run_actor(actor) + + assert run_result.status == 'SUCCEEDED' + # The parent run and its two child runs. + assert (await actor.runs().list()).total == 3 From a809cd06e00f40e285669a7e6c14cdfc3ab04463 Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Wed, 30 Sep 2026 11:31:48 +0200 Subject: [PATCH 3/7] docs: document the child run concurrency limit --- .../06_interacting_with_other_actors.mdx | 18 +++++++++++ .../code/06_interacting_child_run_limit.py | 30 +++++++++++++++++++ 2 files changed, 48 insertions(+) create mode 100644 docs/02_concepts/code/06_interacting_child_run_limit.py diff --git a/docs/02_concepts/06_interacting_with_other_actors.mdx b/docs/02_concepts/06_interacting_with_other_actors.mdx index 7d73ec82..1070aa0b 100644 --- a/docs/02_concepts/06_interacting_with_other_actors.mdx +++ b/docs/02_concepts/06_interacting_with_other_actors.mdx @@ -11,6 +11,7 @@ import InteractingCallExample from '!!raw-loader!roa-loader!./code/06_interactin import InteractingNamedCallExample from '!!raw-loader!roa-loader!./code/06_interacting_named_call.py'; import InteractingChildRunsExample from '!!raw-loader!roa-loader!./code/06_interacting_child_runs.py'; import InteractingAbortWithParentExample from '!!raw-loader!roa-loader!./code/06_interacting_abort_with_parent.py'; +import InteractingChildRunLimitExample from '!!raw-loader!roa-loader!./code/06_interacting_child_run_limit.py'; import InteractingCallTaskExample from '!!raw-loader!roa-loader!./code/06_interacting_call_task.py'; import InteractingMetamorphExample from '!!raw-loader!roa-loader!./code/06_interacting_metamorph.py'; import InteractingAbortExample from '!!raw-loader!roa-loader!./code/06_interacting_abort.py'; @@ -86,6 +87,23 @@ Note that: - Child runs started after your Actor run received `ABORTING` aren't aborted, so don't start new ones while it's shutting down. - A child run started with its own `token` is aborted with that token. After a migration or resurrection, the SDK uses your Actor's token for it until the same named call runs again. If that token can't access the child run, the abort fails and the error is logged. +### Limiting concurrent child runs + +An Actor that starts many child runs at once can hit the concurrency or memory limit of your account. To cap how many named child runs are active at once, call `Actor.set_child_run_limits`. While the limit is reached, a named `Actor.start` or `Actor.call` that would start or resurrect a run waits until one of the active child runs finishes. The SDK counts the child runs from the registry, so child runs started before a migration or resurrection count too. + + + {InteractingChildRunLimitExample} + + +Note that: + +- A child run counts as active while it's `READY`, `RUNNING`, `ABORTING` or `TIMING-OUT`. +- Only named child runs count, and only named calls wait. A call without `name` starts its run right away. +- Reattaching to an active child run never waits, since the run already holds a slot. +- `Actor.call` frees the slot as soon as its run finishes. The SDK doesn't learn right away about a child run that nothing waits for, so it fetches the run again before counting it, once its status is more than 10 seconds old. +- The limit isn't persisted. After a migration or resurrection, call `Actor.set_child_run_limits` again before you start child runs. +- Once your Actor run receives the `ABORTING` event, a call waiting for a slot raises a `RuntimeError` without starting its run. + ## Actor call task The `Actor.call_task` method starts an [Actor task](https://docs.apify.com/platform/actors/tasks) on the Apify platform, and waits for the started Actor run to finish. diff --git a/docs/02_concepts/code/06_interacting_child_run_limit.py b/docs/02_concepts/code/06_interacting_child_run_limit.py new file mode 100644 index 00000000..3da343b6 --- /dev/null +++ b/docs/02_concepts/code/06_interacting_child_run_limit.py @@ -0,0 +1,30 @@ +import asyncio + +from apify import Actor + + +async def main() -> None: + async with Actor: + # Keep at most 3 named child runs active at once. + Actor.set_child_run_limits(max_concurrent_runs=3) + + urls = [f'https://www.apify.com/?page={page}' for page in range(6)] + + # Each call waits for a free slot before it starts its child run. + actor_runs = await asyncio.gather( + *( + Actor.call( + actor_id='apify/screenshot-url', + run_input={'urls': [{'url': url}]}, + name=f'screenshot-{index}', + ) + for index, url in enumerate(urls) + ) + ) + + for actor_run in actor_runs: + Actor.log.info(f'Child run {actor_run.id} finished as {actor_run.status}') + + +if __name__ == '__main__': + asyncio.run(main()) From 75babea68000d4c5380b4e9ab35ca75ee9e43028 Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Wed, 30 Sep 2026 12:07:01 +0200 Subject: [PATCH 4/7] fix: skip a child run that can't be fetched when counting toward the limit --- src/apify/_child_runs.py | 14 +++++++++++--- tests/unit/actor/test_actor_child_runs.py | 23 +++++++++++++++++++++++ 2 files changed, 34 insertions(+), 3 deletions(-) diff --git a/src/apify/_child_runs.py b/src/apify/_child_runs.py index 7ec951da..026c8fd1 100644 --- a/src/apify/_child_runs.py +++ b/src/apify/_child_runs.py @@ -322,15 +322,23 @@ async def _count_active(self, client: ApifyClientAsync, *, exclude: str) -> int: ) ] runs = await asyncio.gather( - *(self._clients.get(name, client).run(records[name].run_id).get() for name in stale) + *(self._clients.get(name, client).run(records[name].run_id).get() for name in stale), + return_exceptions=True, ) for name, run in zip(stale, runs, strict=True): - if run is None: + if isinstance(run, BaseException): + logger.warning( + f'Failed to fetch child run "{name}" to count it toward the concurrency limit', + extra={'run_id': records[name].run_id}, + exc_info=run, + ) + elif run is None: self._observed[name] = ('MISSING', now) else: self._observe(name, run) - active = sum(1 for name in names if self._observed[name][0] in _ACTIVE_STATUSES) + # A run that could not be fetched counts only when an earlier observation saw it active. + active = sum(1 for name in names if name in self._observed and self._observed[name][0] in _ACTIVE_STATUSES) return active + len(self._starting) def _observe(self, name: str, run: Run) -> None: diff --git a/tests/unit/actor/test_actor_child_runs.py b/tests/unit/actor/test_actor_child_runs.py index 1f2b71ff..dc6c4340 100644 --- a/tests/unit/actor/test_actor_child_runs.py +++ b/tests/unit/actor/test_actor_child_runs.py @@ -760,6 +760,29 @@ async def test_child_run_recorded_by_an_earlier_attempt_counts_toward_the_limit( second_task.cancel() +async def test_child_run_that_cannot_be_fetched_does_not_block_the_limit() -> None: + """A recorded child run whose fetch fails is not counted, so it does not fail or block a named start.""" + statuses: dict[str, str] = {} + client = make_client(statuses) + run_client_factory = client.run.side_effect + + def run(run_id: str) -> Mock: + run_client = run_client_factory(run_id) + if run_id == 'old-run': + run_client.get = AsyncMock(side_effect=RuntimeError('forbidden')) + return run_client + + client.run.side_effect = run + + async with Actor: + await record_child_run('first', 'old-run') + registry = ChildRunRegistry(Actor.open_key_value_store) + registry.set_max_concurrent_runs(1) + second = await asyncio.wait_for(start_child(registry, client, 'second', statuses), timeout=1) + + assert second.id == 'second-run' + + async def test_reattach_does_not_wait_for_a_slot() -> None: """Reattaching to an active recorded child run returns it even while the limit is reached.""" statuses = {'old-run': 'RUNNING'} From 45d69a3238b7e8713766fdf00f07d219d41e24c0 Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Wed, 30 Sep 2026 12:07:19 +0200 Subject: [PATCH 5/7] fix: release waiting starts when the child run limit is removed --- src/apify/_child_runs.py | 5 ++++- tests/unit/actor/test_actor_child_runs.py | 19 +++++++++++++++++++ 2 files changed, 23 insertions(+), 1 deletion(-) diff --git a/src/apify/_child_runs.py b/src/apify/_child_runs.py index 026c8fd1..f98bdb5a 100644 --- a/src/apify/_child_runs.py +++ b/src/apify/_child_runs.py @@ -289,7 +289,10 @@ async def _slot(self, name: str, client: ApifyClientAsync) -> AsyncIterator[None return async with self._slots: - while await self._count_active(client, exclude=name) >= self._max_concurrent_runs: + while ( + self._max_concurrent_runs is not None + and await self._count_active(client, exclude=name) >= self._max_concurrent_runs + ): if self._parent_aborting: raise RuntimeError( f'Child run "{name}" was not started, since this Actor run is being aborted and the limit ' diff --git a/tests/unit/actor/test_actor_child_runs.py b/tests/unit/actor/test_actor_child_runs.py index dc6c4340..f6d232c5 100644 --- a/tests/unit/actor/test_actor_child_runs.py +++ b/tests/unit/actor/test_actor_child_runs.py @@ -727,6 +727,25 @@ async def test_named_start_waits_while_the_limit_is_reached() -> None: assert second.id == 'second-run' +async def test_removing_the_limit_releases_waiting_starts(monkeypatch: pytest.MonkeyPatch) -> None: + """A named start waiting for a slot proceeds once the limit is removed.""" + monkeypatch.setattr('apify._child_runs._STATUS_MAX_AGE', timedelta(seconds=0.05)) + statuses: dict[str, str] = {} + client = make_client(statuses) + + async with Actor: + registry = ChildRunRegistry(Actor.open_key_value_store) + registry.set_max_concurrent_runs(1) + await start_child(registry, client, 'first', statuses) + second_task = asyncio.create_task(start_child(registry, client, 'second', statuses)) + await assert_waiting(second_task) + + registry.set_max_concurrent_runs(None) + second = await asyncio.wait_for(second_task, timeout=1) + + assert second.id == 'second-run' + + async def test_stale_active_status_is_refreshed_before_counting(monkeypatch: pytest.MonkeyPatch) -> None: """A child run nobody awaited is fetched again once its status is stale, so a finished one frees its slot.""" monkeypatch.setattr('apify._child_runs._STATUS_MAX_AGE', timedelta(0)) From ae4a0f675cf3cf9b6abe65d3af399e29ae61e2a4 Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Wed, 30 Sep 2026 12:07:28 +0200 Subject: [PATCH 6/7] fix: keep the status of a run that replaced a listed child run --- src/apify/_child_runs.py | 6 +++-- tests/unit/actor/test_actor_child_runs.py | 32 +++++++++++++++++++++++ 2 files changed, 36 insertions(+), 2 deletions(-) diff --git a/src/apify/_child_runs.py b/src/apify/_child_runs.py index f98bdb5a..7309cc15 100644 --- a/src/apify/_child_runs.py +++ b/src/apify/_child_runs.py @@ -215,8 +215,10 @@ async def list_runs(self, client: ApifyClientAsync) -> dict[str, ChildRunInfo]: # Copy the records, since a named start can add one while the runs are fetched. records = dict(await self._load()) runs = await asyncio.gather(*(client.run(record.run_id).get() for record in records.values())) - for (name, record), run in zip(records.items(), runs, strict=True): - if run is not None and run.id == record.run_id: + # A name whose run was replaced during the fetch keeps the status observed for its new run. + current = await self._load() + for name, run in zip(records, runs, strict=True): + if run is not None and name in current and current[name].run_id == run.id: self._observe(name, run) return { name: ChildRunInfo( diff --git a/tests/unit/actor/test_actor_child_runs.py b/tests/unit/actor/test_actor_child_runs.py index f6d232c5..a1b35bfe 100644 --- a/tests/unit/actor/test_actor_child_runs.py +++ b/tests/unit/actor/test_actor_child_runs.py @@ -802,6 +802,38 @@ def run(run_id: str) -> Mock: assert second.id == 'second-run' +async def test_listing_keeps_the_status_of_a_run_that_replaced_the_listed_one() -> None: + """A run that replaces the listed one under a name keeps its slot after the listing observes the old run.""" + statuses = {'old-run': 'FAILED'} + client = make_client(statuses) + fetch_started = asyncio.Event() + release_fetch = asyncio.Event() + + async def get_old_run() -> Run: + fetch_started.set() + await release_fetch.wait() + return make_run('old-run', 'FAILED') + + list_client = Mock() + list_client.run.return_value.get = AsyncMock(side_effect=get_old_run) + + async with Actor: + await record_child_run('first', 'old-run') + registry = ChildRunRegistry(Actor.open_key_value_store) + registry.set_max_concurrent_runs(1) + list_task = asyncio.create_task(registry.list_runs(list_client)) + await fetch_started.wait() + await start_child(registry, client, 'first', statuses) + release_fetch.set() + await list_task + + second_task = asyncio.create_task(start_child(registry, client, 'second', statuses)) + await assert_waiting(second_task) + second_task.cancel() + with pytest.raises(asyncio.CancelledError): + await second_task + + async def test_reattach_does_not_wait_for_a_slot() -> None: """Reattaching to an active recorded child run returns it even while the limit is reached.""" statuses = {'old-run': 'RUNNING'} From 5bbadc2d101f9eca71abde908e99787dcdc93964 Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Wed, 30 Sep 2026 12:07:37 +0200 Subject: [PATCH 7/7] test: await cancelled tasks in child run limit tests --- tests/unit/actor/test_actor_child_runs.py | 3 +++ 1 file changed, 3 insertions(+) diff --git a/tests/unit/actor/test_actor_child_runs.py b/tests/unit/actor/test_actor_child_runs.py index a1b35bfe..dbd14dc6 100644 --- a/tests/unit/actor/test_actor_child_runs.py +++ b/tests/unit/actor/test_actor_child_runs.py @@ -777,6 +777,8 @@ async def test_child_run_recorded_by_an_earlier_attempt_counts_toward_the_limit( second_task = asyncio.create_task(start_child(registry, client, 'second', statuses)) await assert_waiting(second_task) second_task.cancel() + with pytest.raises(asyncio.CancelledError): + await second_task async def test_child_run_that_cannot_be_fetched_does_not_block_the_limit() -> None: @@ -881,6 +883,7 @@ async def test_concurrent_named_starts_respect_the_limit() -> None: await asyncio.sleep(0.05) for task in tasks: task.cancel() + await asyncio.gather(*tasks, return_exceptions=True) assert sum(task.cancelled() for task in tasks) == 1 assert len(statuses) == 2