diff --git a/docs/02_concepts/05_retries.mdx b/docs/02_concepts/05_retries.mdx
index dfdf0917..6b290f46 100644
--- a/docs/02_concepts/05_retries.mdx
+++ b/docs/02_concepts/05_retries.mdx
@@ -12,6 +12,8 @@ import ApiLink from '@theme/ApiLink';
import RetriesAsyncExample from '!!raw-loader!./code/05_retries_async.py';
import RetriesSyncExample from '!!raw-loader!./code/05_retries_sync.py';
+import WaitForResourcesAsyncExample from '!!raw-loader!./code/05_wait_for_resources_async.py';
+import WaitForResourcesSyncExample from '!!raw-loader!./code/05_wait_for_resources_sync.py';
The Apify client automatically retries requests that fail due to:
@@ -43,3 +45,27 @@ Retries with exponential backoff help reduce the load on the server and increase
+
+## Wait for resources to start a run
+
+Starting a run fails with an HTTP 402 error when the account doesn't have enough free memory for the run, or when it already runs as many Actors as its plan allows. The error `type` is `actor-memory-limit-exceeded` or `concurrent-runs-limit-exceeded`. The client doesn't retry these errors on its own, since they clear only after other runs or builds of the account finish.
+
+To keep retrying the start until the resources free up, set the `wait_for_resources` argument of `ActorClient.start`, `ActorClient.call`, or the same methods of `TaskClient`. The client then retries the start every 10 seconds, and the argument value sets how long:
+
+- `True` retries until the run starts.
+- A `timedelta` stops retrying after that time and raises the last error.
+
+Any other error raises right away. A run that asks for more memory than the account's whole memory limit gets the same `actor-memory-limit-exceeded` error and never starts, so with `True` the client retries it forever. In `call`, the time spent retrying doesn't count toward `wait_duration`.
+
+
+
+
+ {WaitForResourcesAsyncExample}
+
+
+
+
+ {WaitForResourcesSyncExample}
+
+
+
diff --git a/docs/02_concepts/code/05_wait_for_resources_async.py b/docs/02_concepts/code/05_wait_for_resources_async.py
new file mode 100644
index 00000000..615e50df
--- /dev/null
+++ b/docs/02_concepts/code/05_wait_for_resources_async.py
@@ -0,0 +1,17 @@
+from datetime import timedelta
+
+from apify_client import ApifyClientAsync
+
+TOKEN = 'MY-APIFY-TOKEN'
+
+
+async def main() -> None:
+ apify_client = ApifyClientAsync(TOKEN)
+
+ # Retry the start until the account has the resources for the run.
+ run = await apify_client.actor('username/actor-name').call(wait_for_resources=True)
+
+ # Stop retrying after 10 minutes and raise the last error.
+ started_run = await apify_client.task('username~task-name').start(
+ wait_for_resources=timedelta(minutes=10),
+ )
diff --git a/docs/02_concepts/code/05_wait_for_resources_sync.py b/docs/02_concepts/code/05_wait_for_resources_sync.py
new file mode 100644
index 00000000..045381ad
--- /dev/null
+++ b/docs/02_concepts/code/05_wait_for_resources_sync.py
@@ -0,0 +1,17 @@
+from datetime import timedelta
+
+from apify_client import ApifyClient
+
+TOKEN = 'MY-APIFY-TOKEN'
+
+
+def main() -> None:
+ apify_client = ApifyClient(TOKEN)
+
+ # Retry the start until the account has the resources for the run.
+ run = apify_client.actor('username/actor-name').call(wait_for_resources=True)
+
+ # Stop retrying after 10 minutes and raise the last error.
+ started_run = apify_client.task('username~task-name').start(
+ wait_for_resources=timedelta(minutes=10),
+ )
diff --git a/src/apify_client/_resource_clients/actor.py b/src/apify_client/_resource_clients/actor.py
index 4ad6c6f8..a4be92f3 100644
--- a/src/apify_client/_resource_clients/actor.py
+++ b/src/apify_client/_resource_clients/actor.py
@@ -26,6 +26,11 @@
from apify_client._utils.encoding import encode_key_value_store_record_value, encode_webhooks_to_base64
from apify_client._utils.http import response_to_dict
from apify_client._utils.time import to_seconds
+from apify_client._utils.wait_for_resources import (
+ prepare_resendable_body,
+ start_waiting_for_resources,
+ start_waiting_for_resources_async,
+)
if TYPE_CHECKING:
from datetime import timedelta
@@ -229,6 +234,7 @@ def start(
force_permission_level: ActorPermissionLevel | None = None,
wait_for_finish: int | None = None,
webhooks: WebhooksList | None = None,
+ wait_for_resources: bool | timedelta = False,
timeout: Timeout = 'medium',
) -> Run:
"""Start the Actor and immediately return the Run object.
@@ -263,12 +269,21 @@ def start(
* `event_types`: List of `WebhookEventType` values which trigger the webhook.
* `request_url`: URL to which to send the webhook HTTP request.
* `payload_template`: Optional template for the request payload.
+ wait_for_resources: Retry the start while the account lacks the memory or a concurrent-run slot for the run,
+ that is while the API rejects it with an `ApifyApiError` of type `actor-memory-limit-exceeded` or
+ `concurrent-runs-limit-exceeded`. Both clear as other runs or builds finish. The start is retried
+ every 10 seconds, and any other error is raised right away. `True` retries until the run starts, a
+ `timedelta` stops retrying after that long and raises the last error. A run that requests more memory
+ than the whole memory limit of the account is rejected with `actor-memory-limit-exceeded` as well and
+ never starts, so `True` retries it forever. A streamed `run_input` that cannot be rewound, such as a
+ generator, is sent only once, so its start is not retried.
timeout: Timeout for the API HTTP request.
Returns:
The run object.
"""
run_input, content_type = encode_key_value_store_record_value(run_input, content_type=content_type)
+ run_input, wait_for_resources = prepare_resendable_body(run_input, wait_for_resources=wait_for_resources)
request_params = self._build_params(
build=build,
@@ -282,13 +297,16 @@ def start(
webhooks=encode_webhooks_to_base64(webhooks),
)
- response = self._http_client.call(
- url=self._build_url('runs'),
- method='POST',
- headers={'content-type': content_type},
- data=run_input,
- params=request_params,
- timeout=timeout,
+ response = start_waiting_for_resources(
+ lambda: self._http_client.call(
+ url=self._build_url('runs'),
+ method='POST',
+ headers={'content-type': content_type},
+ data=run_input,
+ params=request_params,
+ timeout=timeout,
+ ),
+ wait_for_resources=wait_for_resources,
)
result = response_to_dict(response)
@@ -308,6 +326,7 @@ def call(
webhooks: WebhooksList | None = None,
force_permission_level: ActorPermissionLevel | None = None,
wait_duration: timedelta | None = None,
+ wait_for_resources: bool | timedelta = False,
logger: Logger | Literal['default'] | None = 'default',
timeout: Timeout = 'no_timeout',
) -> Run | None:
@@ -341,6 +360,14 @@ def call(
a webhook set up for the Actor, you do not have to add it again here.
wait_duration: The maximum time the server waits for the run to finish. If not provided,
waits indefinitely.
+ wait_for_resources: Retry the start while the account lacks the memory or a concurrent-run slot for the run,
+ that is while the API rejects it with an `ApifyApiError` of type `actor-memory-limit-exceeded` or
+ `concurrent-runs-limit-exceeded`. Both clear as other runs or builds finish. The start is retried
+ every 10 seconds, and any other error is raised right away. `True` retries until the run starts, a
+ `timedelta` stops retrying after that long and raises the last error. A run that requests more memory
+ than the whole memory limit of the account is rejected with `actor-memory-limit-exceeded` as well and
+ never starts, so `True` retries it forever. The time spent retrying doesn't count toward
+ `wait_duration`.
logger: Logger used to redirect logs from the Actor run. Using "default" literal means that a predefined
default logger will be used. Setting `None` will disable any log propagation. Passing custom logger
will redirect logs to the provided logger. The logger is also used to capture status and status message
@@ -361,6 +388,7 @@ def call(
run_timeout=run_timeout,
webhooks=webhooks,
force_permission_level=force_permission_level,
+ wait_for_resources=wait_for_resources,
timeout=timeout,
)
run_client = self._client_registry.run_client(
@@ -740,6 +768,7 @@ async def start(
force_permission_level: ActorPermissionLevel | None = None,
wait_for_finish: int | None = None,
webhooks: WebhooksList | None = None,
+ wait_for_resources: bool | timedelta = False,
timeout: Timeout = 'medium',
) -> Run:
"""Start the Actor and immediately return the Run object.
@@ -774,12 +803,21 @@ async def start(
* `event_types`: List of `WebhookEventType` values which trigger the webhook.
* `request_url`: URL to which to send the webhook HTTP request.
* `payload_template`: Optional template for the request payload.
+ wait_for_resources: Retry the start while the account lacks the memory or a concurrent-run slot for the run,
+ that is while the API rejects it with an `ApifyApiError` of type `actor-memory-limit-exceeded` or
+ `concurrent-runs-limit-exceeded`. Both clear as other runs or builds finish. The start is retried
+ every 10 seconds, and any other error is raised right away. `True` retries until the run starts, a
+ `timedelta` stops retrying after that long and raises the last error. A run that requests more memory
+ than the whole memory limit of the account is rejected with `actor-memory-limit-exceeded` as well and
+ never starts, so `True` retries it forever. A streamed `run_input` that cannot be rewound, such as a
+ generator, is sent only once, so its start is not retried.
timeout: Timeout for the API HTTP request.
Returns:
The run object.
"""
run_input, content_type = encode_key_value_store_record_value(run_input, content_type=content_type)
+ run_input, wait_for_resources = prepare_resendable_body(run_input, wait_for_resources=wait_for_resources)
request_params = self._build_params(
build=build,
@@ -793,13 +831,16 @@ async def start(
webhooks=encode_webhooks_to_base64(webhooks),
)
- response = await self._http_client.call(
- url=self._build_url('runs'),
- method='POST',
- headers={'content-type': content_type},
- data=run_input,
- params=request_params,
- timeout=timeout,
+ response = await start_waiting_for_resources_async(
+ lambda: self._http_client.call(
+ url=self._build_url('runs'),
+ method='POST',
+ headers={'content-type': content_type},
+ data=run_input,
+ params=request_params,
+ timeout=timeout,
+ ),
+ wait_for_resources=wait_for_resources,
)
result = response_to_dict(response)
@@ -819,6 +860,7 @@ async def call(
webhooks: WebhooksList | None = None,
force_permission_level: ActorPermissionLevel | None = None,
wait_duration: timedelta | None = None,
+ wait_for_resources: bool | timedelta = False,
logger: Logger | Literal['default'] | None = 'default',
timeout: Timeout = 'no_timeout',
) -> Run | None:
@@ -852,6 +894,14 @@ async def call(
a webhook set up for the Actor, you do not have to add it again here.
wait_duration: The maximum time the server waits for the run to finish. If not provided,
waits indefinitely.
+ wait_for_resources: Retry the start while the account lacks the memory or a concurrent-run slot for the run,
+ that is while the API rejects it with an `ApifyApiError` of type `actor-memory-limit-exceeded` or
+ `concurrent-runs-limit-exceeded`. Both clear as other runs or builds finish. The start is retried
+ every 10 seconds, and any other error is raised right away. `True` retries until the run starts, a
+ `timedelta` stops retrying after that long and raises the last error. A run that requests more memory
+ than the whole memory limit of the account is rejected with `actor-memory-limit-exceeded` as well and
+ never starts, so `True` retries it forever. The time spent retrying doesn't count toward
+ `wait_duration`.
logger: Logger used to redirect logs from the Actor run. Using "default" literal means that a predefined
default logger will be used. Setting `None` will disable any log propagation. Passing custom logger
will redirect logs to the provided logger. The logger is also used to capture status and status message
@@ -872,6 +922,7 @@ async def call(
run_timeout=run_timeout,
webhooks=webhooks,
force_permission_level=force_permission_level,
+ wait_for_resources=wait_for_resources,
timeout=timeout,
)
diff --git a/src/apify_client/_resource_clients/task.py b/src/apify_client/_resource_clients/task.py
index a25fac88..90f8e57e 100644
--- a/src/apify_client/_resource_clients/task.py
+++ b/src/apify_client/_resource_clients/task.py
@@ -18,6 +18,7 @@
from apify_client._utils.encoding import encode_webhooks_to_base64
from apify_client._utils.http import response_to_dict
from apify_client._utils.time import to_seconds
+from apify_client._utils.wait_for_resources import start_waiting_for_resources, start_waiting_for_resources_async
if TYPE_CHECKING:
from datetime import timedelta
@@ -221,6 +222,7 @@ def start(
restart_on_error: bool | None = None,
wait_for_finish: int | None = None,
webhooks: WebhooksList | None = None,
+ wait_for_resources: bool | timedelta = False,
timeout: Timeout = 'medium',
) -> Run:
"""Start the task and immediately return the Run object.
@@ -248,6 +250,13 @@ def start(
* `event_types`: List of `WebhookEventType` values which trigger the webhook.
* `request_url`: URL to which to send the webhook HTTP request.
* `payload_template`: Optional template for the request payload.
+ wait_for_resources: Retry the start while the account lacks the memory or a concurrent-run slot for the run,
+ that is while the API rejects it with an `ApifyApiError` of type `actor-memory-limit-exceeded` or
+ `concurrent-runs-limit-exceeded`. Both clear as other runs or builds finish. The start is retried
+ every 10 seconds, and any other error is raised right away. `True` retries until the run starts, a
+ `timedelta` stops retrying after that long and raises the last error. A run that requests more memory
+ than the whole memory limit of the account is rejected with `actor-memory-limit-exceeded` as well and
+ never starts, so `True` retries it forever.
timeout: Timeout for the API HTTP request.
Returns:
@@ -266,13 +275,16 @@ def start(
webhooks=encode_webhooks_to_base64(webhooks),
)
- response = self._http_client.call(
- url=self._build_url('runs'),
- method='POST',
- headers={'content-type': 'application/json; charset=utf-8'},
- json=task_input.model_dump() if task_input is not None else None,
- params=request_params,
- timeout=timeout,
+ response = start_waiting_for_resources(
+ lambda: self._http_client.call(
+ url=self._build_url('runs'),
+ method='POST',
+ headers={'content-type': 'application/json; charset=utf-8'},
+ json=task_input.model_dump() if task_input is not None else None,
+ params=request_params,
+ timeout=timeout,
+ ),
+ wait_for_resources=wait_for_resources,
)
result = response_to_dict(response)
@@ -289,6 +301,7 @@ def call(
restart_on_error: bool | None = None,
webhooks: WebhooksList | None = None,
wait_duration: timedelta | None = None,
+ wait_for_resources: bool | timedelta = False,
timeout: Timeout = 'no_timeout',
) -> Run | None:
"""Start a task and wait for it to finish before returning the Run object.
@@ -314,6 +327,14 @@ def call(
the Actor or task, you do not have to add it again here.
wait_duration: The maximum time the server waits for the task run to finish. If not provided,
waits indefinitely.
+ wait_for_resources: Retry the start while the account lacks the memory or a concurrent-run slot for the run,
+ that is while the API rejects it with an `ApifyApiError` of type `actor-memory-limit-exceeded` or
+ `concurrent-runs-limit-exceeded`. Both clear as other runs or builds finish. The start is retried
+ every 10 seconds, and any other error is raised right away. `True` retries until the run starts, a
+ `timedelta` stops retrying after that long and raises the last error. A run that requests more memory
+ than the whole memory limit of the account is rejected with `actor-memory-limit-exceeded` as well and
+ never starts, so `True` retries it forever. The time spent retrying doesn't count toward
+ `wait_duration`.
timeout: Timeout for the API HTTP request.
Returns:
@@ -327,6 +348,7 @@ def call(
run_timeout=run_timeout,
restart_on_error=restart_on_error,
webhooks=webhooks,
+ wait_for_resources=wait_for_resources,
timeout=timeout,
)
@@ -602,6 +624,7 @@ async def start(
restart_on_error: bool | None = None,
wait_for_finish: int | None = None,
webhooks: WebhooksList | None = None,
+ wait_for_resources: bool | timedelta = False,
timeout: Timeout = 'medium',
) -> Run:
"""Start the task and immediately return the Run object.
@@ -629,6 +652,13 @@ async def start(
* `event_types`: List of `WebhookEventType` values which trigger the webhook.
* `request_url`: URL to which to send the webhook HTTP request.
* `payload_template`: Optional template for the request payload.
+ wait_for_resources: Retry the start while the account lacks the memory or a concurrent-run slot for the run,
+ that is while the API rejects it with an `ApifyApiError` of type `actor-memory-limit-exceeded` or
+ `concurrent-runs-limit-exceeded`. Both clear as other runs or builds finish. The start is retried
+ every 10 seconds, and any other error is raised right away. `True` retries until the run starts, a
+ `timedelta` stops retrying after that long and raises the last error. A run that requests more memory
+ than the whole memory limit of the account is rejected with `actor-memory-limit-exceeded` as well and
+ never starts, so `True` retries it forever.
timeout: Timeout for the API HTTP request.
Returns:
@@ -647,13 +677,16 @@ async def start(
webhooks=encode_webhooks_to_base64(webhooks),
)
- response = await self._http_client.call(
- url=self._build_url('runs'),
- method='POST',
- headers={'content-type': 'application/json; charset=utf-8'},
- json=task_input.model_dump() if task_input is not None else None,
- params=request_params,
- timeout=timeout,
+ response = await start_waiting_for_resources_async(
+ lambda: self._http_client.call(
+ url=self._build_url('runs'),
+ method='POST',
+ headers={'content-type': 'application/json; charset=utf-8'},
+ json=task_input.model_dump() if task_input is not None else None,
+ params=request_params,
+ timeout=timeout,
+ ),
+ wait_for_resources=wait_for_resources,
)
result = response_to_dict(response)
@@ -670,6 +703,7 @@ async def call(
restart_on_error: bool | None = None,
webhooks: WebhooksList | None = None,
wait_duration: timedelta | None = None,
+ wait_for_resources: bool | timedelta = False,
timeout: Timeout = 'no_timeout',
) -> Run | None:
"""Start a task and wait for it to finish before returning the Run object.
@@ -695,6 +729,14 @@ async def call(
the Actor or task, you do not have to add it again here.
wait_duration: The maximum time the server waits for the task run to finish. If not provided,
waits indefinitely.
+ wait_for_resources: Retry the start while the account lacks the memory or a concurrent-run slot for the run,
+ that is while the API rejects it with an `ApifyApiError` of type `actor-memory-limit-exceeded` or
+ `concurrent-runs-limit-exceeded`. Both clear as other runs or builds finish. The start is retried
+ every 10 seconds, and any other error is raised right away. `True` retries until the run starts, a
+ `timedelta` stops retrying after that long and raises the last error. A run that requests more memory
+ than the whole memory limit of the account is rejected with `actor-memory-limit-exceeded` as well and
+ never starts, so `True` retries it forever. The time spent retrying doesn't count toward
+ `wait_duration`.
timeout: Timeout for the API HTTP request.
Returns:
@@ -708,6 +750,7 @@ async def call(
run_timeout=run_timeout,
restart_on_error=restart_on_error,
webhooks=webhooks,
+ wait_for_resources=wait_for_resources,
timeout=timeout,
)
run_client = self._client_registry.run_client(
diff --git a/src/apify_client/_utils/wait_for_resources.py b/src/apify_client/_utils/wait_for_resources.py
new file mode 100644
index 00000000..7966fae4
--- /dev/null
+++ b/src/apify_client/_utils/wait_for_resources.py
@@ -0,0 +1,82 @@
+from __future__ import annotations
+
+import asyncio
+import time
+from datetime import timedelta
+from typing import TYPE_CHECKING, Any, TypeVar
+
+from apify_client._logging import logger
+from apify_client.errors import ApifyApiError
+from apify_client.http_clients._streamed_body import StreamedRequestBody
+
+if TYPE_CHECKING:
+ from collections.abc import Awaitable, Callable
+
+T = TypeVar('T')
+
+RESOURCE_LIMIT_ERROR_TYPES = frozenset({'actor-memory-limit-exceeded', 'concurrent-runs-limit-exceeded'})
+"""Error types the API rejects a run start with while the account has no free memory or concurrent-run slot for it.
+
+Both clear as other runs or builds finish.
+"""
+
+WAIT_FOR_RESOURCES_COOLDOWN = timedelta(seconds=10)
+"""Cooldown between two attempts to start a run that was rejected for lack of resources."""
+
+
+def prepare_resendable_body(data: Any, *, wait_for_resources: bool | timedelta) -> tuple[Any, bool | timedelta]:
+ """Wrap a streamed request body once, so every start attempt sends the same body.
+
+ A seekable `io.IOBase` source is rewound before each attempt. Any other streamed source is used up by the first
+ attempt, so the returned `wait_for_resources` is `False` and a start rejected for lack of resources is not retried.
+ """
+ if wait_for_resources is False or not StreamedRequestBody.is_streamable(data):
+ return data, wait_for_resources
+ body = data if isinstance(data, StreamedRequestBody) else StreamedRequestBody(data)
+ return body, wait_for_resources if body.rewindable else False
+
+
+def _next_delay(exc: ApifyApiError, deadline: float | None) -> float:
+ """Return the seconds to sleep before the next attempt, or re-raise `exc` if no attempt should follow."""
+ if exc.type not in RESOURCE_LIMIT_ERROR_TYPES:
+ raise exc
+ delay = WAIT_FOR_RESOURCES_COOLDOWN.total_seconds()
+ if deadline is not None:
+ remaining = deadline - time.monotonic()
+ if remaining <= 0:
+ raise exc
+ delay = min(delay, remaining)
+ logger.info('Not enough resources to start the run, retrying in %.3gs: %s', delay, exc.message)
+ return delay
+
+
+def start_waiting_for_resources(start: Callable[[], T], *, wait_for_resources: bool | timedelta) -> T:
+ """Make the `start` request, retrying it every `WAIT_FOR_RESOURCES_COOLDOWN` while it fails for lack of resources.
+
+ `True` retries until the request succeeds, a `timedelta` bounds the retrying, after which the last error is raised.
+ Any other error is raised right away.
+ """
+ if wait_for_resources is False:
+ return start()
+ deadline = None if wait_for_resources is True else time.monotonic() + wait_for_resources.total_seconds()
+ while True:
+ try:
+ return start()
+ except ApifyApiError as exc:
+ time.sleep(_next_delay(exc, deadline))
+
+
+async def start_waiting_for_resources_async(
+ start: Callable[[], Awaitable[T]],
+ *,
+ wait_for_resources: bool | timedelta,
+) -> T:
+ """Async variant of `start_waiting_for_resources`."""
+ if wait_for_resources is False:
+ return await start()
+ deadline = None if wait_for_resources is True else time.monotonic() + wait_for_resources.total_seconds()
+ while True:
+ try:
+ return await start()
+ except ApifyApiError as exc:
+ await asyncio.sleep(_next_delay(exc, deadline))
diff --git a/tests/unit/test_wait_for_resources.py b/tests/unit/test_wait_for_resources.py
new file mode 100644
index 00000000..d3aa0050
--- /dev/null
+++ b/tests/unit/test_wait_for_resources.py
@@ -0,0 +1,260 @@
+from __future__ import annotations
+
+import inspect
+import io
+import json
+import logging
+from datetime import timedelta
+from types import SimpleNamespace
+from typing import TYPE_CHECKING, Any
+
+import pytest
+from werkzeug import Request, Response
+
+from apify_client import ApifyClient, ApifyClientAsync
+from apify_client._logging import logger
+from apify_client._utils import wait_for_resources
+from apify_client.errors import ApifyApiError
+
+if TYPE_CHECKING:
+ from collections.abc import Callable
+
+ from pytest_httpserver import HTTPServer
+
+RUN = {
+ 'id': 'run-id',
+ 'actId': 'actor-id',
+ 'userId': 'user-id',
+ 'startedAt': '2019-11-30T07:34:24.202Z',
+ 'status': 'SUCCEEDED',
+ 'meta': {'origin': 'API'},
+ 'stats': {'restartCount': 0, 'resurrectCount': 0, 'computeUnits': 0.1},
+ 'options': {'build': 'latest', 'timeoutSecs': 300, 'memoryMbytes': 1024, 'diskMbytes': 2048},
+ 'buildId': 'build-id',
+ 'generalAccess': 'RESTRICTED',
+ 'defaultKeyValueStoreId': 'kvs-id',
+ 'defaultDatasetId': 'dataset-id',
+ 'defaultRequestQueueId': 'rq-id',
+ 'containerUrl': 'https://run.apify.net',
+}
+
+
+ERROR_MESSAGES = {
+ 'actor-memory-limit-exceeded': (
+ 'By launching this job you will exceed the memory limit of 8192MB for all your Actor runs and builds '
+ '(currently used: 4096MB, requested: 8192MB). Please consider upgrading or purchasing extra memory as an '
+ 'add-on at https://console.apify.com/billing/subscription to increase your Actor memory limit.'
+ ),
+ 'concurrent-runs-limit-exceeded': (
+ 'By launching this job you will exceed your limit of 25 concurrent Actor runs. Please consider upgrading or '
+ 'purchasing an increase to concurrent Actor runs as an add-on at '
+ 'https://console.apify.com/billing/subscription to increase your limit.'
+ ),
+ 'invalid-input': 'Input is not valid.',
+}
+"""Error messages the API rejects a run start with, by error type."""
+
+
+def start_actor(client: ApifyClient | ApifyClientAsync, **kwargs: Any) -> Any:
+ return client.actor('actor-id').start(**kwargs)
+
+
+STARTERS = [
+ pytest.param(start_actor, '/v2/actors/actor-id/runs', id='actor start'),
+ pytest.param(
+ lambda c, **kw: c.actor('actor-id').call(logger=None, **kw), '/v2/actors/actor-id/runs', id='actor call'
+ ),
+ pytest.param(lambda c, **kw: c.task('task-id').start(**kw), '/v2/actor-tasks/task-id/runs', id='task start'),
+ pytest.param(lambda c, **kw: c.task('task-id').call(**kw), '/v2/actor-tasks/task-id/runs', id='task call'),
+]
+
+
+class StartServer:
+ """Rejects the start requests with the queued error types, one per request, then starts the run."""
+
+ def __init__(self, httpserver: HTTPServer, start_path: str) -> None:
+ self.rejections: list[str] = []
+ self.bodies: list[bytes] = []
+ httpserver.expect_request(start_path, method='POST').respond_with_handler(self.handle_start)
+ httpserver.expect_request('/v2/actor-runs/run-id').respond_with_json({'data': RUN})
+
+ def handle_start(self, request: Request) -> Response:
+ self.bodies.append(request.get_data())
+ if self.rejections:
+ error_type = self.rejections.pop(0)
+ body = {'error': {'type': error_type, 'message': ERROR_MESSAGES[error_type]}}
+ return Response(json.dumps(body), status=400 if error_type == 'invalid-input' else 402)
+ return Response(json.dumps({'data': RUN}), status=201, mimetype='application/json')
+
+
+@pytest.fixture
+def sleeps(monkeypatch: pytest.MonkeyPatch) -> list[float]:
+ """Make the cooldown sleeps return at once and move a fake clock forward instead, recording each delay."""
+ now = 0.0
+ recorded: list[float] = []
+
+ def sleep(seconds: float) -> None:
+ nonlocal now
+ recorded.append(seconds)
+ now += seconds
+
+ async def async_sleep(seconds: float) -> None:
+ sleep(seconds)
+
+ monkeypatch.setattr(wait_for_resources, 'time', SimpleNamespace(monotonic=lambda: now, sleep=sleep))
+ monkeypatch.setattr(wait_for_resources, 'asyncio', SimpleNamespace(sleep=async_sleep))
+ return recorded
+
+
+@pytest.fixture(params=[pytest.param(ApifyClient, id='sync'), pytest.param(ApifyClientAsync, id='async')])
+def client(request: pytest.FixtureRequest, httpserver: HTTPServer) -> ApifyClient | ApifyClientAsync:
+ return request.param(token='test', api_url=httpserver.url_for('/').removesuffix('/'))
+
+
+async def run(starter: Callable[..., Any], client: ApifyClient | ApifyClientAsync, **kwargs: Any) -> Any:
+ result = starter(client, **kwargs)
+ return await result if inspect.isawaitable(result) else result
+
+
+@pytest.mark.parametrize(('starter', 'start_path'), STARTERS)
+@pytest.mark.parametrize(
+ 'error_type',
+ [
+ pytest.param('actor-memory-limit-exceeded', id='memory limit'),
+ pytest.param('concurrent-runs-limit-exceeded', id='concurrent runs limit'),
+ ],
+)
+async def test_retries_start_every_10_seconds_until_it_succeeds(
+ *,
+ httpserver: HTTPServer,
+ client: ApifyClient | ApifyClientAsync,
+ sleeps: list[float],
+ starter: Callable[..., Any],
+ start_path: str,
+ error_type: str,
+) -> None:
+ """A start rejected for lack of resources is retried every 10 seconds until the run starts."""
+ server = StartServer(httpserver, start_path)
+ server.rejections = [error_type, error_type]
+
+ started_run = await run(starter, client, wait_for_resources=True)
+
+ assert started_run is not None
+ assert started_run.id == 'run-id'
+ assert len(server.bodies) == 3
+ assert sleeps == [10, 10]
+
+
+@pytest.mark.parametrize(('starter', 'start_path'), STARTERS)
+async def test_raises_first_rejection_without_the_option(
+ httpserver: HTTPServer,
+ client: ApifyClient | ApifyClientAsync,
+ sleeps: list[float],
+ starter: Callable[..., Any],
+ start_path: str,
+) -> None:
+ """Without `wait_for_resources`, a start rejected for lack of resources raises right away."""
+ server = StartServer(httpserver, start_path)
+ server.rejections = ['actor-memory-limit-exceeded']
+
+ with pytest.raises(ApifyApiError) as exc_info:
+ await run(starter, client)
+
+ assert exc_info.value.type == 'actor-memory-limit-exceeded'
+ assert len(server.bodies) == 1
+ assert sleeps == []
+
+
+async def test_raises_any_other_error_right_away(
+ httpserver: HTTPServer, client: ApifyClient | ApifyClientAsync, sleeps: list[float]
+) -> None:
+ """An error other than a lack of resources raises without a retry."""
+ server = StartServer(httpserver, '/v2/actors/actor-id/runs')
+ server.rejections = ['invalid-input']
+
+ with pytest.raises(ApifyApiError) as exc_info:
+ await run(start_actor, client, wait_for_resources=True)
+
+ assert exc_info.value.type == 'invalid-input'
+ assert len(server.bodies) == 1
+ assert sleeps == []
+
+
+async def test_timedelta_bounds_the_retrying(
+ *,
+ httpserver: HTTPServer,
+ client: ApifyClient | ApifyClientAsync,
+ sleeps: list[float],
+ caplog: pytest.LogCaptureFixture,
+ monkeypatch: pytest.MonkeyPatch,
+) -> None:
+ """A `timedelta` bounds the retrying, after which the last rejection is raised."""
+ server = StartServer(httpserver, '/v2/actors/actor-id/runs')
+ server.rejections = ['actor-memory-limit-exceeded'] * 10
+ monkeypatch.setattr(logger, 'propagate', True)
+
+ with caplog.at_level(logging.INFO, logger=logger.name), pytest.raises(ApifyApiError) as exc_info:
+ await run(start_actor, client, wait_for_resources=timedelta(seconds=25))
+
+ assert exc_info.value.type == 'actor-memory-limit-exceeded'
+ # Attempts at 0, 10, 20 and 25 seconds, the last cooldown cut short by the bound.
+ assert len(server.bodies) == 4
+ assert sleeps == [10, 10, 5]
+ assert [record.getMessage() for record in caplog.records] == [
+ f'Not enough resources to start the run, retrying in 10s: {ERROR_MESSAGES["actor-memory-limit-exceeded"]}',
+ f'Not enough resources to start the run, retrying in 10s: {ERROR_MESSAGES["actor-memory-limit-exceeded"]}',
+ f'Not enough resources to start the run, retrying in 5s: {ERROR_MESSAGES["actor-memory-limit-exceeded"]}',
+ ]
+
+
+async def test_zero_timedelta_makes_a_single_attempt(
+ httpserver: HTTPServer, client: ApifyClient | ApifyClientAsync, sleeps: list[float]
+) -> None:
+ """A zero `timedelta` raises the first rejection without a retry."""
+ server = StartServer(httpserver, '/v2/actors/actor-id/runs')
+ server.rejections = ['actor-memory-limit-exceeded']
+
+ with pytest.raises(ApifyApiError):
+ await run(start_actor, client, wait_for_resources=timedelta(0))
+
+ assert len(server.bodies) == 1
+ assert sleeps == []
+
+
+async def test_retry_resends_a_file_like_input(
+ httpserver: HTTPServer, client: ApifyClient | ApifyClientAsync, sleeps: list[float]
+) -> None:
+ """A file-like run input is sent whole again on every retry of the start."""
+ server = StartServer(httpserver, '/v2/actors/actor-id/runs')
+ server.rejections = ['actor-memory-limit-exceeded']
+
+ await run(
+ start_actor,
+ client,
+ run_input=io.BytesIO(b'{"a": 1}'),
+ content_type='application/json',
+ wait_for_resources=True,
+ )
+
+ assert server.bodies == [b'{"a": 1}', b'{"a": 1}']
+ assert sleeps == [10]
+
+
+async def test_unrewindable_streamed_input_is_not_retried(
+ httpserver: HTTPServer, client: ApifyClient | ApifyClientAsync, sleeps: list[float]
+) -> None:
+ """A streamed run input that cannot be rewound is sent once, and its rejection is raised without a retry."""
+ server = StartServer(httpserver, '/v2/actors/actor-id/runs')
+ server.rejections = ['actor-memory-limit-exceeded']
+
+ with pytest.raises(ApifyApiError):
+ await run(
+ start_actor,
+ client,
+ run_input=iter([b'{"a": ', b'1}']),
+ content_type='application/json',
+ wait_for_resources=True,
+ )
+
+ assert server.bodies == [b'{"a": 1}']
+ assert sleeps == []