Skip to content
Open
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
2 changes: 2 additions & 0 deletions src/apify/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
)

from apify._actor import Actor
from apify._child_runs import ChildRunInfo
from apify._configuration import Configuration
from apify._consts import ActorEnvVars, ApifyEnvVars
from apify._proxy_configuration import ProxyConfiguration, ProxyInfo
Expand All @@ -26,6 +27,7 @@
'ActorEnvVars',
'ActorEventTypes',
'ApifyEnvVars',
'ChildRunInfo',
'Configuration',
'Event',
'EventAbortingData',
Expand Down
15 changes: 14 additions & 1 deletion src/apify/_actor.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@
ChargingManagerImplementation,
charge_lock_if_charging,
)
from apify._child_runs import ChildRunRegistry
from apify._child_runs import ChildRunInfo, 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 Down Expand Up @@ -1229,6 +1229,19 @@ async def _wait_for_child_run(
async with status_redirector, streamed_log:
return await run_client.wait_for_finish(wait_duration=wait)

@_ensure_context
async def child_runs(self) -> dict[str, ChildRunInfo]:
"""Get the named child runs of this Actor run, with their current state.

Every run started by `Actor.start`, `Actor.call` or `Actor.call_task` with a `run_name` is included, even one
started before a migration or resurrection of this Actor run. Runs started without a `run_name` are not
tracked. Each run is fetched from the API when this method is called, so the result is a snapshot.

Returns:
The child runs by name.
"""
return await self._child_run_registry.list_runs(self.apify_client)

@_ensure_context
async def call_task(
self,
Expand Down
46 changes: 45 additions & 1 deletion src/apify/_child_runs.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,12 +2,15 @@

import asyncio
from collections import defaultdict
from dataclasses import dataclass
from logging import getLogger
from typing import TYPE_CHECKING, Self

from pydantic import BaseModel, ConfigDict, Field, TypeAdapter, ValidationError, model_validator
from pydantic.alias_generators import to_camel

from apify._utils import docs_group

if TYPE_CHECKING:
from collections.abc import Awaitable, Callable

Expand Down Expand Up @@ -43,7 +46,7 @@ class ChildRunRecord(BaseModel):
"""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."""
"""IDs of earlier runs under this name that failed or went missing and were replaced by a new run, oldest first."""

@model_validator(mode='after')
def _check_started_from(self) -> Self:
Expand All @@ -56,6 +59,27 @@ def _describe_started_from(actor_id: str | None, task_id: str | None) -> str:
return f'Actor "{actor_id}"' if actor_id is not None else f'task "{task_id}"'


@docs_group('Actor')
@dataclass(frozen=True)
class ChildRunInfo:
"""A named child run of this Actor run, as returned by `Actor.child_runs`."""

actor_id: str | None
"""The Actor ID or name the child was started with, as the caller passed it, or `None` for a task run."""

task_id: str | None
"""The task ID or name the child was started with, as the caller passed it, or `None` for an Actor run."""

run_id: str
"""ID of the current run under this name."""

run: Run | None
"""The current run as the API returns it now, or `None` when the platform no longer knows it."""

previous_run_ids: list[str]
"""IDs of earlier runs under this name that failed or went missing and were replaced by a new run, oldest first."""


_records_adapter = TypeAdapter(dict[str, ChildRunRecord])


Expand Down Expand Up @@ -137,6 +161,26 @@ async def find_or_start(
logger.info(f'Reattaching to child run "{name}"', extra={'run_id': run.id, 'status': run.status})
return run, False

async def list_runs(self, client: ApifyClientAsync) -> dict[str, ChildRunInfo]:
"""Return every recorded child run by name, with its current state fetched from the API.

Args:
client: Client used to fetch the recorded runs.
"""
# Copy the records, since a named start can add one while the runs are fetched.
records = dict(await self._load())
runs = await asyncio.gather(*(client.run(record.run_id).get() for record in records.values()))
return {
name: ChildRunInfo(
actor_id=record.actor_id,
task_id=record.task_id,
run_id=record.run_id,
run=run,
previous_run_ids=list(record.previous_run_ids),
)
for (name, record), run in zip(records.items(), runs, strict=True)
}

async def _start(
self,
name: str,
Expand Down
9 changes: 8 additions & 1 deletion tests/e2e/test_actor_child_runs.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ 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."""
"""A named child run started before a reboot is reattached, awaited by a named call, and listed after it."""

async def main() -> None:
async with Actor:
Expand All @@ -36,6 +36,13 @@ async def main() -> 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}'

child_runs = await Actor.child_runs()
assert child_runs.keys() == {'child'}, f'child_runs={child_runs}'
assert child_runs['child'].run_id == child_run_id, f'child_runs={child_runs}'
child_run = child_runs['child'].run
assert child_run is not None, 'child_run is None'
assert child_run.status == 'SUCCEEDED', f'child_run.status={child_run.status}'

actor = await make_actor(label='child-run-reattach', main_func=main)
run_result = await run_actor(actor)

Expand Down
62 changes: 62 additions & 0 deletions tests/unit/actor/test_actor_child_runs.py
Original file line number Diff line number Diff line change
Expand Up @@ -425,3 +425,65 @@ async def test_registry_rejects_record_without_actor_or_task(
await Actor.start('some-actor', run_name='scrape-eu')

assert apify_client_async_patcher.calls['actor']['start'] == []


async def test_child_runs_is_empty_without_named_runs(apify_client_async_patcher: ApifyClientAsyncPatcher) -> None:
"""`Actor.child_runs` returns an empty dict and calls no API when nothing is recorded."""
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('new-run', 'READY'))

async with Actor:
await Actor.start('some-actor')
child_runs = await Actor.child_runs()

assert child_runs == {}
assert apify_client_async_patcher.calls['run']['get'] == []


async def test_child_runs_returns_recorded_runs_with_current_state(
apify_client_async_patcher: ApifyClientAsyncPatcher,
) -> None:
"""`Actor.child_runs` returns each recorded run with its fetched state and history, `None` for a missing run."""
runs = {'eu-run': make_run('eu-run', 'RUNNING'), 'us-run': None}

async def get_run(run_client: Any, *_args: Any, **_kwargs: Any) -> Run | None:
return runs[run_client.resource_id]

apify_client_async_patcher.patch('run', 'get', replacement_method=get_run)

async with Actor:
kvs = await Actor.open_key_value_store()
await kvs.set_value(
CHILD_RUNS_KEY,
{
'scrape-eu': {'actorId': 'some-actor', 'runId': 'eu-run', 'previousRunIds': ['failed-run']},
'scrape-us': {'taskId': 'some-task', 'runId': 'us-run', 'previousRunIds': []},
},
)
child_runs = await Actor.child_runs()

assert child_runs.keys() == {'scrape-eu', 'scrape-us'}
assert child_runs['scrape-eu'].actor_id == 'some-actor'
assert child_runs['scrape-eu'].task_id is None
assert child_runs['scrape-eu'].run_id == 'eu-run'
assert child_runs['scrape-eu'].run == runs['eu-run']
assert child_runs['scrape-eu'].previous_run_ids == ['failed-run']
assert child_runs['scrape-us'].actor_id is None
assert child_runs['scrape-us'].task_id == 'some-task'
assert child_runs['scrape-us'].run is None


async def test_child_runs_includes_run_started_in_this_attempt(
apify_client_async_patcher: ApifyClientAsyncPatcher,
) -> None:
"""A run started by a named start in the same attempt shows up in `Actor.child_runs` right away."""
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('new-run', 'RUNNING'))

async with Actor:
await Actor.start('some-actor', run_name='scrape-eu')
child_runs = await Actor.child_runs()

assert child_runs['scrape-eu'].run_id == 'new-run'
assert child_runs['scrape-eu'].run is not None
assert child_runs['scrape-eu'].run.status == 'RUNNING'
Loading