mikebridge commented on code in PR #44835:
URL: https://github.com/apache/superset/pull/44835#discussion_r4190806730


##########
superset/tasks/async_queries.py:
##########
@@ -178,30 +181,63 @@ def _inject_contribution_totals(
     cached dataframe and inject the sums into this query's contribution
     post-processing before it runs — the same result the synchronous
     ``ensure_totals_available`` produces, but reading the cache the 
prerequisite
-    populated instead of re-running the totals query. ``contribution_totals`` 
is
+    populated. Participating semantic queries resolve the totals context under
+    the dependent operation so an obsolete denominator is recomputed.
+    ``contribution_totals`` is
     stripped from the cache key, so this affects only the result, not the key.
     """
     from superset.common.query_context_processor import is_summable
     from superset.common.utils.query_cache_manager import QueryCacheManager
 
-    cache = QueryCacheManager.get(key=totals_cache_key, 
region=CacheRegion.DATA)
-    if not cache.is_loaded or cache.df is None:
-        # The depends_on prerequisite guarantees the totals task succeeded and 
wrote
-        # this cache entry, so a miss is unexpected (e.g. it was evicted 
between the
-        # totals task finishing and this task reading). Fail loudly rather than
-        # caching a silently un-normalized result the client would then 
re-request:
-        # this task's single query cannot reproduce the synchronous path's
-        # ensure_totals_available (it has no totals query to run).
-        raise SupersetException(
-            f"Contribution totals not found in cache under {totals_cache_key}"
+    df: pd.DataFrame
+    if totals_context is not None:
+        # The ordinary key captures this operation's catalog. A matching totals
+        # entry is reused; a refreshed catalog recomputes before normalization.
+        totals_context.force = (
+            totals_context.force
+            and _query_task_cache_key(totals_context, 0) != totals_cache_key
         )
-    df = cache.df
-    totals = {col: df[col].sum() for col in df.columns if is_summable(df[col])}
+        totals_context.is_async_execution = True
+        with _capture_query_cancellation(totals_context):
+            df = totals_context.get_df_payload_result(
+                totals_context.queries[0]
+            ).payload["df"]

Review Comment:
   Confirmed. Commit 9aed7aafd7 checks acquisition status and error as well as 
the dataframe before injecting totals. A failed recomputation stops the 
dependent query before execution or caching; the regression includes a failed 
payload with an empty dataframe.



##########
superset/semantic_layers/metadata_cache.py:
##########
@@ -0,0 +1,109 @@
+# 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.
+
+"""Captured catalog identity for derived caches and read-only inspection."""
+
+from __future__ import annotations
+
+import hashlib
+from dataclasses import dataclass, field
+from typing import TYPE_CHECKING
+from uuid import uuid4
+
+from superset_core.semantic_layers.metadata import CatalogSnapshot, 
MetadataRefreshError
+
+from superset.semantic_layers import metadata_binding
+from superset.semantic_layers.metadata import ScopedMetadataStore
+from superset.semantic_layers.metadata_binding import connection_store
+from superset.utils import json
+
+if TYPE_CHECKING:
+    from superset.semantic_layers.models import SemanticView
+
+
+@dataclass(frozen=True)
+class CompatibilityIdentity:
+    key: str = field(repr=False)
+    source_observed_at: str | None
+
+
+def view_cache_token(view: SemanticView, token: str) -> str:
+    """Include host view configuration without revealing it in cache keys."""
+    identity: str = json.dumps(
+        [token, str(view.uuid), view.name, json.loads(view.configuration)],
+        sort_keys=True,
+        separators=(",", ":"),
+        allow_nan=False,
+    )
+    return hashlib.sha256(identity.encode()).hexdigest()
+
+
+def annotation_cache_token(view: SemanticView) -> str | None:
+    """Key a host chart without discovering its annotation source's 
metadata."""
+    if not metadata_binding.participates(view.semantic_layer):
+        return None
+    try:
+        token: str | None = metadata_binding.peek_view_metadata_token(view)

Review Comment:
   Confirmed: peeking on lookup did not capture a view for the later miss path. 
Commit 9c963e57a9 prepares and authorizes the participating annotation context 
before the parent warehouse query, preserving its observation through the slow 
query without extending the discovery budget. Result-cache hits still perform 
no provider discovery for the annotation source. The regression advances the 
parent beyond 30 seconds and verifies annotation success, with a denial control 
that performs neither provider nor parent execution.



##########
superset/coordination/deadline_backend.py:
##########
@@ -0,0 +1,343 @@
+# 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 os
+import time
+from _thread import LockType
+from contextlib import AsyncExitStack
+from typing import Any, TYPE_CHECKING
+
+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
+
+if TYPE_CHECKING:
+    from threading import local
+
+    from gevent.threadpool import ThreadPool
+
+_METADATA_POOL_SIZE: int = 4
+_metadata_pool_state: local | None = None
+
+
+def _metadata_threadpool() -> ThreadPool:
+    """Reuse a lazy metadata-only pool in its owning native thread and hub."""
+    from gevent import get_hub
+    from gevent.hub import Hub
+    from gevent.monkey import get_original
+    from gevent.threadpool import ThreadPool
+
+    global _metadata_pool_state  # pylint: disable=global-statement
+    if _metadata_pool_state is None:
+        # Patched threading.local is greenlet-local; the pool belongs to a hub.
+        _metadata_pool_state = get_original("threading", "local")()
+    hub: Hub = get_hub()
+    pool: ThreadPool | None = getattr(_metadata_pool_state, "pool", None)
+    if pool is None or pool.hub is not hub or pool.pid != os.getpid():
+        pool = ThreadPool(_METADATA_POOL_SIZE, hub=hub, idle_task_timeout=30)
+        _metadata_pool_state.pool = pool
+    return pool
+
+
+_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 _CommandCancellation:
+    """Transfer caller cancellation to its private native-thread asyncio 
task."""
+
+    def __init__(self, lock: LockType) -> None:
+        self._lock: LockType = lock
+        self._cancelled: bool = False
+        self._loop: asyncio.AbstractEventLoop | None = None
+        self._task: asyncio.Task[Any] | None = None
+
+    def bind(self, loop: asyncio.AbstractEventLoop, task: asyncio.Task[Any]) 
-> None:
+        """A cancelled queued command must not begin Redis work."""
+        with self._lock:
+            self._loop, self._task = loop, task
+            if self._cancelled:
+                task.cancel()
+
+    def cancel(self) -> None:
+        """Schedule cancellation without calling task methods across 
threads."""
+        with self._lock:
+            self._cancelled = True
+            if self._loop is not None and self._task is not None:
+                self._loop.call_soon_threadsafe(self._task.cancel)
+
+    def detach(self) -> None:
+        """Prevent racing cancellation from targeting a closed loop."""
+        with self._lock:
+            self._loop, self._task = None, None
+
+
+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
+
+    def _socket_timeout(self, key: str, remaining: float) -> float:
+        """Preserve shorter node timeouts inside the shared operation 
deadline."""
+        configured: object = self._config.get(key)
+        if (
+            isinstance(configured, (int, float))
+            and not isinstance(configured, bool)
+            and 0 < configured < math.inf
+        ):
+            return min(configured, remaining)
+        return remaining
+
+    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": self._socket_timeout(
+                "CACHE_REDIS_SOCKET_TIMEOUT", remaining
+            ),
+            "socket_connect_timeout": self._socket_timeout(
+                "CACHE_REDIS_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": options["socket_timeout"],
+                        "socket_connect_timeout": 
options["socket_connect_timeout"],
+                        "retry": Retry(NoBackoff(), 0),
+                        "protocol": 2,
+                    },
+                    **options,
+                )
+                sentinel_client: Redis
+                for sentinel_client in sentinel.sentinels:
+                    stack.push_async_callback(
+                        getattr(sentinel_client, "aclose", None)
+                        or sentinel_client.close
+                    )
+                client = sentinel.master_for(
+                    self._config.get("CACHE_REDIS_SENTINEL_MASTER", "mymaster")
+                )
+            else:
+                client = Redis(
+                    host=self._config.get("CACHE_REDIS_HOST", "localhost"),
+                    port=self._config.get("CACHE_REDIS_PORT", 6379),
+                    **options,
+                )
+            await stack.enter_async_context(client)

Review Comment:
   Commit 9aed7aafd7 explicitly disconnects the privately owned Sentinel master 
pool. The committed regression asserts the disconnect callback runs for 
success, error and timeout; separately, a local probe against a redis-py 5.0.0 
wheel confirmed the pool is actually disconnected on that version. A non-mocked 
version-floor test is not in this PR; I can add one to the Redis integration 
lane if you'd like it here rather than as a follow-up.



##########
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.loads(payload, use_decimal=True)
+            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, use_decimal=True),
+                sort_keys=True,
+                separators=(",", ":"),
+                allow_nan=False,
+            )
+        except (TypeError, ValueError, UnicodeError, RecursionError):
+            raise MetadataRefreshError("invalid_payload") from None

Review Comment:
   Confirmed at both parsing boundaries. Commit 9c963e57a9 translates the 
provider decimal conversion failure to invalid_payload and treats the same 
failure in a stored envelope as a cache miss. The regressions use the oversized 
exponent directly, verify a failed refresh preserves the existing snapshot, and 
verify a corrupt stored snapshot can be refilled.



##########
superset/coordination/deadline_backend.py:
##########
@@ -0,0 +1,345 @@
+# 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 os
+import time
+from _thread import LockType
+from contextlib import AsyncExitStack
+from typing import Any, TYPE_CHECKING
+
+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
+from superset.coordination.metadata_resolver import MetadataEventLoop
+
+if TYPE_CHECKING:
+    from threading import local
+
+    from gevent.threadpool import ThreadPool
+
+_METADATA_POOL_SIZE: int = 4
+_metadata_pool_state: local | None = None
+
+
+def _metadata_threadpool() -> ThreadPool:
+    """Reuse a lazy metadata-only pool in its owning native thread and hub."""
+    from gevent import get_hub
+    from gevent.hub import Hub
+    from gevent.monkey import get_original
+    from gevent.threadpool import ThreadPool
+
+    global _metadata_pool_state  # pylint: disable=global-statement
+    if _metadata_pool_state is None:
+        # Patched threading.local is greenlet-local; the pool belongs to a hub.
+        _metadata_pool_state = get_original("threading", "local")()
+    hub: Hub = get_hub()
+    pool: ThreadPool | None = getattr(_metadata_pool_state, "pool", None)
+    if pool is None or pool.hub is not hub or pool.pid != os.getpid():
+        pool = ThreadPool(_METADATA_POOL_SIZE, hub=hub, idle_task_timeout=30)
+        _metadata_pool_state.pool = pool
+    return pool
+
+
+_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 _CommandCancellation:
+    """Transfer caller cancellation to its private native-thread asyncio 
task."""
+
+    def __init__(self, lock: LockType) -> None:
+        self._lock: LockType = lock
+        self._cancelled: bool = False
+        self._loop: asyncio.AbstractEventLoop | None = None
+        self._task: asyncio.Task[Any] | None = None
+
+    def bind(self, loop: asyncio.AbstractEventLoop, task: asyncio.Task[Any]) 
-> None:
+        """A cancelled queued command must not begin Redis work."""
+        with self._lock:
+            self._loop, self._task = loop, task
+            if self._cancelled:
+                task.cancel()
+
+    def cancel(self) -> None:
+        """Schedule cancellation without calling task methods across 
threads."""
+        with self._lock:
+            self._cancelled = True
+            if self._loop is not None and self._task is not None:
+                self._loop.call_soon_threadsafe(self._task.cancel)
+
+    def detach(self) -> None:
+        """Prevent racing cancellation from targeting a closed loop."""
+        with self._lock:
+            self._loop, self._task = None, None
+
+
+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
+
+    def _socket_timeout(self, key: str, remaining: float) -> float:
+        """Preserve shorter node timeouts inside the shared operation 
deadline."""
+        configured: object = self._config.get(key)
+        if (
+            isinstance(configured, (int, float))
+            and not isinstance(configured, bool)
+            and 0 < configured < math.inf
+        ):
+            return min(configured, remaining)
+        return remaining
+
+    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": self._socket_timeout(
+                "CACHE_REDIS_SOCKET_TIMEOUT", remaining
+            ),
+            "socket_connect_timeout": self._socket_timeout(
+                "CACHE_REDIS_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": options["socket_timeout"],
+                        "socket_connect_timeout": 
options["socket_connect_timeout"],
+                        "retry": Retry(NoBackoff(), 0),
+                        "protocol": 2,
+                    },
+                    **options,
+                )
+                sentinel_client: Redis
+                for sentinel_client in sentinel.sentinels:
+                    stack.push_async_callback(
+                        getattr(sentinel_client, "aclose", None)
+                        or sentinel_client.close
+                    )
+                client = sentinel.master_for(
+                    self._config.get("CACHE_REDIS_SENTINEL_MASTER", "mymaster")
+                )
+                # redis-py 5.0 clients do not own the pool supplied by 
Sentinel.
+                # This command does: release it even on errors and 
cancellation.
+                stack.push_async_callback(client.connection_pool.disconnect)
+            else:
+                client = Redis(
+                    host=self._config.get("CACHE_REDIS_HOST", "localhost"),
+                    port=self._config.get("CACHE_REDIS_PORT", 6379),
+                    **options,
+                )
+            await stack.enter_async_context(client)
+            return await client.execute_command(*args)
+
+    def execute(self, *args: str | int) -> Any:
+        """Keep request greenlets from sharing asyncio's native-thread loop 
state."""
+        self._remaining()
+        try:
+            from gevent import getcurrent, Greenlet, Timeout
+            from gevent.event import AsyncResult
+            from gevent.monkey import get_original
+        except ImportError:
+            return self._execute_sync(*args)
+        if not isinstance(getcurrent(), Greenlet):
+            return self._execute_sync(*args)
+        cancellation: _CommandCancellation = _CommandCancellation(
+            get_original("_thread", "allocate_lock")()
+        )
+        try:
+            # The metadata-only native pool leaves the hub DNS pool available.
+            # Its queue wait consumes the original deadline too.
+            with Timeout(
+                self._remaining(), RedisTimeoutError("Metadata deadline 
expired")
+            ):
+                result: AsyncResult = _metadata_threadpool().spawn(
+                    self._execute_sync, *args, cancellation=cancellation
+                )
+                return result.get()
+        finally:
+            cancellation.cancel()
+
+    def _execute_sync(
+        self,
+        *args: str | int,
+        cancellation: _CommandCancellation | None = None,
+    ) -> Any:
+        """Run a private loop and cancel its socket work within the call 
budget."""
+        self._remaining()
+        try:
+            asyncio.get_running_loop()
+        except RuntimeError:
+            pass
+        else:
+            raise RedisError("Metadata backend requires a synchronous caller")
+        loop: asyncio.AbstractEventLoop = MetadataEventLoop()
+        task: asyncio.Task[Any] = loop.create_task(self._command(*args))
+        if cancellation is not None:
+            cancellation.bind(loop, task)
+        try:
+            try:
+                return loop.run_until_complete(task)
+            except asyncio.CancelledError:
+                raise RedisTimeoutError("Metadata command cancelled") from None
+            except TimeoutError:
+                raise RedisTimeoutError("Metadata deadline expired") from None
+        finally:
+            # Outstanding system DNS retains its process-wide admission slot.
+            # Closing this loop cannot resume a cancelled command or publish.
+            if cancellation is not None:
+                cancellation.detach()
+            loop.close()

Review Comment:
   Good catch, thanks. When the deadline cancels disconnect() the TLS transport 
can be left holding its socket until GC. I'll abort any unfinished private 
transports before the loop is closed, with a stalled-TLS-shutdown regression, 
in a follow-up commit on this PR.



-- 
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]

Reply via email to