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


##########
superset/semantic_layers/metadata.py:
##########
@@ -0,0 +1,470 @@
+# 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 decimal import InvalidOperation
+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
+MAX_SNAPSHOT_TTL_SECONDS: int = 2**31 - 1
+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,
+            lease_ttl_ms: int,
+            snapshot_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,
+        snapshot_ttl_seconds: int = CATALOG_TTL_SECONDS,
+        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")
+        if (
+            isinstance(snapshot_ttl_seconds, bool)
+            or not isinstance(snapshot_ttl_seconds, int)
+            or not 1 <= snapshot_ttl_seconds <= MAX_SNAPSHOT_TTL_SECONDS
+        ):
+            raise MetadataRefreshError("configuration")
+        self._snapshot_ttl_seconds: int = snapshot_ttl_seconds
+        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, InvalidOperation):
+            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)
+            if token.startswith(f"{self._scope}:")
+            else None
+        )
+
+    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,
+            snapshot_ttl_seconds=self._snapshot_ttl_seconds,
+            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
+                lease_ttl_ms: int = max(
+                    1, math.ceil(min(REFRESH_LEASE_SECONDS, self._remaining()) 
* 1000)
+                )
+                if self._backend.set(
+                    self._lease_key,
+                    attempt,
+                    px=lease_ttl_ms,
+                    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, 
lease_ttl_ms).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
+        lease_ttl_ms: int = max(
+            1, math.ceil(min(REFRESH_LEASE_SECONDS, self._remaining()) * 1000)
+        )
+        try:
+            if not self._backend.set(
+                self._lease_key,
+                attempt,
+                px=lease_ttl_ms,
+                nx=True,
+            ):
+                raise MetadataRefreshError("in_progress")
+            try:
+                return self._acquire(fetch, attempt, lease_ttl_ms)
+            finally:
+                self._release(attempt)
+        except RedisError:
+            self._remaining()
+            raise MetadataRefreshError("unavailable") from None
+
+    def _acquire(
+        self, fetch: CatalogLoader, attempt: str, lease_ttl_ms: int
+    ) -> 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, 
InvalidOperation):
+            raise MetadataRefreshError("invalid_payload") from None
+        digest: str = hashlib.sha256(payload.encode()).hexdigest()
+        status: Literal["changed", "unchanged"] = (
+            "unchanged"
+            if previous is not None and previous.digest == digest
+            else "changed"
+        )
+        observed_at: str = datetime.now(timezone.utc).isoformat()
+        snapshot: CatalogSnapshot = CatalogSnapshot(
+            payload, f"{self._scope}:{uuid4().hex}", observed_at
+        )
+        if self._before_publish is not None:
+            self._before_publish()
+        self._remaining()
+        envelope: str = json.dumps(
+            {
+                "version": SNAPSHOT_FORMAT_VERSION,
+                "payload": payload,
+                "cache_token": snapshot.cache_token,
+                "observed_at": observed_at,
+                "digest": digest,
+                "attempt": attempt,
+                "created_at": datetime.now(timezone.utc).isoformat(),
+            },
+            separators=(",", ":"),
+        )
+        if len(envelope.encode()) > MAX_CATALOG_BYTES:
+            raise MetadataRefreshError("invalid_payload")
+        ttl_ms: int = math.floor(
+            (self._snapshot_ttl_seconds - (self._clock() - started)) * 1000
+        )
+        self._remaining()
+        if ttl_ms <= 0:
+            raise MetadataRefreshError("deadline")
+        return MetadataRefreshResult(
+            status, self._publish(attempt, envelope, ttl_ms, lease_ttl_ms, 
snapshot)
+        )
+
+    def _publish(
+        self,
+        attempt: str,
+        envelope: str,
+        ttl_ms: int,
+        lease_ttl_ms: int,
+        snapshot: CatalogSnapshot,
+    ) -> CatalogSnapshot:
+        try:
+            accepted: bool = self._backend.compare_and_publish(
+                self._lease_key,
+                attempt,
+                self._snapshot_key,
+                envelope,
+                ttl_ms,
+                lease_ttl_ms,
+                self._snapshot_ttl_seconds * 1000,
+            )
+        except RedisError:
+            try:
+                confirmed: StoredCatalog | None = self._load()
+            except (RedisError, MetadataRefreshError):
+                confirmed = None
+            if confirmed is None or confirmed.attempt != attempt:
+                raise MetadataRefreshError("indeterminate") from None
+            return self._remember(confirmed.snapshot)
+        if not accepted:
+            raise MetadataRefreshError("configuration_changed")
+        return self._remember(snapshot)
+
+    def _release(self, attempt: str) -> None:
+        # Do not perform cleanup I/O after the request budget is exhausted.
+        if self._clock() >= self._deadline:
+            return
+        try:
+            self._backend.compare_and_delete(self._lease_key, attempt)
+        except RedisError:
+            # A failed cleanup leaves only an expiring lease, preserving the 
error.
+            pass
+
+    def invalidate_catalog(self) -> None:
+        """Atomically retire visibility and old writer authority without 
fetching."""
+        self._remaining()
+        try:
+            self._backend.delete(self._snapshot_key, self._lease_key)
+        except RedisError:
+            raise MetadataRefreshError("indeterminate") from None
+
+    def inspect_catalog(self) -> CacheEntryInfo:
+        """Observe timing without filling an empty catalog or renewing its 
lifetime."""
+        self._remaining()
+        try:
+            raw: bytes | None
+            ttl_ms: int
+            raw, ttl_ms = self._backend.get_with_ttl(self._snapshot_key)
+        except RedisError:
+            return CacheEntryInfo("catalog", "unavailable", 
datetime.now(timezone.utc))
+        stored: StoredCatalog | None = self._decode(raw)
+        if raw is not None and stored is None:
+            return CacheEntryInfo("catalog", "unsupported", 
datetime.now(timezone.utc))
+        return describe_entry(
+            "catalog",
+            {
+                "created_at": stored.created_at,
+                "observed_at": stored.snapshot.observed_at,
+            }
+            if stored is not None
+            else None,
+            ttl_ms,
+        )
+
+    def compatibility_generation(self) -> str:
+        """Capture a non-reusable generation, including after eviction."""
+        self._remaining()
+        try:
+            return self._backend.get_or_create(
+                self._generation_key, uuid4().hex, self._snapshot_ttl_seconds
+            ).decode()
+        except RedisError:

Review Comment:
   `compatibility_generation()` maps every `RedisError` straight to 
`unavailable`. When the generation lookup consumes the remaining budget, the 
deadline backend raises a Redis timeout, so an authorized `/compatible` request 
that ran out of time reports 503 `unavailable` rather than the 504 `deadline` 
that `_read()` and `_refresh()` produce by calling `self._remaining()` in their 
handlers. `peek_compatibility_generation` (line 458) has the same shape. Should 
these handlers re-check the budget before raising, as the catalog paths do?



##########
superset/mcp_service/semantic_layer/tool/get_table.py:
##########
@@ -653,9 +658,14 @@ async def get_table(
         datasource_id = request.view_id
 
     try:
-        return await _run_get_table_query(
-            request, ctx, is_builtin, datasource_id, datasource_type
-        )
+        with (
+            metadata_operation()
+            if not is_builtin and metadata_refresh_enabled()
+            else nullcontext()
+        ):
+            return await _run_get_table_query(

Review Comment:
   With refresh enabled and an opted-in provider, `get_table(view_id=…)` runs 
`_run_get_table_query` (a coroutine) and reaches `view.implementation` / 
`view.columns` synchronously on the running event loop. 
`DeadlineRedisBackend._execute_sync()` raises `RedisError("Metadata backend 
requires a synchronous caller")` whenever `asyncio.get_running_loop()` 
succeeds, so even a warm catalog on a healthy Redis comes back as 
`MetadataRefreshError("unavailable")` and the tool returns `InternalError`. 
`get_compatible_metrics` and `list_metrics` use the same `async def` + 
`metadata_operation()` pattern, and `test_view_tool_enters_metadata_operation` 
patches `implementation`, so the real backend is never exercised under the MCP 
loop. Should the discovery calls run off the loop thread (with the operation 
context carried over), or should the backend tolerate a running loop?



##########
tests/unit_tests/semantic_layers/metadata_result_inspection_test.py:
##########
@@ -0,0 +1,313 @@
+# 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.
+
+from __future__ import annotations
+
+from typing import Any
+from unittest.mock import Mock, patch
+
+import numpy as np
+import pytest
+from flask import Flask, g
+
+from superset.common.query_context import QueryContext
+from superset.common.query_object import QueryObject
+from superset.semantic_layers.cache_inspection import CacheEntryInfo
+from superset.semantic_layers.models import SemanticView
+from superset.semantic_layers.registry import registry
+from tests.unit_tests.semantic_layers.metadata_identity_test import (
+    context_for,
+    RefreshLayer,
+    ResultView,
+    view_for,
+)
+
+
[email protected]("total", [np.float32(12.5), np.int32(12)])
+def test_result_capture_excludes_runtime_contribution_totals(
+    app: Flask, monkeypatch: pytest.MonkeyPatch, total: Any
+) -> None:
+    """Diagnostic capture shares the query key's runtime-total exclusion."""
+    from superset.semantic_layers.result_inspection import captured_result_key
+
+    provider: ResultView = ResultView("captured", 17)
+    context: QueryContext
+    query: QueryObject
+    context, query = context_for(view_for(provider))
+    options: dict[str, Any] = {"columns": ["orders"], "contribution_totals": 
total}
+    rename_options: dict[str, Any] = {"columns": {"orders": "Order count"}}
+    query.post_processing = [
+        {"operation": "contribution", "options": options},
+        {"operation": "rename", "options": rename_options},
+    ]
+    monkeypatch.setitem(app.config, "SEMANTIC_LAYER_METADATA_REFRESH_ENABLED", 
True)
+    manager: Mock = Mock()
+    manager.get_rls_cache_key.return_value = []
+    with (
+        app.test_request_context(),
+        patch(
+            "superset.semantic_layers.metadata_binding.participates", 
return_value=True
+        ),
+        patch(
+            
"superset.semantic_layers.result_inspection.metadata_refresh_enabled",
+            return_value=True,
+        ),
+        patch(
+            "superset.semantic_layers.metadata_binding.view_implementation",
+            return_value=provider,
+        ),
+        patch("superset.semantic_layers.result_inspection.security_manager", 
manager),
+        patch("superset.common.query_context_processor.security_manager", 
manager),
+    ):
+        key: str | None = context.query_cache_key(query)
+        assert key is not None
+        assert captured_result_key(context, query) == key
+        assert options["contribution_totals"] is total
+        options["contribution_totals"] = np.float32(99.5)
+        assert captured_result_key(context, query) == key
+        assert rename_options == {"columns": {"orders": "Order count"}}
+        rename_options["columns"] = {"orders": "Renamed count"}
+        assert captured_result_key(context, query) is None
+        rename_options["columns"] = {"orders": "Order count"}
+        assert captured_result_key(context, query) == key
+        options["columns"] = ["revenue"]
+        assert captured_result_key(context, query) is None
+
+
[email protected]("changed", ["subject", "query"])
+def test_result_capture_rejects_changed_subject_or_query(
+    app: Flask, changed: str
+) -> None:
+    """A captured key never survives a subject or query edit with unchanged 
RLS."""
+    from superset.semantic_layers.result_inspection import (
+        capture_result_identity,
+        captured_result_key,
+    )
+
+    context: QueryContext
+    query: QueryObject
+    context, query = context_for(view_for(ResultView("captured", 17)))
+    manager: Mock = Mock()
+    manager.get_rls_cache_key.return_value = ["unchanged-rule"]
+    with (
+        app.test_request_context(),
+        patch(
+            
"superset.semantic_layers.result_inspection.metadata_refresh_enabled",
+            return_value=True,
+        ),
+        patch("superset.semantic_layers.result_inspection.security_manager", 
manager),
+    ):
+        g.user = Mock(id=1)
+        capture_result_identity(context, query, "existing-key")
+        assert captured_result_key(context, query) == "existing-key"
+        if changed == "subject":
+            g.user = Mock(id=2)
+        else:
+            query.metrics = ["revenue"]
+        assert captured_result_key(context, query) is None
+
+
+def test_result_inspection_keeps_query_rls_and_never_constructs_provider(
+    app: Flask, monkeypatch: pytest.MonkeyPatch
+) -> None:
+    from superset.commands.semantic_layer.inspect_query_result import (
+        InspectQueryResultCommand,
+    )
+
+    provider: ResultView = ResultView("unused", 17)
+    view: SemanticView = view_for(provider)
+    monkeypatch.setitem(app.config, "SEMANTIC_LAYER_METADATA_REFRESH_ENABLED", 
True)
+    monkeypatch.setitem(registry, "cache-test", RefreshLayer)
+    context: QueryContext
+    query: QueryObject
+    context, query = context_for(view)
+    manager: Mock = Mock()
+    keys: list[str] = []
+    inspect_entry: Mock
+    construct: Mock
+    rls: list[str]
+    with (
+        app.test_request_context(),
+        patch(
+            "superset.semantic_layers.metadata_binding.is_feature_enabled",
+            return_value=True,
+        ),
+        patch(
+            
"superset.commands.semantic_layer.inspect_query_result.security_manager",
+            manager,
+        ),
+        patch("superset.semantic_layers.result_inspection.security_manager", 
manager),
+        patch("superset.common.query_context_processor.security_manager", 
manager),
+        patch(
+            
"superset.commands.semantic_layer.inspect_query_result.inspect_derived_entry"
+        ) as inspect_entry,
+        patch(
+            "superset.semantic_layers.metadata_binding.view_implementation",
+            return_value=provider,
+        ) as construct,
+    ):
+        assert InspectQueryResultCommand(context, 0).run().state == 
"unsupported"
+        construct.assert_not_called()
+        for rls in (["rule-a"], ["rule-b"]):
+            manager.get_rls_cache_key.return_value = rls
+            key: str | None = context.query_cache_key(query)
+            construct.reset_mock()
+            InspectQueryResultCommand(context, 0).run()
+            assert inspect_entry.call_args.args == (key, "query_result")
+            keys.append(inspect_entry.call_args.args[0])
+            construct.assert_not_called()
+        assert keys[0] != keys[1]
+        manager.get_rls_cache_key.return_value = ["revoked"]
+        assert InspectQueryResultCommand(context, 0).run().state == 
"unsupported"
+        manager.raise_for_access.side_effect = PermissionError("denied")
+        inspect_entry.reset_mock()
+        with pytest.raises(PermissionError):
+            InspectQueryResultCommand(context, 0).run()
+        inspect_entry.assert_not_called()
+
+        # Keep the feature and the same subject/RLS enabled in the next 
request.
+        # Only request-local capture expiry should prevent reading the old key.
+        manager.raise_for_access.side_effect = None
+        manager.get_rls_cache_key.return_value = ["rule-b"]
+        inspect_entry.reset_mock()
+        with app.test_request_context():
+            assert InspectQueryResultCommand(context, 0).run().state == 
"unsupported"
+            inspect_entry.assert_not_called()
+
+
+def test_result_inspection_does_not_refill_after_concurrent_invalidation(
+    app: Flask,
+    monkeypatch: pytest.MonkeyPatch,
+) -> None:
+    from superset.commands.semantic_layer.inspect_query_result import (
+        InspectQueryResultCommand,
+    )
+    from superset.semantic_layers.metadata import ScopedMetadataStore
+    from superset.semantic_layers.metadata_binding import (
+        operation_deadline,
+        request_metadata_budget,
+    )
+    from tests.unit_tests.semantic_layers.metadata_contract_test import 
OptedInLayer
+    from tests.unit_tests.semantic_layers.metadata_store_test import 
MemoryBackend
+
+    manager: Mock = Mock()
+    provider: OptedInLayer = OptedInLayer()
+    resolve_provider: Mock
+    inspect_entry: Mock
+    view: SemanticView = view_for(ResultView("unused", 17))
+    context: QueryContext
+    query: QueryObject
+    context, query = context_for(view)
+    monkeypatch.setitem(app.config, "SEMANTIC_LAYER_METADATA_REFRESH_ENABLED", 
True)
+    monkeypatch.setitem(
+        app.config, "SEMANTIC_LAYER_METADATA_NAMESPACE", "inspection-test"
+    )
+    monkeypatch.setitem(registry, "cache-test", OptedInLayer)
+    with (
+        app.test_request_context(),
+        patch(
+            "superset.semantic_layers.metadata_binding.is_feature_enabled",
+            return_value=True,
+        ),
+        patch(
+            
"superset.commands.semantic_layer.inspect_query_result.security_manager",
+            manager,
+        ),
+        patch(
+            "superset.semantic_layers.metadata_binding.connection_store"
+        ) as resolve_provider,
+        patch(
+            
"superset.commands.semantic_layer.inspect_query_result.inspect_derived_entry"
+        ) as inspect_entry,
+        patch("superset.semantic_layers.result_inspection.security_manager", 
manager),
+        patch("superset.common.query_context_processor.security_manager", 
manager),
+        patch.object(OptedInLayer, "from_configuration", 
return_value=provider),
+    ):
+        manager.get_rls_cache_key.return_value = []
+        request_metadata_budget()
+        deadline: float = operation_deadline()
+        store: ScopedMetadataStore = ScopedMetadataStore(
+            MemoryBackend(), "inspection", deadline=deadline
+        )
+        store.read(lambda budget: '["orders"]', deadline=deadline)
+        resolve_provider.return_value = store
+        key: str | None = context.query_cache_key(query)
+        assert key is not None
+        resolve_provider.assert_called()
+        resolve_provider.reset_mock()
+        # The normal query captured its key before another caller retired the
+        # catalog. Inspection must use that key without resolving a new uid.
+        store.invalidate_catalog()
+        result: CacheEntryInfo = InspectQueryResultCommand(context, 0).run()
+        assert result is inspect_entry.return_value
+        inspect_entry.assert_called_once_with(key, "query_result")
+        assert provider.adapter.fetches == 0
+        resolve_provider.assert_not_called()
+
+
+def test_result_capture_is_disabled_without_http_operation(app: Flask) -> None:

Review Comment:
   This test runs with the default app config, where refresh is disabled, so 
`capture_result_identity` returns at the `metadata_refresh_enabled()` check 
even if the `has_request_context()` guard it is named for were deleted. With 
that guard gone, a Celery-run semantic query with refresh enabled would reach 
`request.environ` outside a request and raise `RuntimeError`, and this suite 
would stay green. Could the test enable 
`SEMANTIC_LAYER_METADATA_REFRESH_ENABLED` (and the feature flag), call capture 
with no request context, and assert nothing is stored and `captured_result_key` 
is `None`?



##########
superset/common/query_context_processor.py:
##########
@@ -428,28 +448,100 @@ def query_cache_key(self, query_obj: QueryObject, 
**kwargs: Any) -> str | None:
         """
         Returns a QueryObject cache key for objects in self.queries
         """
+        return self._query_cache_key(query_obj, **kwargs)[0]
+
+    def _query_cache_key(
+        self,
+        query_obj: QueryObject,
+        *,
+        parent_cache_context: dict[str, Any] | None = None,
+        **kwargs: Any,
+    ) -> tuple[str | None, bool]:
+        """Keep annotation cacheability alongside the opaque hashed result 
key."""
         datasource: Explorable = self._qc_datasource
-        # Reject unenforceable restrictions before provider identity or cache 
reads.
-        rls: list[str] = security_manager.get_rls_cache_key(datasource)
-        extra_cache_keys = datasource.get_extra_cache_keys(query_obj.to_dict())
+        if parent_cache_context is None:
+            parent_cache_context = {}
+        if not parent_cache_context:
+            # Capture parent/Jinja inputs once; annotation discovery may re-key
+            # the result after execution but must not evaluate new parent 
inputs.
+            rls: list[str] = security_manager.get_rls_cache_key(datasource)
+            parent_cache_context.update(
+                datasource=datasource.uid,
+                
extra_cache_keys=datasource.get_extra_cache_keys(query_obj.to_dict()),
+                rls=rls,
+                changed_on=datasource.changed_on,
+            )
+        cacheable: bool = True
 
         # Annotation data is cached on the same entry as the dataframe, so the
         # key must also bind the annotation sources' security context.
         if query_obj and query_obj.annotation_layers:
-            kwargs["annotation_context"] = 
self._annotation_cache_context(query_obj)
+            annotation_context: dict[str, Any] = 
self._annotation_cache_context(
+                query_obj
+            )
+            kwargs["annotation_context"] = annotation_context
+            source_metadata: dict[str, str] = annotation_context.get(
+                "source_metadata", {}
+            )
+            cacheable = not any(
+                token.startswith("uncaptured:") for token in 
source_metadata.values()
+            )
 
-        cache_key = (
+        cache_key: str | None = (
             query_obj.cache_key(
-                datasource=datasource.uid,
-                extra_cache_keys=extra_cache_keys,
-                rls=rls,
-                changed_on=datasource.changed_on,
+                **parent_cache_context,
                 **kwargs,
             )
             if query_obj
             else None
         )
-        return cache_key
+        if cache_key is not None:
+            capture_result_identity(self._query_context, query_obj, cache_key)
+        return cache_key, cacheable
+
+    def _capture_annotation_metadata(self, query_obj: QueryObject) -> None:
+        """Capture authorized annotation views on a miss before a slow parent 
query."""
+        if not query_obj.annotation_layers:
+            return
+        from superset_core.semantic_layers.metadata import MetadataRefreshError
+
+        from superset.semantic_layers.metadata_binding import (
+            metadata_refresh_enabled,
+            participates,
+        )
+        from superset.semantic_layers.models import SemanticView
+
+        if not metadata_refresh_enabled():
+            return
+        layer: dict[str, Any]
+        for layer in query_obj.annotation_layers:
+            if (
+                layer.get("sourceType")
+                not in ANNOTATION_SOURCE_TYPES_WITH_CHART_REFERENCE
+            ):
+                continue
+            value: int | str | None = layer.get("value")
+            chart: Slice | None = (
+                ChartDAO.find_by_id(value) if value is not None else None
+            )
+            if chart is None:
+                continue
+            try:
+                context: QueryContext | None = chart.get_query_context()
+                if context is None:
+                    # The annotation executor reports its missing-context 
error.
+                    continue
+                source: Explorable = context.datasource
+                if not isinstance(source, SemanticView) or not participates(
+                    source.semantic_layer
+                ):
+                    continue
+                context.raise_for_access()
+                # Reuse the annotation command's canonical query authority.
+                # Later execution retains this view without renewing its 
budget.
+                _captured: str | None = source.metadata_cache_token
+            except (SupersetException, MetadataRefreshError) as ex:
+                raise QueryObjectValidationError(error_msg_from_exception(ex)) 
from ex

Review Comment:
   When a chart's annotation layer references a participating semantic-view 
chart and the result cache misses, a typed `MetadataRefreshError` (storage 
unavailable, deadline expired) is rewrapped here as 
`QueryObjectValidationError`. `get_df_payload` turns that into a failed query 
payload, so chart-data returns 400 instead of the 503/504 that 
`metadata_api_errors` gives the other chart-data discovery paths, and 
`semantic-metadata-operations.md` describes the typed mapping for chart-data 
requests. A client that retries on 503/504 would treat a transient Redis outage 
as a bad request. Should `MetadataRefreshError` propagate here (wrapping only 
`SupersetException`) so the existing mapper handles it?



##########
superset/semantic_layers/models.py:
##########
@@ -365,8 +376,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:
   For a participating layer, `implementation` now goes through 
`view_implementation`, which raises `MetadataRefreshError("configuration")` 
unless a metadata operation is active. The MCP `get_compatible_dimensions` tool 
still calls `view.get_compatible_dimensions(...)` and `view.columns` outside 
any `metadata_operation()` (unlike `get_table`, `get_compatible_metrics` and 
`list_metrics`, which this PR wraps). With refresh enabled and an opted-in 
provider, an authorized call for that view therefore fails with a configuration 
error before any provider work. Should that tool enter `metadata_operation()` 
like its siblings?



##########
superset/commands/semantic_layer/refresh_metadata.py:
##########
@@ -0,0 +1,344 @@
+# 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.
+from __future__ import annotations
+
+from collections.abc import Callable, Iterator
+from contextlib import contextmanager
+from datetime import datetime, timezone
+from typing import Any, Literal
+from uuid import UUID
+
+from flask import current_app, g, has_request_context, Request, request
+from flask_appbuilder.security.sqla.models import User
+from sqlalchemy.orm import Session
+from superset_core.semantic_layers.layer import SemanticLayer as 
SemanticLayerABC
+from superset_core.semantic_layers.metadata import (
+    MetadataRefreshError,
+    MetadataRefreshResult,
+)
+from superset_core.semantic_layers.view import SemanticView as SemanticViewABC
+
+from superset import cache_manager, security_manager
+from superset.commands.base import BaseCommand
+from superset.commands.semantic_layer.exceptions import (
+    SemanticLayerForbiddenError,
+    SemanticLayerNotFoundError,
+    SemanticViewNotFoundError,
+)
+from superset.commands.utils import current_user_can_modify_object
+from superset.coordination.deadline_backend import DeadlineRedisBackend
+from superset.daos.semantic_layer import SemanticViewDAO
+from superset.exceptions import SupersetSecurityException
+from superset.extensions import db
+from superset.semantic_layers.cache_inspection import CacheEntryInfo, 
inspect_data_cache
+from superset.semantic_layers.metadata import ScopedMetadataStore
+from superset.semantic_layers.metadata_binding import (
+    connection_metadata_scope,
+    metadata_refresh_enabled,
+    operation_deadline,
+    participates,
+)
+from superset.semantic_layers.metadata_cache import (
+    compatibility_identity,
+    CompatibilityIdentity,
+)
+from superset.semantic_layers.models import SemanticLayer, SemanticView
+from superset.semantic_layers.registry import registry
+from superset.utils import json
+
+
+def authorize_metadata_refresh(view: SemanticView) -> None:
+    """Share one server policy between the affordance and direct mutation."""
+    if not metadata_refresh_enabled():
+        raise SemanticViewNotFoundError()
+    user: User | None = getattr(g, "user", None)
+    if (
+        user is None
+        or user.is_anonymous
+        or getattr(user, "is_guest_user", False)
+        or not user.is_active
+    ):
+        raise SemanticLayerForbiddenError()
+    if not all(
+        security_manager.can_access(action, resource)
+        for action, resource in (
+            ("can_read", "SemanticView"),
+            ("can_read", "SemanticLayer"),
+            ("can_write", "SemanticLayer"),
+        )
+    ):
+        raise SemanticLayerForbiddenError()
+    layer: SemanticLayer | None = view.semantic_layer
+    if layer is None:
+        raise SemanticLayerNotFoundError()
+    try:
+        view.raise_for_access()
+        layer.raise_for_access()
+    except SupersetSecurityException:
+        raise SemanticLayerForbiddenError() from None
+    if not current_user_can_modify_object(layer):
+        raise SemanticLayerForbiddenError()
+    if layer.type not in registry:
+        raise MetadataRefreshError("unsupported")
+    if not participates(layer):
+        raise MetadataRefreshError("unsupported")
+    connection_metadata_scope(layer)
+
+
+def can_refresh_metadata(view: SemanticView) -> bool:
+    """Project policy without constructing a provider or consulting its 
catalog."""
+    try:
+        authorize_metadata_refresh(view)
+    except (
+        SemanticViewNotFoundError,
+        SemanticLayerNotFoundError,
+        SemanticLayerForbiddenError,
+        MetadataRefreshError,
+    ):
+        return False
+    return True
+
+
+def view_binding(view: SemanticView) -> tuple[str, str, str]:
+    """Capture immutable provider selection, independent of Details drafts."""
+    return (
+        str(view.semantic_layer_uuid),
+        view.name,
+        json.dumps(json.loads(view.configuration), sort_keys=True),
+    )
+
+
+@contextmanager
+def fresh_refresh_authority(session: Session) -> Iterator[None]:
+    """Run existing policy against persisted authority without ending request 
work.
+
+    Security-manager and subject helpers use the request-scoped session and
+    principal. Rebind those only for this guard so they cannot reuse an earlier
+    repeatable-read snapshot or cached role membership. The supplied session
+    owns no writes; the caller closes it. Always restore the request's objects.
+    """
+    original_user: User = g.user
+    original_session: Session = db.session()
+    had_login_user: bool = hasattr(g, "_login_user")
+    original_login_user: Any = getattr(g, "_login_user", None)
+    if original_user.is_anonymous or getattr(original_user, "is_guest_user", 
False):
+        raise SemanticLayerForbiddenError()
+    user: User | None = session.get(security_manager.user_model, 
original_user.id)
+    if user is None or not user.is_active:
+        raise SemanticLayerForbiddenError()
+    current_request: Request | None = (
+        request._get_current_object() if has_request_context() else None
+    )
+    had_subject_cache: bool = current_request is not None and hasattr(
+        current_request, "_user_subject_ids"
+    )
+    original_subject_cache: dict[int, list[int]] | None = getattr(
+        current_request, "_user_subject_ids", None
+    )
+    try:
+        if current_request is not None:
+            current_request._user_subject_ids = {}
+        db.session.registry.set(session)
+        g.user = user
+        g._login_user = user
+        with session.no_autoflush:
+            yield
+    finally:
+        if current_request is not None:
+            if had_subject_cache:
+                current_request._user_subject_ids = original_subject_cache
+            else:
+                current_request.__dict__.pop("_user_subject_ids", None)
+        g.user = original_user
+        if had_login_user:
+            g._login_user = original_login_user
+        else:
+            g.pop("_login_user", None)
+        db.session.registry.set(original_session)
+
+
+def guarded_store(
+    view: SemanticView, *, before_publish: Callable[[], None]
+) -> ScopedMetadataStore:
+    """Bind command-specific fresh authority to the unchanged host store 
contract."""
+    deadline: float = operation_deadline()
+    config: Any = current_app.config.get("DISTRIBUTED_COORDINATION_CONFIG")
+    if not isinstance(config, dict):
+        raise MetadataRefreshError("unavailable")
+    try:
+        backend: DeadlineRedisBackend = DeadlineRedisBackend(config, 
deadline=deadline)
+    except ValueError:
+        raise MetadataRefreshError("configuration") from None
+    return ScopedMetadataStore(
+        backend,
+        connection_metadata_scope(view.semantic_layer),
+        deadline=deadline,
+        before_publish=before_publish,
+        snapshot_ttl_seconds=current_app.config[
+            "SEMANTIC_LAYER_METADATA_SNAPSHOT_TTL_SECONDS"
+        ],
+    )
+
+
+class MetadataCommand(BaseCommand):
+    """Resolve stored scope and connection-management authority in fresh 
reads."""
+
+    def __init__(self, view_uuid: UUID) -> None:
+        if not isinstance(view_uuid, UUID):
+            raise SemanticViewNotFoundError()
+        self._view_uuid: UUID = view_uuid
+        self._view: SemanticView | None = None
+        self._binding: tuple[str, str, str] | None = None
+        self._scope: str | None = None
+
+    def validate(self) -> None:
+        """Reject unavailable targets and authority before provider/cache 
work."""
+        if not metadata_refresh_enabled():
+            raise SemanticViewNotFoundError()
+        operation_deadline()
+        session: Session
+        with (
+            Session(
+                bind=db.session.get_bind(mapper=SemanticLayer), autoflush=False
+            ) as session,
+            fresh_refresh_authority(session),
+        ):
+            self._view = SemanticViewDAO.find_by_uuid(str(self._view_uuid))
+            if self._view is None:
+                raise SemanticViewNotFoundError()
+            authorize_metadata_refresh(self._view)
+            self._binding = view_binding(self._view)
+            self._scope = connection_metadata_scope(self._view.semantic_layer)
+
+    def _revalidate(self) -> None:
+        """Observe committed binding, configuration and authority before 
mutation."""
+        operation_deadline()
+        session: Session
+        with (
+            Session(
+                bind=db.session.get_bind(mapper=SemanticLayer), autoflush=False
+            ) as session,
+            fresh_refresh_authority(session),
+        ):
+            fresh: SemanticView | None = SemanticViewDAO.find_by_uuid(
+                str(self._view_uuid)
+            )
+            if (
+                fresh is None
+                or view_binding(fresh) != self._binding
+                or fresh.semantic_layer is None
+                or connection_metadata_scope(fresh.semantic_layer) != 
self._scope
+            ):
+                raise MetadataRefreshError("configuration_changed")
+            authorize_metadata_refresh(fresh)
+
+
+class RefreshMetadataCommand(MetadataCommand):
+    """Refresh the authorized view's stored connection without ORM 
mutations."""
+
+    def run(self) -> MetadataRefreshResult:
+        self.validate()
+        assert self._view is not None
+        layer: SemanticLayer = self._view.semantic_layer
+        try:
+            implementation: SemanticLayerABC[Any, SemanticViewABC] = registry[
+                layer.type
+            ].from_configuration(json.loads(layer.configuration))
+        except (ValueError, TypeError):
+            raise MetadataRefreshError("configuration") from None
+        if implementation.metadata_refresh is None:
+            raise MetadataRefreshError("unsupported")
+        store: ScopedMetadataStore = guarded_store(
+            self._view, before_publish=self._revalidate
+        )
+        deadline: float = operation_deadline()
+        implementation.metadata_refresh.bind(store, deadline=deadline)
+        return implementation.metadata_refresh.refresh(deadline=deadline)
+
+
+class InvalidateCatalogCommand(MetadataCommand):
+    """Retire the connection catalog and any older writer without acquiring 
it."""
+
+    def run(self) -> None:
+        self.validate()
+        assert self._view is not None
+        self._revalidate()
+        guarded_store(self._view, 
before_publish=self._revalidate).invalidate_catalog()
+
+
+class InvalidateCompatibilityCommand(MetadataCommand):
+    """Retire compatibility alone, preserving catalog and result identities."""
+
+    def run(self) -> None:
+        self.validate()
+        assert self._view is not None
+        self._revalidate()
+        guarded_store(
+            self._view, before_publish=self._revalidate
+        ).invalidate_compatibility()
+
+
+class InspectCatalogCommand(MetadataCommand):
+    """Read one connection's entry timing without acquiring or renewing 
metadata."""
+
+    def run(self) -> CacheEntryInfo:
+        self.validate()
+        assert self._view is not None
+        return guarded_store(
+            self._view, before_publish=self._revalidate
+        ).inspect_catalog()
+
+
+def inspect_derived_entry(
+    key: str, kind: Literal["compatibility", "query_result"]
+) -> CacheEntryInfo:
+    """Use the configured data cache, declining unsupported transport 
overrides."""
+    config: Any = current_app.config.get("DATA_CACHE_CONFIG")

Review Comment:
   Thanks for adding the case. As written it cannot catch the swap this thread 
describes: `SEMANTIC_METADATA_TEST_REDIS_URL` is not set in any workflow under 
`.github/workflows`, so the test is skipped in every CI lane, and exchanging 
the `DISTRIBUTED_COORDINATION_CONFIG` and `DATA_CACHE_CONFIG` lookups in the 
factory would leave CI green. The real-Redis tests in `metadata_redis_test.py`, 
including the Lua freshness fence, skip for the same reason. Could the 
two-database check also run somewhere CI exercises it, for example by wiring 
that variable to a Redis service in an existing Python lane, or with a 
non-Redis test that configures two distinct caches and asserts which one each 
reader consults?



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