diff --git a/docs/02_concepts/07_convenience_methods.mdx b/docs/02_concepts/07_convenience_methods.mdx index 24a25593..87c913c8 100644 --- a/docs/02_concepts/07_convenience_methods.mdx +++ b/docs/02_concepts/07_convenience_methods.mdx @@ -18,6 +18,7 @@ The Apify client provides several convenience methods to handle actions that the - `ActorClient.call` - Starts an Actor and waits for it to finish, handling network timeouts internally. Waits indefinitely by default, or up to the specified `wait_duration`. - `ActorClient.start` - Starts an Actor and immediately returns the Run object without waiting for it to finish. - `RunClient.wait_for_finish` - Waits for an already-started run to reach a terminal status. +- `RunClient.iterate_dataset_items` - Yields the items of the run's default dataset as the run pushes them, and returns once the run has finished and every item is read. Additionally, storage-related resources offer flexible options for data retrieval: diff --git a/docs/02_concepts/08_pagination.mdx b/docs/02_concepts/08_pagination.mdx index 757f51c2..86951f66 100644 --- a/docs/02_concepts/08_pagination.mdx +++ b/docs/02_concepts/08_pagination.mdx @@ -83,3 +83,5 @@ The next example uses `iterate_items` on a dataset client to stream items past a + +To read the items of a run that's still going, use `RunClient.iterate_dataset_items`. Its `offset` and `limit` also count dataset rows. For details, see [Retrieve Actor data](/api/client/python/docs/guides/retrieve-actor-data#read-items-while-the-run-is-going). diff --git a/docs/03_guides/03_retrieve_actor_data.mdx b/docs/03_guides/03_retrieve_actor_data.mdx index 7490f114..e430ef03 100644 --- a/docs/03_guides/03_retrieve_actor_data.mdx +++ b/docs/03_guides/03_retrieve_actor_data.mdx @@ -8,6 +8,10 @@ import Tabs from '@theme/Tabs'; import TabItem from '@theme/TabItem'; import CodeBlock from '@theme/CodeBlock'; +import ApiLink from '@theme/ApiLink'; + +import LiveItemsAsyncExample from '!!raw-loader!./code/03_live_items_async.py'; +import LiveItemsSyncExample from '!!raw-loader!./code/03_live_items_sync.py'; import RetrieveAsyncExample from '!!raw-loader!./code/03_retrieve_async.py'; import RetrieveSyncExample from '!!raw-loader!./code/03_retrieve_sync.py'; @@ -27,3 +31,29 @@ The following example shows how to fetch datasets from an Actor's runs, paginate + +## Read items while the run is going + +To process a run's output before the run finishes, use `RunClient.iterate_dataset_items`. It yields the items of the run's default dataset shortly after the run pushes them, and it returns once the run has finished and every item is read. + +Note that: + +- Between polls, the iterator waits up to `poll_interval` (5 seconds by default) for the run to finish. Once it finishes, the iterator reads the remaining items right away. An explicit `timeout` has to leave room for that wait. +- The item options are the same as in `DatasetClient.iterate_items`, except for `desc` and `signature`. +- `offset` and `limit` count dataset rows, not the items you get back. With `clean`, `skip_empty` or `unwind`, the iterator can yield fewer or more items than `limit`. +- A run that's `ABORTING` or `TIMING-OUT` can still push items, so the iterator keeps polling until the run reaches a terminal status. + +The following example starts an Actor and prints its items as the run produces them: + + + + + {LiveItemsAsyncExample} + + + + + {LiveItemsSyncExample} + + + diff --git a/docs/03_guides/code/03_live_items_async.py b/docs/03_guides/code/03_live_items_async.py new file mode 100644 index 00000000..11030d77 --- /dev/null +++ b/docs/03_guides/code/03_live_items_async.py @@ -0,0 +1,22 @@ +import asyncio + +from apify_client import ApifyClientAsync + +TOKEN = 'MY-APIFY-TOKEN' + + +async def main() -> None: + apify_client = ApifyClientAsync(TOKEN) + + # Start the Actor without waiting for it to finish + actor_client = apify_client.actor('username/actor-name') + run = await actor_client.start(run_input={'query': 'web scraping'}) + + # Each item arrives shortly after the run pushes it. The loop ends once the run + # has finished and every item is read. + async for item in apify_client.run(run.id).iterate_dataset_items(skip_empty=True): + print(item) + + +if __name__ == '__main__': + asyncio.run(main()) diff --git a/docs/03_guides/code/03_live_items_sync.py b/docs/03_guides/code/03_live_items_sync.py new file mode 100644 index 00000000..117c71ee --- /dev/null +++ b/docs/03_guides/code/03_live_items_sync.py @@ -0,0 +1,20 @@ +from apify_client import ApifyClient + +TOKEN = 'MY-APIFY-TOKEN' + + +def main() -> None: + apify_client = ApifyClient(TOKEN) + + # Start the Actor without waiting for it to finish + actor_client = apify_client.actor('username/actor-name') + run = actor_client.start(run_input={'query': 'web scraping'}) + + # Each item arrives shortly after the run pushes it. The loop ends once the run + # has finished and every item is read. + for item in apify_client.run(run.id).iterate_dataset_items(skip_empty=True): + print(item) + + +if __name__ == '__main__': + main() diff --git a/src/apify_client/_resource_clients/_resource_client.py b/src/apify_client/_resource_clients/_resource_client.py index ecced233..bb5d85fd 100644 --- a/src/apify_client/_resource_clients/_resource_client.py +++ b/src/apify_client/_resource_clients/_resource_client.py @@ -42,6 +42,7 @@ def __init__( client_registry: Any, resource_id: str | None = None, params: dict | None = None, + api_base_url: str | None = None, ) -> None: """Initialize the resource client. @@ -53,11 +54,13 @@ def __init__( client_registry: Bundle of client classes for dependency injection. resource_id: Optional resource ID for single-resource clients. params: Optional default parameters for all requests. + api_base_url: Base URL of the API itself, for clients of top-level resources. Defaults to `base_url`. """ if resource_path.endswith('/'): raise ValueError('resource_path must not end with "/"') self._base_url = base_url + self._api_base_url = api_base_url or base_url self._public_base_url = public_base_url self._http_client = http_client self._default_params = params or {} @@ -82,11 +85,12 @@ def _resource_url(self) -> str: def _base_client_kwargs(self) -> dict[str, Any]: """Base kwargs for creating nested/child clients. - Returns dict with base_url, public_base_url, http_client, and client_registry. Caller adds + Returns dict with base_url, api_base_url, public_base_url, http_client, and client_registry. Caller adds resource_path, resource_id, and params as needed. """ return { 'base_url': self._resource_url, + 'api_base_url': self._api_base_url, 'public_base_url': self._public_base_url, 'http_client': self._http_client, 'client_registry': self._client_registry, @@ -197,6 +201,7 @@ def __init__( client_registry: ClientRegistry, resource_id: str | None = None, params: dict | None = None, + api_base_url: str | None = None, ) -> None: """Initialize the resource client. @@ -208,6 +213,7 @@ def __init__( client_registry: Bundle of client classes for dependency injection. resource_id: Optional resource ID for single-resource clients. params: Optional default parameters for all requests. + api_base_url: Base URL of the API itself, for clients of top-level resources. Defaults to `base_url`. """ super().__init__( base_url=base_url, @@ -217,6 +223,7 @@ def __init__( client_registry=client_registry, resource_id=resource_id, params=params, + api_base_url=api_base_url, ) def _get(self, *, timeout: Timeout) -> dict | None: @@ -389,6 +396,7 @@ def __init__( client_registry: ClientRegistryAsync, resource_id: str | None = None, params: dict | None = None, + api_base_url: str | None = None, ) -> None: """Initialize the resource client. @@ -400,6 +408,7 @@ def __init__( client_registry: Bundle of client classes for dependency injection. resource_id: Optional resource ID for single-resource clients. params: Optional default parameters for all requests. + api_base_url: Base URL of the API itself, for clients of top-level resources. Defaults to `base_url`. """ super().__init__( base_url=base_url, @@ -409,6 +418,7 @@ def __init__( client_registry=client_registry, resource_id=resource_id, params=params, + api_base_url=api_base_url, ) async def _get(self, *, timeout: Timeout) -> dict | None: diff --git a/src/apify_client/_resource_clients/run.py b/src/apify_client/_resource_clients/run.py index 7ddd6b30..22b34be6 100644 --- a/src/apify_client/_resource_clients/run.py +++ b/src/apify_client/_resource_clients/run.py @@ -10,7 +10,8 @@ from apify_client._docs import docs_group from apify_client._logging import create_redirect_logger from apify_client._models import Run, RunResponse -from apify_client._resource_clients._resource_client import ResourceClient, ResourceClientAsync +from apify_client._pagination import DEFAULT_CHUNK_SIZE +from apify_client._resource_clients._resource_client import _TERMINAL_STATUSES, ResourceClient, ResourceClientAsync from apify_client._status_message_watcher import StatusMessageWatcher, StatusMessageWatcherAsync from apify_client._streamed_log import StreamedLog, StreamedLogAsync from apify_client._utils.encoding import encode_key_value_store_record_value @@ -19,6 +20,7 @@ if TYPE_CHECKING: import logging + from collections.abc import AsyncIterator, Iterator from decimal import Decimal from apify_client._literals import GeneralAccess @@ -32,6 +34,7 @@ RequestQueueClient, RequestQueueClientAsync, ) + from apify_client._resource_clients.dataset import DatasetItemsPage from apify_client.types import Timeout @@ -466,6 +469,122 @@ def get_status_message_watcher( return StatusMessageWatcher(run_client=self, to_logger=to_logger, check_period=check_period) + def iterate_dataset_items( + self, + *, + offset: int | None = None, + limit: int | None = None, + clean: bool | None = None, + fields: list[str] | None = None, + omit: list[str] | None = None, + unwind: list[str] | None = None, + skip_empty: bool | None = None, + skip_hidden: bool | None = None, + chunk_size: int | None = None, + poll_interval: timedelta = timedelta(seconds=5), + timeout: Timeout = 'long', + ) -> Iterator[dict]: + """Iterate over the items of the run's default dataset while the run is still producing them. + + While the run has not finished, each poll yields the rows below the dataset's `item_count` and then waits up to + `poll_interval` for the run to finish, so the last rows are read as soon as it does. Each page is requested with + a `limit` that ends at `item_count`, so it covers exactly the rows it asks for, whatever the filters or `unwind` + do to the items. `item_count` lags a few seconds behind the pushed items, so once the run reaches a terminal + status, the rows past it are read a page at a time until none are left, and the iterator returns. On a + `last_run()` client, the iterator sticks to the run that its first request resolves to. + + https://docs.apify.com/api/v2#/reference/datasets/item-collection/get-items + + Args: + offset: Number of items that should be skipped at the start. The default value is 0. + limit: Maximum number of dataset rows to scan. Fewer items are yielded when filters drop some, more + when `unwind` splits a row into several. By default there is no limit. + clean: If True, returns only non-empty items and skips hidden fields (i.e. fields starting with + the # character). The clean parameter is just a shortcut for skip_hidden=True and skip_empty=True + parameters. + fields: A list of fields which should be picked from the items, only these fields will remain in + the resulting record objects. + omit: A list of fields which should be omitted from the items. + unwind: A list of fields which should be unwound, in order which they should be processed. Each field + should be either an array or an object. If the field is an array then every element of the array + will become a separate record and merged with parent object. If the unwound field is an object then + it is merged with the parent object. + skip_empty: If True, then empty items are skipped from the output. + skip_hidden: If True, then hidden fields are skipped from the output, i.e. fields starting with + the # character. + chunk_size: Maximum number of dataset rows requested per API call. + poll_interval: How long to wait for the run to finish between polls. + timeout: Timeout for each API HTTP request. + + Yields: + An item from the dataset. + """ + page_size = chunk_size or DEFAULT_CHUNK_SIZE + position = offset or 0 + end = position + limit if limit else None + + run = self.get(timeout=timeout) + # A `last_run()` client resolves `runs/last` per request, so a newer run would swap the dataset mid-iteration. + run_client = ( + self._client_registry.run_client( + resource_id=run.id, + base_url=self._api_base_url, + public_base_url=self._public_base_url, + http_client=self._http_client, + client_registry=self._client_registry, + ) + if run is not None and run.id != self._resource_id + else self + ) + dataset_client = run_client.dataset() + + def list_page(page_offset: int, page_limit: int) -> DatasetItemsPage: + return dataset_client.list_items( + offset=page_offset, + limit=page_limit, + clean=clean, + fields=fields, + omit=omit, + unwind=unwind, + skip_empty=skip_empty, + skip_hidden=skip_hidden, + timeout=timeout, + ) + + while True: + is_finished = run is None or run.status in _TERMINAL_STATUSES + dataset = dataset_client.get(timeout=timeout) + item_count = dataset.item_count if dataset else 0 + if end is not None: + item_count = min(item_count, end) + + while position < item_count: + page_limit = min(page_size, item_count - position) + page = list_page(position, page_limit) + yield from page.items + position += page_limit + + if end is not None and position >= end: + return + if is_finished: + break + run = run_client.wait_for_finish(wait_duration=poll_interval, timeout=timeout) + + while True: + page_limit = min(page_size, end - position) if end is not None else page_size + page = list_page(position, page_limit) + yield from page.items + # Only an empty page marks the end, as filters can shorten a full one. A page that `clean`, `skip_empty` or + # `unwind` emptied past a lagging `item_count` reports no scanned rows either, so a plain read checks. + if not page.count and ( + not (clean or skip_empty or unwind) + or not dataset_client.list_items(offset=position, limit=1, timeout=timeout).items + ): + return + position += page_limit + if end is not None and position >= end: + return + @docs_group('Resource clients') class RunClientAsync(ResourceClientAsync): @@ -897,3 +1016,121 @@ async def get_status_message_watcher( to_logger = create_redirect_logger(f'apify.{name}') return StatusMessageWatcherAsync(run_client=self, to_logger=to_logger, check_period=check_period) + + async def iterate_dataset_items( + self, + *, + offset: int | None = None, + limit: int | None = None, + clean: bool | None = None, + fields: list[str] | None = None, + omit: list[str] | None = None, + unwind: list[str] | None = None, + skip_empty: bool | None = None, + skip_hidden: bool | None = None, + chunk_size: int | None = None, + poll_interval: timedelta = timedelta(seconds=5), + timeout: Timeout = 'long', + ) -> AsyncIterator[dict]: + """Iterate over the items of the run's default dataset while the run is still producing them. + + While the run has not finished, each poll yields the rows below the dataset's `item_count` and then waits up to + `poll_interval` for the run to finish, so the last rows are read as soon as it does. Each page is requested with + a `limit` that ends at `item_count`, so it covers exactly the rows it asks for, whatever the filters or `unwind` + do to the items. `item_count` lags a few seconds behind the pushed items, so once the run reaches a terminal + status, the rows past it are read a page at a time until none are left, and the iterator returns. On a + `last_run()` client, the iterator sticks to the run that its first request resolves to. + + https://docs.apify.com/api/v2#/reference/datasets/item-collection/get-items + + Args: + offset: Number of items that should be skipped at the start. The default value is 0. + limit: Maximum number of dataset rows to scan. Fewer items are yielded when filters drop some, more + when `unwind` splits a row into several. By default there is no limit. + clean: If True, returns only non-empty items and skips hidden fields (i.e. fields starting with + the # character). The clean parameter is just a shortcut for skip_hidden=True and skip_empty=True + parameters. + fields: A list of fields which should be picked from the items, only these fields will remain in + the resulting record objects. + omit: A list of fields which should be omitted from the items. + unwind: A list of fields which should be unwound, in order which they should be processed. Each field + should be either an array or an object. If the field is an array then every element of the array + will become a separate record and merged with parent object. If the unwound field is an object then + it is merged with the parent object. + skip_empty: If True, then empty items are skipped from the output. + skip_hidden: If True, then hidden fields are skipped from the output, i.e. fields starting with + the # character. + chunk_size: Maximum number of dataset rows requested per API call. + poll_interval: How long to wait for the run to finish between polls. + timeout: Timeout for each API HTTP request. + + Yields: + An item from the dataset. + """ + page_size = chunk_size or DEFAULT_CHUNK_SIZE + position = offset or 0 + end = position + limit if limit else None + + run = await self.get(timeout=timeout) + # A `last_run()` client resolves `runs/last` per request, so a newer run would swap the dataset mid-iteration. + run_client = ( + self._client_registry.run_client( + resource_id=run.id, + base_url=self._api_base_url, + public_base_url=self._public_base_url, + http_client=self._http_client, + client_registry=self._client_registry, + ) + if run is not None and run.id != self._resource_id + else self + ) + dataset_client = run_client.dataset() + + async def list_page(page_offset: int, page_limit: int) -> DatasetItemsPage: + return await dataset_client.list_items( + offset=page_offset, + limit=page_limit, + clean=clean, + fields=fields, + omit=omit, + unwind=unwind, + skip_empty=skip_empty, + skip_hidden=skip_hidden, + timeout=timeout, + ) + + while True: + is_finished = run is None or run.status in _TERMINAL_STATUSES + dataset = await dataset_client.get(timeout=timeout) + item_count = dataset.item_count if dataset else 0 + if end is not None: + item_count = min(item_count, end) + + while position < item_count: + page_limit = min(page_size, item_count - position) + page = await list_page(position, page_limit) + for item in page.items: + yield item + position += page_limit + + if end is not None and position >= end: + return + if is_finished: + break + run = await run_client.wait_for_finish(wait_duration=poll_interval, timeout=timeout) + + while True: + page_limit = min(page_size, end - position) if end is not None else page_size + page = await list_page(position, page_limit) + for item in page.items: + yield item + # Only an empty page marks the end, as filters can shorten a full one. A page that `clean`, `skip_empty` or + # `unwind` emptied past a lagging `item_count` reports no scanned rows either, so a plain read checks. + if not page.count and ( + not (clean or skip_empty or unwind) + or not (await dataset_client.list_items(offset=position, limit=1, timeout=timeout)).items + ): + return + position += page_limit + if end is not None and position >= end: + return diff --git a/tests/integration/test_run.py b/tests/integration/test_run.py index 811c5794..1cea548d 100644 --- a/tests/integration/test_run.py +++ b/tests/integration/test_run.py @@ -6,8 +6,8 @@ from datetime import UTC, datetime, timedelta from typing import TYPE_CHECKING -from .._utils import maybe_await, poll_until_condition -from apify_client._models import Dataset, KeyValueStore, ListOfRuns, RequestQueue, Run, RunShort +from .._utils import get_random_resource_name, maybe_await, poll_until_condition +from apify_client._models import Actor, Build, Dataset, KeyValueStore, ListOfRuns, RequestQueue, Run, RunShort from apify_client.errors import ApifyApiError if TYPE_CHECKING: @@ -15,6 +15,34 @@ HELLO_WORLD_ACTOR = 'apify/hello-world' +LIVE_ITEM_COUNT = 10 + +LIVE_ITEMS_SOURCE_FILES = [ + { + 'name': 'Dockerfile', + 'format': 'TEXT', + 'content': 'FROM apify/actor-node:22\nCOPY . ./\nCMD ["node", "main.mjs"]\n', + }, + { + 'name': 'main.mjs', + 'format': 'TEXT', + 'content': f""" +const apiUrl = (process.env.APIFY_API_BASE_URL || 'https://api.apify.com').replace(/\\/$/, ''); +const itemsUrl = `${{apiUrl}}/v2/datasets/${{process.env.ACTOR_DEFAULT_DATASET_ID}}/items`; +for (let index = 0; index < {LIVE_ITEM_COUNT}; index++) {{ + const response = await fetch(itemsUrl, {{ + method: 'POST', + headers: {{ 'Content-Type': 'application/json', Authorization: `Bearer ${{process.env.APIFY_TOKEN}}` }}, + body: JSON.stringify({{ runId: process.env.ACTOR_RUN_ID, index }}), + }}); + if (!response.ok) throw new Error(`Pushing item ${{index}} failed with ${{response.status}}`); + await new Promise((resolve) => setTimeout(resolve, 1000)); +}} +""", + }, +] +"""Source of an Actor that pushes `LIVE_ITEM_COUNT` items tagged with its run ID to its dataset, one per second.""" + async def test_run_collection_list_multiple_statuses(client: ApifyClient | ApifyClientAsync) -> None: """Test listing runs with multiple statuses.""" @@ -452,3 +480,63 @@ async def test_run_collection_iterate_actor_runs(client: ApifyClient | ApifyClie assert all(r.act_id == run.act_id for r in collected) finally: await maybe_await(client.run(run.id).delete()) + + +async def test_run_iterate_dataset_items_of_last_run(client: ApifyClient | ApifyClientAsync, *, is_async: bool) -> None: + """`iterate_dataset_items` on `last_run()` reads the run it resolved to, even after a newer run starts.""" + created_actor = await maybe_await( + client.actors().create( + name=get_random_resource_name('actor'), + versions=[ + { + 'versionNumber': '0.0', + 'sourceType': 'SOURCE_FILES', + 'buildTag': 'latest', + 'sourceFiles': LIVE_ITEMS_SOURCE_FILES, + } + ], + ) + ) + assert isinstance(created_actor, Actor) + actor_client = client.actor(created_actor.id) + run_ids: list[str] = [] + + try: + started_build = await maybe_await(actor_client.build(version_number='0.0')) + assert isinstance(started_build, Build) + build = await maybe_await(client.build(started_build.id).wait_for_finish()) + assert isinstance(build, Build) + assert build.status == 'SUCCEEDED' + + async def start_run() -> str: + run = await maybe_await(actor_client.start(memory_mbytes=256, run_timeout=timedelta(seconds=120))) + assert isinstance(run, Run) + return run.id + + first_run_id = await start_run() + run_ids.append(first_run_id) + + items: list[dict] = [] + + async def collect(item: dict) -> None: + items.append(item) + if len(run_ids) == 1: + run_ids.append(await start_run()) + + # A page of one row makes every further read a fresh request, which could land on the newer run. + iterator = actor_client.last_run().iterate_dataset_items(chunk_size=1, poll_interval=timedelta(seconds=1)) + if is_async: + assert isinstance(iterator, AsyncIterator) + async for item in iterator: + await collect(item) + else: + assert isinstance(iterator, Iterator) + for item in iterator: + await collect(item) + + assert items == [{'runId': first_run_id, 'index': index} for index in range(LIVE_ITEM_COUNT)] + finally: + for run_id in run_ids: + await maybe_await(client.run(run_id).wait_for_finish()) + await maybe_await(client.run(run_id).delete()) + await maybe_await(actor_client.delete()) diff --git a/tests/unit/test_run_iterate_dataset_items.py b/tests/unit/test_run_iterate_dataset_items.py new file mode 100644 index 00000000..3e2cd900 --- /dev/null +++ b/tests/unit/test_run_iterate_dataset_items.py @@ -0,0 +1,442 @@ +from __future__ import annotations + +import json +from dataclasses import dataclass, field +from datetime import timedelta +from typing import TYPE_CHECKING, Any +from unittest.mock import AsyncMock, Mock, call + +import pytest +from werkzeug import Response + +if TYPE_CHECKING: + from pytest_httpserver import HTTPServer + from werkzeug import Request + + from apify_client import ApifyClient, ApifyClientAsync + from apify_client._literals import ActorJobStatus + from apify_client._models import Run + +pytestmark = pytest.mark.usefixtures('http_client_classes') + +RUN_ID = 'test-run-id' +RUN_PATH = f'/v2/actor-runs/{RUN_ID}' +ACTOR_ID = 'test-actor-id' +UNWIND_PARTS = 3 +NO_WAIT = timedelta(0) + + +@dataclass +class Step: + """State of the run that the fake API switches to on one run status read.""" + + pushed_rows: int + """Total number of rows in the dataset, readable right away.""" + + item_count: int + """The dataset's `itemCount`, which lags behind the pushed rows.""" + + status: ActorJobStatus + + +@dataclass +class FakeRunApi: + """Fake run, run dataset and dataset items endpoints, advanced by one `Step` per run status read.""" + + steps: list[Step] + step_index: int = -1 + run_requests: list[dict[str, str]] = field(default_factory=list) + last_run_requests: int = 0 + items_requests: list[dict[str, str]] = field(default_factory=list) + + @property + def step(self) -> Step: + return self.steps[max(self.step_index, 0)] + + def register(self, httpserver: HTTPServer) -> None: + httpserver.expect_request(RUN_PATH, method='GET').respond_with_handler(self.handle_run) + httpserver.expect_request(f'{RUN_PATH}/dataset', method='GET').respond_with_handler(self.handle_dataset) + httpserver.expect_request(f'{RUN_PATH}/dataset/items', method='GET').respond_with_handler(self.handle_items) + httpserver.expect_request(f'/v2/actors/{ACTOR_ID}/runs/last', method='GET').respond_with_handler( + self.handle_last_run + ) + + def handle_run(self, request: Request) -> Response: + self.run_requests.append(dict(request.args)) + self.step_index = min(self.step_index + 1, len(self.steps) - 1) + run = { + 'id': RUN_ID, + 'actId': ACTOR_ID, + 'userId': 'test-user-id', + 'startedAt': '2019-11-30T07:34:24.202Z', + 'status': self.step.status, + 'meta': {'origin': 'WEB'}, + 'stats': {'restartCount': 0, 'resurrectCount': 0, 'computeUnits': 0.1}, + 'options': {'build': 'latest', 'timeoutSecs': 300, 'memoryMbytes': 1024, 'diskMbytes': 2048}, + 'buildId': 'test-build-id', + 'defaultKeyValueStoreId': 'test-kvs-id', + 'defaultDatasetId': 'test-dataset-id', + 'defaultRequestQueueId': 'test-rq-id', + 'containerUrl': 'https://test.runs.apify.net', + } + return Response(json.dumps({'data': run}), status=200, mimetype='application/json') + + def handle_last_run(self, request: Request) -> Response: + self.last_run_requests += 1 + return self.handle_run(request) + + def handle_dataset(self, _request: Request) -> Response: + dataset = { + 'id': 'test-dataset-id', + 'userId': 'test-user-id', + 'createdAt': '2019-12-12T07:34:14.202Z', + 'modifiedAt': '2019-12-13T08:36:13.202Z', + 'accessedAt': '2019-12-14T08:36:13.202Z', + 'itemCount': self.step.item_count, + 'cleanItemCount': self.step.item_count, + 'consoleUrl': 'https://console.apify.com/storage/datasets/test-dataset-id', + } + return Response(json.dumps({'data': dataset}), status=200, mimetype='application/json') + + def handle_items(self, request: Request) -> Response: + """Scan exactly `limit` rows from `offset`, like the real endpoint, and derive the count from `itemCount`.""" + offset = int(request.args.get('offset', 0)) + limit = int(request.args.get('limit', 0)) or 999_999_999_999 + self.items_requests.append(dict(request.args)) + rows = range(offset, min(offset + limit, self.step.pushed_rows)) + items = shape_items( + rows, + clean=request.args.get('clean') == 'true', + skip_empty=request.args.get('skipEmpty') == 'true', + unwind=bool(request.args.get('unwind')), + ) + headers = { + 'x-apify-pagination-total': str(self.step.item_count), + 'x-apify-pagination-offset': str(offset), + 'x-apify-pagination-count': str(max(min(self.step.item_count - offset, limit), 0)), + 'x-apify-pagination-limit': str(limit), + 'x-apify-pagination-desc': 'false', + } + return Response(json.dumps(items), status=200, headers=headers, mimetype='application/json') + + +def shape_items( + rows: range, *, clean: bool = False, skip_empty: bool = False, unwind: bool = False +) -> list[dict[str, Any]]: + """Turn dataset rows into items: `clean` and `skip_empty` drop every odd row, `unwind` splits a row into items. + + Under `unwind`, every third row holds an empty array and so unwinds into no items at all. + """ + kept_rows = [row for row in rows if not ((clean or skip_empty) and row % 2)] + if unwind: + return [{'row': row, 'part': part} for row in kept_rows if row % 3 != 2 for part in range(UNWIND_PARTS)] + return [{'row': row} for row in kept_rows] + + +LAGGING_RUN_STEPS = [ + Step(pushed_rows=3, item_count=0, status='READY'), + Step(pushed_rows=40, item_count=3, status='RUNNING'), + Step(pushed_rows=52, item_count=25, status='RUNNING'), + # The run has finished, but `itemCount` still lags more than one page behind the readable rows. + Step(pushed_rows=75, item_count=52, status='SUCCEEDED'), +] + + +def test_iterate_dataset_items_yields_every_row_once_sync(httpserver: HTTPServer, sync_client: ApifyClient) -> None: + """Rows pushed across polls and past a lagging item count after the run finished are all yielded, in order.""" + api = FakeRunApi(LAGGING_RUN_STEPS) + api.register(httpserver) + + items = list(sync_client.run(RUN_ID).iterate_dataset_items(chunk_size=10, poll_interval=NO_WAIT)) + + assert items == shape_items(range(75)) + + +async def test_iterate_dataset_items_yields_every_row_once_async( + httpserver: HTTPServer, async_client: ApifyClientAsync +) -> None: + """Rows pushed across polls and past a lagging item count after the run finished are all yielded, in order.""" + api = FakeRunApi(LAGGING_RUN_STEPS) + api.register(httpserver) + + items = [ + item async for item in async_client.run(RUN_ID).iterate_dataset_items(chunk_size=10, poll_interval=NO_WAIT) + ] + + assert items == shape_items(range(75)) + + +CHUNK_SIZES = [ + pytest.param(10, id='partly filtered pages'), + pytest.param(1, id='fully filtered pages'), +] + + +@pytest.mark.parametrize('chunk_size', CHUNK_SIZES) +@pytest.mark.parametrize( + 'shaping', + [ + pytest.param({'clean': True}, id='clean drops items'), + pytest.param({'skip_empty': True}, id='skip_empty drops items'), + pytest.param({'unwind': True}, id='unwind multiplies or drops items'), + ], +) +def test_iterate_dataset_items_with_shaped_items_sync( + httpserver: HTTPServer, sync_client: ApifyClient, shaping: dict[str, bool], chunk_size: int +) -> None: + """Filters and `unwind` reshape or empty pages without duplicating or skipping rows, also past a lagging count.""" + api = FakeRunApi(LAGGING_RUN_STEPS) + api.register(httpserver) + + items = list( + sync_client.run(RUN_ID).iterate_dataset_items( + clean=shaping.get('clean'), + skip_empty=shaping.get('skip_empty'), + unwind=['parts'] if shaping.get('unwind') else None, + chunk_size=chunk_size, + poll_interval=NO_WAIT, + ) + ) + + assert items == shape_items(range(75), **shaping) + + +@pytest.mark.parametrize('chunk_size', CHUNK_SIZES) +@pytest.mark.parametrize( + 'shaping', + [ + pytest.param({'clean': True}, id='clean drops items'), + pytest.param({'skip_empty': True}, id='skip_empty drops items'), + pytest.param({'unwind': True}, id='unwind multiplies or drops items'), + ], +) +async def test_iterate_dataset_items_with_shaped_items_async( + httpserver: HTTPServer, async_client: ApifyClientAsync, shaping: dict[str, bool], chunk_size: int +) -> None: + """Filters and `unwind` reshape or empty pages without duplicating or skipping rows, also past a lagging count.""" + api = FakeRunApi(LAGGING_RUN_STEPS) + api.register(httpserver) + + items = [ + item + async for item in async_client.run(RUN_ID).iterate_dataset_items( + clean=shaping.get('clean'), + skip_empty=shaping.get('skip_empty'), + unwind=['parts'] if shaping.get('unwind') else None, + chunk_size=chunk_size, + poll_interval=NO_WAIT, + ) + ] + + assert items == shape_items(range(75), **shaping) + + +@pytest.mark.parametrize( + 'status', + [ + pytest.param('ABORTING', id='aborting'), + pytest.param('TIMING-OUT', id='timing out'), + ], +) +def test_iterate_dataset_items_keeps_polling_until_terminal_sync( + httpserver: HTTPServer, sync_client: ApifyClient, status: ActorJobStatus +) -> None: + """A run that is aborting or timing out can still push items, so polling goes on until a terminal status.""" + api = FakeRunApi( + [ + Step(pushed_rows=5, item_count=5, status=status), + Step(pushed_rows=8, item_count=8, status='ABORTED'), + ] + ) + api.register(httpserver) + + items = list(sync_client.run(RUN_ID).iterate_dataset_items(poll_interval=NO_WAIT)) + + assert items == shape_items(range(8)) + + +@pytest.mark.parametrize( + 'status', + [ + pytest.param('ABORTING', id='aborting'), + pytest.param('TIMING-OUT', id='timing out'), + ], +) +async def test_iterate_dataset_items_keeps_polling_until_terminal_async( + httpserver: HTTPServer, async_client: ApifyClientAsync, status: ActorJobStatus +) -> None: + """A run that is aborting or timing out can still push items, so polling goes on until a terminal status.""" + api = FakeRunApi( + [ + Step(pushed_rows=5, item_count=5, status=status), + Step(pushed_rows=8, item_count=8, status='ABORTED'), + ] + ) + api.register(httpserver) + + items = [item async for item in async_client.run(RUN_ID).iterate_dataset_items(poll_interval=NO_WAIT)] + + assert items == shape_items(range(8)) + + +def test_iterate_dataset_items_respects_offset_and_limit_sync(httpserver: HTTPServer, sync_client: ApifyClient) -> None: + """Iteration starts at `offset` and stops once `limit` rows are scanned, without waiting for the run to finish.""" + api = FakeRunApi( + [Step(pushed_rows=4, item_count=4, status='RUNNING'), Step(pushed_rows=20, item_count=20, status='RUNNING')] + ) + api.register(httpserver) + + items = list(sync_client.run(RUN_ID).iterate_dataset_items(offset=2, limit=7, chunk_size=3, poll_interval=NO_WAIT)) + + assert items == shape_items(range(2, 9)) + assert api.step.status == 'RUNNING' + + +async def test_iterate_dataset_items_respects_offset_and_limit_async( + httpserver: HTTPServer, async_client: ApifyClientAsync +) -> None: + """Iteration starts at `offset` and stops once `limit` rows are scanned, without waiting for the run to finish.""" + api = FakeRunApi( + [Step(pushed_rows=4, item_count=4, status='RUNNING'), Step(pushed_rows=20, item_count=20, status='RUNNING')] + ) + api.register(httpserver) + + items = [ + item + async for item in async_client.run(RUN_ID).iterate_dataset_items( + offset=2, limit=7, chunk_size=3, poll_interval=NO_WAIT + ) + ] + + assert items == shape_items(range(2, 9)) + assert api.step.status == 'RUNNING' + + +def test_iterate_dataset_items_limit_ends_reading_past_item_count_sync( + httpserver: HTTPServer, sync_client: ApifyClient +) -> None: + """A `limit` beyond the lagging item count of a finished run stops the reads past it at the limit.""" + api = FakeRunApi(LAGGING_RUN_STEPS) + api.register(httpserver) + + items = list(sync_client.run(RUN_ID).iterate_dataset_items(limit=60, chunk_size=10, poll_interval=NO_WAIT)) + + assert items == shape_items(range(60)) + + +async def test_iterate_dataset_items_limit_ends_reading_past_item_count_async( + httpserver: HTTPServer, async_client: ApifyClientAsync +) -> None: + """A `limit` beyond the lagging item count of a finished run stops the reads past it at the limit.""" + api = FakeRunApi(LAGGING_RUN_STEPS) + api.register(httpserver) + + items = [ + item + async for item in async_client.run(RUN_ID).iterate_dataset_items(limit=60, chunk_size=10, poll_interval=NO_WAIT) + ] + + assert items == shape_items(range(60)) + + +def test_iterate_dataset_items_forwards_options_without_extra_read_sync( + httpserver: HTTPServer, sync_client: ApifyClient +) -> None: + """Item options reach every page request, and without a dropping filter an empty page ends with no plain read.""" + api = FakeRunApi([Step(pushed_rows=5, item_count=3, status='SUCCEEDED')]) + api.register(httpserver) + + items = list(sync_client.run(RUN_ID).iterate_dataset_items(fields=['row'], chunk_size=10, poll_interval=NO_WAIT)) + + assert items == shape_items(range(5)) + assert [request.get('fields') for request in api.items_requests] == ['row', 'row', 'row'] + + +async def test_iterate_dataset_items_forwards_options_without_extra_read_async( + httpserver: HTTPServer, async_client: ApifyClientAsync +) -> None: + """Item options reach every page request, and without a dropping filter an empty page ends with no plain read.""" + api = FakeRunApi([Step(pushed_rows=5, item_count=3, status='SUCCEEDED')]) + api.register(httpserver) + + items = [ + item + async for item in async_client.run(RUN_ID).iterate_dataset_items( + fields=['row'], chunk_size=10, poll_interval=NO_WAIT + ) + ] + + assert items == shape_items(range(5)) + assert [request.get('fields') for request in api.items_requests] == ['row', 'row', 'row'] + + +def test_iterate_dataset_items_waits_for_finish_between_polls_sync( + httpserver: HTTPServer, sync_client: ApifyClient, monkeypatch: pytest.MonkeyPatch +) -> None: + """Each poll of an unfinished run waits up to `poll_interval` for the run to finish, and the final one does not.""" + api = FakeRunApi(LAGGING_RUN_STEPS) + api.register(httpserver) + run_client = sync_client.run(RUN_ID) + wait_with_holding = run_client.wait_for_finish + + # The fake API answers at once, so a real wait would re-read the run until the deadline and skip steps. + def wait_without_holding(**kwargs: Any) -> Run | None: + return wait_with_holding(**{**kwargs, 'wait_duration': NO_WAIT}) + + wait_for_finish = Mock(side_effect=wait_without_holding) + monkeypatch.setattr(run_client, 'wait_for_finish', wait_for_finish) + + items = list(run_client.iterate_dataset_items(poll_interval=timedelta(seconds=2))) + + assert items == shape_items(range(75)) + assert wait_for_finish.call_args_list == [call(wait_duration=timedelta(seconds=2), timeout='long')] * 3 + assert api.run_requests == [{}] + [{'waitForFinish': '0'}] * 3 + + +async def test_iterate_dataset_items_waits_for_finish_between_polls_async( + httpserver: HTTPServer, async_client: ApifyClientAsync, monkeypatch: pytest.MonkeyPatch +) -> None: + """Each poll of an unfinished run waits up to `poll_interval` for the run to finish, and the final one does not.""" + api = FakeRunApi(LAGGING_RUN_STEPS) + api.register(httpserver) + run_client = async_client.run(RUN_ID) + wait_with_holding = run_client.wait_for_finish + + # The fake API answers at once, so a real wait would re-read the run until the deadline and skip steps. + async def wait_without_holding(**kwargs: Any) -> Run | None: + return await wait_with_holding(**{**kwargs, 'wait_duration': NO_WAIT}) + + wait_for_finish = AsyncMock(side_effect=wait_without_holding) + monkeypatch.setattr(run_client, 'wait_for_finish', wait_for_finish) + + items = [item async for item in run_client.iterate_dataset_items(poll_interval=timedelta(seconds=2))] + + assert items == shape_items(range(75)) + assert wait_for_finish.call_args_list == [call(wait_duration=timedelta(seconds=2), timeout='long')] * 3 + assert api.run_requests == [{}] + [{'waitForFinish': '0'}] * 3 + + +def test_iterate_dataset_items_pins_the_last_run_sync(httpserver: HTTPServer, sync_client: ApifyClient) -> None: + """A `last_run()` client resolves `runs/last` once and reads that run's dataset to the end.""" + api = FakeRunApi(LAGGING_RUN_STEPS) + api.register(httpserver) + + run_client = sync_client.actor(ACTOR_ID).last_run() + items = list(run_client.iterate_dataset_items(chunk_size=10, poll_interval=NO_WAIT)) + + assert items == shape_items(range(75)) + assert api.last_run_requests == 1 + + +async def test_iterate_dataset_items_pins_the_last_run_async( + httpserver: HTTPServer, async_client: ApifyClientAsync +) -> None: + """A `last_run()` client resolves `runs/last` once and reads that run's dataset to the end.""" + api = FakeRunApi(LAGGING_RUN_STEPS) + api.register(httpserver) + + run_client = async_client.actor(ACTOR_ID).last_run() + items = [item async for item in run_client.iterate_dataset_items(chunk_size=10, poll_interval=NO_WAIT)] + + assert items == shape_items(range(75)) + assert api.last_run_requests == 1