mikebridge commented on code in PR #44849:
URL: https://github.com/apache/superset/pull/44849#discussion_r4169300085
##########
superset/common/query_context_processor.py:
##########
@@ -472,7 +479,18 @@ def _annotation_cache_context(self, query_obj:
QueryObject) -> dict[str, Any]:
if annotation_datasource
else None
)
- return {"user_id": get_user_id(), "source_rls": source_rls}
+ metadata_datasource: Datasource | None = (
+ chart.resolved_datasource if chart else None
+ )
+ if isinstance(metadata_datasource, SemanticView):
+ token: str | None = metadata_datasource.metadata_cache_token
Review Comment:
Confirmed, thanks: the annotation key could run provider discovery before
reading a warm host result. Fixed in the store layer (#44835) and merged up to
this branch at `27907475`:
- `8cf8697d` builds the key from stored metadata with no provider call,
keeps identities already captured in the request, and forces a unique miss when
freshness is unknown, so a provider failure no longer takes down a host chart
that can be served and nothing stale is served as fresh. On a miss the write
key is bound after acquisition, so a concurrent refresh cannot cache under a
superseded token.
- `b6c26325` skips the persistent result write when the keyed view still has
no captured identity after acquisition; the dataframe and annotation payload
are still returned.
Tests cover both annotation sources, all three error categories, flag-off
(keys unchanged), snapshot changes and the acquisition race.
##########
superset/semantic_layers/api.py:
##########
@@ -196,8 +233,22 @@ class SemanticViewRestApi(BaseSupersetModelRestApi):
method_permission_name = {
**MODEL_API_RW_METHOD_PERMISSION_MAP,
"structure": "read",
+ "refresh_metadata": "read",
Review Comment:
The `read` route gate is deliberate, and I should have documented it. The
routes require SemanticView read; the commands then require write on the owning
SemanticLayer, access to the object, and connection-modify authority. Mapping
the routes to `write` would add a SemanticView-write requirement that
connection managers do not otherwise need, and view-edit alone should never
grant a connection-wide sync.
`d8c74284` documents the split and tests all three POST routes: a view-read
principal reaches the command and is refused before any provider or cache work,
view-edit alone is also refused, and an authorized connection manager succeeds
without view-write. The tests fail if the mapping changes or either command
check is removed. The mapping itself is unchanged. Happy to revisit if you'd
prefer the stricter route gate as defence in depth.
##########
superset/security/manager.py:
##########
@@ -5242,7 +5248,12 @@ def has_promiscuous_chart_access() -> bool:
and dashboard_.json_metadata
and (json_metadata :=
json.loads(dashboard_.json_metadata))
and any(
- target.get("datasetId") ==
datasource.data["id"]
+ target.get("datasetId")
+ == (
+ datasource.id
Review Comment:
Added the allowed native-filter case in `9d6a1d34`, using the real payload
and target validation methods with a stored semantic-view ID. Each of the two
`id` → `-1` mutations now fails the allowed case; the denial tests stay green.
Production authorization logic is unchanged.
##########
superset/common/query_context_processor.py:
##########
@@ -446,6 +448,8 @@ def query_cache_key(self, query_obj: QueryObject, **kwargs:
Any) -> str | None:
if query_obj
else None
)
+ if cache_key is not None:
+ capture_result_identity(self._query_context, query_obj, cache_key)
Review Comment:
Fair point. Capture is an internal hook with no production caller yet; it is
kept so the first inspection caller can use the captured identity without
adding a second capture path. I'm happy to remove it from this PR if you'd
prefer it land with its caller.
`9d6a1d34` pins the flag-off and annotation exclusions with tests that fail
when either guard is removed, so those paths do no work. For eligible keys the
cost is serialization, an RLS-key read, one SHA-256 and a request-local record.
I measured 6.8–62.9 µs locally for serialization and hashing; that is a lower
bound, since it excludes the per-key RLS read and query serialization
preparation.
##########
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:
Review Comment:
`9d6a1d34` adds the opted-in / no-adapter case: the command returns
unsupported before any store work. With the guard removed the test fails on
`None.bind`, as you described. Existing HTTP mapping coverage keeps unsupported
→ 422. No production change.
--
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]