diff --git a/sentry-options/schemas/snuba/schema.json b/sentry-options/schemas/snuba/schema.json index dcec5f5e46..c3c42a920a 100644 --- a/sentry-options/schemas/snuba/schema.json +++ b/sentry-options/schemas/snuba/schema.json @@ -274,6 +274,11 @@ "default": false, "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", "properties": { diff --git a/snuba/web/rpc/storage_routing/load_retriever.py b/snuba/clusters/load_info.py similarity index 93% rename from snuba/web/rpc/storage_routing/load_retriever.py rename to snuba/clusters/load_info.py index e9af18700c..d3d837dc29 100644 --- a/snuba/web/rpc/storage_routing/load_retriever.py +++ b/snuba/clusters/load_info.py @@ -15,11 +15,12 @@ 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( environment.metrics, - "snuba.web.rpc.storage_routing.load_retriever", + "snuba.clusters.load_info", ) @@ -47,6 +48,9 @@ def from_dict(cls, load_info_dict: dict[str, float | int | None]) -> LoadInfo: } ) + def is_idle(self) -> bool: + return get_option("storage_routing.enable_dynamic_allocation_policy", False) + def cache( ttl_secs: int = 60, @@ -86,9 +90,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/query/allocation_policies/__init__.py b/snuba/query/allocation_policies/__init__.py index dbcb7783aa..8bc2e67243 100644 --- a/snuba/query/allocation_policies/__init__.py +++ b/snuba/query/allocation_policies/__init__.py @@ -4,13 +4,14 @@ import os from abc import ABC, abstractmethod from dataclasses import asdict, dataclass, field -from enum import Enum -from typing import Any, cast +from enum import Enum, IntEnum +from typing import Any, assert_never, 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, @@ -26,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" @@ -54,6 +56,13 @@ class AllocationPolicyConfig(Configuration): pass +class QuotaAllowanceDecision(IntEnum): + REJECTED = 0 + PARDONED = 1 + ALLOWED = 2 + THROTTLED = 3 + + @dataclass(frozen=True) class QuotaAllowance: can_run: bool @@ -78,6 +87,23 @@ class QuotaAllowance: def to_dict(self) -> dict[str, Any]: return asdict(self) + def decision( + self, + policy: AllocationPolicy, + 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() + 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: if not isinstance(other, QuotaAllowance): return False @@ -221,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. @@ -342,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.", @@ -375,6 +409,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 bool(self.get_config_value(IS_PARDONABLE)) + @property def max_threads(self) -> int: """Maximum number of threads run a single query on ClickHouse with.""" @@ -412,7 +450,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__, @@ -420,6 +461,8 @@ def get_quota_allowance( ) as span: 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: @@ -442,7 +485,7 @@ 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 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( @@ -450,23 +493,41 @@ 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"))}, - ) - elif allowance.max_threads < self.max_threads: - # NOTE: The elif is very intentional here. Don't count the throttling - # if the request was rejected. - self.metrics.increment( - "db_request_throttled", - tags={ - "referrer": str(tenant_ids.get("referrer", "no_referrer")), - "max_threads": str(allowance.max_threads), - }, - ) - span.set_attribute("db_request_throttled", True) + 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")) + 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( can_run=True, @@ -479,6 +540,7 @@ def get_quota_allowance( quota_unit=allowance.quota_unit, suggestion=allowance.suggestion, ) + # 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..2df6f9d7ea 100644 --- a/snuba/web/db_query.py +++ b/snuba/web/db_query.py @@ -26,7 +26,9 @@ 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 from snuba.downsampled_storage_tiers import Tier from snuba.query import ProcessableQuery @@ -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/routing_strategies/storage_routing.py b/snuba/web/rpc/storage_routing/routing_strategies/storage_routing.py index 05e0e51f11..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" @@ -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..54dcb949b5 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 ( @@ -308,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 @@ -317,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 [ @@ -337,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", @@ -354,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", @@ -448,6 +463,37 @@ 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: + 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 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 == reject_policy.max_threads + assert pardoned.max_bytes_to_read == 0 + assert pardoned.explanation["idle_pardon"] == idle.to_dict() + + 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..0b32505ee0 100644 --- a/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py +++ b/tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py @@ -1,8 +1,11 @@ +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 +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. @@ -46,10 +49,30 @@ def test_from_dict_ignores_unknown_keys() -> None: assert load_info.concurrent_queries == -1 +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: + with override_options("snuba", {"storage_routing.enable_get_cluster_loadinfo": False}): + assert get_cluster_loadinfo() is None + with override_options("snuba", {"storage_routing.enable_get_cluster_loadinfo": True}): + assert get_cluster_loadinfo() is not 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 @@ -58,9 +81,11 @@ def test_get_cluster_load_from_cache() -> None: with 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() @@ -71,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 @@ -79,4 +106,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 load_info is not None + _assert_probe_failed(load_info) + + +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 diff --git a/tests/web/rpc/v1/test_storage_routing.py b/tests/web/rpc/v1/test_storage_routing.py index ee805484a0..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 @@ -553,6 +554,71 @@ def _update_quota_balance( assert not exc.details["can_run"] +@pytest.mark.redis_db +def test_routing_strategy_idle_pardon_allows_rejected_query() -> None: + 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}, + ), + 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( + in_msg=_get_in_msg(), + timer=Timer("test"), + query_id="abc", + ) + ) + assert decision.can_run is True + assert decision.clickhouse_settings["max_threads"] == 10 + pardoned = decision.routing_context.allocation_policies_recommendations[ + "IdlePardonRejectionPolicy" + ] + assert ( + pardoned.explanation["idle_pardon"] + == LoadInfo(cluster_load=1.0, concurrent_queries=1).to_dict() + ) + + @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..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 @@ -1360,3 +1361,94 @@ 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: + 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}, + ), + mock.patch( + "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( + 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"] == 10 + assert details["explanation"]["idle_pardon"] == idle.to_dict() + quota = query_settings.get_resource_quota() + assert quota is not None + assert quota.max_threads == 10 + + +def test_idle_pardon_still_rejects_when_not_idle() -> None: + 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}, + ), + mock.patch( + "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( + 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