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
19 changes: 19 additions & 0 deletions docs/02_concepts/06_interacting_with_other_actors.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import InteractingNamedCallExample from '!!raw-loader!roa-loader!./code/06_inter
import InteractingChildRunsExample from '!!raw-loader!roa-loader!./code/06_interacting_child_runs.py';
import InteractingAbortWithParentExample from '!!raw-loader!roa-loader!./code/06_interacting_abort_with_parent.py';
import InteractingChildRunLimitExample from '!!raw-loader!roa-loader!./code/06_interacting_child_run_limit.py';
import InteractingChildRunBudgetExample from '!!raw-loader!roa-loader!./code/06_interacting_child_run_budget.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 Down Expand Up @@ -104,6 +105,24 @@ Note that:
- The limit isn't persisted. After a migration or resurrection, call `Actor.set_child_run_limits` again before you start child runs.
- Once your Actor run receives the `ABORTING` event, a call waiting for a slot raises a `RuntimeError` without starting its run.

### Sharing the charge budget with child runs

When your Actor run is started with a maximum total charge (`max_total_charge_usd`), its named child runs share that budget. Each named `Actor.start` or `Actor.call` reserves a charge limit for its child run from the part of the budget your Actor run hasn't charged yet and hasn't reserved for other child runs. Without `max_total_charge_usd`, the child run gets all of that part. A higher value is lowered to it. Your Actor run can't charge the reserved part itself, so the whole tree of runs stays within the budget, however many child runs start at once.

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

Note that:

- Pass `max_total_charge_usd` when several child runs run at once. Otherwise the first one reserves the whole budget, and the next named start raises a `RuntimeError`.
- When a child run finishes, the SDK keeps only its charge (`usage_total_usd`) reserved and releases the rest. The platform can add to that charge for about 3 minutes after the run finishes, so the SDK fetches the run again until then.
- The charges of a failed child run stay reserved after a new run replaces it under the same name.
- A reattached child run keeps the limit it was started with. A resurrected one gets a new limit, which can include the part it reserved before.
- The reservations are stored in the registry, so they survive a migration or resurrection of your Actor run.
- Child runs started without `name` aren't tracked, so they don't reserve any part of the budget.
- A limit that the platform sets by default, which it does for pay-per-event Actors, isn't shared. Only a limit set for the run counts.

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

from apify import Actor


async def main() -> None:
async with Actor:
# The budget this Actor run was started with, shared with its named child runs.
budget = Actor.get_charging_manager().get_pricing_info().max_total_charge_usd
Actor.log.info(f'Budget of this run: {budget} USD')

# Give each of the three child runs a quarter of the budget, so they can all
# start at once and this run keeps the rest for its own charges.
per_child = Decimal(1) if budget.is_infinite() else budget / 4

actor_runs = await asyncio.gather(
*(
Actor.call(
actor_id='apify/screenshot-url',
run_input={'urls': [{'url': f'https://www.apify.com/?page={page}'}]},
name=f'screenshot-{page}',
max_total_charge_usd=per_child,
)
for page in range(3)
)
)

for actor_run in actor_runs:
cost = actor_run.usage_total_usd
Actor.log.info(f'Child run {actor_run.id} cost {cost} USD')


if __name__ == '__main__':
asyncio.run(main())
34 changes: 24 additions & 10 deletions 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 ChildRunInfo, ChildRunRegistry
from apify._child_runs import ChildRunInfo, ChildRunRegistry, StartRun
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 @@ -50,7 +50,7 @@

if TYPE_CHECKING:
import logging
from collections.abc import Awaitable, Callable, MutableMapping
from collections.abc import Callable, MutableMapping
from decimal import Decimal
from types import TracebackType
from typing import Self
Expand Down Expand Up @@ -153,7 +153,9 @@ 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._child_run_registry = ChildRunRegistry(
self.open_key_value_store, lambda: self._charging_manager_implementation
)

self._active = False
"""Whether the Actor instance is currently active (initialized and within context)."""
Expand Down Expand Up @@ -212,6 +214,7 @@ async def __aenter__(self) -> Self:
self.log.debug('Event manager initialized')

# Initialize the charging manager.
self._charging_manager_implementation.child_run_reservations = self._child_run_registry.reserved_usd
try:
await self._charging_manager_implementation.__aenter__()
except BaseException:
Expand All @@ -224,6 +227,10 @@ async def __aenter__(self) -> Self:
# Mark initialization as complete and update global state.
self._active = True

# Child runs recorded by an earlier attempt of this run keep their part of the budget reserved.
if self._charging_manager_implementation.get_max_total_charge_usd().is_finite():
await self._child_run_registry.load()

if not Actor.is_at_home():
# Make sure that the input related KVS is initialized to ensure that the input aware client is used
await self.open_key_value_store()
Expand Down Expand Up @@ -966,7 +973,11 @@ async def start(
content_type: The content type of the input.
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_total_charge_usd: A limit on the total charged amount for pay-per-event Actors.
max_total_charge_usd: A limit on the total charged amount for pay-per-event Actors. When `name` is set
and this Actor run was started with a `max_total_charge_usd` set by the user, the limit defaults to
the part of that budget not charged by this Actor run nor reserved for its other named child runs,
and a higher value is lowered to it. The limit stays reserved until the child run finishes and its
charge is known.
restart_on_error: If true, the Actor 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
Expand Down Expand Up @@ -1012,7 +1023,6 @@ async def 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_start_timeout,
Expand All @@ -1021,7 +1031,7 @@ async def start(
)

if name is None:
return await start_run()
return await start_run(max_total_charge_usd=max_total_charge_usd)

run, _ = await self._find_or_start_child_run(
name,
Expand Down Expand Up @@ -1105,7 +1115,11 @@ async def call(
content_type: The content type of the input.
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_total_charge_usd: A limit on the total charged amount for pay-per-event Actors.
max_total_charge_usd: A limit on the total charged amount for pay-per-event Actors. When `name` is set
and this Actor run was started with a `max_total_charge_usd` set by the user, the limit defaults to
the part of that budget not charged by this Actor run nor reserved for its other named child runs,
and a higher value is lowered to it. The limit stays reserved until the child run finishes and its
charge is known.
restart_on_error: If true, the Actor 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
Expand Down Expand Up @@ -1175,7 +1189,6 @@ async def 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,
Expand Down Expand Up @@ -1207,7 +1220,7 @@ async def _find_or_start_child_run(
*,
actor_id: str,
client: ApifyClientAsync,
start_run: Callable[[], Awaitable[Run]],
start_run: StartRun,
build: str | None,
max_total_charge_usd: Decimal | None,
restart_on_error: bool | None,
Expand All @@ -1220,14 +1233,15 @@ async def _find_or_start_child_run(
actor_id=actor_id,
client=client,
start_run=start_run,
resurrect_run=lambda run_client: run_client.resurrect(
resurrect_run=lambda run_client, max_total_charge_usd: 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,
),
abort_with_parent=abort_with_parent,
max_total_charge_usd=max_total_charge_usd,
)

def _remove_internal_listeners(self) -> None:
Expand Down
34 changes: 31 additions & 3 deletions src/apify/_charging.py
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@
from apify.storages import Dataset

if TYPE_CHECKING:
from collections.abc import AsyncIterator
from collections.abc import AsyncIterator, Callable
from types import TracebackType

from apify_client import ApifyClientAsync
Expand Down Expand Up @@ -343,6 +343,10 @@ def __init__(self, configuration: Configuration, client: ApifyClientAsync) -> No

self.charge_lock = ReentrantLock()

self.child_run_reservations: Callable[[], Decimal] = Decimal
"""Returns the part of `max_total_charge_usd` reserved for child runs of this Actor run."""
self._is_max_total_charge_usd_set_by_user: bool | None = None

async def __aenter__(self) -> None:
"""Initialize the charging manager - this is called by the `Actor` class and shouldn't be invoked manually."""
# Validate config
Expand Down Expand Up @@ -563,9 +567,33 @@ def calculate_max_event_charge_count_within_limit(self, event_name: str) -> int
if not price:
return None

result = (self._max_total_charge_usd - self.calculate_total_charged_amount()) / price
result = self.calculate_remaining_budget() / price
return max(0, math.floor(result)) if result.is_finite() else None

@_ensure_context
def calculate_remaining_budget(self) -> Decimal:
"""Return the part of `max_total_charge_usd` not charged by this Actor run nor reserved for its child runs."""
return self._max_total_charge_usd - self.calculate_total_charged_amount() - self.child_run_reservations()

@_ensure_context
async def is_max_total_charge_usd_set_by_user(self) -> bool:
"""Return whether `max_total_charge_usd` was set for this Actor run, not defaulted by the platform.

The platform gives pay-per-event runs a limit even when nobody set one, and marks the run options when the
limit was set. A run that does not say so is treated as having a default limit.
"""
if not self._max_total_charge_usd.is_finite():
return False
if not self._is_at_home:
return True
if self._is_max_total_charge_usd_set_by_user is None:
if self._actor_run_id is None:
raise RuntimeError('Actor run ID not configured')
run = await self._client.run(self._actor_run_id).get()
extra = (run.options.model_extra or {}) if run is not None else {}
self._is_max_total_charge_usd_set_by_user = extra.get('isMaxTotalChargeUsdSetByUser') is True
return self._is_max_total_charge_usd_set_by_user

@_ensure_context
def get_pricing_info(self) -> ActorPricingInfo:
return ActorPricingInfo(
Expand Down Expand Up @@ -603,7 +631,7 @@ def compute_push_data_limit(
if not combined_price:
return items_count

result = (self._max_total_charge_usd - self.calculate_total_charged_amount()) / combined_price
result = self.calculate_remaining_budget() / combined_price
max_count = max(0, math.floor(result)) if result.is_finite() else items_count
return min(items_count, max_count)

Expand Down
Loading
Loading