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
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
20 changes: 16 additions & 4 deletions src/apify/_actor.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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.

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand All @@ -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
Expand All @@ -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),
Expand All @@ -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,
Expand All @@ -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
Expand All @@ -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),
Expand Down
6 changes: 3 additions & 3 deletions src/apify/storage_clients/_apify/_utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand Down
4 changes: 2 additions & 2 deletions tests/integration/test_request_queue.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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'],
Expand Down
26 changes: 26 additions & 0 deletions tests/unit/actor/test_actor_helpers.py
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down
22 changes: 11 additions & 11 deletions tests/unit/storage_clients/test_apify_request_queue_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand All @@ -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:
Expand Down Expand Up @@ -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,
Expand All @@ -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,
Expand Down
8 changes: 4 additions & 4 deletions uv.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading