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())
diff --git a/src/apify/_actor.py b/src/apify/_actor.py
index 3c5845ee5..8aa532dfc 100644
--- a/src/apify/_actor.py
+++ b/src/apify/_actor.py
@@ -7,7 +7,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 pathlib import Path
from typing import TYPE_CHECKING, Any, Literal, TypeVar, cast, overload
@@ -35,6 +35,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
@@ -49,13 +50,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
@@ -151,6 +153,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(self.open_key_value_store)
+
self._active = False
"""Whether the Actor instance is currently active (initialized and within context)."""
@@ -942,6 +946,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.
@@ -968,6 +973,12 @@ 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` 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
@@ -984,7 +995,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,
@@ -996,6 +1008,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,
@@ -1050,6 +1078,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.
@@ -1079,6 +1108,12 @@ 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` 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.
@@ -1095,25 +1130,103 @@ 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]:
+ 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..05427df0a
--- /dev/null
+++ b/src/apify/_child_runs.py
@@ -0,0 +1,152 @@
+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, ValidationError
+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, 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()
+ 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:
+ key_value_store = await self._open_key_value_store()
+ stored = await key_value_store.get_value(CHILD_RUNS_KEY)
+ 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:
+ records = await self._load()
+ key_value_store = await self._open_key_value_store()
+ async with self._write_lock:
+ records[name] = record
+ await key_value_store.set_value(
+ CHILD_RUNS_KEY, _records_adapter.dump_python(records, by_alias=True, mode='json')
+ )
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
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..b5d297172
--- /dev/null
+++ b/tests/unit/actor/test_actor_child_runs.py
@@ -0,0 +1,332 @@
+from __future__ import annotations
+
+import asyncio
+from datetime import timedelta
+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._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:
+ 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_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()
+ 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'] == []
+
+
+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'] == []