From c593fde45127641b09193b2862d42cc5de4e5b59 Mon Sep 17 00:00:00 2001 From: Prajjwal Bhandari Date: Fri, 28 Aug 2026 14:29:21 -0400 Subject: [PATCH 01/12] feat(allocation_policies): pardon rejections when the cluster is idle --- sentry-options/schemas/snuba/schema.json | 12 ++- snuba/query/allocation_policies/__init__.py | 53 +++++++-- snuba/web/db_query.py | 12 ++- .../web/rpc/storage_routing/load_retriever.py | 22 +++- .../routing_strategies/storage_routing.py | 8 +- .../test_allocation_policy_base.py | 43 ++++++++ .../test_cluster_loadinfo.py | 24 ++++- tests/web/rpc/v1/test_storage_routing.py | 70 ++++++++++++ tests/web/test_db_query.py | 102 ++++++++++++++++++ 9 files changed, 323 insertions(+), 23 deletions(-) diff --git a/sentry-options/schemas/snuba/schema.json b/sentry-options/schemas/snuba/schema.json index dcec5f5e46..b92f4975f4 100644 --- a/sentry-options/schemas/snuba/schema.json +++ b/sentry-options/schemas/snuba/schema.json @@ -272,7 +272,17 @@ "storage_routing.enable_get_cluster_loadinfo": { "type": "boolean", "default": false, - "description": "When true, the storage routing strategy fetches ClickHouse cluster load info to inform routing decisions." + "description": "When true, the storage routing strategy fetches ClickHouse cluster load info to inform routing decisions. Also the kill switch for idle allocation-policy pardons." + }, + "storage_routing.idle_cluster_load_threshold": { + "type": "number", + "default": 0, + "description": "Pardon allocation-policy rejections when cluster_load is strictly below this value (and concurrent queries are below storage_routing.idle_concurrent_queries_threshold). 0 means never idle." + }, + "storage_routing.idle_concurrent_queries_threshold": { + "type": "integer", + "default": 0, + "description": "Pardon allocation-policy rejections when concurrent_queries is strictly below this value (and cluster load is below storage_routing.idle_cluster_load_threshold). 0 means never idle." }, "retention_days": { "type": "object", diff --git a/snuba/query/allocation_policies/__init__.py b/snuba/query/allocation_policies/__init__.py index dbcb7783aa..bee3260912 100644 --- a/snuba/query/allocation_policies/__init__.py +++ b/snuba/query/allocation_policies/__init__.py @@ -3,9 +3,9 @@ import json import os from abc import ABC, abstractmethod -from dataclasses import asdict, dataclass, field +from dataclasses import asdict, dataclass, field, replace from enum import Enum -from typing import Any, cast +from typing import TYPE_CHECKING, Any, cast from redis.exceptions import TimeoutError as RedisTimeoutError from sentry_sdk import traces @@ -25,6 +25,9 @@ from snuba.utils.serializable_exception import JsonSerializable, SerializableException from snuba.web import QueryResult +if TYPE_CHECKING: + from snuba.web.rpc.storage_routing.load_retriever import LoadInfo + IS_ENFORCED = "is_enforced" MAX_THREADS = "max_threads" NO_UNITS = "no_units" @@ -375,6 +378,10 @@ def metrics(self) -> MetricsWrapper: def is_enforced(self) -> bool: return bool(self.get_config_value(IS_ENFORCED)) + @property + def is_pardonable(self) -> bool: + return True + @property def max_threads(self) -> int: """Maximum number of threads run a single query on ClickHouse with.""" @@ -412,7 +419,10 @@ def _get_default_config_definitions(self) -> list[Configuration]: return cast(list[Configuration], self._default_config_definitions) def get_quota_allowance( - self, tenant_ids: dict[str, str | int], query_id: str + self, + tenant_ids: dict[str, str | int], + query_id: str, + load_info: LoadInfo | None = None, ) -> QuotaAllowance: with traces.start_span( name=self.__class__.__name__, @@ -442,7 +452,9 @@ def get_quota_allowance( 1, tags={"method": "get_quota_allowance", "reason": type(e).__name__}, ) - return DEFAULT_PASSTHROUGH_POLICY.get_quota_allowance(tenant_ids, query_id) + return DEFAULT_PASSTHROUGH_POLICY.get_quota_allowance( + tenant_ids, query_id, load_info + ) except Exception: self.metrics.increment("fail_open", 1, tags={"method": "get_quota_allowance"}) logger.exception( @@ -450,12 +462,26 @@ def get_quota_allowance( ) if settings.RAISE_ON_ALLOCATION_POLICY_FAILURES: raise - return DEFAULT_PASSTHROUGH_POLICY.get_quota_allowance(tenant_ids, query_id) - if not allowance.can_run: - self.metrics.increment( - "db_request_rejected", - tags={"referrer": str(tenant_ids.get("referrer", "no_referrer"))}, + return DEFAULT_PASSTHROUGH_POLICY.get_quota_allowance( + tenant_ids, query_id, load_info ) + idle_pardon = ( + not allowance.can_run + and self.is_pardonable + and load_info is not None + and load_info.is_idle() + ) + if not allowance.can_run: + if idle_pardon: + self.metrics.increment( + "db_request_pardoned", + tags={"referrer": str(tenant_ids.get("referrer", "no_referrer"))}, + ) + else: + self.metrics.increment( + "db_request_rejected", + tags={"referrer": str(tenant_ids.get("referrer", "no_referrer"))}, + ) elif allowance.max_threads < self.max_threads: # NOTE: The elif is very intentional here. Don't count the throttling # if the request was rejected. @@ -479,6 +505,15 @@ def get_quota_allowance( quota_unit=allowance.quota_unit, suggestion=allowance.suggestion, ) + elif idle_pardon: + assert load_info is not None + allowance = replace( + allowance, + can_run=True, + max_threads=MAX_THRESHOLD, + max_bytes_to_read=0, + ) + allowance.explanation["idle_pardon"] = load_info.to_dict() # make sure we always know which storage key we rejected a query from allowance.explanation["storage_key"] = self._resource_identifier.value for k, v in allowance.to_dict().items(): diff --git a/snuba/web/db_query.py b/snuba/web/db_query.py index 2dc0464eab..d130f10b6b 100644 --- a/snuba/web/db_query.py +++ b/snuba/web/db_query.py @@ -27,6 +27,7 @@ from snuba.clickhouse.query_dsl.accessors import get_time_range_estimate from snuba.clickhouse.query_profiler import generate_profile from snuba.configs.configuration import ResourceIdentifier +from snuba.datasets.storages.factory import get_storage from snuba.datasets.storages.storage_key import StorageKey from snuba.downsampled_storage_tiers import Tier from snuba.query import ProcessableQuery @@ -76,6 +77,7 @@ SerializableExceptionDict, ) from snuba.web import QueryException, QueryResult, constants +from snuba.web.rpc.storage_routing.load_retriever import get_cluster_loadinfo metrics = MetricsWrapper(environment.metrics, "db_query") @@ -878,15 +880,21 @@ def _apply_allocation_policies_quota( rejection_quota_and_policy = None throttle_quota_and_policy = None min_threads_across_policies = MAX_THRESHOLD + load_info = get_cluster_loadinfo( + get_storage( + StorageKey(allocation_policies[0].resource_identifier.value) + ).get_storage_set_key() + ) with traces.start_span( name="_apply_allocation_policies_quota", attributes={SENTRY_OP: "allocation_policy"}, ) as span: for allocation_policy in allocation_policies: - allowance = allocation_policy.get_quota_allowance(attribution_info.tenant_ids, query_id) + allowance = allocation_policy.get_quota_allowance( + attribution_info.tenant_ids, query_id, load_info + ) can_run &= allowance.can_run quota_allowances[allocation_policy.class_name()] = allowance - # QuotaAllowance isn't a valid attribute value; serialize it. span.set_attribute( "quota_allowance", json.dumps( diff --git a/snuba/web/rpc/storage_routing/load_retriever.py b/snuba/web/rpc/storage_routing/load_retriever.py index e9af18700c..4b3f3543b0 100644 --- a/snuba/web/rpc/storage_routing/load_retriever.py +++ b/snuba/web/rpc/storage_routing/load_retriever.py @@ -15,6 +15,7 @@ from snuba.clusters.cluster import ClickhouseClientSettings, get_cluster from snuba.clusters.storage_sets import StorageSetKey from snuba.redis import RedisClientKey, get_redis_client +from snuba.state.sentry_options import get_option from snuba.utils.metrics.wrapper import MetricsWrapper metrics = MetricsWrapper( @@ -47,6 +48,17 @@ def from_dict(cls, load_info_dict: dict[str, float | int | None]) -> LoadInfo: } ) + def is_idle(self) -> bool: + if self.cluster_load == -1.0 or self.concurrent_queries == -1: + return False + idle_load = self.cluster_load < get_option( + "storage_routing.idle_cluster_load_threshold", 0.0 + ) + idle_conc_queries = self.concurrent_queries < get_option( + "storage_routing.idle_concurrent_queries_threshold", 0 + ) + return idle_load and idle_conc_queries + def cache( ttl_secs: int = 60, @@ -86,9 +98,17 @@ def wrapper(*args: Any, **kwargs: Any) -> LoadInfo: return decorator -@cache(ttl_secs=60) def get_cluster_loadinfo( storage_set_key: StorageSetKey = StorageSetKey.EVENTS_ANALYTICS_PLATFORM, +) -> LoadInfo | None: + if not get_option("storage_routing.enable_get_cluster_loadinfo", False): + return None + return _get_cluster_loadinfo(storage_set_key) + + +@cache(ttl_secs=60) +def _get_cluster_loadinfo( + storage_set_key: StorageSetKey = StorageSetKey.EVENTS_ANALYTICS_PLATFORM, ) -> LoadInfo: cluster_name = None try: diff --git a/snuba/web/rpc/storage_routing/routing_strategies/storage_routing.py b/snuba/web/rpc/storage_routing/routing_strategies/storage_routing.py index 05e0e51f11..3a4ecf385d 100644 --- a/snuba/web/rpc/storage_routing/routing_strategies/storage_routing.py +++ b/snuba/web/rpc/storage_routing/routing_strategies/storage_routing.py @@ -462,6 +462,7 @@ def _get_recommendations_from_allocation_policies( recommendations[allocation_policy_name] = allocation_policy.get_quota_allowance( routing_context.tenant_ids, routing_context.query_id, + routing_context.cluster_load_info, ) # QuotaAllowance isn't a valid attribute value; serialize it. span.set_attribute( @@ -482,6 +483,7 @@ def get_routing_decision(self, routing_context: RoutingContext) -> RoutingDecisi try: routing_context.timer.mark(_START_ESTIMATION_MARK) + routing_context.cluster_load_info = get_cluster_loadinfo() routing_context.allocation_policies_recommendations = ( self._get_recommendations_from_allocation_policies(routing_context) ) @@ -500,12 +502,6 @@ def get_routing_decision(self, routing_context: RoutingContext) -> RoutingDecisi is_throttled=combined_allocation_policies_recommendations["is_throttled"], ) - routing_context.cluster_load_info = ( - get_cluster_loadinfo() - if get_option("storage_routing.enable_get_cluster_loadinfo", False) - else None - ) - self._update_routing_decision(routing_decision) org_overrides = self._get_org_clickhouse_setting_overrides( diff --git a/tests/query/allocation_policies/test_allocation_policy_base.py b/tests/query/allocation_policies/test_allocation_policy_base.py index 93fa00b24d..862cfbb0c6 100644 --- a/tests/query/allocation_policies/test_allocation_policy_base.py +++ b/tests/query/allocation_policies/test_allocation_policy_base.py @@ -3,6 +3,7 @@ from unittest import TestCase, mock import pytest +from sentry_options.testing import override_options from snuba.configs.configuration import Configuration, InvalidConfig from snuba.datasets.storages.storage_key import StorageKey @@ -448,6 +449,48 @@ def test_is_not_enforced() -> None: assert throttled_metrics[1].tags["is_enforced"] == "False" +@pytest.mark.redis_db +def test_is_pardonable_overrides_can_run_when_idle() -> None: + from snuba.web.rpc.storage_routing.load_retriever import LoadInfo + + reject_policy = RejectingEverythingAllocationPolicy( + StorageKey("some_storage"), + is_enforced=1, + ) + tenant_ids: dict[str, int | str] = { + "organization_id": 123, + "referrer": "some_referrer", + } + idle = LoadInfo(cluster_load=1.0, concurrent_queries=1) + with override_options( + "snuba", + { + "storage_routing.idle_cluster_load_threshold": 10.0, + "storage_routing.idle_concurrent_queries_threshold": 5, + }, + ): + assert reject_policy.get_quota_allowance(tenant_ids, "deadbeef").can_run is False + pardoned = reject_policy.get_quota_allowance(tenant_ids, "deadbeef", idle) + assert pardoned.can_run is True + assert pardoned.max_threads == MAX_THRESHOLD + assert pardoned.max_bytes_to_read == 0 + assert pardoned.explanation["idle_pardon"] == { + "cluster_load": 1.0, + "concurrent_queries": 1, + } + + pardoned_metrics = get_recorded_metric_calls( + "increment", "allocation_policy.db_request_pardoned" + ) + assert pardoned_metrics is not None + assert len(pardoned_metrics) == 1 + rejected_metrics = get_recorded_metric_calls( + "increment", "allocation_policy.db_request_rejected" + ) + assert rejected_metrics is not None + assert len(rejected_metrics) == 1 + + class TestComponentNameBackwardsCompatibility: """Test that component_name() keeps the historical ``{storage}.{ClassName}`` shape.""" diff --git a/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py b/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py index 78027ae9bb..ed08c51a18 100644 --- a/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py +++ b/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py @@ -1,6 +1,7 @@ from unittest.mock import Mock, patch import pytest +from sentry_options.testing import override_options from snuba.web.rpc.storage_routing.load_retriever import LoadInfo, get_cluster_loadinfo @@ -45,22 +46,31 @@ def test_from_dict_ignores_unknown_keys() -> None: assert load_info.cluster_load == 2.0 assert load_info.concurrent_queries == -1 +ENABLE_LOADINFO = {"storage_routing.enable_get_cluster_loadinfo": True} + + +def test_get_cluster_loadinfo_disabled() -> None: + assert get_cluster_loadinfo() is None + @pytest.mark.redis_db @pytest.mark.clickhouse_db def test_get_cluster_load() -> None: - _assert_probe_ok(get_cluster_loadinfo()) + with override_options("snuba", ENABLE_LOADINFO): + _assert_probe_ok(get_cluster_loadinfo()) @pytest.mark.redis_db @pytest.mark.clickhouse_db def test_get_cluster_load_from_cache() -> None: - with patch("time.time") as mock_time: + with override_options("snuba", ENABLE_LOADINFO), patch("time.time") as mock_time: mock_time.return_value = 0 load_info = get_cluster_loadinfo() + assert load_info is not None mock_time.return_value = 59 second_load_info = get_cluster_loadinfo() + assert second_load_info is not None assert load_info.to_dict() == second_load_info.to_dict() @@ -69,7 +79,10 @@ def test_get_cluster_load_from_cache() -> None: def test_get_cluster_loadinfo_if_cache_fails() -> None: mock_redis = Mock() mock_redis.side_effect = Exception("Test error") - with patch("snuba.redis.get_redis_client") as mock_redis_client: + with ( + override_options("snuba", ENABLE_LOADINFO), + patch("snuba.redis.get_redis_client") as mock_redis_client, + ): mock_redis_client.return_value = mock_redis _assert_probe_ok(get_cluster_loadinfo()) @@ -77,6 +90,9 @@ def test_get_cluster_loadinfo_if_cache_fails() -> None: @pytest.mark.redis_db @pytest.mark.clickhouse_db def test_get_cluster_load_error_handling() -> None: - with patch("snuba.clickhouse.connect.ClickhouseConnectPool.execute") as mock_execute: + with ( + override_options("snuba", ENABLE_LOADINFO), + patch("snuba.clickhouse.connect.ClickhouseConnectPool.execute") as mock_execute, + ): mock_execute.side_effect = Exception("Test error") _assert_probe_failed(get_cluster_loadinfo()) diff --git a/tests/web/rpc/v1/test_storage_routing.py b/tests/web/rpc/v1/test_storage_routing.py index ee805484a0..6201b58b75 100644 --- a/tests/web/rpc/v1/test_storage_routing.py +++ b/tests/web/rpc/v1/test_storage_routing.py @@ -553,6 +553,76 @@ def _update_quota_balance( assert not exc.details["can_run"] +@pytest.mark.redis_db +def test_routing_strategy_idle_pardon_allows_rejected_query() -> None: + from snuba.web.rpc.storage_routing.load_retriever import LoadInfo + + class IdlePardonRejectionPolicy(AllocationPolicy): + def _additional_config_definitions(self) -> list[Configuration]: + return [] + + def _get_quota_allowance( + self, tenant_ids: dict[str, str | int], query_id: str + ) -> QuotaAllowance: + return QuotaAllowance( + can_run=False, + max_threads=0, + explanation={"reason": "policy rejects all queries"}, + is_throttled=False, + throttle_threshold=MAX_THRESHOLD, + rejection_threshold=MAX_THRESHOLD, + quota_used=0, + quota_unit=NO_UNITS, + suggestion=NO_SUGGESTION, + ) + + def _update_quota_balance( + self, + tenant_ids: dict[str, str | int], + query_id: str, + result_or_error: QueryResultOrError, + ) -> None: + return + + with ( + mock.patch.object( + BaseRoutingStrategy, + "get_allocation_policies", + return_value=[ + IdlePardonRejectionPolicy(ResourceIdentifier(StorageKey("doesntmatter"))) + ], + ), + override_options( + "snuba", + { + "storage_routing.enable_get_cluster_loadinfo": True, + "storage_routing.idle_cluster_load_threshold": 10.0, + "storage_routing.idle_concurrent_queries_threshold": 5, + }, + ), + mock.patch( + "snuba.web.rpc.storage_routing.routing_strategies.storage_routing.get_cluster_loadinfo", + return_value=LoadInfo(cluster_load=1.0, concurrent_queries=1), + ), + ): + decision = OutcomesBasedRoutingStrategy().get_routing_decision( + RoutingContext( + in_msg=_get_in_msg(), + timer=Timer("test"), + query_id="abc", + ) + ) + assert decision.can_run is True + assert decision.clickhouse_settings["max_threads"] == MAX_THRESHOLD + pardoned = decision.routing_context.allocation_policies_recommendations[ + "IdlePardonRejectionPolicy" + ] + assert pardoned.explanation["idle_pardon"] == { + "cluster_load": 1.0, + "concurrent_queries": 1, + } + + @pytest.mark.redis_db def test_routing_strategy_with_throttling_allocation_policy() -> None: POLICY_THREADS = 4 diff --git a/tests/web/test_db_query.py b/tests/web/test_db_query.py index 72013a0472..177217e701 100644 --- a/tests/web/test_db_query.py +++ b/tests/web/test_db_query.py @@ -1360,3 +1360,105 @@ def _update_quota_balance( }, }, } + + +class _RejectAllPolicy(AllocationPolicy): + def _additional_config_definitions(self) -> list[Configuration]: + return [] + + def _get_quota_allowance( + self, tenant_ids: dict[str, str | int], query_id: str + ) -> QuotaAllowance: + return QuotaAllowance( + can_run=False, + max_threads=0, + explanation={"reason": "reject"}, + is_throttled=False, + throttle_threshold=MAX_THRESHOLD, + rejection_threshold=MAX_THRESHOLD, + quota_used=0, + quota_unit=NO_UNITS, + suggestion=NO_SUGGESTION, + ) + + def _update_quota_balance( + self, + tenant_ids: dict[str, str | int], + query_id: str, + result_or_error: QueryResultOrError, + ) -> None: + return + + +def test_idle_pardon_allows_rejected_query() -> None: + from snuba.web.rpc.storage_routing.load_retriever import LoadInfo + + attribution_info = mock.Mock() + attribution_info.tenant_ids = {"referrer": "test_referrer", "organization_id": 1} + stats: MutableMapping[str, Any] = {} + with ( + override_options( + "snuba", + { + "storage_routing.enable_get_cluster_loadinfo": True, + "storage_routing.idle_cluster_load_threshold": 10.0, + "storage_routing.idle_concurrent_queries_threshold": 5, + }, + ), + mock.patch( + "snuba.web.rpc.storage_routing.load_retriever.get_cluster_loadinfo", + return_value=LoadInfo(cluster_load=1.0, concurrent_queries=1), + ), + ): + query_settings = HTTPQuerySettings() + _apply_allocation_policies_quota( + query_settings=query_settings, + attribution_info=attribution_info, + formatted_query=mock.Mock(), + stats=stats, + allocation_policies=[_RejectAllPolicy(ResourceIdentifier(StorageKey("errors_ro")))], + query_id="pardon_query", + ) + assert stats["quota_allowance"]["summary"]["is_rejected"] is False + details = stats["quota_allowance"]["details"]["_RejectAllPolicy"] + assert details["can_run"] is True + assert details["max_threads"] == MAX_THRESHOLD + assert details["explanation"]["idle_pardon"] == { + "cluster_load": 1.0, + "concurrent_queries": 1, + } + quota = query_settings.get_resource_quota() + assert quota is not None + assert quota.max_threads == MAX_THRESHOLD + + +def test_idle_pardon_still_rejects_when_not_idle() -> None: + from snuba.web.rpc.storage_routing.load_retriever import LoadInfo + + attribution_info = mock.Mock() + attribution_info.tenant_ids = {"referrer": "test_referrer", "organization_id": 1} + stats: MutableMapping[str, Any] = {} + with ( + override_options( + "snuba", + { + "storage_routing.enable_get_cluster_loadinfo": True, + "storage_routing.idle_cluster_load_threshold": 10.0, + "storage_routing.idle_concurrent_queries_threshold": 5, + }, + ), + mock.patch( + "snuba.web.rpc.storage_routing.load_retriever.get_cluster_loadinfo", + return_value=LoadInfo(cluster_load=50.0, concurrent_queries=1), + ), + pytest.raises(AllocationPolicyViolations), + ): + _apply_allocation_policies_quota( + query_settings=HTTPQuerySettings(), + attribution_info=attribution_info, + formatted_query=mock.Mock(), + stats=stats, + allocation_policies=[_RejectAllPolicy(ResourceIdentifier(StorageKey("errors_ro")))], + query_id="no_pardon_query", + ) + assert stats["quota_allowance"]["summary"]["is_rejected"] is True From ecffa7be71bdc421319f89951286babc72e2b02b Mon Sep 17 00:00:00 2001 From: Prajjwal Bhandari Date: Thu, 3 Sep 2026 13:25:31 -0400 Subject: [PATCH 02/12] ref(allocation_policies): address idle-pardon review comments --- sentry-options/schemas/snuba/schema.json | 10 --- snuba/query/allocation_policies/__init__.py | 89 +++++++++++-------- .../test_allocation_policy_base.py | 26 ++---- .../test_cluster_loadinfo.py | 32 ++++--- tests/web/rpc/v1/test_storage_routing.py | 10 +-- tests/web/test_db_query.py | 29 +++--- 6 files changed, 96 insertions(+), 100 deletions(-) diff --git a/sentry-options/schemas/snuba/schema.json b/sentry-options/schemas/snuba/schema.json index b92f4975f4..1a3d358639 100644 --- a/sentry-options/schemas/snuba/schema.json +++ b/sentry-options/schemas/snuba/schema.json @@ -274,16 +274,6 @@ "default": false, "description": "When true, the storage routing strategy fetches ClickHouse cluster load info to inform routing decisions. Also the kill switch for idle allocation-policy pardons." }, - "storage_routing.idle_cluster_load_threshold": { - "type": "number", - "default": 0, - "description": "Pardon allocation-policy rejections when cluster_load is strictly below this value (and concurrent queries are below storage_routing.idle_concurrent_queries_threshold). 0 means never idle." - }, - "storage_routing.idle_concurrent_queries_threshold": { - "type": "integer", - "default": 0, - "description": "Pardon allocation-policy rejections when concurrent_queries is strictly below this value (and cluster load is below storage_routing.idle_cluster_load_threshold). 0 means never idle." - }, "retention_days": { "type": "object", "properties": { diff --git a/snuba/query/allocation_policies/__init__.py b/snuba/query/allocation_policies/__init__.py index bee3260912..75cf713e8f 100644 --- a/snuba/query/allocation_policies/__init__.py +++ b/snuba/query/allocation_policies/__init__.py @@ -3,8 +3,8 @@ import json import os from abc import ABC, abstractmethod -from dataclasses import asdict, dataclass, field, replace -from enum import Enum +from dataclasses import asdict, dataclass, field +from enum import Enum, IntEnum from typing import TYPE_CHECKING, Any, cast from redis.exceptions import TimeoutError as RedisTimeoutError @@ -26,6 +26,8 @@ from snuba.web import QueryResult if TYPE_CHECKING: + # Importing load_retriever at runtime pulls in snuba.web.rpc, which + # eventually imports db_query → AllocationPolicy (cycle). from snuba.web.rpc.storage_routing.load_retriever import LoadInfo IS_ENFORCED = "is_enforced" @@ -57,6 +59,13 @@ class AllocationPolicyConfig(Configuration): pass +class QuotaAllowanceDecision(IntEnum): + REJECTED = 0 + PARDONED = 1 + ALLOWED = 2 + THROTTLED = 3 + + @dataclass(frozen=True) class QuotaAllowance: can_run: bool @@ -81,6 +90,21 @@ class QuotaAllowance: def to_dict(self) -> dict[str, Any]: return asdict(self) + def decision( + self, + *, + policy_max_threads: int, + idle_pardon: bool, + ) -> QuotaAllowanceDecision: + if not self.can_run: + if idle_pardon: + return QuotaAllowanceDecision.PARDONED + return QuotaAllowanceDecision.REJECTED + max_threads = min(self.max_threads, policy_max_threads) + if max_threads < policy_max_threads: + return QuotaAllowanceDecision.THROTTLED + return QuotaAllowanceDecision.ALLOWED + def __eq__(self, other: Any) -> bool: if not isinstance(other, QuotaAllowance): return False @@ -452,9 +476,9 @@ def get_quota_allowance( 1, tags={"method": "get_quota_allowance", "reason": type(e).__name__}, ) - return DEFAULT_PASSTHROUGH_POLICY.get_quota_allowance( - tenant_ids, query_id, load_info - ) + return _default_passthough_policy( + self._resource_identifier.value + ).get_quota_allowance(tenant_ids, query_id, load_info) except Exception: self.metrics.increment("fail_open", 1, tags={"method": "get_quota_allowance"}) logger.exception( @@ -462,35 +486,22 @@ def get_quota_allowance( ) if settings.RAISE_ON_ALLOCATION_POLICY_FAILURES: raise - return DEFAULT_PASSTHROUGH_POLICY.get_quota_allowance( - tenant_ids, query_id, load_info - ) - idle_pardon = ( - not allowance.can_run - and self.is_pardonable - and load_info is not None - and load_info.is_idle() + return _default_passthough_policy( + self._resource_identifier.value + ).get_quota_allowance(tenant_ids, query_id, load_info) + idle_pardon = self.is_pardonable and load_info is not None and load_info.is_idle() + decision = allowance.decision( + policy_max_threads=self.max_threads, idle_pardon=idle_pardon ) - if not allowance.can_run: - if idle_pardon: - self.metrics.increment( - "db_request_pardoned", - tags={"referrer": str(tenant_ids.get("referrer", "no_referrer"))}, - ) - else: - self.metrics.increment( - "db_request_rejected", - tags={"referrer": str(tenant_ids.get("referrer", "no_referrer"))}, - ) - elif allowance.max_threads < self.max_threads: - # NOTE: The elif is very intentional here. Don't count the throttling - # if the request was rejected. + referrer = str(tenant_ids.get("referrer", "no_referrer")) + if decision == QuotaAllowanceDecision.PARDONED: + self.metrics.increment("db_request_pardoned", tags={"referrer": referrer}) + elif decision == QuotaAllowanceDecision.REJECTED: + self.metrics.increment("db_request_rejected", tags={"referrer": referrer}) + elif decision == QuotaAllowanceDecision.THROTTLED: self.metrics.increment( "db_request_throttled", - tags={ - "referrer": str(tenant_ids.get("referrer", "no_referrer")), - "max_threads": str(allowance.max_threads), - }, + tags={"referrer": referrer, "max_threads": str(allowance.max_threads)}, ) span.set_attribute("db_request_throttled", True) if not self.is_enforced: @@ -505,15 +516,23 @@ def get_quota_allowance( quota_unit=allowance.quota_unit, suggestion=allowance.suggestion, ) - elif idle_pardon: + elif decision == QuotaAllowanceDecision.PARDONED: assert load_info is not None - allowance = replace( - allowance, + allowance = QuotaAllowance( can_run=True, max_threads=MAX_THRESHOLD, + explanation={ + **allowance.explanation, + "idle_pardon": load_info.to_dict(), + }, + is_throttled=allowance.is_throttled, + throttle_threshold=allowance.throttle_threshold, + rejection_threshold=allowance.rejection_threshold, + quota_used=allowance.quota_used, + quota_unit=allowance.quota_unit, + suggestion=allowance.suggestion, max_bytes_to_read=0, ) - allowance.explanation["idle_pardon"] = load_info.to_dict() # make sure we always know which storage key we rejected a query from allowance.explanation["storage_key"] = self._resource_identifier.value for k, v in allowance.to_dict().items(): diff --git a/tests/query/allocation_policies/test_allocation_policy_base.py b/tests/query/allocation_policies/test_allocation_policy_base.py index 862cfbb0c6..18a5b1141f 100644 --- a/tests/query/allocation_policies/test_allocation_policy_base.py +++ b/tests/query/allocation_policies/test_allocation_policy_base.py @@ -3,7 +3,6 @@ from unittest import TestCase, mock import pytest -from sentry_options.testing import override_options from snuba.configs.configuration import Configuration, InvalidConfig from snuba.datasets.storages.storage_key import StorageKey @@ -31,6 +30,7 @@ from snuba.query.allocation_policies.cross_org import CrossOrgQueryAllocationPolicy from snuba.query.allocation_policies.per_referrer import ReferrerGuardRailPolicy from snuba.utils.metrics.backends.testing import get_recorded_metric_calls +from snuba.web.rpc.storage_routing.load_retriever import LoadInfo from tests.configs.component_config import ( delete_component_config, set_component_config, @@ -451,8 +451,6 @@ def test_is_not_enforced() -> None: @pytest.mark.redis_db def test_is_pardonable_overrides_can_run_when_idle() -> None: - from snuba.web.rpc.storage_routing.load_retriever import LoadInfo - reject_policy = RejectingEverythingAllocationPolicy( StorageKey("some_storage"), is_enforced=1, @@ -462,22 +460,16 @@ def test_is_pardonable_overrides_can_run_when_idle() -> None: "referrer": "some_referrer", } idle = LoadInfo(cluster_load=1.0, concurrent_queries=1) - with override_options( - "snuba", - { - "storage_routing.idle_cluster_load_threshold": 10.0, - "storage_routing.idle_concurrent_queries_threshold": 5, - }, - ): + with mock.patch.object(LoadInfo, "is_idle", return_value=True): assert reject_policy.get_quota_allowance(tenant_ids, "deadbeef").can_run is False pardoned = reject_policy.get_quota_allowance(tenant_ids, "deadbeef", idle) - assert pardoned.can_run is True - assert pardoned.max_threads == MAX_THRESHOLD - assert pardoned.max_bytes_to_read == 0 - assert pardoned.explanation["idle_pardon"] == { - "cluster_load": 1.0, - "concurrent_queries": 1, - } + assert pardoned.can_run is True + assert pardoned.max_threads == MAX_THRESHOLD + assert pardoned.max_bytes_to_read == 0 + assert pardoned.explanation["idle_pardon"] == { + "cluster_load": 1.0, + "concurrent_queries": 1, + } pardoned_metrics = get_recorded_metric_calls( "increment", "allocation_policy.db_request_pardoned" diff --git a/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py b/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py index ed08c51a18..b6c30db7d1 100644 --- a/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py +++ b/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py @@ -1,6 +1,8 @@ +from collections.abc import Generator from unittest.mock import Mock, patch import pytest +from sentry_options import OptionValue from sentry_options.testing import override_options from snuba.web.rpc.storage_routing.load_retriever import LoadInfo, get_cluster_loadinfo @@ -46,24 +48,34 @@ def test_from_dict_ignores_unknown_keys() -> None: assert load_info.cluster_load == 2.0 assert load_info.concurrent_queries == -1 -ENABLE_LOADINFO = {"storage_routing.enable_get_cluster_loadinfo": True} +ENABLE_LOADINFO: dict[str, OptionValue] = {"storage_routing.enable_get_cluster_loadinfo": True} +@pytest.fixture(autouse=True) +def enable_get_cluster_loadinfo() -> Generator[None]: + with override_options("snuba", ENABLE_LOADINFO): + yield + + +@pytest.mark.redis_db +@pytest.mark.clickhouse_db def test_get_cluster_loadinfo_disabled() -> None: - assert get_cluster_loadinfo() is None + with override_options("snuba", {"storage_routing.enable_get_cluster_loadinfo": False}): + assert get_cluster_loadinfo() is None + with override_options("snuba", ENABLE_LOADINFO): + assert get_cluster_loadinfo() is not None @pytest.mark.redis_db @pytest.mark.clickhouse_db def test_get_cluster_load() -> None: - with override_options("snuba", ENABLE_LOADINFO): - _assert_probe_ok(get_cluster_loadinfo()) + _assert_probe_ok(get_cluster_loadinfo()) @pytest.mark.redis_db @pytest.mark.clickhouse_db def test_get_cluster_load_from_cache() -> None: - with override_options("snuba", ENABLE_LOADINFO), patch("time.time") as mock_time: + with patch("time.time") as mock_time: mock_time.return_value = 0 load_info = get_cluster_loadinfo() assert load_info is not None @@ -79,10 +91,7 @@ def test_get_cluster_load_from_cache() -> None: def test_get_cluster_loadinfo_if_cache_fails() -> None: mock_redis = Mock() mock_redis.side_effect = Exception("Test error") - with ( - override_options("snuba", ENABLE_LOADINFO), - patch("snuba.redis.get_redis_client") as mock_redis_client, - ): + with patch("snuba.redis.get_redis_client") as mock_redis_client: mock_redis_client.return_value = mock_redis _assert_probe_ok(get_cluster_loadinfo()) @@ -90,9 +99,6 @@ def test_get_cluster_loadinfo_if_cache_fails() -> None: @pytest.mark.redis_db @pytest.mark.clickhouse_db def test_get_cluster_load_error_handling() -> None: - with ( - override_options("snuba", ENABLE_LOADINFO), - patch("snuba.clickhouse.connect.ClickhouseConnectPool.execute") as mock_execute, - ): + with patch("snuba.clickhouse.connect.ClickhouseConnectPool.execute") as mock_execute: mock_execute.side_effect = Exception("Test error") _assert_probe_failed(get_cluster_loadinfo()) diff --git a/tests/web/rpc/v1/test_storage_routing.py b/tests/web/rpc/v1/test_storage_routing.py index 6201b58b75..038cc027e1 100644 --- a/tests/web/rpc/v1/test_storage_routing.py +++ b/tests/web/rpc/v1/test_storage_routing.py @@ -35,6 +35,7 @@ encode_routing_hint, extract_message_meta, ) +from snuba.web.rpc.storage_routing.load_retriever import LoadInfo from snuba.web.rpc.storage_routing.routing_strategies.outcomes_based import ( OutcomesBasedRoutingStrategy, ) @@ -555,8 +556,6 @@ def _update_quota_balance( @pytest.mark.redis_db def test_routing_strategy_idle_pardon_allows_rejected_query() -> None: - from snuba.web.rpc.storage_routing.load_retriever import LoadInfo - class IdlePardonRejectionPolicy(AllocationPolicy): def _additional_config_definitions(self) -> list[Configuration]: return [] @@ -594,16 +593,13 @@ def _update_quota_balance( ), override_options( "snuba", - { - "storage_routing.enable_get_cluster_loadinfo": True, - "storage_routing.idle_cluster_load_threshold": 10.0, - "storage_routing.idle_concurrent_queries_threshold": 5, - }, + {"storage_routing.enable_get_cluster_loadinfo": True}, ), mock.patch( "snuba.web.rpc.storage_routing.routing_strategies.storage_routing.get_cluster_loadinfo", return_value=LoadInfo(cluster_load=1.0, concurrent_queries=1), ), + mock.patch.object(LoadInfo, "is_idle", return_value=True), ): decision = OutcomesBasedRoutingStrategy().get_routing_decision( RoutingContext( diff --git a/tests/web/test_db_query.py b/tests/web/test_db_query.py index 177217e701..f7b3ff4c9a 100644 --- a/tests/web/test_db_query.py +++ b/tests/web/test_db_query.py @@ -41,6 +41,7 @@ db_query, execute_query, ) +from snuba.web.rpc.storage_routing.load_retriever import LoadInfo from tests.query.allocation_policies.attachment import ( match_block, override_allocation_policy, @@ -1391,24 +1392,20 @@ def _update_quota_balance( def test_idle_pardon_allows_rejected_query() -> None: - from snuba.web.rpc.storage_routing.load_retriever import LoadInfo - attribution_info = mock.Mock() attribution_info.tenant_ids = {"referrer": "test_referrer", "organization_id": 1} stats: MutableMapping[str, Any] = {} + idle = LoadInfo(cluster_load=1.0, concurrent_queries=1) with ( override_options( "snuba", - { - "storage_routing.enable_get_cluster_loadinfo": True, - "storage_routing.idle_cluster_load_threshold": 10.0, - "storage_routing.idle_concurrent_queries_threshold": 5, - }, + {"storage_routing.enable_get_cluster_loadinfo": True}, ), mock.patch( - "snuba.web.rpc.storage_routing.load_retriever.get_cluster_loadinfo", - return_value=LoadInfo(cluster_load=1.0, concurrent_queries=1), + "snuba.web.db_query.get_cluster_loadinfo", + return_value=idle, ), + mock.patch.object(LoadInfo, "is_idle", return_value=True), ): query_settings = HTTPQuerySettings() _apply_allocation_policies_quota( @@ -1433,24 +1430,20 @@ def test_idle_pardon_allows_rejected_query() -> None: def test_idle_pardon_still_rejects_when_not_idle() -> None: - from snuba.web.rpc.storage_routing.load_retriever import LoadInfo - attribution_info = mock.Mock() attribution_info.tenant_ids = {"referrer": "test_referrer", "organization_id": 1} stats: MutableMapping[str, Any] = {} + busy = LoadInfo(cluster_load=50.0, concurrent_queries=1) with ( override_options( "snuba", - { - "storage_routing.enable_get_cluster_loadinfo": True, - "storage_routing.idle_cluster_load_threshold": 10.0, - "storage_routing.idle_concurrent_queries_threshold": 5, - }, + {"storage_routing.enable_get_cluster_loadinfo": True}, ), mock.patch( - "snuba.web.rpc.storage_routing.load_retriever.get_cluster_loadinfo", - return_value=LoadInfo(cluster_load=50.0, concurrent_queries=1), + "snuba.web.db_query.get_cluster_loadinfo", + return_value=busy, ), + mock.patch.object(LoadInfo, "is_idle", return_value=False), pytest.raises(AllocationPolicyViolations), ): _apply_allocation_policies_quota( From 5729ba10c2ea86235678466cb5d91133685ecffd Mon Sep 17 00:00:00 2001 From: Prajjwal Bhandari Date: Thu, 3 Sep 2026 13:41:12 -0400 Subject: [PATCH 03/12] ref(allocation_policies): fold idle check into QuotaAllowance.decision --- snuba/query/allocation_policies/__init__.py | 27 +++++++++---------- .../test_allocation_policy_base.py | 2 +- .../test_cluster_loadinfo.py | 2 +- tests/web/rpc/v1/test_storage_routing.py | 2 +- tests/web/test_db_query.py | 4 +-- 5 files changed, 17 insertions(+), 20 deletions(-) diff --git a/snuba/query/allocation_policies/__init__.py b/snuba/query/allocation_policies/__init__.py index 75cf713e8f..030bad8e53 100644 --- a/snuba/query/allocation_policies/__init__.py +++ b/snuba/query/allocation_policies/__init__.py @@ -92,17 +92,18 @@ def to_dict(self) -> dict[str, Any]: def decision( self, - *, - policy_max_threads: int, - idle_pardon: bool, + policy: AllocationPolicy, + load_info: LoadInfo | None, ) -> QuotaAllowanceDecision: if not self.can_run: - if idle_pardon: - return QuotaAllowanceDecision.PARDONED - return QuotaAllowanceDecision.REJECTED - max_threads = min(self.max_threads, policy_max_threads) - if max_threads < policy_max_threads: + idle_pardon = policy.is_pardonable and load_info is not None and load_info.is_idle() + return ( + QuotaAllowanceDecision.PARDONED if idle_pardon else QuotaAllowanceDecision.REJECTED + ) + + if self.max_threads < policy.max_threads: return QuotaAllowanceDecision.THROTTLED + return QuotaAllowanceDecision.ALLOWED def __eq__(self, other: Any) -> bool: @@ -489,10 +490,7 @@ def get_quota_allowance( return _default_passthough_policy( self._resource_identifier.value ).get_quota_allowance(tenant_ids, query_id, load_info) - idle_pardon = self.is_pardonable and load_info is not None and load_info.is_idle() - decision = allowance.decision( - policy_max_threads=self.max_threads, idle_pardon=idle_pardon - ) + decision = allowance.decision(self, load_info) referrer = str(tenant_ids.get("referrer", "no_referrer")) if decision == QuotaAllowanceDecision.PARDONED: self.metrics.increment("db_request_pardoned", tags={"referrer": referrer}) @@ -517,13 +515,12 @@ def get_quota_allowance( suggestion=allowance.suggestion, ) elif decision == QuotaAllowanceDecision.PARDONED: - assert load_info is not None allowance = QuotaAllowance( can_run=True, - max_threads=MAX_THRESHOLD, + max_threads=self.max_threads, explanation={ **allowance.explanation, - "idle_pardon": load_info.to_dict(), + "idle_pardon": load_info.to_dict() if load_info is not None else {}, }, is_throttled=allowance.is_throttled, throttle_threshold=allowance.throttle_threshold, diff --git a/tests/query/allocation_policies/test_allocation_policy_base.py b/tests/query/allocation_policies/test_allocation_policy_base.py index 18a5b1141f..9f00d3875f 100644 --- a/tests/query/allocation_policies/test_allocation_policy_base.py +++ b/tests/query/allocation_policies/test_allocation_policy_base.py @@ -464,7 +464,7 @@ def test_is_pardonable_overrides_can_run_when_idle() -> None: assert reject_policy.get_quota_allowance(tenant_ids, "deadbeef").can_run is False pardoned = reject_policy.get_quota_allowance(tenant_ids, "deadbeef", idle) assert pardoned.can_run is True - assert pardoned.max_threads == MAX_THRESHOLD + assert pardoned.max_threads == reject_policy.max_threads assert pardoned.max_bytes_to_read == 0 assert pardoned.explanation["idle_pardon"] == { "cluster_load": 1.0, diff --git a/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py b/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py index b6c30db7d1..e5263e2512 100644 --- a/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py +++ b/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py @@ -62,7 +62,7 @@ def enable_get_cluster_loadinfo() -> Generator[None]: def test_get_cluster_loadinfo_disabled() -> None: with override_options("snuba", {"storage_routing.enable_get_cluster_loadinfo": False}): assert get_cluster_loadinfo() is None - with override_options("snuba", ENABLE_LOADINFO): + with override_options("snuba", {"storage_routing.enable_get_cluster_loadinfo": True}): assert get_cluster_loadinfo() is not None diff --git a/tests/web/rpc/v1/test_storage_routing.py b/tests/web/rpc/v1/test_storage_routing.py index 038cc027e1..e1aedc4c9b 100644 --- a/tests/web/rpc/v1/test_storage_routing.py +++ b/tests/web/rpc/v1/test_storage_routing.py @@ -609,7 +609,7 @@ def _update_quota_balance( ) ) assert decision.can_run is True - assert decision.clickhouse_settings["max_threads"] == MAX_THRESHOLD + assert decision.clickhouse_settings["max_threads"] == 10 pardoned = decision.routing_context.allocation_policies_recommendations[ "IdlePardonRejectionPolicy" ] diff --git a/tests/web/test_db_query.py b/tests/web/test_db_query.py index f7b3ff4c9a..144b2e34da 100644 --- a/tests/web/test_db_query.py +++ b/tests/web/test_db_query.py @@ -1419,14 +1419,14 @@ def test_idle_pardon_allows_rejected_query() -> None: assert stats["quota_allowance"]["summary"]["is_rejected"] is False details = stats["quota_allowance"]["details"]["_RejectAllPolicy"] assert details["can_run"] is True - assert details["max_threads"] == MAX_THRESHOLD + assert details["max_threads"] == 10 assert details["explanation"]["idle_pardon"] == { "cluster_load": 1.0, "concurrent_queries": 1, } quota = query_settings.get_resource_quota() assert quota is not None - assert quota.max_threads == MAX_THRESHOLD + assert quota.max_threads == 10 def test_idle_pardon_still_rejects_when_not_idle() -> None: From 27af7336cea99d22855953d209074120c70a454a Mon Sep 17 00:00:00 2001 From: Prajjwal Bhandari Date: Thu, 3 Sep 2026 13:47:53 -0400 Subject: [PATCH 04/12] ref(allocation_policies): getattr is_idle; space get_quota_allowance --- snuba/query/allocation_policies/__init__.py | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/snuba/query/allocation_policies/__init__.py b/snuba/query/allocation_policies/__init__.py index 030bad8e53..83ac8b46f7 100644 --- a/snuba/query/allocation_policies/__init__.py +++ b/snuba/query/allocation_policies/__init__.py @@ -96,7 +96,8 @@ def decision( load_info: LoadInfo | None, ) -> QuotaAllowanceDecision: if not self.can_run: - idle_pardon = policy.is_pardonable and load_info is not None and load_info.is_idle() + is_idle_load = getattr(load_info, "is_idle", lambda: False) + idle_pardon = policy.is_pardonable and is_idle_load() return ( QuotaAllowanceDecision.PARDONED if idle_pardon else QuotaAllowanceDecision.REJECTED ) @@ -455,6 +456,7 @@ def get_quota_allowance( ) as span: for t, tid in tenant_ids.items(): span.set_attribute(f"tenant_ids.{t}", str(tid)) + try: allowance = self._get_quota_allowance(tenant_ids, query_id) except InvalidTenantsForAllocationPolicy as e: @@ -490,6 +492,7 @@ def get_quota_allowance( return _default_passthough_policy( self._resource_identifier.value ).get_quota_allowance(tenant_ids, query_id, load_info) + decision = allowance.decision(self, load_info) referrer = str(tenant_ids.get("referrer", "no_referrer")) if decision == QuotaAllowanceDecision.PARDONED: @@ -502,6 +505,7 @@ def get_quota_allowance( tags={"referrer": referrer, "max_threads": str(allowance.max_threads)}, ) span.set_attribute("db_request_throttled", True) + if not self.is_enforced: allowance = QuotaAllowance( can_run=True, @@ -530,6 +534,7 @@ def get_quota_allowance( suggestion=allowance.suggestion, max_bytes_to_read=0, ) + # make sure we always know which storage key we rejected a query from allowance.explanation["storage_key"] = self._resource_identifier.value for k, v in allowance.to_dict().items(): From ac4654ccc64d2baa85b4b6a64027c3a2f71c6b09 Mon Sep 17 00:00:00 2001 From: Prajjwal Bhandari Date: Fri, 2 Oct 2026 15:08:57 -0700 Subject: [PATCH 05/12] Simplify when we do and don't pardon allocation policies --- sentry-options/schemas/snuba/schema.json | 7 ++++++- snuba/web/rpc/storage_routing/load_retriever.py | 8 +------- .../v1/routing_strategies/test_cluster_loadinfo.py | 11 ++++++++++- 3 files changed, 17 insertions(+), 9 deletions(-) diff --git a/sentry-options/schemas/snuba/schema.json b/sentry-options/schemas/snuba/schema.json index 1a3d358639..c3c42a920a 100644 --- a/sentry-options/schemas/snuba/schema.json +++ b/sentry-options/schemas/snuba/schema.json @@ -272,7 +272,12 @@ "storage_routing.enable_get_cluster_loadinfo": { "type": "boolean", "default": false, - "description": "When true, the storage routing strategy fetches ClickHouse cluster load info to inform routing decisions. Also the kill switch for idle allocation-policy pardons." + "description": "When true, the storage routing strategy fetches ClickHouse cluster load info to inform routing decisions." + }, + "storage_routing.enable_dynamic_allocation_policy": { + "type": "boolean", + "default": false, + "description": "When true, allocation-policy rejections are pardoned. Load thresholds are not applied yet." }, "retention_days": { "type": "object", diff --git a/snuba/web/rpc/storage_routing/load_retriever.py b/snuba/web/rpc/storage_routing/load_retriever.py index 4b3f3543b0..c6ef390bd6 100644 --- a/snuba/web/rpc/storage_routing/load_retriever.py +++ b/snuba/web/rpc/storage_routing/load_retriever.py @@ -51,13 +51,7 @@ def from_dict(cls, load_info_dict: dict[str, float | int | None]) -> LoadInfo: def is_idle(self) -> bool: if self.cluster_load == -1.0 or self.concurrent_queries == -1: return False - idle_load = self.cluster_load < get_option( - "storage_routing.idle_cluster_load_threshold", 0.0 - ) - idle_conc_queries = self.concurrent_queries < get_option( - "storage_routing.idle_concurrent_queries_threshold", 0 - ) - return idle_load and idle_conc_queries + return get_option("storage_routing.enable_dynamic_allocation_policy", False) def cache( diff --git a/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py b/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py index e5263e2512..9fd942d1cc 100644 --- a/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py +++ b/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py @@ -101,4 +101,13 @@ def test_get_cluster_loadinfo_if_cache_fails() -> None: def test_get_cluster_load_error_handling() -> None: with patch("snuba.clickhouse.connect.ClickhouseConnectPool.execute") as mock_execute: mock_execute.side_effect = Exception("Test error") - _assert_probe_failed(get_cluster_loadinfo()) + load_info = get_cluster_loadinfo() + _assert_probe_failed(load_info) + assert load_info.is_idle() is False + + +def test_is_idle_reads_dynamic_allocation_policy_flag() -> None: + load_info = LoadInfo(cluster_load=1.0, concurrent_queries=1) + assert load_info.is_idle() is False + with override_options("snuba", {"storage_routing.enable_dynamic_allocation_policy": True}): + assert load_info.is_idle() is True From d464149d2ebc87e9a3d910444c7c457b8cbbb805 Mon Sep 17 00:00:00 2001 From: Prajjwal Bhandari Date: Fri, 2 Oct 2026 15:18:18 -0700 Subject: [PATCH 06/12] Fix tests post rebase --- .../allocation_policies/test_allocation_policy_base.py | 5 +---- tests/web/rpc/v1/test_storage_routing.py | 8 ++++---- tests/web/test_db_query.py | 5 +---- 3 files changed, 6 insertions(+), 12 deletions(-) diff --git a/tests/query/allocation_policies/test_allocation_policy_base.py b/tests/query/allocation_policies/test_allocation_policy_base.py index 9f00d3875f..1c861d94d9 100644 --- a/tests/query/allocation_policies/test_allocation_policy_base.py +++ b/tests/query/allocation_policies/test_allocation_policy_base.py @@ -466,10 +466,7 @@ def test_is_pardonable_overrides_can_run_when_idle() -> None: assert pardoned.can_run is True assert pardoned.max_threads == reject_policy.max_threads assert pardoned.max_bytes_to_read == 0 - assert pardoned.explanation["idle_pardon"] == { - "cluster_load": 1.0, - "concurrent_queries": 1, - } + assert pardoned.explanation["idle_pardon"] == idle.to_dict() pardoned_metrics = get_recorded_metric_calls( "increment", "allocation_policy.db_request_pardoned" diff --git a/tests/web/rpc/v1/test_storage_routing.py b/tests/web/rpc/v1/test_storage_routing.py index e1aedc4c9b..b68859e7fe 100644 --- a/tests/web/rpc/v1/test_storage_routing.py +++ b/tests/web/rpc/v1/test_storage_routing.py @@ -613,10 +613,10 @@ def _update_quota_balance( pardoned = decision.routing_context.allocation_policies_recommendations[ "IdlePardonRejectionPolicy" ] - assert pardoned.explanation["idle_pardon"] == { - "cluster_load": 1.0, - "concurrent_queries": 1, - } + assert ( + pardoned.explanation["idle_pardon"] + == LoadInfo(cluster_load=1.0, concurrent_queries=1).to_dict() + ) @pytest.mark.redis_db diff --git a/tests/web/test_db_query.py b/tests/web/test_db_query.py index 144b2e34da..0d467867d6 100644 --- a/tests/web/test_db_query.py +++ b/tests/web/test_db_query.py @@ -1420,10 +1420,7 @@ def test_idle_pardon_allows_rejected_query() -> None: details = stats["quota_allowance"]["details"]["_RejectAllPolicy"] assert details["can_run"] is True assert details["max_threads"] == 10 - assert details["explanation"]["idle_pardon"] == { - "cluster_load": 1.0, - "concurrent_queries": 1, - } + assert details["explanation"]["idle_pardon"] == idle.to_dict() quota = query_settings.get_resource_quota() assert quota is not None assert quota.max_threads == 10 From 264d4bdf6048d8bbdec1f4d7f47c52ed8443d548 Mon Sep 17 00:00:00 2001 From: Prajjwal Bhandari Date: Fri, 2 Oct 2026 15:30:34 -0700 Subject: [PATCH 07/12] ref(allocation_policies): move LoadInfo out of snuba.web.rpc --- .../load_retriever.py => clusters/load_info.py} | 2 +- snuba/query/allocation_policies/__init__.py | 8 ++------ snuba/web/db_query.py | 2 +- .../routing_strategies/storage_routing.py | 2 +- .../test_allocation_policy_base.py | 2 +- .../v1/routing_strategies/test_cluster_loadinfo.py | 12 +++++++++--- tests/web/rpc/v1/test_storage_routing.py | 2 +- tests/web/test_db_query.py | 2 +- 8 files changed, 17 insertions(+), 15 deletions(-) rename snuba/{web/rpc/storage_routing/load_retriever.py => clusters/load_info.py} (99%) diff --git a/snuba/web/rpc/storage_routing/load_retriever.py b/snuba/clusters/load_info.py similarity index 99% rename from snuba/web/rpc/storage_routing/load_retriever.py rename to snuba/clusters/load_info.py index c6ef390bd6..8d8e7523cf 100644 --- a/snuba/web/rpc/storage_routing/load_retriever.py +++ b/snuba/clusters/load_info.py @@ -20,7 +20,7 @@ metrics = MetricsWrapper( environment.metrics, - "snuba.web.rpc.storage_routing.load_retriever", + "snuba.clusters.load_info", ) diff --git a/snuba/query/allocation_policies/__init__.py b/snuba/query/allocation_policies/__init__.py index 83ac8b46f7..6f52bc8834 100644 --- a/snuba/query/allocation_policies/__init__.py +++ b/snuba/query/allocation_policies/__init__.py @@ -5,12 +5,13 @@ from abc import ABC, abstractmethod from dataclasses import asdict, dataclass, field from enum import Enum, IntEnum -from typing import TYPE_CHECKING, Any, cast +from typing import Any, cast from redis.exceptions import TimeoutError as RedisTimeoutError from sentry_sdk import traces from snuba import environment, settings +from snuba.clusters.load_info import LoadInfo from snuba.configs.configuration import ( ConfigurableComponent, ConfigurableComponentData, @@ -25,11 +26,6 @@ from snuba.utils.serializable_exception import JsonSerializable, SerializableException from snuba.web import QueryResult -if TYPE_CHECKING: - # Importing load_retriever at runtime pulls in snuba.web.rpc, which - # eventually imports db_query → AllocationPolicy (cycle). - from snuba.web.rpc.storage_routing.load_retriever import LoadInfo - IS_ENFORCED = "is_enforced" MAX_THREADS = "max_threads" NO_UNITS = "no_units" diff --git a/snuba/web/db_query.py b/snuba/web/db_query.py index d130f10b6b..2df6f9d7ea 100644 --- a/snuba/web/db_query.py +++ b/snuba/web/db_query.py @@ -26,6 +26,7 @@ from snuba.clickhouse.query import Query from snuba.clickhouse.query_dsl.accessors import get_time_range_estimate from snuba.clickhouse.query_profiler import generate_profile +from snuba.clusters.load_info import get_cluster_loadinfo from snuba.configs.configuration import ResourceIdentifier from snuba.datasets.storages.factory import get_storage from snuba.datasets.storages.storage_key import StorageKey @@ -77,7 +78,6 @@ SerializableExceptionDict, ) from snuba.web import QueryException, QueryResult, constants -from snuba.web.rpc.storage_routing.load_retriever import get_cluster_loadinfo metrics = MetricsWrapper(environment.metrics, "db_query") diff --git a/snuba/web/rpc/storage_routing/routing_strategies/storage_routing.py b/snuba/web/rpc/storage_routing/routing_strategies/storage_routing.py index 3a4ecf385d..74e2780379 100644 --- a/snuba/web/rpc/storage_routing/routing_strategies/storage_routing.py +++ b/snuba/web/rpc/storage_routing/routing_strategies/storage_routing.py @@ -26,6 +26,7 @@ from sentry_sdk import traces from snuba import environment, settings +from snuba.clusters.load_info import LoadInfo, get_cluster_loadinfo from snuba.configs.configuration import ( ConfigurableComponent, ConfigurableComponentData, @@ -54,7 +55,6 @@ from snuba.web.rpc.common.exceptions import RPCAllocationPolicyException from snuba.web.rpc.common.query_info import extract_query_info, extract_query_info_tags from snuba.web.rpc.storage_routing.common import extract_message_meta -from snuba.web.rpc.storage_routing.load_retriever import LoadInfo, get_cluster_loadinfo _SAMPLING_IN_STORAGE_PREFIX = "sampling_in_storage_" _START_ESTIMATION_MARK = "start_sampling_in_storage_estimation" diff --git a/tests/query/allocation_policies/test_allocation_policy_base.py b/tests/query/allocation_policies/test_allocation_policy_base.py index 1c861d94d9..22ddcb15a9 100644 --- a/tests/query/allocation_policies/test_allocation_policy_base.py +++ b/tests/query/allocation_policies/test_allocation_policy_base.py @@ -4,6 +4,7 @@ import pytest +from snuba.clusters.load_info import LoadInfo from snuba.configs.configuration import Configuration, InvalidConfig from snuba.datasets.storages.storage_key import StorageKey from snuba.query.allocation_policies import ( @@ -30,7 +31,6 @@ from snuba.query.allocation_policies.cross_org import CrossOrgQueryAllocationPolicy from snuba.query.allocation_policies.per_referrer import ReferrerGuardRailPolicy from snuba.utils.metrics.backends.testing import get_recorded_metric_calls -from snuba.web.rpc.storage_routing.load_retriever import LoadInfo from tests.configs.component_config import ( delete_component_config, set_component_config, diff --git a/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py b/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py index 9fd942d1cc..e01e196f54 100644 --- a/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py +++ b/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py @@ -5,7 +5,7 @@ from sentry_options import OptionValue from sentry_options.testing import override_options -from snuba.web.rpc.storage_routing.load_retriever import LoadInfo, get_cluster_loadinfo +from snuba.clusters.load_info import LoadInfo, get_cluster_loadinfo # Always present on CH. CGroupUserTimeNormalized / disk_inflight_ops are # host-dependent and may be -1 without failing the probe. @@ -48,6 +48,7 @@ def test_from_dict_ignores_unknown_keys() -> None: assert load_info.cluster_load == 2.0 assert load_info.concurrent_queries == -1 + ENABLE_LOADINFO: dict[str, OptionValue] = {"storage_routing.enable_get_cluster_loadinfo": True} @@ -69,7 +70,9 @@ def test_get_cluster_loadinfo_disabled() -> None: @pytest.mark.redis_db @pytest.mark.clickhouse_db def test_get_cluster_load() -> None: - _assert_probe_ok(get_cluster_loadinfo()) + load_info = get_cluster_loadinfo() + assert load_info is not None + _assert_probe_ok(load_info) @pytest.mark.redis_db @@ -93,7 +96,9 @@ def test_get_cluster_loadinfo_if_cache_fails() -> None: mock_redis.side_effect = Exception("Test error") with patch("snuba.redis.get_redis_client") as mock_redis_client: mock_redis_client.return_value = mock_redis - _assert_probe_ok(get_cluster_loadinfo()) + load_info = get_cluster_loadinfo() + assert load_info is not None + _assert_probe_ok(load_info) @pytest.mark.redis_db @@ -102,6 +107,7 @@ def test_get_cluster_load_error_handling() -> None: with patch("snuba.clickhouse.connect.ClickhouseConnectPool.execute") as mock_execute: mock_execute.side_effect = Exception("Test error") load_info = get_cluster_loadinfo() + assert load_info is not None _assert_probe_failed(load_info) assert load_info.is_idle() is False diff --git a/tests/web/rpc/v1/test_storage_routing.py b/tests/web/rpc/v1/test_storage_routing.py index b68859e7fe..b6b77ff3ea 100644 --- a/tests/web/rpc/v1/test_storage_routing.py +++ b/tests/web/rpc/v1/test_storage_routing.py @@ -13,6 +13,7 @@ from sentry_protos.snuba.v1.endpoint_time_series_pb2 import TimeSeriesRequest from sentry_protos.snuba.v1.request_common_pb2 import RequestMeta, TraceItemType +from snuba.clusters.load_info import LoadInfo from snuba.configs.configuration import Configuration, ResourceIdentifier from snuba.datasets.storages.storage_key import StorageKey from snuba.downsampled_storage_tiers import Tier @@ -35,7 +36,6 @@ encode_routing_hint, extract_message_meta, ) -from snuba.web.rpc.storage_routing.load_retriever import LoadInfo from snuba.web.rpc.storage_routing.routing_strategies.outcomes_based import ( OutcomesBasedRoutingStrategy, ) diff --git a/tests/web/test_db_query.py b/tests/web/test_db_query.py index 0d467867d6..bc80acab7d 100644 --- a/tests/web/test_db_query.py +++ b/tests/web/test_db_query.py @@ -12,6 +12,7 @@ from snuba.attribution.attribution_info import AttributionInfo from snuba.clickhouse.formatter.query import format_query from snuba.clickhouse.query import Query as ClickhouseQuery +from snuba.clusters.load_info import LoadInfo from snuba.configs.configuration import Configuration, ResourceIdentifier from snuba.datasets.schemas.tables import TableSource from snuba.datasets.storage import Storage @@ -41,7 +42,6 @@ db_query, execute_query, ) -from snuba.web.rpc.storage_routing.load_retriever import LoadInfo from tests.query.allocation_policies.attachment import ( match_block, override_allocation_policy, From a407e0643059e1112d7e1f4f85703173f356e7ce Mon Sep 17 00:00:00 2001 From: Prajjwal Bhandari Date: Fri, 2 Oct 2026 16:00:39 -0700 Subject: [PATCH 08/12] ref(allocation_policies): match on QuotaAllowanceDecision --- snuba/query/allocation_policies/__init__.py | 58 +++++++++++---------- 1 file changed, 31 insertions(+), 27 deletions(-) diff --git a/snuba/query/allocation_policies/__init__.py b/snuba/query/allocation_policies/__init__.py index 6f52bc8834..1fbc0c8c32 100644 --- a/snuba/query/allocation_policies/__init__.py +++ b/snuba/query/allocation_policies/__init__.py @@ -5,7 +5,7 @@ from abc import ABC, abstractmethod from dataclasses import asdict, dataclass, field from enum import Enum, IntEnum -from typing import Any, cast +from typing import Any, assert_never, cast from redis.exceptions import TimeoutError as RedisTimeoutError from sentry_sdk import traces @@ -491,16 +491,36 @@ def get_quota_allowance( decision = allowance.decision(self, load_info) referrer = str(tenant_ids.get("referrer", "no_referrer")) - if decision == QuotaAllowanceDecision.PARDONED: - self.metrics.increment("db_request_pardoned", tags={"referrer": referrer}) - elif decision == QuotaAllowanceDecision.REJECTED: - self.metrics.increment("db_request_rejected", tags={"referrer": referrer}) - elif decision == QuotaAllowanceDecision.THROTTLED: - self.metrics.increment( - "db_request_throttled", - tags={"referrer": referrer, "max_threads": str(allowance.max_threads)}, - ) - span.set_attribute("db_request_throttled", True) + match decision: + case QuotaAllowanceDecision.PARDONED: + self.metrics.increment("db_request_pardoned", tags={"referrer": referrer}) + allowance = QuotaAllowance( + can_run=True, + max_threads=self.max_threads, + explanation={ + **allowance.explanation, + "idle_pardon": load_info.to_dict() if load_info is not None else {}, + }, + is_throttled=allowance.is_throttled, + throttle_threshold=allowance.throttle_threshold, + rejection_threshold=allowance.rejection_threshold, + quota_used=allowance.quota_used, + quota_unit=allowance.quota_unit, + suggestion=allowance.suggestion, + max_bytes_to_read=0, + ) + case QuotaAllowanceDecision.REJECTED: + self.metrics.increment("db_request_rejected", tags={"referrer": referrer}) + case QuotaAllowanceDecision.THROTTLED: + self.metrics.increment( + "db_request_throttled", + tags={"referrer": referrer, "max_threads": str(allowance.max_threads)}, + ) + span.set_attribute("db_request_throttled", True) + case QuotaAllowanceDecision.ALLOWED: + pass + case unreachable: + assert_never(unreachable) if not self.is_enforced: allowance = QuotaAllowance( @@ -514,22 +534,6 @@ def get_quota_allowance( quota_unit=allowance.quota_unit, suggestion=allowance.suggestion, ) - elif decision == QuotaAllowanceDecision.PARDONED: - allowance = QuotaAllowance( - can_run=True, - max_threads=self.max_threads, - explanation={ - **allowance.explanation, - "idle_pardon": load_info.to_dict() if load_info is not None else {}, - }, - is_throttled=allowance.is_throttled, - throttle_threshold=allowance.throttle_threshold, - rejection_threshold=allowance.rejection_threshold, - quota_used=allowance.quota_used, - quota_unit=allowance.quota_unit, - suggestion=allowance.suggestion, - max_bytes_to_read=0, - ) # make sure we always know which storage key we rejected a query from allowance.explanation["storage_key"] = self._resource_identifier.value From 696639abfbb5f35d31e4a2dad938c57972babdca Mon Sep 17 00:00:00 2001 From: Prajjwal Bhandari Date: Fri, 2 Oct 2026 16:05:42 -0700 Subject: [PATCH 09/12] ref(allocation_policies): read is_pardonable from config like is_enforced --- snuba/clusters/load_info.py | 2 -- snuba/query/allocation_policies/__init__.py | 22 ++++++++++++------- .../test_allocation_policy_base.py | 18 +++++++++++++-- .../test_cluster_loadinfo.py | 1 - 4 files changed, 30 insertions(+), 13 deletions(-) diff --git a/snuba/clusters/load_info.py b/snuba/clusters/load_info.py index 8d8e7523cf..d3d837dc29 100644 --- a/snuba/clusters/load_info.py +++ b/snuba/clusters/load_info.py @@ -49,8 +49,6 @@ def from_dict(cls, load_info_dict: dict[str, float | int | None]) -> LoadInfo: ) def is_idle(self) -> bool: - if self.cluster_load == -1.0 or self.concurrent_queries == -1: - return False return get_option("storage_routing.enable_dynamic_allocation_policy", False) diff --git a/snuba/query/allocation_policies/__init__.py b/snuba/query/allocation_policies/__init__.py index 1fbc0c8c32..8bc2e67243 100644 --- a/snuba/query/allocation_policies/__init__.py +++ b/snuba/query/allocation_policies/__init__.py @@ -27,6 +27,7 @@ from snuba.web import QueryResult IS_ENFORCED = "is_enforced" +IS_PARDONABLE = "is_pardonable" MAX_THREADS = "max_threads" NO_UNITS = "no_units" NO_SUGGESTION = "no_suggestion" @@ -246,9 +247,11 @@ class AllocationPolicy(ConfigurableComponent, ABC): Any configuration definition that exists in your sub class' `_additional_config_definitions()` will appear in the Capacity Management Snuba Admin UI for the policy. From there you can modify the live values to alter how your policy works. - The base class comes with a built in config accessible as a property of the class itself: + The base class comes with built in configs accessible as properties of the class itself: - is_enforced - Use this to throttle/reject queries OR just log stuff. A configured policy is always active. + - is_pardonable + - When true, rejections from this policy can be pardoned if the cluster is idle. Eg. @@ -367,6 +370,12 @@ def __init__( value_type=int, default=kwargs.get(IS_ENFORCED, 1), ), + AllocationPolicyConfig( + name=IS_PARDONABLE, + description="Toggles whether rejections from this policy can be pardoned when the cluster is idle.", + value_type=int, + default=kwargs.get(IS_PARDONABLE, 1), + ), AllocationPolicyConfig( name=MAX_THREADS, description="The max threads Clickhouse can use for the query.", @@ -402,7 +411,7 @@ def is_enforced(self) -> bool: @property def is_pardonable(self) -> bool: - return True + return bool(self.get_config_value(IS_PARDONABLE)) @property def max_threads(self) -> int: @@ -453,6 +462,7 @@ def get_quota_allowance( for t, tid in tenant_ids.items(): span.set_attribute(f"tenant_ids.{t}", str(tid)) + passthrough = _default_passthough_policy(self._resource_identifier.value) try: allowance = self._get_quota_allowance(tenant_ids, query_id) except InvalidTenantsForAllocationPolicy as e: @@ -475,9 +485,7 @@ def get_quota_allowance( 1, tags={"method": "get_quota_allowance", "reason": type(e).__name__}, ) - return _default_passthough_policy( - self._resource_identifier.value - ).get_quota_allowance(tenant_ids, query_id, load_info) + return passthrough.get_quota_allowance(tenant_ids, query_id, load_info) except Exception: self.metrics.increment("fail_open", 1, tags={"method": "get_quota_allowance"}) logger.exception( @@ -485,9 +493,7 @@ def get_quota_allowance( ) if settings.RAISE_ON_ALLOCATION_POLICY_FAILURES: raise - return _default_passthough_policy( - self._resource_identifier.value - ).get_quota_allowance(tenant_ids, query_id, load_info) + return passthrough.get_quota_allowance(tenant_ids, query_id, load_info) decision = allowance.decision(self, load_info) referrer = str(tenant_ids.get("referrer", "no_referrer")) diff --git a/tests/query/allocation_policies/test_allocation_policy_base.py b/tests/query/allocation_policies/test_allocation_policy_base.py index 22ddcb15a9..54dcb949b5 100644 --- a/tests/query/allocation_policies/test_allocation_policy_base.py +++ b/tests/query/allocation_policies/test_allocation_policy_base.py @@ -309,6 +309,12 @@ def test_default_config_overrides(policy: AllocationPolicy) -> None: set_component_config(policy, config_key="is_enforced", value=1) assert policy.is_enforced == 1 + assert policy.is_pardonable == 1 + set_component_config(policy, config_key="is_pardonable", value=0) + assert policy.is_pardonable == 0 + set_component_config(policy, config_key="is_pardonable", value=1) + assert policy.is_pardonable == 1 + assert policy.max_threads == 10 set_component_config(policy, config_key="max_threads", value=4) assert policy.max_threads == 4 @@ -318,7 +324,7 @@ def test_default_config_overrides(policy: AllocationPolicy) -> None: @pytest.mark.redis_db def test_get_current_configs(policy: AllocationPolicy) -> None: - assert len(policy_configs := policy.get_current_configs()) == 3 + assert len(policy_configs := policy.get_current_configs()) == 4 assert all( config in policy_configs for config in [ @@ -338,6 +344,14 @@ def test_get_current_configs(policy: AllocationPolicy) -> None: "value": 1, "params": {}, }, + { + "name": "is_pardonable", + "type": "int", + "default": 1, + "description": "Toggles whether rejections from this policy can be pardoned when the cluster is idle.", + "value": 1, + "params": {}, + }, { "name": "max_threads", "type": "int", @@ -355,7 +369,7 @@ def test_get_current_configs(policy: AllocationPolicy) -> None: ) set_component_config(policy, config_key="is_enforced", value=0) set_component_config(policy, config_key="max_threads", value=4) - assert len(policy_configs := policy.get_current_configs()) == 4 + assert len(policy_configs := policy.get_current_configs()) == 5 assert { "name": "my_param_config", "type": "int", diff --git a/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py b/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py index e01e196f54..0b32505ee0 100644 --- a/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py +++ b/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py @@ -109,7 +109,6 @@ def test_get_cluster_load_error_handling() -> None: load_info = get_cluster_loadinfo() assert load_info is not None _assert_probe_failed(load_info) - assert load_info.is_idle() is False def test_is_idle_reads_dynamic_allocation_policy_flag() -> None: From 155432be42ee9178e60ff04bc7a8a9546b7b50fc Mon Sep 17 00:00:00 2001 From: Prajjwal Bhandari Date: Tue, 6 Oct 2026 11:24:57 -0400 Subject: [PATCH 10/12] ref(allocation_policies): rename LoadInfo.is_idle to should_pardon --- snuba/clusters/load_info.py | 2 +- snuba/query/allocation_policies/__init__.py | 4 ++-- .../allocation_policies/test_allocation_policy_base.py | 2 +- .../web/rpc/v1/routing_strategies/test_cluster_loadinfo.py | 6 +++--- tests/web/rpc/v1/test_storage_routing.py | 2 +- tests/web/test_db_query.py | 4 ++-- 6 files changed, 10 insertions(+), 10 deletions(-) diff --git a/snuba/clusters/load_info.py b/snuba/clusters/load_info.py index d3d837dc29..920c131b69 100644 --- a/snuba/clusters/load_info.py +++ b/snuba/clusters/load_info.py @@ -48,7 +48,7 @@ def from_dict(cls, load_info_dict: dict[str, float | int | None]) -> LoadInfo: } ) - def is_idle(self) -> bool: + def should_pardon(self) -> bool: return get_option("storage_routing.enable_dynamic_allocation_policy", False) diff --git a/snuba/query/allocation_policies/__init__.py b/snuba/query/allocation_policies/__init__.py index 8bc2e67243..928831a738 100644 --- a/snuba/query/allocation_policies/__init__.py +++ b/snuba/query/allocation_policies/__init__.py @@ -93,8 +93,8 @@ def decision( load_info: LoadInfo | None, ) -> QuotaAllowanceDecision: if not self.can_run: - is_idle_load = getattr(load_info, "is_idle", lambda: False) - idle_pardon = policy.is_pardonable and is_idle_load() + should_pardon = getattr(load_info, "should_pardon", lambda: False) + idle_pardon = policy.is_pardonable and should_pardon() return ( QuotaAllowanceDecision.PARDONED if idle_pardon else QuotaAllowanceDecision.REJECTED ) diff --git a/tests/query/allocation_policies/test_allocation_policy_base.py b/tests/query/allocation_policies/test_allocation_policy_base.py index 54dcb949b5..90ff9fb145 100644 --- a/tests/query/allocation_policies/test_allocation_policy_base.py +++ b/tests/query/allocation_policies/test_allocation_policy_base.py @@ -474,7 +474,7 @@ def test_is_pardonable_overrides_can_run_when_idle() -> None: "referrer": "some_referrer", } idle = LoadInfo(cluster_load=1.0, concurrent_queries=1) - with mock.patch.object(LoadInfo, "is_idle", return_value=True): + with mock.patch.object(LoadInfo, "should_pardon", return_value=True): assert reject_policy.get_quota_allowance(tenant_ids, "deadbeef").can_run is False pardoned = reject_policy.get_quota_allowance(tenant_ids, "deadbeef", idle) assert pardoned.can_run is True diff --git a/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py b/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py index 0b32505ee0..958df3a1a2 100644 --- a/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py +++ b/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py @@ -111,8 +111,8 @@ def test_get_cluster_load_error_handling() -> None: _assert_probe_failed(load_info) -def test_is_idle_reads_dynamic_allocation_policy_flag() -> None: +def test_should_pardon_reads_dynamic_allocation_policy_flag() -> None: load_info = LoadInfo(cluster_load=1.0, concurrent_queries=1) - assert load_info.is_idle() is False + assert load_info.should_pardon() is False with override_options("snuba", {"storage_routing.enable_dynamic_allocation_policy": True}): - assert load_info.is_idle() is True + assert load_info.should_pardon() is True diff --git a/tests/web/rpc/v1/test_storage_routing.py b/tests/web/rpc/v1/test_storage_routing.py index b6b77ff3ea..cacbba1bbd 100644 --- a/tests/web/rpc/v1/test_storage_routing.py +++ b/tests/web/rpc/v1/test_storage_routing.py @@ -599,7 +599,7 @@ def _update_quota_balance( "snuba.web.rpc.storage_routing.routing_strategies.storage_routing.get_cluster_loadinfo", return_value=LoadInfo(cluster_load=1.0, concurrent_queries=1), ), - mock.patch.object(LoadInfo, "is_idle", return_value=True), + mock.patch.object(LoadInfo, "should_pardon", return_value=True), ): decision = OutcomesBasedRoutingStrategy().get_routing_decision( RoutingContext( diff --git a/tests/web/test_db_query.py b/tests/web/test_db_query.py index bc80acab7d..2448f6d372 100644 --- a/tests/web/test_db_query.py +++ b/tests/web/test_db_query.py @@ -1405,7 +1405,7 @@ def test_idle_pardon_allows_rejected_query() -> None: "snuba.web.db_query.get_cluster_loadinfo", return_value=idle, ), - mock.patch.object(LoadInfo, "is_idle", return_value=True), + mock.patch.object(LoadInfo, "should_pardon", return_value=True), ): query_settings = HTTPQuerySettings() _apply_allocation_policies_quota( @@ -1440,7 +1440,7 @@ def test_idle_pardon_still_rejects_when_not_idle() -> None: "snuba.web.db_query.get_cluster_loadinfo", return_value=busy, ), - mock.patch.object(LoadInfo, "is_idle", return_value=False), + mock.patch.object(LoadInfo, "should_pardon", return_value=False), pytest.raises(AllocationPolicyViolations), ): _apply_allocation_policies_quota( From a03c7421ac6adfd07e06ef406a558243e99aeff1 Mon Sep 17 00:00:00 2001 From: Prajjwal Bhandari Date: Tue, 6 Oct 2026 11:25:02 -0400 Subject: [PATCH 11/12] Fix allocation-policy tests broken by is_pardonable config --- snuba/web/db_query.py | 12 +++++++----- .../test_allocation_policy_base.py | 2 +- 2 files changed, 8 insertions(+), 6 deletions(-) diff --git a/snuba/web/db_query.py b/snuba/web/db_query.py index 2df6f9d7ea..a959b130b3 100644 --- a/snuba/web/db_query.py +++ b/snuba/web/db_query.py @@ -880,11 +880,13 @@ def _apply_allocation_policies_quota( rejection_quota_and_policy = None throttle_quota_and_policy = None min_threads_across_policies = MAX_THRESHOLD - load_info = get_cluster_loadinfo( - get_storage( - StorageKey(allocation_policies[0].resource_identifier.value) - ).get_storage_set_key() - ) + load_info = None + if get_option("storage_routing.enable_get_cluster_loadinfo", False): + load_info = get_cluster_loadinfo( + get_storage( + StorageKey(allocation_policies[0].resource_identifier.value) + ).get_storage_set_key() + ) with traces.start_span( name="_apply_allocation_policies_quota", attributes={SENTRY_OP: "allocation_policy"}, diff --git a/tests/query/allocation_policies/test_allocation_policy_base.py b/tests/query/allocation_policies/test_allocation_policy_base.py index 90ff9fb145..9dc9d2df9e 100644 --- a/tests/query/allocation_policies/test_allocation_policy_base.py +++ b/tests/query/allocation_policies/test_allocation_policy_base.py @@ -229,7 +229,7 @@ def test_bad_config_key_in_option(self) -> None: configs = policy.get_current_configs() # the bad configs are not returned - assert len(configs) == 3 + assert len(configs) == 4 # the bad configs are logged assert len(captured.records) == 3 From 77075b571d6dd73946cd7aae28675bb61e6fa878 Mon Sep 17 00:00:00 2001 From: Prajjwal Bhandari Date: Tue, 6 Oct 2026 11:52:57 -0400 Subject: [PATCH 12/12] Fix test_simple_config_values for is_pardonable config --- .../test_bytes_scanned_window_allocation_policy.py | 1 + 1 file changed, 1 insertion(+) diff --git a/tests/query/allocation_policies/test_bytes_scanned_window_allocation_policy.py b/tests/query/allocation_policies/test_bytes_scanned_window_allocation_policy.py index 80f90684d7..dcbc91fa38 100644 --- a/tests/query/allocation_policies/test_bytes_scanned_window_allocation_policy.py +++ b/tests/query/allocation_policies/test_bytes_scanned_window_allocation_policy.py @@ -159,6 +159,7 @@ def test_simple_config_values(policy: AllocationPolicy) -> None: "org_limit_bytes_scanned_override", "throttled_thread_number", "is_enforced", + "is_pardonable", "max_threads", } assert policy.get_config_value("org_limit_bytes_scanned") == ORG_SCAN_LIMIT