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
60 changes: 42 additions & 18 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 `run_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 run_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(
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 `run_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 @@ -1208,7 +1221,7 @@ async def _find_or_start_child_run(
actor_id: str | None = None,
task_id: str | None = None,
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 @@ -1222,14 +1235,15 @@ async def _find_or_start_child_run(
task_id=task_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 Expand Up @@ -1335,7 +1349,9 @@ async def call_task(
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 `task_id` exactly as passed, so reusing it with
any other value, or for an Actor, raises a `ValueError`.
any other value, or for an Actor, raises a `ValueError`. When this Actor run was started with
a `max_total_charge_usd` set by the user, a named call that would start a new run raises
a `RuntimeError`, since the task's run cannot be given a part of that budget.
abort_with_parent: If true, the child run is gracefully aborted when this Actor run is gracefully
aborted. It requires `run_name`, and the value is recorded under it, replacing the one from an earlier
call. A hard abort, a timeout or a crash of this Actor run leaves the child running.
Expand Down Expand Up @@ -1370,19 +1386,27 @@ async def call_task(
wait_duration=wait,
)
else:
started_run, _ = await self._find_or_start_child_run(
run_name,
task_id=task_id,
client=client,
start_run=partial(
task_client.start,

async def start_task_run(*, max_total_charge_usd: Decimal | None) -> Run:
if max_total_charge_usd is not None:
raise RuntimeError(
f'Child run "{run_name}" was not started, since a task run cannot be given a part of the '
'budget of this Actor run.'
)
return await task_client.start(
task_input=task_input,
build=build,
restart_on_error=restart_on_error,
memory_mbytes=memory_mbytes,
run_timeout=task_call_timeout,
webhooks=to_client_representations(webhooks),
),
)

started_run, _ = await self._find_or_start_child_run(
run_name,
task_id=task_id,
client=client,
start_run=start_task_run,
build=build,
max_total_charge_usd=None,
restart_on_error=restart_on_error,
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