Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 24 additions & 0 deletions docs/02_concepts/06_interacting_with_other_actors.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand All @@ -34,6 +35,29 @@ The <ApiLink to="class/Actor#call">`Actor.call`</ApiLink> method starts another
{InteractingCallExample}
</RunnableCodeBlock>

## 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 <ApiLink to="class/Actor#start">`Actor.start`</ApiLink> or <ApiLink to="class/Actor#call">`Actor.call`</ApiLink>. 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.

<RunnableCodeBlock className="language-python" language="python">
{InteractingNamedCallExample}
</RunnableCodeBlock>

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 <ApiLink to="class/Actor#call_task">`Actor.call_task`</ApiLink> 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.
Expand Down
21 changes: 21 additions & 0 deletions docs/02_concepts/code/06_interacting_named_call.py
Original file line number Diff line number Diff line change
@@ -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())
145 changes: 129 additions & 16 deletions src/apify/_actor.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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
Expand All @@ -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

Expand Down Expand Up @@ -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)."""

Expand Down Expand Up @@ -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.

Expand All @@ -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
Expand All @@ -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,
Expand All @@ -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,
Expand Down Expand Up @@ -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.

Expand Down Expand Up @@ -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.
Expand All @@ -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,
Expand Down
Loading
Loading