Skip to content
Merged
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
9 changes: 9 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 InteractingStartTaskExample from '!!raw-loader!roa-loader!./code/06_interacting_start_task.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,14 @@ The <ApiLink to="class/Actor#call">`Actor.call`</ApiLink> method starts another
{InteractingCallExample}
</RunnableCodeBlock>

## Actor start task

The <ApiLink to="class/Actor#start_task">`Actor.start_task`</ApiLink> method starts an [Actor task](https://docs.apify.com/platform/actors/tasks) on the Apify platform, and immediately returns the details of the started Actor run.

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

## 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
16 changes: 16 additions & 0 deletions docs/02_concepts/code/06_interacting_start_task.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
import asyncio

from apify import Actor


async def main() -> None:
async with Actor:
# Start the Actor task by its ID without waiting for it to finish.
actor_run = await Actor.start_task(task_id='Z3m6FPSj0GYZ25rQc')

# Log the run ID, which you can use to check on the run later.
Actor.log.info(f'Started task run: {actor_run.id}')


if __name__ == '__main__':
asyncio.run(main())
56 changes: 56 additions & 0 deletions src/apify/_actor.py
Original file line number Diff line number Diff line change
Expand Up @@ -1153,6 +1153,62 @@ async def call(

return run

@_ensure_context
async def start_task(
self,
task_id: str,
task_input: dict | None = None,
*,
build: str | None = None,
max_items: int | None = None,
restart_on_error: bool | None = None,
memory_mbytes: int | None = None,
timeout: timedelta | Literal['inherit'] | None = None,
webhooks: list[Webhook] | None = None,
token: str | None = None,
) -> Run:
"""Start an Actor task on the Apify Platform.

Unlike `Actor.call_task`, this method just starts the run without waiting for finish. To wait for the run to
finish, use `Actor.call_task` instead.

Note that an Actor task is a saved input configuration and options for an Actor. If you want to run an Actor
directly rather than an Actor task, use `Actor.start`.

Args:
task_id: The ID of the Actor task to be run.
task_input: Overrides the input to pass to the Actor run.
token: The Apify API token to use for this request (defaults to the `APIFY_TOKEN` environment variable).
build: Specifies the Actor build to run. It can be either a build tag or build number. By default,
the run uses the build specified in the default run configuration for the Actor (typically latest).
max_items: Maximum number of dataset items you are charged for, for pay-per-result Actors. It caps the
charge, not the output, so the run can return fewer or more items than this.
restart_on_error: If true, the Task run process will be restarted whenever it exits with
a non-zero status code.
memory_mbytes: Memory limit for the run, in megabytes. By default, the run uses a memory limit specified
in the default run configuration for the Actor.
timeout: Optional timeout for the run. By default, the run uses timeout specified in
the default run configuration for the Actor. Using `inherit` will set timeout of the other Actor to the
time remaining from this Actor timeout.
webhooks: Optional webhooks (https://docs.apify.com/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, you do not have to add it again here.

Returns:
Info about the started Actor run.
"""
client = self.new_client(token=token) if token else self.apify_client
task_client = client.task(task_id)
return await task_client.start(
task_input=task_input,
build=build,
max_items=max_items,
restart_on_error=restart_on_error,
memory_mbytes=memory_mbytes,
run_timeout=self._resolve_run_timeout(timeout),
webhooks=to_client_representations(webhooks),
)

@_ensure_context
async def call_task(
self,
Expand Down
59 changes: 59 additions & 0 deletions tests/e2e/test_actor_api_helpers.py
Original file line number Diff line number Diff line change
Expand Up @@ -202,6 +202,65 @@ async def main_outer() -> None:
assert inner_output_record['value'] == f'{test_value}_XXX_{test_value}'


async def test_actor_starts_task(
make_actor: MakeActorFunction,
run_actor: RunActorFunction,
apify_client_async: ApifyClientAsync,
) -> None:
"""`Actor.start_task` returns the run before it finishes, and the started task run completes with its output."""

async def main_inner() -> None:
async with Actor:
await asyncio.sleep(5)
actor_input = await Actor.get_input() or {}
test_value = actor_input.get('test_value')
await Actor.set_value('OUTPUT', f'{test_value}_XXX_{test_value}')

async def main_outer() -> None:
async with Actor:
actor_input = await Actor.get_input() or {}
inner_task_id = actor_input.get('inner_task_id')

assert inner_task_id is not None

inner_run = await Actor.start_task(inner_task_id)

assert inner_run.actor_task_id == inner_task_id
assert inner_run.status in {'READY', 'RUNNING'}

inner_actor = await make_actor(label='start-task-inner', main_func=main_inner)
outer_actor = await make_actor(label='start-task-outer', main_func=main_outer)

inner_actor_get_result = await inner_actor.get()
assert inner_actor_get_result is not None, 'Failed to get inner actor ID'

inner_actor_id = inner_actor_get_result.id
test_value = crypto_random_object_id()

task = await apify_client_async.tasks().create(
actor_id=inner_actor_id,
name=generate_unique_resource_name('actor-start-task'),
task_input={'test_value': test_value},
)

try:
run_result_outer = await run_actor(
outer_actor,
run_input={'inner_task_id': task.id},
force_permission_level='FULL_PERMISSIONS',
)

assert run_result_outer.status == 'SUCCEEDED'

await inner_actor.last_run().wait_for_finish(wait_duration=timedelta(seconds=600))

inner_output_record = await inner_actor.last_run().key_value_store().get_record('OUTPUT')
assert inner_output_record is not None
assert inner_output_record['value'] == f'{test_value}_XXX_{test_value}'
finally:
await apify_client_async.task(task.id).delete()


async def test_actor_calls_task(
make_actor: MakeActorFunction,
run_actor: RunActorFunction,
Expand Down
46 changes: 36 additions & 10 deletions tests/unit/actor/test_actor_helpers.py
Original file line number Diff line number Diff line change
Expand Up @@ -107,6 +107,21 @@ async def test_call_actor_task(apify_client_async_patcher: ApifyClientAsyncPatch
assert apify_client_async_patcher.calls['task']['call'][0][0][0].resource_id == task_id


async def test_start_actor_task(apify_client_async_patcher: ApifyClientAsyncPatcher, fake_actor_run: Run) -> None:
"""`Actor.start_task` starts the task through the client's `task.start` and returns the run."""
apify_client_async_patcher.patch('task', 'start', return_value=fake_actor_run)
task_id = 'some-task-id'

async with Actor:
run = await Actor.start_task(task_id, {'foo': 'bar'})

assert run is fake_actor_run
calls = apify_client_async_patcher.calls['task']['start']
assert len(calls) == 1
assert calls[0][0][0].resource_id == task_id
assert calls[0][1]['task_input'] == {'foo': 'bar'}


async def test_start_actor(apify_client_async_patcher: ApifyClientAsyncPatcher, fake_actor_run: Run) -> None:
apify_client_async_patcher.patch('actor', 'start', return_value=fake_actor_run)
actor_id = 'some-id'
Expand All @@ -123,6 +138,7 @@ async def test_start_actor(apify_client_async_patcher: ApifyClientAsyncPatcher,
[
pytest.param('actor', 'start', 'start', id='start'),
pytest.param('actor', 'call', 'call', id='call'),
pytest.param('task', 'start', 'start_task', id='start_task'),
pytest.param('task', 'call', 'call_task', id='call_task'),
],
)
Expand All @@ -133,7 +149,7 @@ async def test_max_items_forwarded_to_client(
client_method: str,
sdk_method: str,
) -> None:
"""`max_items` passed to `Actor.start`, `Actor.call` or `Actor.call_task` reaches the API client."""
"""`max_items` passed to any of the run-starting helpers reaches the API client."""
apify_client_async_patcher.patch(client_resource, client_method, return_value=fake_actor_run)

async with Actor:
Expand Down Expand Up @@ -300,6 +316,7 @@ async def listener(_data: object) -> None:
_ACTOR_REMOTE_METHODS = [
pytest.param('actor', 'start', 'start', 'some-actor-id', id='start'),
pytest.param('actor', 'call', 'call', 'some-actor-id', id='call'),
pytest.param('task', 'start', 'start_task', 'some-task-id', id='start_task'),
pytest.param('task', 'call', 'call_task', 'some-task-id', id='call_task'),
]

Expand All @@ -314,7 +331,7 @@ async def test_remote_method_with_webhooks(
actor_method_name: str,
entity_id: str,
) -> None:
"""Test that start/call/call_task correctly serialize webhooks."""
"""Test that start/call/start_task/call_task correctly serialize webhooks."""
apify_client_async_patcher.patch(client_resource, client_method, return_value=fake_actor_run)

async with Actor:
Expand All @@ -341,7 +358,7 @@ async def test_remote_method_with_timedelta_timeout(
actor_method_name: str,
entity_id: str,
) -> None:
"""Test that start/call/call_task accept a timedelta timeout."""
"""Test that start/call/start_task/call_task accept a timedelta timeout."""
apify_client_async_patcher.patch(client_resource, client_method, return_value=fake_actor_run)

async with Actor:
Expand All @@ -364,7 +381,7 @@ async def test_remote_method_with_invalid_timeout(
actor_method_name: str,
entity_id: str,
) -> None:
"""Test that start/call/call_task raise ValueError for invalid timeout."""
"""Test that start/call/start_task/call_task raise ValueError for invalid timeout."""
apify_client_async_patcher.patch(client_resource, client_method, return_value=fake_actor_run)

async with Actor:
Expand Down Expand Up @@ -437,19 +454,28 @@ async def test_actor_start_and_call_skipped_when_no_inherited_time_remains(
assert apify_client_async_patcher.calls['actor'][method_name][0][1]['run_timeout'] == timedelta(seconds=1)


async def test_actor_call_task_skipped_when_no_inherited_time_remains(
@pytest.mark.parametrize(
('client_method', 'actor_method_name'),
[
pytest.param('start', 'start_task', id='start_task'),
pytest.param('call', 'call_task', id='call_task'),
],
)
async def test_actor_task_methods_skipped_when_no_inherited_time_remains(
apify_client_async_patcher: ApifyClientAsyncPatcher,
client_method: str,
actor_method_name: str,
) -> None:
"""Test that Actor.call_task with `timeout='inherit'` is skipped when the run is past its timeout."""
apify_client_async_patcher.patch('task', 'call', return_value=Mock())
"""Test that start_task/call_task with `timeout='inherit'` is skipped when the run is past its timeout."""
apify_client_async_patcher.patch('task', client_method, return_value=Mock())

async with Actor:
Actor.configuration.is_at_home = True
Actor.configuration.timeout_at = datetime.now(tz=UTC) - timedelta(minutes=5)
await Actor.call_task('some-task-id', timeout='inherit')
await getattr(Actor, actor_method_name)('some-task-id', timeout='inherit')

assert len(apify_client_async_patcher.calls['task']['call']) == 1
assert apify_client_async_patcher.calls['task']['call'][0][1]['run_timeout'] == timedelta(seconds=1)
assert len(apify_client_async_patcher.calls['task'][client_method]) == 1
assert apify_client_async_patcher.calls['task'][client_method][0][1]['run_timeout'] == timedelta(seconds=1)


async def test_reboot_runs_all_listeners_even_when_one_fails(
Expand Down
Loading