diff --git a/packages/evo-objects/src/evo/objects/endpoints/api/objects_api.py b/packages/evo-objects/src/evo/objects/endpoints/api/objects_api.py index 00ac48c1..f04cfa7f 100644 --- a/packages/evo-objects/src/evo/objects/endpoints/api/objects_api.py +++ b/packages/evo-objects/src/evo/objects/endpoints/api/objects_api.py @@ -24,14 +24,19 @@ API version: 1.21.0 """ +import logging + from evo.common.connector import APIConnector from evo.common.data import EmptyResponse, RequestMethod -from evo.common.utils import get_header_metadata +from evo.common.utils import EvoAPIRetry, get_header_metadata +from evo.common.utils.retry import BackoffExponential from ..models import * # noqa: F403 __all__ = ["ObjectsApi"] +logger = logging.getLogger("objects.endpoints.api") + class ObjectsApi: """API client for the Objects endpoint. @@ -46,6 +51,12 @@ class ObjectsApi: def __init__(self, connector: APIConnector): self.connector = connector + self.api_retry = EvoAPIRetry( + logger, + max_attempts=3, + backoff_method=BackoffExponential(backoff_factor=1, max_delay=15), + statuses={429, 503}, + ) async def delete_object_by_path( self, @@ -103,15 +114,19 @@ async def delete_object_by_path( "204": EmptyResponse, } - return await self.connector.call_api( - method=RequestMethod.DELETE, - resource_path="/geoscience-object/orgs/{org_id}/workspaces/{workspace_id}/objects/path/{objects_path}", - path_params=_path_params, - header_params=_header_params, - collection_formats=_collection_formats, - response_types_map=_response_types_map, - request_timeout=request_timeout, - ) + async for handle in self.api_retry(): + with handle.suppress_errors(): + return await self.connector.call_api( + method=RequestMethod.DELETE, + resource_path="/geoscience-object/orgs/{org_id}/workspaces/{workspace_id}/objects/path/{objects_path}", + path_params=_path_params, + header_params=_header_params, + collection_formats=_collection_formats, + response_types_map=_response_types_map, + request_timeout=request_timeout, + ) + + raise RuntimeError("EvoAPIRetry neither yielded nor raised any errors") async def delete_objects_by_id( self, @@ -170,15 +185,19 @@ async def delete_objects_by_id( "204": EmptyResponse, } - return await self.connector.call_api( - method=RequestMethod.DELETE, - resource_path="/geoscience-object/orgs/{org_id}/workspaces/{workspace_id}/objects/{object_id}", - path_params=_path_params, - header_params=_header_params, - collection_formats=_collection_formats, - response_types_map=_response_types_map, - request_timeout=request_timeout, - ) + async for handle in self.api_retry(): + with handle.suppress_errors(): + return await self.connector.call_api( + method=RequestMethod.DELETE, + resource_path="/geoscience-object/orgs/{org_id}/workspaces/{workspace_id}/objects/{object_id}", + path_params=_path_params, + header_params=_header_params, + collection_formats=_collection_formats, + response_types_map=_response_types_map, + request_timeout=request_timeout, + ) + + raise RuntimeError("EvoAPIRetry neither yielded nor raised any errors") async def get_object( self, @@ -257,16 +276,20 @@ async def get_object( "304": EmptyResponse, } - return await self.connector.call_api( - method=RequestMethod.GET, - resource_path="/geoscience-object/orgs/{org_id}/workspaces/{workspace_id}/objects/path/{objects_path}", - path_params=_path_params, - query_params=_query_params, - header_params=_header_params, - collection_formats=_collection_formats, - response_types_map=_response_types_map, - request_timeout=request_timeout, - ) + async for handle in self.api_retry(): + with handle.suppress_errors(): + return await self.connector.call_api( + method=RequestMethod.GET, + resource_path="/geoscience-object/orgs/{org_id}/workspaces/{workspace_id}/objects/path/{objects_path}", + path_params=_path_params, + query_params=_query_params, + header_params=_header_params, + collection_formats=_collection_formats, + response_types_map=_response_types_map, + request_timeout=request_timeout, + ) + + raise RuntimeError("EvoAPIRetry neither yielded nor raised any errors") async def get_object_by_id( self, @@ -351,16 +374,20 @@ async def get_object_by_id( "304": EmptyResponse, } - return await self.connector.call_api( - method=RequestMethod.GET, - resource_path="/geoscience-object/orgs/{org_id}/workspaces/{workspace_id}/objects/{object_id}", - path_params=_path_params, - query_params=_query_params, - header_params=_header_params, - collection_formats=_collection_formats, - response_types_map=_response_types_map, - request_timeout=request_timeout, - ) + async for handle in self.api_retry(): + with handle.suppress_errors(): + return await self.connector.call_api( + method=RequestMethod.GET, + resource_path="/geoscience-object/orgs/{org_id}/workspaces/{workspace_id}/objects/{object_id}", + path_params=_path_params, + query_params=_query_params, + header_params=_header_params, + collection_formats=_collection_formats, + response_types_map=_response_types_map, + request_timeout=request_timeout, + ) + + raise RuntimeError("EvoAPIRetry neither yielded nor raised any errors") async def list_object_version_ids( self, @@ -645,16 +672,18 @@ async def list_objects( "200": ListObjectsResponse, # noqa: F405 } - return await self.connector.call_api( - method=RequestMethod.GET, - resource_path="/geoscience-object/orgs/{org_id}/workspaces/{workspace_id}/objects", - path_params=_path_params, - query_params=_query_params, - header_params=_header_params, - collection_formats=_collection_formats, - response_types_map=_response_types_map, - request_timeout=request_timeout, - ) + async for handle in self.api_retry(): + with handle.suppress_errors(): + return await self.connector.call_api( + method=RequestMethod.GET, + resource_path="/geoscience-object/orgs/{org_id}/workspaces/{workspace_id}/objects", + path_params=_path_params, + query_params=_query_params, + header_params=_header_params, + collection_formats=_collection_formats, + response_types_map=_response_types_map, + request_timeout=request_timeout, + ) async def list_objects_by_org( self, @@ -788,16 +817,18 @@ async def list_objects_by_org( "200": ListOrgObjectsResponse, # noqa: F405 } - return await self.connector.call_api( - method=RequestMethod.GET, - resource_path="/geoscience-object/orgs/{org_id}/objects", - path_params=_path_params, - query_params=_query_params, - header_params=_header_params, - collection_formats=_collection_formats, - response_types_map=_response_types_map, - request_timeout=request_timeout, - ) + async for handle in self.api_retry(): + with handle.suppress_errors(): + return await self.connector.call_api( + method=RequestMethod.GET, + resource_path="/geoscience-object/orgs/{org_id}/objects", + path_params=_path_params, + query_params=_query_params, + header_params=_header_params, + collection_formats=_collection_formats, + response_types_map=_response_types_map, + request_timeout=request_timeout, + ) async def post_objects( self, diff --git a/packages/evo-sdk-common/src/evo/common/utils/__init__.py b/packages/evo-sdk-common/src/evo/common/utils/__init__.py index ae63ec94..14355580 100644 --- a/packages/evo-sdk-common/src/evo/common/utils/__init__.py +++ b/packages/evo-sdk-common/src/evo/common/utils/__init__.py @@ -11,6 +11,7 @@ from .cache import Cache from .data import parse_order_by +from .evo_api_retry import EvoAPIRetry from .feedback import ( NoFeedback, PartialFeedback, @@ -30,6 +31,7 @@ "BackoffLinear", "BackoffMethod", "Cache", + "EvoAPIRetry", "NoFeedback", "PartialFeedback", "Retry", diff --git a/packages/evo-sdk-common/src/evo/common/utils/evo_api_retry.py b/packages/evo-sdk-common/src/evo/common/utils/evo_api_retry.py new file mode 100644 index 00000000..8f1e0d34 --- /dev/null +++ b/packages/evo-sdk-common/src/evo/common/utils/evo_api_retry.py @@ -0,0 +1,158 @@ +# Copyright © 2026 Bentley Systems, Incorporated +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# http://www.apache.org/licenses/LICENSE-2.0 +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +import asyncio +import contextlib +import logging +import typing as tp +from collections.abc import Set +from datetime import datetime, timezone +from email.utils import parsedate_to_datetime +from random import random + +from evo.common.exceptions import EvoAPIException, RetryError +from evo.common.utils.retry import BackoffMethod + + +class EvoAPIRetryHandler: + """Handler for a single retry attempt in EvoAPIRetry.""" + + def __init__(self, logger: logging.Logger, attempt: int): + self.logger = logger + self.attempt = attempt + self.exception: Exception | None = None + + @contextlib.contextmanager + def suppress_errors( + self, excs: type[BaseException] | tuple[type[BaseException], ...] | None = None + ) -> tp.Generator[None, tp.Any, tp.Any]: + """Suppress errors raised during a retry attempt. + + :param excs: Additional exception types to suppress + """ + try: + yield + except Exception as exc: + if isinstance(exc, EvoAPIException) or (excs is not None and isinstance(exc, excs)): + self.exception = exc + else: + raise + + def set_exception(self, exc: Exception) -> None: + """Set the exception handled during the retry attempt. + + :param exc: The exception to set. + """ + self.exception = exc + + @property + def succeeded(self) -> bool: + return self.exception is None + + @property + def failed(self) -> bool: + return self.exception is not None + + +def _parse_retry_after(logger: logging.Logger, retry_after_str: str | None) -> float | None: + if retry_after_str is None or retry_after_str == "": + return None + try: + return float(retry_after_str) + except ValueError: + try: + retry_after = parsedate_to_datetime(retry_after_str) + if retry_after.tzinfo is None: + retry_after = retry_after.replace(tzinfo=timezone.utc) + return (retry_after - datetime.now(timezone.utc)).total_seconds() + except (TypeError, ValueError, IndexError): + logger.info("Failed to parse Retry-After header: %s", repr(retry_after_str)) + return None + + +class EvoAPIRetry: + """EvoAPIException-aware retry implementation + + .. note:: Retrying requests that have different outcomes each time they are called can lead to unexpected results such + as duplicate transactions or data corruption. Although operations such as GET, PUT and DELETE are generally safe to retry, + it is the responsibility of the caller to ensure that retrying is safe. Consult the API documentation or contact + the API provider for guidance on which operations are safe to retry. + + Usage:: + retry = EvoAPIRetry(logger=logging.getLogger(__name__), max_attempts=3, backoff_method=BackoffLinear(1)) + async for handler in retry(): # mandatory to call the retry object + # do some things + ... + with handler.suppress_errors(): # mandatory to suppress EvoAPIException + # make request + ... + if handler.failed: + # do some cleanup + ... + """ + + def __init__( + self, + logger: logging.Logger, + max_attempts: int, + backoff_method: BackoffMethod, + statuses: Set[int] = frozenset({429, 503}), + ) -> None: + """Initialise a EvoAPIRetry object used when retrying after failures. + + :param logger: Logger instance for logging retry attempts. + :param max_attempts: Maximum number of times to retry. + :param backoff_method: Backoff method to apply. + :param statuses: HTTP status codes that should trigger a retry. + """ + if max_attempts < 1: + raise ValueError("max_attempts must be greater than 0") + if len(statuses) == 0: + raise ValueError("statuses must contain at least one status code") + + self._statuses = statuses + self._logger = logger + self._max_attempts = max_attempts + self._backoff_method = backoff_method + + async def _recover(self, handler: EvoAPIRetryHandler) -> None: + """Recover from a failed attempt, applying backoff and jitter before the next attempt.""" + + retry_after: float | None = None + if isinstance(handler.exception, EvoAPIException) and handler.exception.status in self._statuses: + headers = handler.exception.headers if handler.exception.headers else {} + retry_after = _parse_retry_after(self._logger, headers.get("Retry-After", "").strip()) + + delay = self._backoff_method.get_backoff_time(handler.attempt) + if retry_after is not None: + delay = max(delay, retry_after) + + delay += random() # jitter + self._logger.debug(f"Waiting {delay}s") + await asyncio.sleep(delay) + + async def __call__(self) -> tp.AsyncGenerator[EvoAPIRetryHandler, None]: + """Returns an async generator that yields EvoAPIRetryHandler objects for each retry attempt.""" + errors: list[Exception] = [] + + for attempt in range(1, self._max_attempts + 1): + handler = EvoAPIRetryHandler(self._logger, attempt) + yield handler + + if handler.succeeded: + break + else: + assert handler.exception is not None + errors.append(handler.exception) + if attempt < self._max_attempts: + await self._recover(handler) + else: + raise RetryError("Retry failed", errors) diff --git a/packages/evo-sdk-common/tests/common/utils/test_evo_api_retry.py b/packages/evo-sdk-common/tests/common/utils/test_evo_api_retry.py new file mode 100644 index 00000000..bcd480ce --- /dev/null +++ b/packages/evo-sdk-common/tests/common/utils/test_evo_api_retry.py @@ -0,0 +1,154 @@ +# Copyright © 2026 Bentley Systems, Incorporated +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# http://www.apache.org/licenses/LICENSE-2.0 +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +import logging +import unittest +from datetime import datetime, timedelta, timezone +from unittest import mock + +from evo.common.data import HTTPHeaderDict +from evo.common.exceptions import EvoAPIException, RetryError +from evo.common.utils.evo_api_retry import EvoAPIRetry, EvoAPIRetryHandler +from evo.common.utils.retry import BackoffIncremental + +logger = logging.getLogger(__name__) + + +class _TestException(Exception): ... + + +class TestEvoAPIRetryHandler(unittest.TestCase): + def test_suppress_evo_api_exception(self) -> None: + handler = EvoAPIRetryHandler(logger, 1) + exception = EvoAPIException(503, "Service unavailable", None, None) + + with handler.suppress_errors(): + raise exception + + self.assertIs(exception, handler.exception) + self.assertFalse(handler.succeeded) + self.assertTrue(handler.failed) + + def test_suppress_additional_exception(self) -> None: + handler = EvoAPIRetryHandler(logger, 1) + exception = _TestException("Expected exception") + + with handler.suppress_errors(_TestException): + raise exception + + self.assertIs(exception, handler.exception) + + def test_unexpected_exception_is_not_suppressed(self) -> None: + handler = EvoAPIRetryHandler(logger, 1) + + with self.assertRaises(_TestException): + with handler.suppress_errors(): + raise _TestException("Unexpected exception") + + self.assertTrue(handler.succeeded) + + def test_set_exception(self) -> None: + handler = EvoAPIRetryHandler(logger, 1) + exception = _TestException("Expected exception") + + handler.set_exception(exception) + + self.assertIs(exception, handler.exception) + self.assertTrue(handler.failed) + + +class TestEvoAPIRetry(unittest.IsolatedAsyncioTestCase): + def setUp(self) -> None: + self.retry = EvoAPIRetry(logger, max_attempts=5, backoff_method=BackoffIncremental(1)) + + async def test_successful_attempt(self) -> None: + with mock.patch("asyncio.sleep", spec_set=True) as mock_sleep: + async for _ in self.retry(): + pass + + mock_sleep.assert_not_called() + + @mock.patch("evo.common.utils.evo_api_retry.random", return_value=0) + @mock.patch("asyncio.sleep", spec_set=True) + async def test_max_attempts(self, mock_sleep: mock.MagicMock, mock_random: mock.MagicMock) -> None: + with self.assertRaises(RetryError): + async for handler in self.retry(): + with handler.suppress_errors(): + raise EvoAPIException(503, "Service unavailable", None, None) + + self.assertEqual(4, mock_sleep.call_count) # 5 attempts == 4 sleeps. + mock_sleep.assert_has_calls([mock.call(1), mock.call(2), mock.call(3), mock.call(4)]) + self.assertEqual(4, mock_random.call_count) + + @mock.patch("evo.common.utils.evo_api_retry.random", return_value=0) + @mock.patch("asyncio.sleep", spec_set=True) + async def test_retry_after_sets_minimum_delay( + self, mock_sleep: mock.MagicMock, mock_random: mock.MagicMock + ) -> None: + retry = EvoAPIRetry(logger, max_attempts=2, backoff_method=BackoffIncremental(1)) + + async for handler in retry(): + with handler.suppress_errors(): + if handler.attempt == 1: + raise EvoAPIException(429, "Too many requests", None, HTTPHeaderDict({"Retry-After": "3"})) + + mock_sleep.assert_called_once_with(3) + mock_random.assert_called_once() + + async def test_retry_after_date_sets_minimum_delay(self) -> None: + retry = EvoAPIRetry(logger, max_attempts=2, backoff_method=BackoffIncremental(1)) + now = datetime(2026, 8, 20, 12, 0, 0, tzinfo=timezone.utc) + retry_at = datetime(2100, 1, 1, 0, 0, 0, tzinfo=timezone(-timedelta(hours=5))) + + with ( + mock.patch("asyncio.sleep", spec_set=True) as mock_sleep, + mock.patch("evo.common.utils.evo_api_retry.datetime", wraps=datetime) as mock_datetime, + mock.patch("evo.common.utils.evo_api_retry.random", return_value=0) as mock_random, + ): + mock_datetime.now.side_effect = lambda timezone_=None: ( + now if timezone_ is None else now.astimezone(timezone_) + ) + + async for handler in retry(): + with handler.suppress_errors(): + if handler.attempt == 1: + raise EvoAPIException( + 429, + "Too many requests", + None, + HTTPHeaderDict({"Retry-After": "Fri, 01 Jan 2100 00:00:00 EST"}), + ) + + mock_sleep.assert_called_once_with((retry_at - now).total_seconds()) + mock_random.assert_called_once() + + @mock.patch("evo.common.utils.evo_api_retry.random", return_value=0) + @mock.patch("asyncio.sleep", spec_set=True) + async def test_retry_after_is_ignored_for_unconfigured_status( + self, mock_sleep: mock.MagicMock, mock_random: mock.MagicMock + ) -> None: + retry = EvoAPIRetry(logger, max_attempts=2, backoff_method=BackoffIncremental(1), statuses=frozenset({503})) + + async for handler in retry(): + with handler.suppress_errors(): + if handler.attempt == 1: + raise EvoAPIException(429, "Too many requests", None, HTTPHeaderDict({"Retry-After": "3"})) + + mock_sleep.assert_called_once_with(1) + mock_random.assert_called_once() + + def test_invalid_max_attempts(self) -> None: + with self.assertRaisesRegex(ValueError, "max_attempts must be greater than 0"): + EvoAPIRetry(logger, max_attempts=0, backoff_method=BackoffIncremental(1)) + + def test_empty_statuses(self) -> None: + with self.assertRaisesRegex(ValueError, "statuses must contain at least one status code"): + EvoAPIRetry(logger, max_attempts=1, backoff_method=BackoffIncremental(1), statuses=frozenset()) diff --git a/uv.lock b/uv.lock index 746a2aae..13c4b50e 100644 --- a/uv.lock +++ b/uv.lock @@ -844,7 +844,7 @@ wheels = [ [[package]] name = "evo-blockmodels" -version = "0.5.1" +version = "0.5.2" source = { editable = "packages/evo-blockmodels" } dependencies = [ { name = "evo-sdk-common" }, @@ -1265,7 +1265,7 @@ test = [ [[package]] name = "evo-sdk-common" -version = "0.5.26" +version = "0.5.27" source = { editable = "packages/evo-sdk-common" } dependencies = [ { name = "pure-interface" },