diff --git a/pyproject.toml b/pyproject.toml index 8c6d53f08..0c7de4148 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -35,7 +35,7 @@ keywords = [ "scraping", ] dependencies = [ - "apify-client[brotli]>=3.1.0,<4.0.0", + "apify-client[brotli]>=3.3.0,<4.0.0", "crawlee>=1.8.0,<2.0.0,!=1.10.1", "cachetools>=5.5.0", "cryptography>=42.0.0", diff --git a/src/apify/_actor.py b/src/apify/_actor.py index 1bb9857c1..11ebe78da 100644 --- a/src/apify/_actor.py +++ b/src/apify/_actor.py @@ -958,7 +958,8 @@ async def start( 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. - 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, 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. 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 @@ -1058,8 +1059,10 @@ async def resurrect( timeout. max_items: Maximum number of items that the resurrected pay-per-result run will return. By default, the resurrected run uses the same limit as before. The limit can only be increased. - max_total_charge_usd: Maximum cost for the resurrected pay-per-event run in USD. By default, the resurrected - run uses the same limit as before. The limit can only be increased. + max_total_charge_usd: A limit on the total charged amount of the resurrected run, 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. By default, the resurrected run uses the same limit as before. The limit can only + be increased. restart_on_error: If true, the resurrected run process will be restarted whenever it exits with a non-zero status code. By default, the resurrected run uses the same setting as before. @@ -1109,7 +1112,8 @@ async def call( 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. - 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, 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. 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 @@ -1161,6 +1165,7 @@ async def start_task( *, build: str | None = None, max_items: int | None = None, + max_total_charge_usd: Decimal | None = None, restart_on_error: bool | None = None, memory_mbytes: int | None = None, timeout: timedelta | Literal['inherit'] | None = None, @@ -1183,6 +1188,8 @@ async def start_task( 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. + 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. 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 @@ -1203,6 +1210,7 @@ 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), @@ -1217,6 +1225,7 @@ async def call_task( *, build: str | None = None, max_items: int | None = None, + max_total_charge_usd: Decimal | None = None, restart_on_error: bool | None = None, memory_mbytes: int | None = None, timeout: timedelta | Literal['inherit'] | None = None, @@ -1239,6 +1248,8 @@ async def call_task( 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. + 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. 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 @@ -1261,6 +1272,7 @@ 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), diff --git a/src/apify/storage_clients/_apify/_utils.py b/src/apify/storage_clients/_apify/_utils.py index 21363431b..e66023fa3 100644 --- a/src/apify/storage_clients/_apify/_utils.py +++ b/src/apify/storage_clients/_apify/_utils.py @@ -14,8 +14,8 @@ if TYPE_CHECKING: from collections.abc import Iterable - from apify_client._models import HeadRequest, LockedHeadRequest - from apify_client._models import Request as ClientRequest + from apify_client._models import LockedRequestQueueHeadItem, RequestQueueHeadItem + from apify_client._models import RequestResource as ClientRequest from crawlee.storage_clients.models import AddRequestsResponse from apify import Configuration @@ -51,7 +51,7 @@ def hash_api_public_base_url_and_token(configuration: Configuration) -> str: return compute_short_hash(f'{configuration.api_public_base_url}{configuration.token}'.encode()) -def to_crawlee_request(client_request: ClientRequest | HeadRequest | LockedHeadRequest) -> Request: +def to_crawlee_request(client_request: ClientRequest | RequestQueueHeadItem | LockedRequestQueueHeadItem) -> Request: """Convert an Apify API client's `Request` model to a Crawlee's `Request` model. Args: diff --git a/tests/integration/test_request_queue.py b/tests/integration/test_request_queue.py index 6702937a3..b80a61b73 100644 --- a/tests/integration/test_request_queue.py +++ b/tests/integration/test_request_queue.py @@ -9,7 +9,7 @@ import pytest -from apify_client._models import BatchAddResult, RequestDraft +from apify_client._models import BatchAddResult, UnprocessedRequest from crawlee import service_locator from crawlee.crawlers import BasicCrawler @@ -1290,7 +1290,7 @@ async def test_request_queue_deduplication_unprocessed_requests( def return_unprocessed_requests(requests: list[dict], *_: Any, **__: Any) -> BatchAddResult: """Simulate API returning unprocessed requests.""" unprocessed_requests = [ - RequestDraft.model_construct( + UnprocessedRequest.model_construct( url=request['url'], unique_key=request['uniqueKey'], method=request['method'], diff --git a/tests/unit/actor/test_actor_helpers.py b/tests/unit/actor/test_actor_helpers.py index a5b678358..a30b2fb73 100644 --- a/tests/unit/actor/test_actor_helpers.py +++ b/tests/unit/actor/test_actor_helpers.py @@ -158,6 +158,32 @@ async def test_max_items_forwarded_to_client( assert apify_client_async_patcher.calls[client_resource][client_method][0][1]['max_items'] == 42 +@pytest.mark.parametrize( + ('client_resource', 'client_method', 'sdk_method'), + [ + 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'), + ], +) +async def test_max_total_charge_usd_forwarded_to_client( + apify_client_async_patcher: ApifyClientAsyncPatcher, + fake_actor_run: Run, + client_resource: str, + client_method: str, + sdk_method: str, +) -> None: + """`max_total_charge_usd` 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: + await getattr(Actor, sdk_method)('some-id', max_total_charge_usd=Decimal('2.5')) + + kwargs = apify_client_async_patcher.calls[client_resource][client_method][0][1] + assert kwargs['max_total_charge_usd'] == Decimal('2.5') + + async def test_abort_actor_run(apify_client_async_patcher: ApifyClientAsyncPatcher, fake_actor_run: Run) -> None: apify_client_async_patcher.patch('run', 'abort', return_value=fake_actor_run) run_id = 'some-run-id' diff --git a/tests/unit/storage_clients/test_apify_request_queue_client.py b/tests/unit/storage_clients/test_apify_request_queue_client.py index 4d086cb10..34d26a5d9 100644 --- a/tests/unit/storage_clients/test_apify_request_queue_client.py +++ b/tests/unit/storage_clients/test_apify_request_queue_client.py @@ -12,16 +12,16 @@ from apify_client._models import ( AddedRequest, BatchAddResult, - HeadRequest, - LockedHeadRequest, LockedRequestQueueHead, - RequestDraft, + LockedRequestQueueHeadItem, RequestLockInfo, RequestQueueHead, + RequestQueueHeadItem, RequestQueueStats, RequestRegistration, + UnprocessedRequest, ) -from apify_client._models import Request as ClientRequest +from apify_client._models import RequestResource as ClientRequest from crawlee.storage_clients.models import AddRequestsResponse, ProcessedRequest, RequestQueueMetadata from apify import Request @@ -66,7 +66,7 @@ def _batch_result( for request in processed ], unprocessed_requests=[ - RequestDraft.model_construct( + UnprocessedRequest.model_construct( unique_key=request.unique_key, url=request.url, method=request.method, @@ -589,9 +589,9 @@ async def test_shared_is_finished_does_not_refetch_requests_confirmed_in_an_unfi api_client.get_request.assert_awaited_once_with(straggler_id) -def _locked_item(request: Request, *, lock_expires_at: datetime) -> LockedHeadRequest: +def _locked_item(request: Request, *, lock_expires_at: datetime) -> LockedRequestQueueHeadItem: """Build a single locked head entry for `request` with the given lock expiry.""" - return LockedHeadRequest( + return LockedRequestQueueHeadItem( id=unique_key_to_request_id(request.unique_key), unique_key=request.unique_key, url=request.url, @@ -602,7 +602,7 @@ def _locked_item(request: Request, *, lock_expires_at: datetime) -> LockedHeadRe def _locked_head( - items: Sequence[LockedHeadRequest], + items: Sequence[LockedRequestQueueHeadItem], *, queue_has_locked_requests: bool = True, ) -> LockedRequestQueueHead: @@ -850,9 +850,9 @@ async def test_reclaim_request_frees_in_progress() -> None: assert second.unique_key == request.unique_key -def _head_item(request: Request) -> HeadRequest: +def _head_item(request: Request) -> RequestQueueHeadItem: """Build a `list_head` item for the given request.""" - return HeadRequest( + return RequestQueueHeadItem( id=unique_key_to_request_id(request.unique_key), unique_key=request.unique_key, url=request.url, @@ -861,7 +861,7 @@ def _head_item(request: Request) -> HeadRequest: ) -def _head(*items: HeadRequest) -> RequestQueueHead: +def _head(*items: RequestQueueHeadItem) -> RequestQueueHead: """Build a `list_head` response wrapping the given items.""" return RequestQueueHead( limit=200, diff --git a/uv.lock b/uv.lock index ee2b192df..d84116f92 100644 --- a/uv.lock +++ b/uv.lock @@ -225,7 +225,7 @@ dev = [ [package.metadata] requires-dist = [ - { name = "apify-client", extras = ["brotli"], specifier = ">=3.1.0,<4.0.0" }, + { name = "apify-client", extras = ["brotli"], specifier = ">=3.3.0,<4.0.0" }, { name = "cachetools", specifier = ">=5.5.0" }, { name = "crawlee", specifier = ">=1.8.0,!=1.10.1,<2.0.0" }, { name = "cryptography", specifier = ">=42.0.0" }, @@ -265,7 +265,7 @@ dev = [ [[package]] name = "apify-client" -version = "3.2.1" +version = "3.3.0" source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "colorama" }, @@ -274,9 +274,9 @@ dependencies = [ { name = "pydantic", extra = ["email"] }, { name = "typing-extensions" }, ] -sdist = { url = "https://files.pythonhosted.org/packages/cb/37/0e7374e97e7cdb0749f8a0b1e9d84b1036b0e08b1ff636fb140463246084/apify_client-3.2.1.tar.gz", hash = "sha256:604561af6eada8ca2ed4f3b6afb62f22e5bf6e00f8f4c46f6f51636301412215", size = 144669, upload-time = "2026-09-25T12:55:14.751Z" } +sdist = { url = "https://files.pythonhosted.org/packages/f8/eb/e99684fd036fd63f9b3f95ebeafa0e706e02b7c7066ac091444683cb0471/apify_client-3.3.0.tar.gz", hash = "sha256:4adb1d2d424ceb72f7382bd51b9372a687e5cc1ea54f1bc252228e0eb9316687", size = 158619, upload-time = "2026-10-08T11:22:35.225Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/d5/05/9dff8e2767335f159fb9506926888a990b3f3855dfbef34dada70be96762/apify_client-3.2.1-py3-none-any.whl", hash = "sha256:b464064ae14e1f96b0fc253529a093bcf4ccb8ff3112647c326d11011b2454ea", size = 160238, upload-time = "2026-09-25T12:55:12.808Z" }, + { url = "https://files.pythonhosted.org/packages/a4/30/0ef4c0a26ef68c491f502dbe149478f97ec68d56bcc097d3505ee6cd660d/apify_client-3.3.0-py3-none-any.whl", hash = "sha256:57edb1382e44b54a58f0ca2f62fe884c09994fcfad94ff327bcbed490022d708", size = 174351, upload-time = "2026-10-08T11:22:33.826Z" }, ] [package.optional-dependencies]