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'] == []