From a472d615de0626ca67e794353ac0a3f6e5a638d6 Mon Sep 17 00:00:00 2001 From: Victor Skvortsov Date: Mon, 17 Aug 2026 15:29:44 +0500 Subject: [PATCH 1/5] Handle try_collect_stats() errors --- .../server/background/scheduled_tasks/gateways.py | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/src/dstack/_internal/server/background/scheduled_tasks/gateways.py b/src/dstack/_internal/server/background/scheduled_tasks/gateways.py index 73d153730..2d7dbb93f 100644 --- a/src/dstack/_internal/server/background/scheduled_tasks/gateways.py +++ b/src/dstack/_internal/server/background/scheduled_tasks/gateways.py @@ -1,5 +1,6 @@ import asyncio +import httpx from sqlalchemy import select from dstack._internal.core.errors import SSHError @@ -60,4 +61,9 @@ async def _process_connection(conn: GatewayConnection): logger.warning("Connection to gateway %s failed: %s", conn.ip_address, e) return - await conn.try_collect_stats() + try: + await conn.try_collect_stats() + except httpx.HTTPError as e: + logger.warning("Failed to collect stats from gateway %s: %r", conn.ip_address, e) + except Exception: + logger.exception("Failed to collect stats from gateway %s", conn.ip_address) From 943cad2c5fe2a6467b9eb3ee695e17c1126030f4 Mon Sep 17 00:00:00 2001 From: Victor Skvortsov Date: Mon, 17 Aug 2026 15:34:17 +0500 Subject: [PATCH 2/5] Turn remove_dangling_tasks_from_instance() error into warining --- .../server/background/pipeline_tasks/instances/check.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/dstack/_internal/server/background/pipeline_tasks/instances/check.py b/src/dstack/_internal/server/background/pipeline_tasks/instances/check.py index 9be54939c..793d0577d 100644 --- a/src/dstack/_internal/server/background/pipeline_tasks/instances/check.py +++ b/src/dstack/_internal/server/background/pipeline_tasks/instances/check.py @@ -432,7 +432,7 @@ def _check_instance_inner( try: remove_dangling_tasks_from_instance(shim_client, instance) except Exception as exc: - logger.exception("%s: error removing dangling tasks: %s", fmt(instance), exc) + logger.warning("%s: error removing dangling tasks: %s", fmt(instance), exc) # There should be no shim API calls after this function call since it can request shim restart. _maybe_install_components(instance, shim_client) From 7c115fa1bdb6123e4284883b02b9ed37e8f76362 Mon Sep 17 00:00:00 2001 From: Victor Skvortsov Date: Mon, 17 Aug 2026 15:50:07 +0500 Subject: [PATCH 3/5] Add tolerable _release_advisory_lock() --- .../_internal/server/services/locking.py | 29 +++++- .../_internal/server/services/test_locking.py | 91 +++++++++++++++++++ 2 files changed, 115 insertions(+), 5 deletions(-) create mode 100644 src/tests/_internal/server/services/test_locking.py diff --git a/src/dstack/_internal/server/services/locking.py b/src/dstack/_internal/server/services/locking.py index 2a2b833c0..405681dc0 100644 --- a/src/dstack/_internal/server/services/locking.py +++ b/src/dstack/_internal/server/services/locking.py @@ -9,6 +9,10 @@ from sqlalchemy import func, select from sqlalchemy.ext.asyncio import AsyncConnection, AsyncSession +from dstack._internal.utils.logging import get_logger + +logger = get_logger(__name__) + KeyT = TypeVar("KeyT") @@ -134,12 +138,12 @@ async def advisory_lock_ctx( To prevent unreleased locks: - 1. When possible, prefer using `pg_advisory_xact_lock` instead of this context manager. + * When possible, prefer using `pg_advisory_xact_lock` instead of this context manager. `pg_advisory_xact_lock` is automatically released at the end of transaction. - 1. Prefer using `AsyncConnection` as `bind`. + * Prefer using `AsyncConnection` as `bind`. - 1. If using `AsyncSession` as `bind`, **do not** commit before exiting from the context manager. + * If using `AsyncSession` as `bind`, **do not** commit before exiting from the context manager. Committing will prompt `AsyncSession` to start a new transaction for releasing the lock, which may be assigned to a different database connection, which will fail to release. """ @@ -150,7 +154,7 @@ async def advisory_lock_ctx( yield finally: if dialect_name == "postgresql": - await bind.execute(select(func.pg_advisory_unlock(string_to_lock_id(resource)))) + await _release_advisory_lock(bind, resource) @asynccontextmanager @@ -165,7 +169,7 @@ async def try_advisory_lock_ctx( yield locked finally: if dialect_name == "postgresql" and locked: - await bind.execute(select(func.pg_advisory_unlock(string_to_lock_id(resource)))) + await _release_advisory_lock(bind, resource) _in_memory_locker = InMemoryResourceLocker() @@ -203,3 +207,18 @@ async def _wait_to_lock_many( if not left_to_lock: return await asyncio.sleep(delay) + + +async def _release_advisory_lock( + bind: Union[AsyncConnection, AsyncSession], resource: str +) -> None: + """ + Release an advisory lock, tolerating failures. + Releasing typically fails with `PendingRollbackError` because the connection has been + invalidated. In this case Postgres has already dropped the lock along with the + session and there is nothing left to release. + """ + try: + await bind.execute(select(func.pg_advisory_unlock(string_to_lock_id(resource)))) + except Exception as e: + logger.warning("Failed to release advisory lock on %s: %r", resource, e) diff --git a/src/tests/_internal/server/services/test_locking.py b/src/tests/_internal/server/services/test_locking.py new file mode 100644 index 000000000..3fe0c7163 --- /dev/null +++ b/src/tests/_internal/server/services/test_locking.py @@ -0,0 +1,91 @@ +from typing import Any, cast + +import pytest +from sqlalchemy.ext.asyncio import AsyncSession + +from dstack._internal.server.services.locking import advisory_lock_ctx, try_advisory_lock_ctx + + +class _Result: + def __init__(self, value: bool) -> None: + self._value = value + + def scalar_one(self) -> bool: + return self._value + + +class _ReleaseFailsBind: + """Acquires the lock successfully, then fails to release it.""" + + def __init__(self, locked: bool = True) -> None: + self._locked = locked + self.executes = 0 + + async def execute(self, statement: Any) -> _Result: + self.executes += 1 + if self.executes > 1: + raise RuntimeError("connection invalidated") + return _Result(self._locked) + + +def _bind(locked: bool = True) -> tuple[_ReleaseFailsBind, AsyncSession]: + bind = _ReleaseFailsBind(locked) + return bind, cast(AsyncSession, bind) + + +@pytest.mark.asyncio +class TestAdvisoryLockCtx: + async def test_tolerates_release_failure(self): + bind, session = _bind() + + async with advisory_lock_ctx( + bind=session, dialect_name="postgresql", resource="test_resource" + ): + pass + + assert bind.executes == 2 + + async def test_release_failure_does_not_mask_error(self): + bind, session = _bind() + + with pytest.raises(ValueError, match="failed inside the lock"): + async with advisory_lock_ctx( + bind=session, dialect_name="postgresql", resource="test_resource" + ): + raise ValueError("failed inside the lock") + + assert bind.executes == 2 + + +@pytest.mark.asyncio +class TestTryAdvisoryLockCtx: + async def test_tolerates_release_failure(self): + bind, session = _bind() + + async with try_advisory_lock_ctx( + bind=session, dialect_name="postgresql", resource="test_resource" + ) as locked: + assert locked + + assert bind.executes == 2 + + async def test_release_failure_does_not_mask_error(self): + bind, session = _bind() + + with pytest.raises(ValueError, match="failed inside the lock"): + async with try_advisory_lock_ctx( + bind=session, dialect_name="postgresql", resource="test_resource" + ): + raise ValueError("failed inside the lock") + + assert bind.executes == 2 + + async def test_does_not_release_when_lock_not_acquired(self): + bind, session = _bind(locked=False) + + async with try_advisory_lock_ctx( + bind=session, dialect_name="postgresql", resource="test_resource" + ) as locked: + assert not locked + + assert bind.executes == 1 From 54a0a8a1fab4152452fd096c41b52e79e00746a2 Mon Sep 17 00:00:00 2001 From: Victor Skvortsov Date: Mon, 17 Aug 2026 16:02:31 +0500 Subject: [PATCH 4/5] Handle all get_service_quota() client errors --- .../_internal/core/backends/aws/compute.py | 21 ++++++-- .../core/backends/aws/test_compute.py | 53 +++++++++++++++++++ 2 files changed, 71 insertions(+), 3 deletions(-) create mode 100644 src/tests/_internal/core/backends/aws/test_compute.py diff --git a/src/dstack/_internal/core/backends/aws/compute.py b/src/dstack/_internal/core/backends/aws/compute.py index 4fa3db952..1ed48b8c0 100644 --- a/src/dstack/_internal/core/backends/aws/compute.py +++ b/src/dstack/_internal/core/backends/aws/compute.py @@ -1291,6 +1291,19 @@ def _get_vpc_id_subnets_ids_by_vpc_name_or_error( "L-DB2E81BA": "G/OnDemand", } +# `GetServiceQuota` errors that say nothing about dstack: the quota stays unknown and +# the offer availability stays `UNKNOWN`. Any other error code is reported, since it +# may mean the request is malformed. +_EXPECTED_QUOTA_ERROR_CODES = { + "408", # request timed out + "TooManyRequestsException", # rate limits + "AccessDeniedException", # no servicequotas:GetServiceQuota permission + "AuthFailure", # invalid, expired, or deactivated credentials + "InvalidClientTokenId", + "RequestExpired", + "UnrecognizedClientException", +} + def _get_regions_to_quotas( session: boto3.Session, regions: List[str] @@ -1302,14 +1315,16 @@ def get_region_quotas(region_name: str, client: botocore.client.BaseClient) -> D resp = client.get_service_quota(ServiceCode="ec2", QuotaCode=quota_code) region_quotas[quota_class] = resp["Quota"]["Value"] except botocore.exceptions.ClientError as e: - if "TooManyRequestsException" in str(e): + error_code = e.response.get("Error", {}).get("Code", "") + if error_code in _EXPECTED_QUOTA_ERROR_CODES: logger.warning( - "Failed to get quota %s in %s due to rate limits", + "Failed to get quota %s in %s: %s", quota_code, region_name, + error_code, ) else: - logger.exception(e) + logger.exception("Failed to get quota %s in %s", quota_code, region_name) return region_quotas regions_to_quotas = {} diff --git a/src/tests/_internal/core/backends/aws/test_compute.py b/src/tests/_internal/core/backends/aws/test_compute.py new file mode 100644 index 000000000..d38f97459 --- /dev/null +++ b/src/tests/_internal/core/backends/aws/test_compute.py @@ -0,0 +1,53 @@ +import logging +from unittest.mock import Mock + +import botocore.exceptions +import pytest + +from dstack._internal.core.backends.aws.compute import _get_regions_to_quotas + + +def _session_raising(error_code: str) -> Mock: + client = Mock() + client.get_service_quota.side_effect = botocore.exceptions.ClientError( + {"Error": {"Code": error_code, "Message": ""}}, "GetServiceQuota" + ) + session = Mock() + session.client.return_value = client + return session + + +class TestGetRegionsToQuotas: + def test_returns_quotas(self): + client = Mock() + client.get_service_quota.return_value = {"Quota": {"Value": 8}} + session = Mock() + session.client.return_value = client + + assert _get_regions_to_quotas(session=session, regions=["eu-west-1"]) == { + "eu-west-1": { + "Standard/OnDemand": 8, + "P/OnDemand": 8, + "G/OnDemand": 8, + } + } + + @pytest.mark.parametrize("error_code", ["408", "UnrecognizedClientException"]) + def test_expected_errors_are_not_reported(self, error_code: str, caplog): + with caplog.at_level(logging.WARNING): + regions_to_quotas = _get_regions_to_quotas( + session=_session_raising(error_code), regions=["eu-west-1"] + ) + + assert regions_to_quotas == {"eu-west-1": {}} + assert not [r for r in caplog.records if r.levelno >= logging.ERROR] + assert error_code in caplog.text + + def test_unexpected_error_is_reported(self, caplog): + with caplog.at_level(logging.WARNING): + regions_to_quotas = _get_regions_to_quotas( + session=_session_raising("NoSuchResourceException"), regions=["eu-west-1"] + ) + + assert regions_to_quotas == {"eu-west-1": {}} + assert [r for r in caplog.records if r.levelno >= logging.ERROR] From c5940e58a30045eb46f4f528de0f1892a9976fc5 Mon Sep 17 00:00:00 2001 From: Victor Skvortsov Date: Mon, 17 Aug 2026 16:06:14 +0500 Subject: [PATCH 5/5] Drop tests --- .../_internal/core/backends/aws/compute.py | 2 +- .../core/backends/aws/test_compute.py | 53 ----------- .../_internal/server/services/test_locking.py | 91 ------------------- 3 files changed, 1 insertion(+), 145 deletions(-) delete mode 100644 src/tests/_internal/core/backends/aws/test_compute.py delete mode 100644 src/tests/_internal/server/services/test_locking.py diff --git a/src/dstack/_internal/core/backends/aws/compute.py b/src/dstack/_internal/core/backends/aws/compute.py index 1ed48b8c0..39521cc7c 100644 --- a/src/dstack/_internal/core/backends/aws/compute.py +++ b/src/dstack/_internal/core/backends/aws/compute.py @@ -1321,7 +1321,7 @@ def get_region_quotas(region_name: str, client: botocore.client.BaseClient) -> D "Failed to get quota %s in %s: %s", quota_code, region_name, - error_code, + e, ) else: logger.exception("Failed to get quota %s in %s", quota_code, region_name) diff --git a/src/tests/_internal/core/backends/aws/test_compute.py b/src/tests/_internal/core/backends/aws/test_compute.py deleted file mode 100644 index d38f97459..000000000 --- a/src/tests/_internal/core/backends/aws/test_compute.py +++ /dev/null @@ -1,53 +0,0 @@ -import logging -from unittest.mock import Mock - -import botocore.exceptions -import pytest - -from dstack._internal.core.backends.aws.compute import _get_regions_to_quotas - - -def _session_raising(error_code: str) -> Mock: - client = Mock() - client.get_service_quota.side_effect = botocore.exceptions.ClientError( - {"Error": {"Code": error_code, "Message": ""}}, "GetServiceQuota" - ) - session = Mock() - session.client.return_value = client - return session - - -class TestGetRegionsToQuotas: - def test_returns_quotas(self): - client = Mock() - client.get_service_quota.return_value = {"Quota": {"Value": 8}} - session = Mock() - session.client.return_value = client - - assert _get_regions_to_quotas(session=session, regions=["eu-west-1"]) == { - "eu-west-1": { - "Standard/OnDemand": 8, - "P/OnDemand": 8, - "G/OnDemand": 8, - } - } - - @pytest.mark.parametrize("error_code", ["408", "UnrecognizedClientException"]) - def test_expected_errors_are_not_reported(self, error_code: str, caplog): - with caplog.at_level(logging.WARNING): - regions_to_quotas = _get_regions_to_quotas( - session=_session_raising(error_code), regions=["eu-west-1"] - ) - - assert regions_to_quotas == {"eu-west-1": {}} - assert not [r for r in caplog.records if r.levelno >= logging.ERROR] - assert error_code in caplog.text - - def test_unexpected_error_is_reported(self, caplog): - with caplog.at_level(logging.WARNING): - regions_to_quotas = _get_regions_to_quotas( - session=_session_raising("NoSuchResourceException"), regions=["eu-west-1"] - ) - - assert regions_to_quotas == {"eu-west-1": {}} - assert [r for r in caplog.records if r.levelno >= logging.ERROR] diff --git a/src/tests/_internal/server/services/test_locking.py b/src/tests/_internal/server/services/test_locking.py deleted file mode 100644 index 3fe0c7163..000000000 --- a/src/tests/_internal/server/services/test_locking.py +++ /dev/null @@ -1,91 +0,0 @@ -from typing import Any, cast - -import pytest -from sqlalchemy.ext.asyncio import AsyncSession - -from dstack._internal.server.services.locking import advisory_lock_ctx, try_advisory_lock_ctx - - -class _Result: - def __init__(self, value: bool) -> None: - self._value = value - - def scalar_one(self) -> bool: - return self._value - - -class _ReleaseFailsBind: - """Acquires the lock successfully, then fails to release it.""" - - def __init__(self, locked: bool = True) -> None: - self._locked = locked - self.executes = 0 - - async def execute(self, statement: Any) -> _Result: - self.executes += 1 - if self.executes > 1: - raise RuntimeError("connection invalidated") - return _Result(self._locked) - - -def _bind(locked: bool = True) -> tuple[_ReleaseFailsBind, AsyncSession]: - bind = _ReleaseFailsBind(locked) - return bind, cast(AsyncSession, bind) - - -@pytest.mark.asyncio -class TestAdvisoryLockCtx: - async def test_tolerates_release_failure(self): - bind, session = _bind() - - async with advisory_lock_ctx( - bind=session, dialect_name="postgresql", resource="test_resource" - ): - pass - - assert bind.executes == 2 - - async def test_release_failure_does_not_mask_error(self): - bind, session = _bind() - - with pytest.raises(ValueError, match="failed inside the lock"): - async with advisory_lock_ctx( - bind=session, dialect_name="postgresql", resource="test_resource" - ): - raise ValueError("failed inside the lock") - - assert bind.executes == 2 - - -@pytest.mark.asyncio -class TestTryAdvisoryLockCtx: - async def test_tolerates_release_failure(self): - bind, session = _bind() - - async with try_advisory_lock_ctx( - bind=session, dialect_name="postgresql", resource="test_resource" - ) as locked: - assert locked - - assert bind.executes == 2 - - async def test_release_failure_does_not_mask_error(self): - bind, session = _bind() - - with pytest.raises(ValueError, match="failed inside the lock"): - async with try_advisory_lock_ctx( - bind=session, dialect_name="postgresql", resource="test_resource" - ): - raise ValueError("failed inside the lock") - - assert bind.executes == 2 - - async def test_does_not_release_when_lock_not_acquired(self): - bind, session = _bind(locked=False) - - async with try_advisory_lock_ctx( - bind=session, dialect_name="postgresql", resource="test_resource" - ) as locked: - assert not locked - - assert bind.executes == 1