From df5f19458c1289f2edafd9bda13a814c73ef423c Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Tue, 29 Sep 2026 13:42:26 +0200 Subject: [PATCH 1/6] feat: reattach named child runs after migration or resurrection --- src/apify/_actor.py | 146 +++++++++++-- src/apify/_child_runs.py | 144 +++++++++++++ tests/unit/actor/test_actor_child_runs.py | 248 ++++++++++++++++++++++ 3 files changed, 522 insertions(+), 16 deletions(-) create mode 100644 src/apify/_child_runs.py create mode 100644 tests/unit/actor/test_actor_child_runs.py diff --git a/src/apify/_actor.py b/src/apify/_actor.py index 23d5cff2b..0f4f380a6 100644 --- a/src/apify/_actor.py +++ b/src/apify/_actor.py @@ -6,7 +6,7 @@ import warnings from dataclasses import asdict from datetime import UTC, datetime, timedelta -from functools import cached_property +from functools import cached_property, partial from typing import TYPE_CHECKING, Any, Literal, TypeVar, cast, overload from lazy_object_proxy import Proxy @@ -33,6 +33,7 @@ ChargingManagerImplementation, charge_lock_if_charging, ) +from apify._child_runs import ChildRunRegistry from apify._configuration import Configuration from apify._consts import EVENT_LISTENERS_TIMEOUT, EXIT_CODE_ERROR_USER_FUNCTION_THREW, ActorEnvVars, ApifyEnvVars from apify._crypto import decrypt_input_secrets, load_private_key @@ -47,13 +48,14 @@ if TYPE_CHECKING: import logging - from collections.abc import Callable, MutableMapping + from collections.abc import Awaitable, Callable, MutableMapping from decimal import Decimal from types import TracebackType from typing import Self from apify_client._literals import ActorPermissionLevel from apify_client._models import Run + from apify_client._resource_clients import RunClientAsync from crawlee._types import JsonSerializable from crawlee.proxy_configuration import _NewUrlFunction @@ -149,6 +151,8 @@ def __init__( # Keep track of all used state stores to persist their values on exit self._use_state_stores: set[str | None] = set() + self._child_run_registry: ChildRunRegistry | None = None + self._active = False """Whether the Actor instance is currently active (initialized and within context).""" @@ -927,6 +931,7 @@ async def start( timeout: timedelta | Literal['inherit'] | None = None, force_permission_level: ActorPermissionLevel | None = None, webhooks: list[Webhook] | None = None, + name: str | None = None, ) -> Run: """Run an Actor on the Apify platform. @@ -953,6 +958,11 @@ async def start( webhooks: Optional ad-hoc webhooks (https://docs.apify.com/webhooks/ad-hoc-webhooks) associated with the Actor run which can be used to receive a notification, e.g. when the Actor finished or failed. If you already have a webhook set up for the Actor or task, you do not have to add it again here. + name: Optional name of the child run, unique within this Actor run. A named run is recorded in the + default key-value store, so after a migration or resurrection of this Actor the same call reattaches + to the recorded run. A `SUCCEEDED` run is returned as is, an `ABORTED` or `TIMED-OUT` one is + resurrected, and a new run is started only when nothing is recorded under the name or the recorded + run `FAILED`. Returns: Info about the started Actor run @@ -969,7 +979,8 @@ async def start( raise ValueError(f'Invalid timeout {timeout!r}: expected `None`, `"inherit"`, or a `timedelta`.') actor_client = client.actor(actor_id) - return await actor_client.start( + start_run = partial( + actor_client.start, run_input=run_input, content_type=content_type, build=build, @@ -981,6 +992,22 @@ async def start( webhooks=to_client_representations(webhooks), ) + if name is None: + return await start_run() + + run, _ = await self._find_or_start_child_run( + name, + actor_id=actor_id, + client=client, + start_run=start_run, + build=build, + max_total_charge_usd=max_total_charge_usd, + restart_on_error=restart_on_error, + memory_mbytes=memory_mbytes, + run_timeout=actor_start_timeout, + ) + return run + @_ensure_context async def abort( self, @@ -1035,6 +1062,7 @@ async def call( webhooks: list[Webhook] | None = None, wait: timedelta | None = None, logger: logging.Logger | Literal['default'] | None = 'default', + name: str | None = None, ) -> Run: """Start an Actor on the Apify Platform and wait for it to finish before returning. @@ -1064,6 +1092,11 @@ async def call( logger: Logger used to redirect logs from the Actor run. Using "default" literal means that a predefined default logger will be used. Setting `None` will disable any log propagation. Passing custom logger will redirect logs to the provided logger. + name: Optional name of the child run, unique within this Actor run. A named run is recorded in the + default key-value store, so after a migration or resurrection of this Actor the same call reattaches + to the recorded run. A `SUCCEEDED` run is returned as is, an `ABORTED` or `TIMED-OUT` one is + resurrected, and a new run is started only when nothing is recorded under the name or the recorded + run `FAILED`. Returns: Info about the started Actor run. @@ -1080,25 +1113,106 @@ async def call( raise ValueError(f'Invalid timeout {timeout!r}: expected `None`, `"inherit"`, or a `timedelta`.') actor_client = client.actor(actor_id) - run = await actor_client.call( - run_input=run_input, - content_type=content_type, - build=build, - max_total_charge_usd=max_total_charge_usd, - restart_on_error=restart_on_error, - memory_mbytes=memory_mbytes, - run_timeout=actor_call_timeout, - force_permission_level=force_permission_level, - webhooks=to_client_representations(webhooks), - wait_duration=wait, - logger=logger, - ) + + if name is None: + run = await actor_client.call( + run_input=run_input, + content_type=content_type, + build=build, + max_total_charge_usd=max_total_charge_usd, + restart_on_error=restart_on_error, + memory_mbytes=memory_mbytes, + run_timeout=actor_call_timeout, + force_permission_level=force_permission_level, + webhooks=to_client_representations(webhooks), + wait_duration=wait, + logger=logger, + ) + else: + started_run, is_new = await self._find_or_start_child_run( + name, + actor_id=actor_id, + client=client, + start_run=partial( + actor_client.start, + run_input=run_input, + content_type=content_type, + build=build, + max_total_charge_usd=max_total_charge_usd, + restart_on_error=restart_on_error, + memory_mbytes=memory_mbytes, + run_timeout=actor_call_timeout, + force_permission_level=force_permission_level, + webhooks=to_client_representations(webhooks), + ), + build=build, + max_total_charge_usd=max_total_charge_usd, + restart_on_error=restart_on_error, + memory_mbytes=memory_mbytes, + run_timeout=actor_call_timeout, + ) + # The earlier attempt of this call already streamed the log of a reattached or resurrected run. + 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 None: raise RuntimeError(f'Failed to call Actor with ID "{actor_id}".') return run + async def _find_or_start_child_run( + self, + name: str, + *, + actor_id: str, + client: ApifyClientAsync, + start_run: Callable[[], Awaitable[Run]], + build: str | None, + max_total_charge_usd: Decimal | None, + restart_on_error: bool | None, + memory_mbytes: int | None, + run_timeout: timedelta | None, + ) -> tuple[Run, bool]: + if self._child_run_registry is None: + self._child_run_registry = ChildRunRegistry(await self.open_key_value_store()) + + return await self._child_run_registry.find_or_start( + name, + actor_id=actor_id, + client=client, + start_run=start_run, + resurrect_run=lambda run_client: run_client.resurrect( + build=build, + max_total_charge_usd=max_total_charge_usd, + restart_on_error=restart_on_error, + memory_mbytes=memory_mbytes, + run_timeout=run_timeout, + ), + ) + + async def _wait_for_child_run( + self, + run_client: RunClientAsync, + run: Run, + *, + wait: timedelta | None, + logger: logging.Logger | Literal['default'] | None, + from_start: bool, + ) -> Run | None: + if run.status == 'SUCCEEDED': + return run + + if not logger: + return await run_client.wait_for_finish(wait_duration=wait) + + to_logger = None if logger == 'default' else logger + status_redirector = await run_client.get_status_message_watcher(to_logger=to_logger) + streamed_log = await run_client.get_streamed_log(to_logger=to_logger, from_start=from_start) + + async with status_redirector, streamed_log: + return await run_client.wait_for_finish(wait_duration=wait) + @_ensure_context async def call_task( self, diff --git a/src/apify/_child_runs.py b/src/apify/_child_runs.py new file mode 100644 index 000000000..63fa9024c --- /dev/null +++ b/src/apify/_child_runs.py @@ -0,0 +1,144 @@ +from __future__ import annotations + +import asyncio +from collections import defaultdict +from logging import getLogger +from typing import TYPE_CHECKING + +from pydantic import BaseModel, ConfigDict, Field, TypeAdapter +from pydantic.alias_generators import to_camel + +if TYPE_CHECKING: + from collections.abc import Awaitable, Callable + + from apify_client import ApifyClientAsync + from apify_client._models import Run + from apify_client._resource_clients import RunClientAsync + + from apify.storages import KeyValueStore + +logger = getLogger(__name__) + +CHILD_RUNS_KEY = 'APIFY_CHILD_RUNS' +"""Key in the default key-value store under which the child run registry is persisted.""" + +_SETTLING_STATUSES = frozenset({'ABORTING', 'TIMING-OUT'}) +"""Statuses that end as `ABORTED` / `TIMED-OUT` shortly, and are resurrectable once they do.""" + +_RESURRECTABLE_STATUSES = frozenset({'ABORTED', 'TIMED-OUT'}) + + +class ChildRunRecord(BaseModel): + """A child run tracked under a name in the child run registry.""" + + model_config = ConfigDict(populate_by_name=True, alias_generator=to_camel) + + actor_id: str + """The Actor ID or name the child was started with, as the caller passed it.""" + + run_id: str + """ID of the current run under this name.""" + + previous_run_ids: list[str] = Field(default_factory=list) + """IDs of earlier runs under this name that failed and were replaced by a new run, oldest first.""" + + +_records_adapter = TypeAdapter(dict[str, ChildRunRecord]) + + +class ChildRunRegistry: + """Persisted name -> run map that lets named child runs survive a migration or resurrection of the parent. + + Every change is written to the key-value store right away. A hard kill of the parent between the platform + starting the child and that write can still orphan the child, since nothing but the platform knows about it. + """ + + def __init__(self, key_value_store: KeyValueStore) -> None: + self._key_value_store = key_value_store + self._records: dict[str, ChildRunRecord] | None = None + self._load_lock = asyncio.Lock() + self._write_lock = asyncio.Lock() + self._name_locks: defaultdict[str, asyncio.Lock] = defaultdict(asyncio.Lock) + + async def find_or_start( + self, + name: str, + *, + actor_id: str, + client: ApifyClientAsync, + start_run: Callable[[], Awaitable[Run]], + resurrect_run: Callable[[RunClientAsync], Awaitable[Run]], + ) -> tuple[Run, bool]: + """Return the run recorded under `name`, or start one when there is none to reuse. + + A recorded run that is `READY` or `RUNNING` is reattached and one that `SUCCEEDED` is returned as is. + 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. + + 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`. + client: Client used to look up and resurrect the recorded run. + start_run: Starts a new run of the Actor. + resurrect_run: Resurrects the recorded run, given its run client. + + Returns: + The run, and whether it was newly started. + """ + async with self._name_locks[name]: + records = await self._load() + record = records.get(name) + + if record is None: + return await self._start(name, actor_id=actor_id, start_run=start_run, previous_run_ids=[]), True + + if record.actor_id != actor_id: + raise ValueError( + f'Child run "{name}" is already recorded for Actor "{record.actor_id}", ' + f'it cannot be reused for Actor "{actor_id}".' + ) + + run_client = client.run(record.run_id) + run = await run_client.get() + + if run is not None and run.status in _SETTLING_STATUSES: + run = await run_client.wait_for_finish() + + if run is None or run.status == 'FAILED': + previous_run_ids = [*record.previous_run_ids, record.run_id] + run = await self._start(name, actor_id=actor_id, start_run=start_run, previous_run_ids=previous_run_ids) + return run, True + + 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 + + logger.info(f'Reattaching to child run "{name}"', extra={'run_id': run.id, 'status': run.status}) + return run, False + + async def _start( + self, + name: str, + *, + actor_id: str, + start_run: Callable[[], Awaitable[Run]], + previous_run_ids: list[str], + ) -> Run: + run = await start_run() + await self._save(name, ChildRunRecord(actor_id=actor_id, run_id=run.id, previous_run_ids=previous_run_ids)) + return run + + async def _load(self) -> dict[str, ChildRunRecord]: + async with self._load_lock: + if self._records is None: + stored = await self._key_value_store.get_value(CHILD_RUNS_KEY) + self._records = _records_adapter.validate_python(stored or {}) + return self._records + + async def _save(self, name: str, record: ChildRunRecord) -> None: + records = await self._load() + async with self._write_lock: + records[name] = record + await self._key_value_store.set_value( + CHILD_RUNS_KEY, _records_adapter.dump_python(records, by_alias=True, mode='json') + ) diff --git a/tests/unit/actor/test_actor_child_runs.py b/tests/unit/actor/test_actor_child_runs.py new file mode 100644 index 000000000..26c162255 --- /dev/null +++ b/tests/unit/actor/test_actor_child_runs.py @@ -0,0 +1,248 @@ +from __future__ import annotations + +import asyncio +from typing import TYPE_CHECKING, Any +from unittest.mock import MagicMock, Mock + +import pytest + +from apify_client._models import Run + +from apify import Actor +from apify._child_runs import CHILD_RUNS_KEY + +if TYPE_CHECKING: + from ..conftest import ApifyClientAsyncPatcher + + +def make_run(run_id: str, status: str) -> Run: + return Run.model_validate( + { + 'id': run_id, + 'actId': 'actor_id', + 'userId': 'user_id', + 'startedAt': '2024-08-08T12:12:44Z', + 'status': status, + 'meta': {'origin': 'API'}, + 'buildId': 'build_id', + 'defaultDatasetId': 'dataset_id', + 'defaultKeyValueStoreId': 'kvs_id', + 'defaultRequestQueueId': 'rq_id', + 'generalAccess': 'RESTRICTED', + 'stats': {'restartCount': 0, 'resurrectCount': 0, 'computeUnits': 0}, + 'options': {'build': '', 'timeoutSecs': 44, 'memoryMbytes': 4096, 'diskMbytes': 16384}, + } + ) + + +async def record_child_run(name: str, run_id: str, *, actor_id: str = 'some-actor') -> None: + """Seed the registry the way an earlier attempt of this Actor run would have left it.""" + kvs = await Actor.open_key_value_store() + await kvs.set_value(CHILD_RUNS_KEY, {name: {'actorId': actor_id, 'runId': run_id, 'previousRunIds': []}}) + + +async def test_named_start_records_run_in_kvs(apify_client_async_patcher: ApifyClientAsyncPatcher) -> None: + """A named start persists the name -> run ID entry to the default KVS as soon as the run starts.""" + apify_client_async_patcher.patch('actor', 'start', return_value=make_run('new-run', 'READY')) + + async with Actor: + run = await Actor.start('some-actor', name='scrape-eu') + kvs = await Actor.open_key_value_store() + stored = await kvs.get_value(CHILD_RUNS_KEY) + + assert run.id == 'new-run' + assert stored == {'scrape-eu': {'actorId': 'some-actor', 'runId': 'new-run', 'previousRunIds': []}} + + +async def test_unnamed_start_is_not_recorded(apify_client_async_patcher: ApifyClientAsyncPatcher) -> None: + """A start without a name leaves the registry untouched.""" + apify_client_async_patcher.patch('actor', 'start', return_value=make_run('new-run', 'READY')) + + async with Actor: + await Actor.start('some-actor') + kvs = await Actor.open_key_value_store() + stored = await kvs.get_value(CHILD_RUNS_KEY) + + assert stored is None + + +@pytest.mark.parametrize( + 'status', + [ + pytest.param('READY', id='ready'), + pytest.param('RUNNING', id='running'), + pytest.param('SUCCEEDED', id='succeeded'), + ], +) +async def test_named_start_reuses_recorded_run( + apify_client_async_patcher: ApifyClientAsyncPatcher, status: str +) -> None: + """A recorded run that is active or succeeded is returned without starting a new one.""" + apify_client_async_patcher.patch('actor', 'start', return_value=make_run('new-run', 'READY')) + apify_client_async_patcher.patch('run', 'get', return_value=make_run('old-run', status)) + + async with Actor: + await record_child_run('scrape-eu', 'old-run') + run = await Actor.start('some-actor', name='scrape-eu') + + assert run.id == 'old-run' + assert run.status == status + assert apify_client_async_patcher.calls['actor']['start'] == [] + + +@pytest.mark.parametrize( + 'status', + [ + pytest.param('ABORTED', id='aborted'), + pytest.param('TIMED-OUT', id='timed out'), + ], +) +async def test_named_start_resurrects_recorded_run( + apify_client_async_patcher: ApifyClientAsyncPatcher, status: str +) -> None: + """A recorded run that was aborted or timed out is resurrected, and its options are passed through.""" + apify_client_async_patcher.patch('actor', 'start', return_value=make_run('new-run', 'READY')) + apify_client_async_patcher.patch('run', 'get', return_value=make_run('old-run', status)) + apify_client_async_patcher.patch('run', 'resurrect', return_value=make_run('old-run', 'RUNNING')) + + async with Actor: + await record_child_run('scrape-eu', 'old-run') + run = await Actor.start('some-actor', name='scrape-eu', memory_mbytes=2048) + + assert run.id == 'old-run' + assert run.status == 'RUNNING' + assert apify_client_async_patcher.calls['actor']['start'] == [] + [(args, kwargs)] = apify_client_async_patcher.calls['run']['resurrect'] + assert args[0].resource_id == 'old-run' + assert kwargs['memory_mbytes'] == 2048 + + +@pytest.mark.parametrize( + 'status', + [ + pytest.param('ABORTING', id='aborting'), + pytest.param('TIMING-OUT', id='timing out'), + ], +) +async def test_named_start_resurrects_settling_run_after_it_finishes( + apify_client_async_patcher: ApifyClientAsyncPatcher, status: str +) -> None: + """A recorded run still aborting or timing out is waited for, then resurrected.""" + finished_status = {'ABORTING': 'ABORTED', 'TIMING-OUT': 'TIMED-OUT'}[status] + apify_client_async_patcher.patch('run', 'get', return_value=make_run('old-run', status)) + apify_client_async_patcher.patch('run', 'wait_for_finish', return_value=make_run('old-run', finished_status)) + apify_client_async_patcher.patch('run', 'resurrect', return_value=make_run('old-run', 'RUNNING')) + + async with Actor: + await record_child_run('scrape-eu', 'old-run') + run = await Actor.start('some-actor', name='scrape-eu') + + assert run.status == 'RUNNING' + assert len(apify_client_async_patcher.calls['run']['wait_for_finish']) == 1 + assert len(apify_client_async_patcher.calls['run']['resurrect']) == 1 + + +@pytest.mark.parametrize( + 'recorded_run', + [ + pytest.param(make_run('old-run', 'FAILED'), id='failed'), + pytest.param(None, id='not found'), + ], +) +async def test_named_start_replaces_failed_or_missing_run( + apify_client_async_patcher: ApifyClientAsyncPatcher, recorded_run: Run | None +) -> None: + """A recorded run that failed or no longer exists is replaced by a new run and kept in the history.""" + apify_client_async_patcher.patch('actor', 'start', return_value=make_run('new-run', 'READY')) + apify_client_async_patcher.patch('run', 'get', return_value=recorded_run) + + async with Actor: + await record_child_run('scrape-eu', 'old-run') + run = await Actor.start('some-actor', name='scrape-eu') + kvs = await Actor.open_key_value_store() + stored = await kvs.get_value(CHILD_RUNS_KEY) + + assert run.id == 'new-run' + assert stored == {'scrape-eu': {'actorId': 'some-actor', 'runId': 'new-run', 'previousRunIds': ['old-run']}} + + +async def test_named_start_rejects_name_recorded_for_another_actor( + apify_client_async_patcher: ApifyClientAsyncPatcher, +) -> None: + """Reusing a name for a different Actor raises instead of attaching to the other Actor's run.""" + apify_client_async_patcher.patch('actor', 'start', return_value=make_run('new-run', 'READY')) + + async with Actor: + await record_child_run('scrape-eu', 'old-run', actor_id='other-actor') + with pytest.raises(ValueError, match='already recorded for Actor "other-actor"'): + await Actor.start('some-actor', name='scrape-eu') + + assert apify_client_async_patcher.calls['actor']['start'] == [] + + +async def test_concurrent_named_starts_start_one_run(apify_client_async_patcher: ApifyClientAsyncPatcher) -> None: + """Concurrent starts under the same name start a single run and share it.""" + started = make_run('new-run', 'RUNNING') + + async def slow_start(*_args: Any, **_kwargs: Any) -> Run: + await asyncio.sleep(0.05) + return started + + apify_client_async_patcher.patch('actor', 'start', replacement_method=slow_start) + apify_client_async_patcher.patch('run', 'get', return_value=started) + + async with Actor: + runs = await asyncio.gather(*(Actor.start('some-actor', name='scrape-eu') for _ in range(3))) + + assert {run.id for run in runs} == {'new-run'} + assert len(apify_client_async_patcher.calls['actor']['start']) == 1 + + +async def test_named_call_waits_for_reattached_run(apify_client_async_patcher: ApifyClientAsyncPatcher) -> None: + """A named call waits for the reattached run and streams only its new log lines.""" + streamed_log = MagicMock() + get_streamed_log = Mock(return_value=streamed_log) + apify_client_async_patcher.patch('run', 'get', return_value=make_run('old-run', 'RUNNING')) + apify_client_async_patcher.patch('run', 'wait_for_finish', return_value=make_run('old-run', 'SUCCEEDED')) + apify_client_async_patcher.patch('run', 'get_status_message_watcher', return_value=MagicMock()) + apify_client_async_patcher.patch('run', 'get_streamed_log', replacement_method=get_streamed_log) + + async with Actor: + await record_child_run('scrape-eu', 'old-run') + run = await Actor.call('some-actor', name='scrape-eu') + + assert run.status == 'SUCCEEDED' + assert get_streamed_log.call_args.kwargs['from_start'] is False + streamed_log.__aenter__.assert_awaited_once() + + +async def test_named_call_streams_new_run_log_from_start(apify_client_async_patcher: ApifyClientAsyncPatcher) -> None: + """A named call that starts a new run streams its log from the start.""" + get_streamed_log = Mock(return_value=MagicMock()) + apify_client_async_patcher.patch('actor', 'start', return_value=make_run('new-run', 'READY')) + apify_client_async_patcher.patch('run', 'wait_for_finish', return_value=make_run('new-run', 'SUCCEEDED')) + apify_client_async_patcher.patch('run', 'get_status_message_watcher', return_value=MagicMock()) + apify_client_async_patcher.patch('run', 'get_streamed_log', replacement_method=get_streamed_log) + + async with Actor: + run = await Actor.call('some-actor', name='scrape-eu') + + assert run.id == 'new-run' + assert run.status == 'SUCCEEDED' + assert get_streamed_log.call_args.kwargs['from_start'] is True + assert apify_client_async_patcher.calls['actor']['call'] == [] + + +async def test_named_call_returns_succeeded_run_without_waiting( + apify_client_async_patcher: ApifyClientAsyncPatcher, +) -> None: + """A named call whose recorded run already succeeded returns it without waiting or streaming logs.""" + apify_client_async_patcher.patch('run', 'get', return_value=make_run('old-run', 'SUCCEEDED')) + apify_client_async_patcher.patch('run', 'wait_for_finish', return_value=None) + + async with Actor: + await record_child_run('scrape-eu', 'old-run') + run = await Actor.call('some-actor', name='scrape-eu') + + assert run.id == 'old-run' + assert apify_client_async_patcher.calls['run']['wait_for_finish'] == [] From 9deab5f9dce7ed8ae0a29a160914fb884ef92763 Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Tue, 29 Sep 2026 14:00:03 +0200 Subject: [PATCH 2/6] docs: mention missing recorded runs in the child run name docstring --- src/apify/_actor.py | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/src/apify/_actor.py b/src/apify/_actor.py index 0f4f380a6..725a9bde0 100644 --- a/src/apify/_actor.py +++ b/src/apify/_actor.py @@ -961,8 +961,8 @@ async def start( name: Optional name of the child run, unique within this Actor run. A named run is recorded in the default key-value store, so after a migration or resurrection of this Actor the same call reattaches to the recorded run. A `SUCCEEDED` run is returned as is, an `ABORTED` or `TIMED-OUT` one is - resurrected, and a new run is started only when nothing is recorded under the name or the recorded - run `FAILED`. + resurrected, and a new run is started only when nothing is recorded under the name, or the recorded + run `FAILED` or no longer exists. Returns: Info about the started Actor run @@ -1095,8 +1095,8 @@ async def call( name: Optional name of the child run, unique within this Actor run. A named run is recorded in the default key-value store, so after a migration or resurrection of this Actor the same call reattaches to the recorded run. A `SUCCEEDED` run is returned as is, an `ABORTED` or `TIMED-OUT` one is - resurrected, and a new run is started only when nothing is recorded under the name or the recorded - run `FAILED`. + resurrected, and a new run is started only when nothing is recorded under the name, or the recorded + run `FAILED` or no longer exists. Returns: Info about the started Actor run. From 58bb12db429dc0aa576c24f0358756be7553c9e9 Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Tue, 29 Sep 2026 14:00:04 +0200 Subject: [PATCH 3/6] fix: share one child run registry across concurrent named starts --- src/apify/_actor.py | 5 +--- src/apify/_child_runs.py | 10 +++++--- tests/unit/actor/test_actor_child_runs.py | 31 +++++++++++++++++++++++ 3 files changed, 38 insertions(+), 8 deletions(-) diff --git a/src/apify/_actor.py b/src/apify/_actor.py index 725a9bde0..ec6ff4a72 100644 --- a/src/apify/_actor.py +++ b/src/apify/_actor.py @@ -151,7 +151,7 @@ def __init__( # Keep track of all used state stores to persist their values on exit self._use_state_stores: set[str | None] = set() - self._child_run_registry: ChildRunRegistry | None = None + self._child_run_registry = ChildRunRegistry(self.open_key_value_store) self._active = False """Whether the Actor instance is currently active (initialized and within context).""" @@ -1174,9 +1174,6 @@ async def _find_or_start_child_run( memory_mbytes: int | None, run_timeout: timedelta | None, ) -> tuple[Run, bool]: - if self._child_run_registry is None: - self._child_run_registry = ChildRunRegistry(await self.open_key_value_store()) - return await self._child_run_registry.find_or_start( name, actor_id=actor_id, diff --git a/src/apify/_child_runs.py b/src/apify/_child_runs.py index 63fa9024c..5c9f092c3 100644 --- a/src/apify/_child_runs.py +++ b/src/apify/_child_runs.py @@ -53,8 +53,8 @@ class ChildRunRegistry: starting the child and that write can still orphan the child, since nothing but the platform knows about it. """ - def __init__(self, key_value_store: KeyValueStore) -> None: - self._key_value_store = key_value_store + def __init__(self, open_key_value_store: Callable[[], Awaitable[KeyValueStore]]) -> None: + self._open_key_value_store = open_key_value_store self._records: dict[str, ChildRunRecord] | None = None self._load_lock = asyncio.Lock() self._write_lock = asyncio.Lock() @@ -131,14 +131,16 @@ async def _start( async def _load(self) -> dict[str, ChildRunRecord]: async with self._load_lock: if self._records is None: - stored = await self._key_value_store.get_value(CHILD_RUNS_KEY) + key_value_store = await self._open_key_value_store() + stored = await key_value_store.get_value(CHILD_RUNS_KEY) self._records = _records_adapter.validate_python(stored or {}) return self._records async def _save(self, name: str, record: ChildRunRecord) -> None: records = await self._load() + key_value_store = await self._open_key_value_store() async with self._write_lock: records[name] = record - await self._key_value_store.set_value( + await key_value_store.set_value( CHILD_RUNS_KEY, _records_adapter.dump_python(records, by_alias=True, mode='json') ) diff --git a/tests/unit/actor/test_actor_child_runs.py b/tests/unit/actor/test_actor_child_runs.py index 26c162255..fd6bfeb9c 100644 --- a/tests/unit/actor/test_actor_child_runs.py +++ b/tests/unit/actor/test_actor_child_runs.py @@ -9,10 +9,12 @@ from apify_client._models import Run from apify import Actor +from apify._actor import _ActorType from apify._child_runs import CHILD_RUNS_KEY if TYPE_CHECKING: from ..conftest import ApifyClientAsyncPatcher + from apify.storages import KeyValueStore def make_run(run_id: str, status: str) -> Run: @@ -198,6 +200,35 @@ async def slow_start(*_args: Any, **_kwargs: Any) -> Run: assert len(apify_client_async_patcher.calls['actor']['start']) == 1 +async def test_concurrent_first_named_starts_share_one_registry( + apify_client_async_patcher: ApifyClientAsyncPatcher, monkeypatch: pytest.MonkeyPatch +) -> None: + """Concurrent first named starts start a single run even when opening the default KVS yields to the event loop.""" + started = make_run('new-run', 'RUNNING') + + async def slow_start(*_args: Any, **_kwargs: Any) -> Run: + await asyncio.sleep(0.05) + return started + + apify_client_async_patcher.patch('actor', 'start', replacement_method=slow_start) + apify_client_async_patcher.patch('run', 'get', return_value=started) + open_key_value_store = _ActorType.open_key_value_store + + async def yielding_open_key_value_store(self: _ActorType, *args: Any, **kwargs: Any) -> KeyValueStore: + # On the platform the default KVS is opened lazily through the API, so opening it suspends. + await asyncio.sleep(0.01) + return await open_key_value_store(self, *args, **kwargs) + + monkeypatch.setattr(_ActorType, 'open_key_value_store', yielding_open_key_value_store) + + # A fresh instance, since the registry binds the opener when the Actor is created. + async with _ActorType() as actor: + runs = await asyncio.gather(*(actor.start('some-actor', name='scrape-eu') for _ in range(3))) + + assert {run.id for run in runs} == {'new-run'} + assert len(apify_client_async_patcher.calls['actor']['start']) == 1 + + async def test_named_call_waits_for_reattached_run(apify_client_async_patcher: ApifyClientAsyncPatcher) -> None: """A named call waits for the reattached run and streams only its new log lines.""" streamed_log = MagicMock() From 8575cf74a44a9f55baf5d7226cbe6e88a3946452 Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Tue, 29 Sep 2026 15:01:06 +0200 Subject: [PATCH 4/6] test: add e2e tests for named child runs --- tests/e2e/test_actor_child_runs.py | 86 ++++++++++++++++++++++++++++++ 1 file changed, 86 insertions(+) create mode 100644 tests/e2e/test_actor_child_runs.py diff --git a/tests/e2e/test_actor_child_runs.py b/tests/e2e/test_actor_child_runs.py new file mode 100644 index 000000000..c9d0d9ccf --- /dev/null +++ b/tests/e2e/test_actor_child_runs.py @@ -0,0 +1,86 @@ +from __future__ import annotations + +import asyncio +from typing import TYPE_CHECKING + +from apify import Actor + +if TYPE_CHECKING: + from .conftest import MakeActorFunction, RunActorFunction + + +async def test_named_child_run_is_reattached_after_reboot( + make_actor: MakeActorFunction, + run_actor: RunActorFunction, +) -> None: + """A named child run started before a reboot is reattached and awaited by a named call after it.""" + + 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(20) + return + + actor_id = Actor.configuration.actor_id or '' + child_run_id = await Actor.get_value('child_run_id') + + if child_run_id is None: + run = await Actor.start(actor_id=actor_id, run_input={'is_child': True}, name='child') + await Actor.set_value('child_run_id', run.id) + await Actor.reboot() + return + + run = await Actor.call(actor_id=actor_id, run_input={'is_child': True}, name='child') + assert run is not None, 'run is None' + assert run.id == child_run_id, f'run.id={run.id}, child_run_id={child_run_id}' + assert run.status == 'SUCCEEDED', f'run.status={run.status}' + + actor = await make_actor(label='child-run-reattach', main_func=main) + run_result = await run_actor(actor) + + assert run_result.status == 'SUCCEEDED' + # The parent run and the one child run it reattached to. + assert (await actor.runs().list()).total == 2 + + +async def test_named_aborted_child_run_is_resurrected_after_reboot( + make_actor: MakeActorFunction, + run_actor: RunActorFunction, +) -> None: + """A named child run aborted before a reboot is resurrected by a named start after it.""" + + 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(300) + return + + actor_id = Actor.configuration.actor_id or '' + child_run_id = await Actor.get_value('child_run_id') + + if child_run_id is None: + run = await Actor.start(actor_id=actor_id, run_input={'is_child': True}, name='child') + await Actor.set_value('child_run_id', run.id) + run_client = Actor.apify_client.run(run.id) + await run_client.abort() + aborted_run = await run_client.wait_for_finish() + assert aborted_run is not None, 'aborted_run is None' + assert aborted_run.status == 'ABORTED', f'aborted_run.status={aborted_run.status}' + await Actor.reboot() + return + + run = await Actor.start(actor_id=actor_id, run_input={'is_child': True}, name='child') + try: + assert run.id == child_run_id, f'run.id={run.id}, child_run_id={child_run_id}' + assert run.status in {'READY', 'RUNNING'}, f'run.status={run.status}' + finally: + await Actor.apify_client.run(run.id).abort() + + actor = await make_actor(label='child-run-resurrect', main_func=main) + run_result = await run_actor(actor) + + assert run_result.status == 'SUCCEEDED' + # The parent run and the one child run it resurrected. + assert (await actor.runs().list()).total == 2 From 59202e15a1c277e9a618a7ff445b834d31af1b2e Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Wed, 30 Sep 2026 09:21:48 +0200 Subject: [PATCH 5/6] docs: document named child runs --- .../06_interacting_with_other_actors.mdx | 24 +++++++++++++++++++ .../code/06_interacting_named_call.py | 21 ++++++++++++++++ 2 files changed, 45 insertions(+) create mode 100644 docs/02_concepts/code/06_interacting_named_call.py diff --git a/docs/02_concepts/06_interacting_with_other_actors.mdx b/docs/02_concepts/06_interacting_with_other_actors.mdx index 1a362d7a1..2f6c195dc 100644 --- a/docs/02_concepts/06_interacting_with_other_actors.mdx +++ b/docs/02_concepts/06_interacting_with_other_actors.mdx @@ -8,6 +8,7 @@ import RunnableCodeBlock from '@site/src/components/RunnableCodeBlock'; import InteractingStartExample from '!!raw-loader!roa-loader!./code/06_interacting_start.py'; import InteractingCallExample from '!!raw-loader!roa-loader!./code/06_interacting_call.py'; +import InteractingNamedCallExample from '!!raw-loader!roa-loader!./code/06_interacting_named_call.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'; @@ -34,6 +35,29 @@ The `Actor.call` method starts another {InteractingCallExample} +## Named child runs + +When your Actor migrates to another server or is resurrected, it starts again from the beginning. An unnamed `Actor.start` or `Actor.call` then starts a second child run, and the first one keeps running with nobody waiting for it. + +To avoid the duplicate, pass a `name` to `Actor.start` or `Actor.call`. The name must be unique within your Actor run. Right after the child run starts, the SDK records its name and run ID under the `APIFY_CHILD_RUNS` key in the default key-value store. After a restart, the same call looks up the recorded run and reuses it based on its status: + +- `READY` or `RUNNING`: the call reattaches to the run. +- `SUCCEEDED`: the call returns the run as is. +- `ABORTED` or `TIMED-OUT`: the call resurrects the run. A run that is still `ABORTING` or `TIMING-OUT` is waited for first. +- `FAILED`, or the run no longer exists: the call starts a new run under the same name. + +A named `Actor.call` that reattaches to a run streams only the log lines the child writes from then on, so the parent log doesn't repeat what the previous attempt already printed. + + + {InteractingNamedCallExample} + + +Note that: + +- The name is bound to the `actor_id` you pass. Using the same name for a different Actor, or for the same Actor referenced by its ID instead of its name, raises a `ValueError`. +- Concurrent calls under one name in the same Actor run share a single child run. +- If your Actor is killed after the platform starts the child but before the SDK records it, the child run isn't recorded and the next attempt starts a new one. + ## 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_named_call.py b/docs/02_concepts/code/06_interacting_named_call.py new file mode 100644 index 000000000..ac7f58f6f --- /dev/null +++ b/docs/02_concepts/code/06_interacting_named_call.py @@ -0,0 +1,21 @@ +import asyncio + +from apify import Actor + + +async def main() -> None: + async with Actor: + # Call the apify/screenshot-url Actor under the name 'screenshot'. If this run + # migrates while the child is running, the same call after the restart waits + # for the recorded child run instead of starting a new one. + actor_run = await Actor.call( + actor_id='apify/screenshot-url', + run_input={'urls': [{'url': 'https://www.apify.com/'}]}, + name='screenshot', + ) + + Actor.log.info(f'Child run {actor_run.id} finished with {actor_run.status}') + + +if __name__ == '__main__': + asyncio.run(main()) From 82366a3d6233c7d3f144af12a1e9a3b26bfdef19 Mon Sep 17 00:00:00 2001 From: Vlada Dusek Date: Wed, 30 Sep 2026 09:21:49 +0200 Subject: [PATCH 6/6] feat: reject a malformed child run registry and cover the remaining call paths --- src/apify/_actor.py | 6 ++- src/apify/_child_runs.py | 10 ++++- tests/unit/actor/test_actor_child_runs.py | 53 +++++++++++++++++++++++ 3 files changed, 65 insertions(+), 4 deletions(-) diff --git a/src/apify/_actor.py b/src/apify/_actor.py index a67bcf65c..8aa532dfc 100644 --- a/src/apify/_actor.py +++ b/src/apify/_actor.py @@ -977,7 +977,8 @@ async def start( default key-value store, so after a migration or resurrection of this Actor the same call reattaches to the recorded run. A `SUCCEEDED` run is returned as is, an `ABORTED` or `TIMED-OUT` one is resurrected, and a new run is started only when nothing is recorded under the name, or the recorded - run `FAILED` or no longer exists. + run `FAILED` or no longer exists. The name is bound to `actor_id` exactly as passed, so reusing it with + any other value raises a `ValueError`. Returns: Info about the started Actor run @@ -1111,7 +1112,8 @@ async def call( default key-value store, so after a migration or resurrection of this Actor the same call reattaches to the recorded run. A `SUCCEEDED` run is returned as is, an `ABORTED` or `TIMED-OUT` one is resurrected, and a new run is started only when nothing is recorded under the name, or the recorded - run `FAILED` or no longer exists. + run `FAILED` or no longer exists. The name is bound to `actor_id` exactly as passed, so reusing it with + any other value raises a `ValueError`. Returns: Info about the started Actor run. diff --git a/src/apify/_child_runs.py b/src/apify/_child_runs.py index 5c9f092c3..05427df0a 100644 --- a/src/apify/_child_runs.py +++ b/src/apify/_child_runs.py @@ -5,7 +5,7 @@ from logging import getLogger from typing import TYPE_CHECKING -from pydantic import BaseModel, ConfigDict, Field, TypeAdapter +from pydantic import BaseModel, ConfigDict, Field, TypeAdapter, ValidationError from pydantic.alias_generators import to_camel if TYPE_CHECKING: @@ -133,7 +133,13 @@ async def _load(self) -> dict[str, ChildRunRecord]: if self._records is None: key_value_store = await self._open_key_value_store() stored = await key_value_store.get_value(CHILD_RUNS_KEY) - self._records = _records_adapter.validate_python(stored or {}) + try: + self._records = _records_adapter.validate_python(stored or {}) + except ValidationError as exc: + raise ValueError( + f'The child run registry under the "{CHILD_RUNS_KEY}" key in the default key-value store ' + 'is malformed.' + ) from exc return self._records async def _save(self, name: str, record: ChildRunRecord) -> None: diff --git a/tests/unit/actor/test_actor_child_runs.py b/tests/unit/actor/test_actor_child_runs.py index fd6bfeb9c..b5d297172 100644 --- a/tests/unit/actor/test_actor_child_runs.py +++ b/tests/unit/actor/test_actor_child_runs.py @@ -1,6 +1,7 @@ from __future__ import annotations import asyncio +from datetime import timedelta from typing import TYPE_CHECKING, Any from unittest.mock import MagicMock, Mock @@ -277,3 +278,55 @@ async def test_named_call_returns_succeeded_run_without_waiting( assert run.id == 'old-run' assert apify_client_async_patcher.calls['run']['wait_for_finish'] == [] + + +async def test_named_call_waits_for_resurrected_run(apify_client_async_patcher: ApifyClientAsyncPatcher) -> None: + """A named call resurrects an aborted run with its own timeout and streams only the new log lines.""" + get_streamed_log = Mock(return_value=MagicMock()) + apify_client_async_patcher.patch('run', 'get', return_value=make_run('old-run', 'ABORTED')) + apify_client_async_patcher.patch('run', 'resurrect', return_value=make_run('old-run', 'RUNNING')) + apify_client_async_patcher.patch('run', 'wait_for_finish', return_value=make_run('old-run', 'SUCCEEDED')) + apify_client_async_patcher.patch('run', 'get_status_message_watcher', return_value=MagicMock()) + apify_client_async_patcher.patch('run', 'get_streamed_log', replacement_method=get_streamed_log) + + async with Actor: + await record_child_run('scrape-eu', 'old-run') + run = await Actor.call('some-actor', name='scrape-eu', timeout=timedelta(minutes=5)) + + assert run.id == 'old-run' + assert run.status == 'SUCCEEDED' + assert apify_client_async_patcher.calls['actor']['start'] == [] + [(_, kwargs)] = apify_client_async_patcher.calls['run']['resurrect'] + assert kwargs['run_timeout'] == timedelta(minutes=5) + assert get_streamed_log.call_args.kwargs['from_start'] is False + + +async def test_named_call_without_logger_only_waits(apify_client_async_patcher: ApifyClientAsyncPatcher) -> None: + """A named call with `logger=None` waits for the run without redirecting its log or status messages.""" + get_streamed_log = Mock(return_value=MagicMock()) + get_status_message_watcher = Mock(return_value=MagicMock()) + apify_client_async_patcher.patch('actor', 'start', return_value=make_run('new-run', 'READY')) + apify_client_async_patcher.patch('run', 'wait_for_finish', return_value=make_run('new-run', 'SUCCEEDED')) + apify_client_async_patcher.patch('run', 'get_status_message_watcher', replacement_method=get_status_message_watcher) + apify_client_async_patcher.patch('run', 'get_streamed_log', replacement_method=get_streamed_log) + + async with Actor: + run = await Actor.call('some-actor', name='scrape-eu', logger=None) + + assert run.status == 'SUCCEEDED' + assert len(apify_client_async_patcher.calls['run']['wait_for_finish']) == 1 + get_streamed_log.assert_not_called() + get_status_message_watcher.assert_not_called() + + +async def test_named_start_rejects_malformed_registry(apify_client_async_patcher: ApifyClientAsyncPatcher) -> None: + """A malformed registry in the default KVS raises a `ValueError` naming the key, without starting a run.""" + apify_client_async_patcher.patch('actor', 'start', return_value=make_run('new-run', 'READY')) + + async with Actor: + kvs = await Actor.open_key_value_store() + await kvs.set_value(CHILD_RUNS_KEY, {'scrape-eu': {'runId': 'old-run'}}) + with pytest.raises(ValueError, match=CHILD_RUNS_KEY): + await Actor.start('some-actor', name='scrape-eu') + + assert apify_client_async_patcher.calls['actor']['start'] == []