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


##########
superset/semantic_layers/metadata_binding.py:
##########
@@ -0,0 +1,273 @@
+# 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.
+
+"""Host construction and operation lifetime for optional shared metadata."""
+
+from __future__ import annotations
+
+import logging
+import math
+import time
+from collections.abc import Iterator
+from contextlib import contextmanager
+from contextvars import ContextVar, Token
+from copy import deepcopy
+from dataclasses import dataclass, field
+from typing import Any, TYPE_CHECKING
+
+from flask import current_app, has_app_context, has_request_context, request
+from sqlalchemy.exc import SQLAlchemyError
+from sqlalchemy.orm import Session
+from superset_core.semantic_layers.layer import SemanticLayer as LayerABC
+from superset_core.semantic_layers.metadata import (
+    CatalogSnapshot,
+    MetadataRefreshError,
+    remaining_budget,
+)
+from superset_core.semantic_layers.view import SemanticView as ViewABC
+
+from superset import db, is_feature_enabled
+from superset.coordination.deadline_backend import DeadlineRedisBackend
+from superset.semantic_layers.metadata import (
+    FETCH_DEADLINE_SECONDS,
+    metadata_scope,
+    ScopedMetadataStore,
+)
+from superset.semantic_layers.registry import registry
+from superset.utils import json
+
+if TYPE_CHECKING:
+    from superset.semantic_layers.models import SemanticLayer, SemanticView
+
+
+logger: logging.Logger = logging.getLogger(__name__)
+
+
+@dataclass
+class MetadataOperation:
+    """One request or worker operation; nested discovery shares its 
deadline."""
+
+    deadline: float
+    layers: dict[str, LayerABC[Any, ViewABC]] = field(default_factory=dict)
+    views: dict[tuple[str, str, str], ViewABC] = field(default_factory=dict)
+    stores: dict[str, ScopedMetadataStore] = field(default_factory=dict)
+    configurations: dict[str, dict[str, Any]] = field(default_factory=dict)
+
+
+_OPERATION_KEY: str = "superset.semantic_metadata.operation"
+_worker_operation: ContextVar[MetadataOperation | None] = ContextVar(
+    _OPERATION_KEY, default=None
+)
+
+
+def request_metadata_budget() -> None:
+    """Register before authentication hooks; this performs no provider or 
cache I/O."""
+    if current_app.config.get("SEMANTIC_LAYER_METADATA_REFRESH_ENABLED") is 
True:
+        request.environ.setdefault(
+            _OPERATION_KEY, MetadataOperation(time.monotonic() + 
FETCH_DEADLINE_SECONDS)
+        )
+
+
+def _current_operation() -> MetadataOperation | None:
+    if has_request_context():
+        return request.environ.get(_OPERATION_KEY) or _worker_operation.get()
+    return _worker_operation.get()
+
+
+def _operation(*, require_budget: bool = True) -> MetadataOperation:
+    state: MetadataOperation | None = _current_operation()
+    if state is None or not math.isfinite(state.deadline):
+        raise MetadataRefreshError("configuration")
+    if require_budget:
+        remaining_budget(state.deadline, now=time.monotonic())
+    return state
+
+
+def operation_deadline() -> float:
+    return _operation().deadline
+
+
+@contextmanager
+def metadata_operation(*, deadline: float | None = None) -> Iterator[None]:
+    """Workers opt in before access checks; nested calls never replenish the 
budget."""
+    if deadline is not None and not math.isfinite(deadline):
+        raise MetadataRefreshError("configuration")
+    if deadline is not None:
+        remaining_budget(deadline, now=time.monotonic())
+    if _current_operation() is not None:
+        _operation(require_budget=False)
+        yield
+        return
+    if has_request_context():
+        # HTTP requests must enter through the registered early request hook.
+        raise MetadataRefreshError("configuration")
+    ceiling: float = time.monotonic() + FETCH_DEADLINE_SECONDS
+    state: MetadataOperation = MetadataOperation(
+        min(deadline, ceiling) if deadline is not None else ceiling
+    )
+    token: Token[MetadataOperation | None] = _worker_operation.set(state)
+    try:
+        operation_deadline()
+        yield
+    finally:
+        _worker_operation.reset(token)
+
+
+def metadata_refresh_enabled() -> bool:
+    return (
+        has_app_context()
+        and current_app.config.get("SEMANTIC_LAYER_METADATA_REFRESH_ENABLED") 
is True
+        and is_feature_enabled("SEMANTIC_LAYERS")
+    )
+
+
+def _configuration(raw: str) -> dict[str, Any]:
+    """Cache parsing by stored bytes; providers receive independent mutable 
copies."""
+    state: MetadataOperation | None = _current_operation()
+    if state is not None and raw in state.configurations:
+        return deepcopy(state.configurations[raw])
+    parsed: Any = json.loads(raw)
+    if not isinstance(parsed, dict):
+        raise MetadataRefreshError("configuration")
+    if state is not None:
+        state.configurations[raw] = parsed
+    return deepcopy(parsed)
+
+
+def participates(layer: SemanticLayer) -> bool:
+    """Classify stored configuration without leaking parser or registry 
errors."""
+    if not metadata_refresh_enabled():
+        return False
+    try:
+        configuration: dict[str, Any] = _configuration(layer.configuration)
+        return registry[layer.type].supports_metadata_refresh(configuration)
+    except (KeyError, TypeError, ValueError):
+        raise MetadataRefreshError("configuration") from None
+
+
+def connection_metadata_scope(layer: SemanticLayer) -> str:
+    namespace: Any = 
current_app.config.get("SEMANTIC_LAYER_METADATA_NAMESPACE")
+    if callable(namespace):
+        namespace = namespace()
+    secret: Any = current_app.config.get("SECRET_KEY")
+    if isinstance(secret, bytes):
+        secret = secret.decode()
+    if (
+        not isinstance(namespace, str)
+        or not isinstance(secret, str)
+        or layer.uuid is None
+    ):
+        raise MetadataRefreshError("configuration")
+    configuration: str = json.dumps(
+        {"provider": layer.type, "configuration": 
_configuration(layer.configuration)},
+        sort_keys=True,
+    )
+    return metadata_scope(secret, namespace, str(layer.uuid), configuration)
+
+
+def connection_store(layer: SemanticLayer) -> ScopedMetadataStore:
+    """Resolve only stored, server-owned scope and recheck it before 
publication."""
+    state: MetadataOperation = _operation()
+    scope: str = connection_metadata_scope(layer)
+    if scope in state.stores:
+        return state.stores[scope]
+    config: Any = current_app.config.get("DISTRIBUTED_COORDINATION_CONFIG")
+    if not isinstance(config, dict):
+        raise MetadataRefreshError("unavailable")
+    try:
+        backend: DeadlineRedisBackend = DeadlineRedisBackend(
+            config, deadline=state.deadline
+        )
+    except ValueError:
+        raise MetadataRefreshError("configuration") from None
+
+    def revalidate() -> None:
+        # Avoid app-init regression: binding loads before encrypted model 
fields.
+        from superset.semantic_layers.models import SemanticLayer
+
+        if not participates(layer):
+            raise MetadataRefreshError("configuration_changed")
+        session: Session
+        try:
+            with Session(
+                bind=db.session.get_bind(mapper=SemanticLayer), autoflush=False
+            ) as session:
+                fresh: SemanticLayer | None = session.get(SemanticLayer, 
layer.uuid)
+                if fresh is None or connection_metadata_scope(fresh) != scope:
+                    raise MetadataRefreshError("configuration_changed")
+        except SQLAlchemyError:
+            logger.warning("Metadata layer revalidation failed", exc_info=True)
+            raise MetadataRefreshError("unavailable") from None
+
+    store: ScopedMetadataStore = ScopedMetadataStore(
+        backend, scope, deadline=state.deadline, before_publish=revalidate
+    )
+    state.stores[scope] = store
+    return store
+
+
+def layer_implementation(layer: SemanticLayer) -> LayerABC[Any, ViewABC]:
+    """Construct once per operation, binding the store before any discovery."""
+    state: MetadataOperation = _operation(require_budget=False)
+    scope: str = connection_metadata_scope(layer)
+    if scope not in state.layers:
+        operation_deadline()
+        store: ScopedMetadataStore = connection_store(layer)
+        implementation: LayerABC[Any, ViewABC] = registry[
+            layer.type
+        ].from_configuration(_configuration(layer.configuration))
+        if implementation.metadata_refresh is None:
+            raise MetadataRefreshError("configuration")
+        implementation.metadata_refresh.bind(store, deadline=state.deadline)
+        state.layers[scope] = implementation
+    return state.layers[scope]
+
+
+def _view_key(view: SemanticView) -> tuple[str, str, str]:
+    """Identify a configured view within one metadata operation."""
+    return (
+        connection_metadata_scope(view.semantic_layer),
+        str(view.uuid),
+        json.dumps([view.name, _configuration(view.configuration)], 
sort_keys=True),
+    )
+
+
+def peek_view_metadata_token(view: SemanticView) -> str | None:
+    """Use a captured view or stored snapshot without constructing a 
provider."""
+    state: MetadataOperation = _operation(require_budget=False)
+    captured: ViewABC | None = state.views.get(_view_key(view))
+    if captured is not None:
+        # An in-flight annotation query must retain its own observation even if
+        # another request publishes a newer catalog before its host key is 
built.
+        return captured.metadata_cache_token
+    snapshot: CatalogSnapshot | None = 
connection_store(view.semantic_layer).peek()
+    return snapshot.cache_token if snapshot is not None else None

Review Comment:
   Clarified in `634ada11`: the provider view must return the exact captured 
`CatalogSnapshot.cache_token`, not a derived per-view token. The host adds 
stored-view identity separately. The binding already rejects missing, unknown, 
forged, expired or other-scope tokens through the operation store’s provenance 
check; the existing parameterized tests exercise those cases.



##########
superset/semantic_layers/result_inspection.py:
##########
@@ -0,0 +1,92 @@
+# 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
+
+import hashlib
+from dataclasses import dataclass, field
+from typing import TYPE_CHECKING
+
+from flask import g, has_request_context, request
+
+from superset import security_manager
+from superset.semantic_layers.metadata_binding import metadata_refresh_enabled
+from superset.utils import json
+from superset.utils.json import json_int_dttm_ser
+
+if TYPE_CHECKING:
+    from superset.common.query_context import QueryContext
+    from superset.common.query_object import QueryObject
+
+_IDENTITY_KEY: str = "superset.semantic_metadata.result_identities"
+
+
+@dataclass(frozen=True)
+class CapturedResultIdentity:
+    """A private result key and its request-local authorization/query 
fingerprint."""
+
+    key: str = field(repr=False)
+    fingerprint: str = field(repr=False)
+
+
+def _fingerprint(context: QueryContext, query: QueryObject) -> str:
+    """Bind an existing key to its subject, stored view, query and RLS 
scope."""
+    return hashlib.sha256(
+        json.dumps(
+            [
+                id(getattr(g, "user", None)),
+                str(getattr(context.datasource, "uuid", None)),
+                context.datasource.changed_on,
+                query.to_dict(),

Review Comment:
   Fixed in `925580bc`. The diagnostic fingerprint excludes runtime 
`contribution_totals`, matching the query cache key’s treatment, without 
mutating the query or its options. The regressions exercise the real query-key 
capture path with NumPy float32 and int32 totals, confirm changing totals does 
not invalidate the captured key, and confirm changing query inputs still does.



##########
superset/semantic_layers/result_inspection.py:
##########
@@ -0,0 +1,92 @@
+# 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
+
+import hashlib
+from dataclasses import dataclass, field
+from typing import TYPE_CHECKING
+
+from flask import g, has_request_context, request
+
+from superset import security_manager
+from superset.semantic_layers.metadata_binding import metadata_refresh_enabled
+from superset.utils import json
+from superset.utils.json import json_int_dttm_ser
+
+if TYPE_CHECKING:
+    from superset.common.query_context import QueryContext
+    from superset.common.query_object import QueryObject
+
+_IDENTITY_KEY: str = "superset.semantic_metadata.result_identities"
+
+
+@dataclass(frozen=True)
+class CapturedResultIdentity:
+    """A private result key and its request-local authorization/query 
fingerprint."""
+
+    key: str = field(repr=False)
+    fingerprint: str = field(repr=False)
+
+
+def _fingerprint(context: QueryContext, query: QueryObject) -> str:
+    """Bind an existing key to its subject, stored view, query and RLS 
scope."""
+    return hashlib.sha256(
+        json.dumps(
+            [
+                id(getattr(g, "user", None)),

Review Comment:
   Added both cases in `925580bc`: after capturing a key, the test replaces 
g.user or changes the retained query’s metrics within the same request while 
keeping RLS unchanged. Both must return no captured key. The existing 
implementation already rejects those changes; the tests now pin the two 
fingerprint terms explicitly.



##########
superset/commands/semantic_layer/refresh_metadata.py:
##########
@@ -0,0 +1,341 @@
+# 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,

Review Comment:
   Fixed in `925580bc`. The command factory passes 
`SEMANTIC_LAYER_METADATA_SNAPSHOT_TTL_SECONDS` to the store just like ordinary 
discovery. Tests build the real factory and check both catalog publication and 
compatibility-generation invalidation at 10, 60 and 3600 seconds; each 
previously used 300 seconds and now respects the configured expiry.



##########
docs/developer_docs/semantic-metadata-store.md:
##########
@@ -0,0 +1,286 @@
+---
+title: Shared semantic metadata storage
+---
+
+<!--
+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.
+-->
+
+## Enablement and scope
+
+This host implementation supports the optional [SDK metadata 
contract](./semantic-metadata-contract.md).
+It provides storage and cache identity; it does not add refresh endpoints or 
UI.

Review Comment:
   Updated in `925580bc`: both passages link to the metadata operations page 
and describe the shipped refresh, invalidation and inspection endpoints. The 
API layer retains only the no-UI limitation; `58e8cbb1` updates that 
description for the UI layer’s controls.



##########
superset/commands/semantic_layer/refresh_metadata.py:
##########
@@ -0,0 +1,341 @@
+# 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,
+    )
+
+
+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")
+    backend: DeadlineRedisBackend | None = None
+    if (
+        isinstance(config, dict)
+        and not config.get("CACHE_REDIS_URL")
+        and not config.get("CACHE_OPTIONS")
+    ):
+        try:
+            backend = DeadlineRedisBackend(config, 
deadline=operation_deadline())
+        except ValueError:
+            pass
+    return inspect_data_cache(cache_manager.data_cache, key, kind, 
backend=backend)
+
+
+class InspectCompatibilityCommand(MetadataCommand):
+    """Inspect one normalized selection with connection-management 
authority."""
+
+    def __init__(
+        self, view_uuid: UUID, metrics: list[str], dimensions: list[str]
+    ) -> None:
+        super().__init__(view_uuid)
+        self._metrics: list[str] = metrics
+        self._dimensions: list[str] = dimensions
+
+    def run(self) -> CacheEntryInfo:
+        self.validate()
+        assert self._view is not None
+        identity: CompatibilityIdentity | None = compatibility_identity(
+            self._view, self._metrics, self._dimensions, inspection=True
+        )
+        if identity is None:
+            return CacheEntryInfo(
+                "compatibility", "missing", datetime.now(timezone.utc)
+            )
+        return inspect_derived_entry(identity.key, "compatibility")

Review Comment:
   Added in `925580bc`. The test populates the real derived-cache key for 
metrics ["orders"] and dimensions ["country"], then calls cache_metadata 
through the real route, inspection command and cache reader. It requires 
present state and stored timestamps, checks that the metadata store is 
unchanged, and verifies provider construction never occurs. Database 
authorization fixtures and the in-memory cache backend remain controlled test 
boundaries.



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