Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion sentry-options/schemas/snuba/schema.json
Original file line number Diff line number Diff line change
Expand Up @@ -264,7 +264,7 @@
"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."

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Idle pardon options missing from schema

High Severity

is_idle reads storage_routing.idle_cluster_load_threshold and storage_routing.idle_concurrent_queries_threshold, but those keys are not declared in the snuba sentry-options schema. get_option then always falls back to 0, so cluster_load < 0 and concurrent_queries < 0 never hold for real loadinfo, and the idle pardon path cannot be enabled.

Additional Locations (1)
Fix in Cursor Fix in Web

Reviewed by Cursor Bugbot for commit ae1a2d8. Configure here.

},
"retention_days": {
"type": "object",
Expand Down
90 changes: 73 additions & 17 deletions snuba/query/allocation_policies/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,8 +4,8 @@
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 TYPE_CHECKING, Any, cast

from redis.exceptions import TimeoutError as RedisTimeoutError
from sentry_sdk import traces
Expand All @@ -25,6 +25,11 @@
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"
Expand Down Expand Up @@ -54,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
Expand All @@ -78,6 +90,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
Expand Down Expand Up @@ -375,6 +404,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."""
Expand Down Expand Up @@ -412,14 +445,18 @@ 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__,
attributes={SENTRY_OP: "allocation_policy.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:
Expand All @@ -442,31 +479,33 @@ 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_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(
"Allocation policy failed to get quota allowance, this is a bug, fix it"
)
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.
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:
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:
allowance = QuotaAllowance(
can_run=True,
Expand All @@ -479,6 +518,23 @@ 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
for k, v in allowance.to_dict().items():
Expand Down
12 changes: 10 additions & 2 deletions snuba/web/db_query.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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")

Expand Down Expand Up @@ -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(
Expand Down
22 changes: 21 additions & 1 deletion snuba/web/rpc/storage_routing/load_retriever.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,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(
Expand Down Expand Up @@ -39,6 +40,17 @@ def from_dict(cls, load_info_dict: dict[str, float | int]) -> "LoadInfo":
concurrent_queries=int(load_info_dict["concurrent_queries"]),
)

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,
Expand Down Expand Up @@ -78,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:
try:
cluster = get_cluster(storage_set_key)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand All @@ -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)
)
Expand All @@ -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(
Expand Down
35 changes: 35 additions & 0 deletions tests/query/allocation_policies/test_allocation_policy_base.py
Original file line number Diff line number Diff line change
Expand Up @@ -30,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,
Expand Down Expand Up @@ -448,6 +449,40 @@ 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"] == {
"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."""

Expand Down
22 changes: 22 additions & 0 deletions tests/web/rpc/v1/routing_strategies/test_cluster_loadinfo.py
Original file line number Diff line number Diff line change
@@ -1,9 +1,29 @@
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 get_cluster_loadinfo

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
Expand All @@ -20,9 +40,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()


Expand Down
Loading
Loading