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


##########
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:
   Fixed in published store commit `5199837fe44894d2ab306f72986d3ba08123cb3b` 
(implementation in `b6b40bd3e843da9a65d3865937ab8e207a3ffa39`), now merged into 
this PR at `2740673693`. Catalog normalization and warm reads use opt-in 
Decimal parsing, preserving the long decimal, `1e400` and `1e-400`; ordinary 
JSON callers retain existing behavior.



##########
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:
   Fixed in published store commit `5199837fe44894d2ab306f72986d3ba08123cb3b` 
(follow-up `3e4e330dc71c175b572a5754aecbc21c82a39756`), now merged into this PR 
at `2740673693`. Finite positive configured node timeouts are clamped to the 
operation budget; invalid/unset values use the remaining budget. Docs recommend 
starting at 1 second per node and tuning. Defaults can still exhaust the budget 
on the first node, and each command can pay that timeout again. Constructor 
configuration is tested; live second-node failover is not claimed.



##########
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:
   Corrected in published store commit 
`5199837fe44894d2ab306f72986d3ba08123cb3b` (implementation in 
`b6b40bd3e843da9a65d3865937ab8e207a3ffa39`), now merged into this PR at 
`2740673693`. Both guards patch `superset.coordination.deadline_backend.Redis`; 
injecting an early constructor call made both assertions fail, and restored 
production code passes.



##########
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:
   Fixed in `2740673693`: datasource metadata/query/column-values and Explore 
GET use the established typed-error mapper. Authorized discovery regressions 
return 503/504 with safe categories; denials stay 403 before provider work, and 
unrelated DB errors retain their handling with refresh on/off. Dashboard 
datasets already isolate provider failures per datasource: a real 
mixed-dashboard HTTP regression returns 200 and preserves healthy table 
metadata. That behaviour is unchanged and now documented.



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