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