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
46 changes: 31 additions & 15 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 ChildRunRegistry
from apify._child_runs import 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 @@ -163,7 +163,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 @@ -223,6 +225,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 Down Expand Up @@ -1027,7 +1030,11 @@ async def start(
max_items: Deprecated, use `max_total_charge_usd` instead. Will be removed in version 5.0.0. Works only with
legacy pay-per-result Actors.
max_total_charge_usd: A limit on the total charged amount, in USD. Once the run exceeds it, the platform
aborts the run, which takes a few seconds, so the final charge can slightly overshoot the limit.
aborts the run, which takes a few seconds, so the final charge can slightly overshoot the limit. 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 @@ -1068,7 +1075,6 @@ async def start(
content_type=content_type,
build=build,
max_items=max_items,
max_total_charge_usd=max_total_charge_usd,
restart_on_error=restart_on_error,
memory_mbytes=memory_mbytes,
run_timeout=self._resolve_run_timeout(timeout),
Expand All @@ -1077,7 +1083,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 @@ -1222,7 +1228,11 @@ async def call(
max_items: Deprecated, use `max_total_charge_usd` instead. Will be removed in version 5.0.0. Works only with
legacy pay-per-result Actors.
max_total_charge_usd: A limit on the total charged amount, in USD. Once the run exceeds it, the platform
aborts the run, which takes a few seconds, so the final charge can slightly overshoot the limit.
aborts the run, which takes a few seconds, so the final charge can slightly overshoot the limit. 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 @@ -1289,7 +1299,6 @@ async def call(
content_type=content_type,
build=build,
max_items=max_items,
max_total_charge_usd=max_total_charge_usd,
restart_on_error=restart_on_error,
memory_mbytes=memory_mbytes,
run_timeout=self._resolve_run_timeout(timeout),
Expand Down Expand Up @@ -1322,7 +1331,7 @@ async def _find_or_start_child_run(
task_id: str | None = None,
run_input: Any,
client: ApifyClientAsync,
start_run: Callable[[], Awaitable[Run]],
start_run: StartRun,
build: str | None,
max_items: int | None,
max_total_charge_usd: Decimal | None,
Expand All @@ -1338,7 +1347,7 @@ async def _find_or_start_child_run(
run_input=run_input,
client=client,
start_run=start_run,
resurrect_run=lambda run_id: client.run(run_id).resurrect(
resurrect_run=lambda run_id, *, max_total_charge_usd: client.run(run_id).resurrect(
build=build,
memory_mbytes=memory_mbytes,
run_timeout=self._resolve_run_timeout(timeout),
Expand All @@ -1347,6 +1356,7 @@ async def _find_or_start_child_run(
restart_on_error=restart_on_error,
),
abort_with_parent=abort_with_parent,
max_total_charge_usd=max_total_charge_usd,
)

def _remove_internal_listeners(self) -> None:
Expand Down Expand Up @@ -1417,7 +1427,11 @@ async def start_task(
max_items: Deprecated, use `max_total_charge_usd` instead. Will be removed in version 5.0.0. Works only with
legacy pay-per-result Actors.
max_total_charge_usd: A limit on the total charged amount, in USD. Once the run exceeds it, the platform
aborts the run, which takes a few seconds, so the final charge can slightly overshoot the limit.
aborts the run, which takes a few seconds, so the final charge can slightly overshoot the limit. 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 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
Expand Down Expand Up @@ -1454,15 +1468,14 @@ async def start_task(
task_input=task_input,
build=build,
max_items=max_items,
max_total_charge_usd=max_total_charge_usd,
restart_on_error=restart_on_error,
memory_mbytes=memory_mbytes,
run_timeout=self._resolve_run_timeout(timeout),
webhooks=to_client_representations(webhooks),
)

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 @@ -1514,7 +1527,11 @@ async def call_task(
max_items: Deprecated, use `max_total_charge_usd` instead. Will be removed in version 5.0.0. Works only with
legacy pay-per-result Actors.
max_total_charge_usd: A limit on the total charged amount, in USD. Once the run exceeds it, the platform
aborts the run, which takes a few seconds, so the final charge can slightly overshoot the limit.
aborts the run, which takes a few seconds, so the final charge can slightly overshoot the limit. 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 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
Expand Down Expand Up @@ -1572,7 +1589,6 @@ async def call_task(
task_input=task_input,
build=build,
max_items=max_items,
max_total_charge_usd=max_total_charge_usd,
restart_on_error=restart_on_error,
memory_mbytes=memory_mbytes,
run_timeout=self._resolve_run_timeout(timeout),
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