sadpandajoe commented on code in PR #44849: URL: https://github.com/apache/superset/pull/44849#discussion_r4161350529
########## superset/semantic_layers/metadata.py: ########## @@ -0,0 +1,442 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +"""Scoped publication and invalidation of provider-owned metadata.""" + +from __future__ import annotations + +import hashlib +import hmac +import math +import time +from collections.abc import Callable +from dataclasses import dataclass, field +from datetime import datetime, timezone +from typing import Literal, Protocol, TYPE_CHECKING +from uuid import uuid4 + +from redis.exceptions import RedisError +from superset_core.semantic_layers.metadata import ( + CatalogLoader, + CatalogSnapshot, + MetadataRefreshError, + MetadataRefreshResult, + remaining_budget, +) + +from superset.semantic_layers.cache_inspection import CacheEntryInfo, describe_entry +from superset.utils import json + +CATALOG_TTL_SECONDS: int = 300 +REFRESH_LEASE_SECONDS: int = 60 +FETCH_DEADLINE_SECONDS: int = 30 +MAX_CATALOG_BYTES: int = 10 * 1024 * 1024 +SNAPSHOT_FORMAT_VERSION: int = 2 +READER_POLL_SECONDS: float = 0.05 + + +if TYPE_CHECKING: + + class PublicationBackend(Protocol): + """The shared coordinator operations used by semantic metadata.""" + + def with_deadline(self, deadline: float) -> PublicationBackend: ... + def get(self, name: str) -> bytes | None: ... + def set( + self, + name: str, + value: str, + ex: int | None = None, + px: int | None = None, + nx: bool = False, + xx: bool = False, + ) -> bool | None: ... + def delete(self, *names: str) -> int: ... + def compare_and_delete(self, name: str, expected: str) -> int: ... + def compare_and_publish( + self, + lease_key: str, + expected: str, + snapshot_key: str, + value: str, + ttl_ms: int, + ) -> bool: ... + def get_with_ttl(self, name: str) -> tuple[bytes | None, int]: ... + def get_or_create(self, name: str, value: str, ttl: int) -> bytes: ... + + +def metadata_scope( + secret: str, namespace: str, connection_uuid: str, configuration: str +) -> str: + """Derive a private identity from trusted deployment, tenant and connection data.""" + if not secret or not namespace or not connection_uuid: + raise MetadataRefreshError("configuration") + try: + canonical: str = json.dumps( + [namespace, connection_uuid, json.loads(configuration)], + sort_keys=True, + separators=(",", ":"), + allow_nan=False, + ) + except (TypeError, ValueError): + raise MetadataRefreshError("configuration") from None + return hmac.new(secret.encode(), canonical.encode(), hashlib.sha256).hexdigest() + + +@dataclass(frozen=True) +class StoredCatalog: + """Internal envelope; publication bookkeeping is never a provider revision.""" + + snapshot: CatalogSnapshot + digest: str = field(repr=False) + attempt: str = field(repr=False) + created_at: str + + +class ScopedMetadataStore: + """One shared observation with a request-wide budget and no local fallback.""" + + def __init__( + self, + backend: PublicationBackend, + scope: str, + *, + deadline: float, + before_publish: Callable[[], None] | None = None, + clock: Callable[[], float] = time.monotonic, + wait: Callable[[float], None] = time.sleep, + ) -> None: + if not math.isfinite(deadline) or not scope or "{" in scope or "}" in scope: + raise MetadataRefreshError("configuration") + self._backend: PublicationBackend = backend + self._scope: str = scope + self._deadline: float = deadline + self._lease_key: str = f"semantic-metadata:{{{scope}}}:lease" + self._snapshot_key: str = f"semantic-metadata:{{{scope}}}:snapshot" + self._generation_key: str = f"semantic-metadata:{{{scope}}}:compatibility" + self._before_publish: Callable[[], None] | None = before_publish + self._clock: Callable[[], float] = clock + self._wait: Callable[[float], None] = wait + self._observations: dict[str, str] = {} + + def _remaining(self) -> float: + return remaining_budget(self._deadline, now=self._clock()) + + def _decode(self, raw: bytes | None) -> StoredCatalog | None: + if raw is None or len(raw) > MAX_CATALOG_BYTES: + return None + try: + envelope: object = json.loads(raw) + if ( + not isinstance(envelope, dict) + or envelope.get("version") != SNAPSHOT_FORMAT_VERSION + ): + return None + if any( + not isinstance(envelope.get(key), str) + for key in ( + "payload", + "cache_token", + "observed_at", + "digest", + "attempt", + "created_at", + ) + ): + return None + if any( + not envelope[key] + for key in ( + "cache_token", + "observed_at", + "digest", + "attempt", + "created_at", + ) + ) or not envelope["cache_token"].startswith(f"{self._scope}:"): + return None + payload: str = envelope["payload"] + json.dumps(json.loads(payload), allow_nan=False) + if hashlib.sha256(payload.encode()).hexdigest() != envelope["digest"]: + return None + return StoredCatalog( + CatalogSnapshot( + payload, envelope["cache_token"], envelope["observed_at"] + ), + envelope["digest"], + envelope["attempt"], + envelope["created_at"], + ) + except (ValueError, UnicodeError, RecursionError): + return None + + def _load(self) -> StoredCatalog | None: + self._remaining() + stored: StoredCatalog | None = self._decode( + self._backend.get(self._snapshot_key) + ) + self._remaining() + return stored + + def _remember(self, snapshot: CatalogSnapshot) -> CatalogSnapshot: + self._observations[snapshot.cache_token] = snapshot.observed_at + return snapshot + + def observed_at(self, token: str) -> str | None: + """Read the timestamp captured with a provider's token, without backend I/O.""" + return self._observations.get(token) + + def peek(self) -> CatalogSnapshot | None: + """Read the current observation without acquiring, filling or renewing it.""" + try: + stored: StoredCatalog | None = self._load() + except RedisError: + raise MetadataRefreshError("unavailable") from None + return stored.snapshot if stored is not None else None + + def _for_deadline(self, deadline: float) -> ScopedMetadataStore: + """Narrow one call without mutating the operation or another call's budget.""" + remaining_budget(deadline, now=self._clock()) + if deadline > self._deadline: + raise MetadataRefreshError("deadline") + scoped: ScopedMetadataStore = ScopedMetadataStore( + self._backend.with_deadline(deadline), + self._scope, + deadline=deadline, + before_publish=self._before_publish, + clock=self._clock, + wait=self._wait, + ) + scoped._observations = self._observations + return scoped + + def read(self, fetch: CatalogLoader, *, deadline: float) -> CatalogSnapshot: + """Honor the explicit caller budget, including cache hits and transport.""" + return self._for_deadline(deadline)._read(fetch) + + def refresh( + self, fetch: CatalogLoader, *, deadline: float + ) -> MetadataRefreshResult: + """Publish within the caller budget, which cannot extend the host operation.""" + return self._for_deadline(deadline)._refresh(fetch) + + def _read(self, fetch: CatalogLoader) -> CatalogSnapshot: + """Wait for an owner or acquire once using the same remaining request budget.""" + try: + while True: + current: StoredCatalog | None = self._load() + if current is not None: + return self._remember(current.snapshot) + attempt: str = uuid4().hex + self._remaining() + if self._backend.set( + self._lease_key, + attempt, + px=max( + 1, + math.ceil(min(REFRESH_LEASE_SECONDS, self._remaining()) * 1000), + ), + nx=True, + ): + try: + # Another owner may have published between our read and SET NX. + current = self._load() + if current is not None: + return self._remember(current.snapshot) + return self._acquire(fetch, attempt).snapshot + finally: + self._release(attempt) + self._wait(min(READER_POLL_SECONDS, self._remaining())) + except RedisError: + self._remaining() + raise MetadataRefreshError("unavailable") from None + + def _refresh(self, fetch: CatalogLoader) -> MetadataRefreshResult: + """Publish a new observation, or report explicit contention without retry.""" + self._remaining() + attempt: str = uuid4().hex + try: + if not self._backend.set( + self._lease_key, + attempt, + px=max( + 1, math.ceil(min(REFRESH_LEASE_SECONDS, self._remaining()) * 1000) + ), + nx=True, + ): + raise MetadataRefreshError("in_progress") + try: + return self._acquire(fetch, attempt) + finally: + self._release(attempt) + except RedisError: + self._remaining() + raise MetadataRefreshError("unavailable") from None + + def _acquire(self, fetch: CatalogLoader, attempt: str) -> MetadataRefreshResult: + started: float = self._clock() + previous: StoredCatalog | None = self._load() + self._remaining() + try: + payload: str = fetch(self._deadline) + except MetadataRefreshError: + raise + except Exception: # pylint: disable=broad-except + # Provider failures cannot transport vendor payloads into host errors. + raise MetadataRefreshError("upstream") from None + self._remaining() + if not isinstance(payload, str): + raise MetadataRefreshError("invalid_payload") + try: + if len(payload.encode()) > MAX_CATALOG_BYTES: + raise MetadataRefreshError("invalid_payload") + payload = json.dumps( + json.loads(payload), Review Comment: This round-trip changes provider-owned numbers: `0.12345678901234567890123456789` becomes `0.12345678901234568`, and `1e400` becomes `null`, so a provider preserving decimal-valued definitions receives altered catalog data on refresh and warm reads. Could we validate or normalize the payload without losing numeric values? ########## tests/unit_tests/coordination/test_deadline_backend.py: ########## @@ -0,0 +1,159 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +"""The native Redis transport must cancel within the original remaining budget.""" + +from __future__ import annotations + +import asyncio +import time +from unittest.mock import AsyncMock, Mock, patch + +import pytest +from redis.exceptions import TimeoutError as RedisTimeoutError + +from superset.coordination.deadline_backend import DeadlineRedisBackend + + +def test_deadline_includes_waiting_for_the_transport() -> None: + async def slow(*args: object, **kwargs: object) -> None: + await asyncio.sleep(2) + + started: float = time.monotonic() + backend: DeadlineRedisBackend = DeadlineRedisBackend( + {"CACHE_TYPE": "RedisCache"}, + deadline=started + 0.05, + ) + with patch("redis.asyncio.Redis.execute_command", slow): + with pytest.raises(RedisTimeoutError): + backend.get("owned-key") + assert time.monotonic() - started < 0.5 + + +def test_expired_budget_never_opens_a_connection() -> None: + backend: DeadlineRedisBackend = DeadlineRedisBackend( + {"CACHE_TYPE": "RedisCache"}, + deadline=time.monotonic() - 1, + ) + client: Mock + with patch("redis.asyncio.Redis") as client: Review Comment: The backend imports `Redis` into its own module, so this patch observes an unused binding and `assert_not_called()` still passes if the real constructor is invoked before the expected rejection. Could this test and the async-caller test patch `superset.coordination.deadline_backend.Redis` so they actually detect premature client construction? ########## superset/semantic_layers/models.py: ########## @@ -355,8 +366,15 @@ def after_delete( security_manager.semantic_view_after_delete(mapper, connection, target) - @cached_property + @property def implementation(self) -> SemanticViewABC: + if metadata_binding.participates(self.semantic_layer): + # Preserve canonical chart/dashboard/guest policy at the caller. + return metadata_binding.view_implementation(self) Review Comment: With refresh enabled, unavailable Redis or an expired discovery deadline now raises `MetadataRefreshError` here, but datasource metadata/query and Explore still lack the mapping used by `compatible()` and the semantic-layer endpoints. Those authorized requests return generic 500s instead of the intended 503/504 category, leaving clients unable to identify the failure; could we map typed metadata errors at these remaining HTTP boundaries too? ########## superset/coordination/deadline_backend.py: ########## @@ -0,0 +1,228 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +"""Private, cancellation-bounded Redis operations for synchronous metadata callers. + +Each command owns its async client and loop. Cancelling the command disconnects +its socket; shared coordinator pools and their retry/timeout policy are untouched. +""" + +from __future__ import annotations + +import asyncio +import math +import time +from contextlib import AsyncExitStack +from typing import Any + +from redis.asyncio import Redis +from redis.asyncio.retry import Retry +from redis.asyncio.sentinel import Sentinel +from redis.backoff import NoBackoff +from redis.exceptions import RedisError, TimeoutError as RedisTimeoutError +from superset_core.semantic_layers.metadata import ( + MetadataRefreshError, + remaining_budget, +) + +from superset.coordination.cache_backend import _COMPARE_AND_DELETE_LUA + +_COMPARE_AND_PUBLISH_LUA: str = """ +if redis.call('get', KEYS[1]) ~= ARGV[1] then + return 0 +end +redis.call('psetex', KEYS[2], ARGV[3], ARGV[2]) +redis.call('del', KEYS[1]) +return 1 +""" + +_GET_WITH_TTL_LUA: str = """ +return {redis.call('get', KEYS[1]), redis.call('pttl', KEYS[1])} +""" + +_GET_OR_CREATE_LUA: str = """ +local value = redis.call('get', KEYS[1]) +if value then + return value +end +redis.call('set', KEYS[1], ARGV[1], 'EX', ARGV[2]) +return ARGV[1] +""" + + +class DeadlineRedisBackend: + """Use the coordinator configuration without sharing mutable connections.""" + + def __init__(self, config: dict[str, Any], *, deadline: float) -> None: + if not math.isfinite(deadline) or config.get("CACHE_TYPE") not in { + "RedisCache", + "RedisSentinelCache", + }: + raise ValueError("Unsupported metadata coordination configuration") + self._config: dict[str, Any] = dict(config) + self._deadline: float = deadline + + def with_deadline(self, deadline: float) -> DeadlineRedisBackend: + """Create a private call budget without extending the operation ceiling.""" + if not math.isfinite(deadline): + raise ValueError("Metadata deadline must be finite") + return DeadlineRedisBackend( + self._config, deadline=min(self._deadline, deadline) + ) + + def _remaining(self) -> float: + """Keep the transport's Redis error boundary while sharing SDK validation.""" + try: + return remaining_budget(self._deadline, now=time.monotonic()) + except MetadataRefreshError: + raise RedisTimeoutError("Metadata deadline invalid or expired") from None + + async def _command(self, *args: str | int) -> Any: + remaining: float = self._remaining() + options: dict[str, Any] = { + "db": self._config.get("CACHE_REDIS_DB", 0), + "username": self._config.get("CACHE_REDIS_USER"), + "password": self._config.get("CACHE_REDIS_PASSWORD"), + "socket_timeout": remaining, + "socket_connect_timeout": remaining, + "retry": Retry(NoBackoff(), 0), + "protocol": 2, + } + if self._config.get("CACHE_REDIS_SSL", False): + options.update( + { + "ssl": True, + "ssl_certfile": self._config.get("CACHE_REDIS_SSL_CERTFILE"), + "ssl_keyfile": self._config.get("CACHE_REDIS_SSL_KEYFILE"), + "ssl_ca_certs": self._config.get("CACHE_REDIS_SSL_CA_CERTS"), + "ssl_cert_reqs": self._config.get( + "CACHE_REDIS_SSL_CERT_REQS", "required" + ), + } + ) + # One cancellation deadline covers DNS, Sentinel discovery, authentication, + # response parsing (including trickled responses) and connection cleanup. + stack: AsyncExitStack + async with asyncio.timeout(remaining), AsyncExitStack() as stack: + client: Redis + if self._config["CACHE_TYPE"] == "RedisSentinelCache": + sentinel: Sentinel = Sentinel( + self._config.get("CACHE_REDIS_SENTINELS", [("127.0.0.1", 26379)]), + sentinel_kwargs={ + "password": self._config.get("CACHE_REDIS_SENTINEL_PASSWORD"), + "socket_timeout": remaining, + "socket_connect_timeout": remaining, Review Comment: Giving the first Sentinel the entire remaining budget prevents fallback when that node blackholes connections: the outer deadline cancels discovery before redis-py can try a healthy second node, and each command recreates the same ordering. Could we preserve shorter configured per-node timeouts, clamped to the operation budget, so Sentinel redundancy still works? -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
