mikebridge commented on code in PR #44849: URL: https://github.com/apache/superset/pull/44849#discussion_r4159903662
########## superset/semantic_layers/metadata_cache.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. + +"""Captured catalog identity for derived caches and read-only inspection.""" + +from __future__ import annotations + +import hashlib +from dataclasses import dataclass, field +from typing import TYPE_CHECKING + +from superset_core.semantic_layers.metadata import CatalogSnapshot, MetadataRefreshError + +from superset.semantic_layers.metadata import ScopedMetadataStore +from superset.semantic_layers.metadata_binding import connection_store +from superset.utils import json + +if TYPE_CHECKING: + from superset.semantic_layers.models import SemanticView + + +@dataclass(frozen=True) +class CompatibilityIdentity: + key: str = field(repr=False) + source_observed_at: str | None + + +def view_cache_token(view: SemanticView, token: str) -> str: + """Include host view configuration without revealing it in cache keys.""" + identity: str = json.dumps( + [token, str(view.uuid), view.name, json.loads(view.configuration)], + sort_keys=True, + separators=(",", ":"), + allow_nan=False, + ) + return hashlib.sha256(identity.encode()).hexdigest() + + +def compatibility_identity( + view: SemanticView, + metrics: list[str], + dimensions: list[str], + *, + inspection: bool = False, +) -> CompatibilityIdentity | None: + """Capture identity before lookup; inspection never discovers or initializes.""" + store: ScopedMetadataStore = connection_store(view.semantic_layer) + token: str + generation: str | None + observed: str | None + if inspection: + snapshot: CatalogSnapshot | None = store.peek() + generation = store.peek_compatibility_generation() + if snapshot is None or generation is None: + return None + token, observed = snapshot.cache_token, snapshot.observed_at + else: + captured: str | None = view.implementation.metadata_cache_token + if not captured: + raise MetadataRefreshError("configuration") + token = captured + observed = store.observed_at(token) + generation = store.compatibility_generation() Review Comment: Addressed in the store layer, 74a604376434379748549c4d7f6e3857e4deedd4, and merged forward: compatibility generation is captured before resolving the provider view. A deterministic interleaving clears compatibility during view resolution and proves the old observation keeps its old namespace. Callers must not pre-resolve the view before this capture. ########## superset/commands/semantic_layer/refresh_metadata.py: ########## @@ -0,0 +1,332 @@ +# 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 +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, +) +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() + provider: type[SemanticLayerABC[Any, SemanticViewABC]] | None = registry.get( + layer.type + ) + if provider is None: + raise MetadataRefreshError("unsupported") + try: + configuration: dict[str, Any] = json.loads(layer.configuration) + supported: bool = provider.supports_metadata_refresh(configuration) + except (ValueError, TypeError): + raise MetadataRefreshError("configuration") from None + if not supported: + 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() + try: + db.session.registry.set(session) Review Comment: Addressed in ac2497f6b92619a707677c571f8c6e83d9c003fe: fresh authority checks isolate the request subject cache alongside the fresh session/principal, then restore the caller's cache and principal. Tests revoke extra-editorship membership during the fetch and prove publication is rejected without changing the catalog or compatibility generation. Both warmed and absent request caches are covered; persisted subject queries are mocked, while the editorship policy, memoized helper and publication guard are real. ########## tests/unit_tests/semantic_layers/refresh_metadata_command_test.py: ########## @@ -0,0 +1,418 @@ +# 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 typing import cast +from unittest.mock import MagicMock, Mock, patch +from uuid import UUID + +import pytest +from flask import Flask, g +from superset_core.semantic_layers.metadata import ( + MetadataRefreshError, + MetadataRefreshResult, +) + +MODULE: str = "superset.commands.semantic_layer.refresh_metadata" +VIEW_UUID: UUID = UUID("bd2f07da-c65e-40da-b75e-c62b7cdd67f1") + + [email protected] +def refresh_context(app: Flask) -> Iterator[tuple[Mock, Mock, Mock]]: + """Authorize synthetic records; provider construction is always observable.""" + from superset.commands.semantic_layer import refresh_metadata as module + + layer: Mock = Mock(uuid="connection", type="test", configuration='{"token":"test"}') + view: Mock = Mock( + uuid=VIEW_UUID, semantic_layer_uuid="connection", configuration="{}" + ) + view.name = "full" + view.semantic_layer = layer + provider: Mock = Mock() + provider.supports_metadata_refresh.return_value = True + manager: Mock = Mock() + with ( + app.app_context(), + patch.object( + g, + "user", + Mock(id=1, is_anonymous=False, is_guest_user=False), + create=True, + ), + patch.object(module, "Session", return_value=MagicMock()), + patch.object(module, "operation_deadline", return_value=130.0), + patch.object(module, "metadata_refresh_enabled", return_value=True), + patch.object(module, "security_manager", manager), + patch.object(module, "current_user_can_modify_object", return_value=True), + patch.object(module, "connection_metadata_scope", return_value="scope"), + patch.dict(module.registry, {"test": provider}), + patch.object(module.SemanticViewDAO, "find_by_uuid", return_value=view), + patch.object(module, "guarded_store"), + patch.object(module.db, "session", MagicMock()), + ): + cast( + Mock, module.Session + ).return_value.__enter__.return_value.get.return_value = Mock( + id=1, is_anonymous=False, is_guest_user=False, is_active=True + ) + yield view, provider, manager + + +def test_refresh_uses_stored_connection_after_all_authority_checks( + refresh_context: tuple[Mock, Mock, Mock], +) -> None: + from superset.commands.semantic_layer import refresh_metadata as module + from superset.commands.semantic_layer.refresh_metadata import RefreshMetadataCommand + + view: Mock + provider: Mock + manager: Mock + view, provider, manager = refresh_context + result: MetadataRefreshResult = RefreshMetadataCommand(VIEW_UUID).run() + provider.from_configuration.assert_called_once_with({"token": "test"}) + adapter: Mock = provider.from_configuration.return_value.metadata_refresh + assert result is adapter.refresh.return_value + adapter.bind.assert_called_once_with( + cast(Mock, module.guarded_store).return_value, + deadline=130.0, + ) + adapter.refresh.assert_called_once_with(deadline=130.0) + manager.can_access.assert_any_call("can_read", "SemanticView") + manager.can_access.assert_any_call("can_read", "SemanticLayer") + manager.can_access.assert_any_call("can_write", "SemanticLayer") + view.raise_for_access.assert_called_once() + view.semantic_layer.raise_for_access.assert_called_once() + + [email protected]( + "permission", + [ + ("can_read", "SemanticView"), + ("can_read", "SemanticLayer"), + ("can_write", "SemanticLayer"), + ], +) +def test_denied_permissions_do_no_provider_or_cache_work( + refresh_context: tuple[Mock, Mock, Mock], + permission: tuple[str, str], +) -> None: + from superset.commands.semantic_layer import refresh_metadata as module + from superset.commands.semantic_layer.exceptions import SemanticLayerForbiddenError + + view: Mock + provider: Mock + manager: Mock + view, provider, manager = refresh_context + manager.can_access.side_effect = ( + lambda action, resource: (action, resource) != permission + ) + with pytest.raises(SemanticLayerForbiddenError): + module.RefreshMetadataCommand(VIEW_UUID).run() + assert module.can_refresh_metadata(view) is False + provider.from_configuration.assert_not_called() + cast(Mock, module.guarded_store).assert_not_called() + + +def test_creator_without_edit_authority_cannot_refresh( + refresh_context: tuple[Mock, Mock, Mock], +) -> None: + from superset.commands.semantic_layer import refresh_metadata as module + from superset.commands.semantic_layer.exceptions import SemanticLayerForbiddenError + + view: Mock + provider: Mock + view, provider, _ = refresh_context + with patch.object(module, "current_user_can_modify_object", return_value=False): + with pytest.raises(SemanticLayerForbiddenError): + module.RefreshMetadataCommand(VIEW_UUID).run() + assert module.can_refresh_metadata(view) is False + provider.from_configuration.assert_not_called() + + +def test_disabled_refresh_is_unavailable_without_discovery( + refresh_context: tuple[Mock, Mock, Mock], +) -> None: + from superset.commands.semantic_layer import refresh_metadata as module + from superset.commands.semantic_layer.exceptions import SemanticViewNotFoundError + + view: Mock + provider: Mock + view, provider, _ = refresh_context + with patch.object(module, "metadata_refresh_enabled", return_value=False): + with pytest.raises(SemanticViewNotFoundError): + module.RefreshMetadataCommand(VIEW_UUID).run() + assert module.can_refresh_metadata(view) is False + provider.from_configuration.assert_not_called() + + +def test_unsupported_provider_is_not_constructed( + refresh_context: tuple[Mock, Mock, Mock], +) -> None: + from superset.commands.semantic_layer import refresh_metadata as module + + view: Mock + provider: Mock + view, provider, _ = refresh_context + provider.supports_metadata_refresh.return_value = False + with pytest.raises(MetadataRefreshError, match="unsupported"): + module.RefreshMetadataCommand(VIEW_UUID).run() + assert module.can_refresh_metadata(view) is False + provider.from_configuration.assert_not_called() + + [email protected]( + "changed", ["missing", "connection", "configuration", "name", "authority"] +) +def test_publication_revalidates_view_in_supplied_fresh_session( + refresh_context: tuple[Mock, Mock, Mock], + changed: str, +) -> None: + from superset.commands.semantic_layer import refresh_metadata as module + from superset.commands.semantic_layer.exceptions import SemanticLayerForbiddenError + + view: Mock + provider: Mock + manager: Mock + view, provider, manager = refresh_context + module.RefreshMetadataCommand(VIEW_UUID).run() + guard: Callable[[], None] = cast(Mock, module.guarded_store).call_args.kwargs[ + "before_publish" + ] + fresh: Mock = Mock( + uuid=VIEW_UUID, semantic_layer_uuid="connection", configuration="{}" + ) + fresh.name = "full" + fresh.semantic_layer = view.semantic_layer + session: Mock = cast(Mock, module.Session).return_value.__enter__.return_value + cast(Mock, module.SemanticViewDAO.find_by_uuid).return_value = fresh + if changed == "missing": + cast(Mock, module.SemanticViewDAO.find_by_uuid).return_value = None + elif changed == "connection": + fresh.semantic_layer_uuid = "another" + elif changed == "configuration": + fresh.configuration = '{"metrics": ["old"]}' + elif changed == "name": + fresh.name = "another" + else: + manager.can_access.return_value = False + with pytest.raises((MetadataRefreshError, SemanticLayerForbiddenError)): + guard() + session.commit.assert_not_called() + session.rollback.assert_not_called() + + +def test_incomplete_configuration_is_sanitized_before_construction( + refresh_context: tuple[Mock, Mock, Mock], +) -> None: + from superset.commands.semantic_layer import refresh_metadata as module + + view: Mock + provider: Mock + manager: Mock + view, provider, manager = refresh_context + provider.supports_metadata_refresh.side_effect = ValueError("private configuration") + with pytest.raises(MetadataRefreshError, match="^configuration$"): + module.RefreshMetadataCommand(VIEW_UUID).run() + assert module.can_refresh_metadata(view) is False + provider.from_configuration.assert_not_called() + cast(Mock, module.guarded_store).assert_not_called() + + [email protected]("resource", ["view", "layer"]) +def test_datasource_access_denial_matches_capability_and_prevents_construction( + refresh_context: tuple[Mock, Mock, Mock], resource: str +) -> None: + from superset.commands.semantic_layer import refresh_metadata as module + from superset.commands.semantic_layer.exceptions import SemanticLayerForbiddenError + from superset.errors import ErrorLevel, SupersetError, SupersetErrorType + from superset.exceptions import SupersetSecurityException + + view: Mock + provider: Mock + manager: Mock + view, provider, manager = refresh_context + target: Mock = view if resource == "view" else view.semantic_layer + target.raise_for_access.side_effect = SupersetSecurityException( + SupersetError( + message="controlled denial", + error_type=SupersetErrorType.DATASOURCE_SECURITY_ACCESS_ERROR, + level=ErrorLevel.ERROR, + ) + ) + with pytest.raises(SemanticLayerForbiddenError): + module.RefreshMetadataCommand(VIEW_UUID).run() + assert module.can_refresh_metadata(view) is False + provider.from_configuration.assert_not_called() + cast(Mock, module.guarded_store).assert_not_called() + + [email protected]("target", ["view", "layer"]) +def test_missing_stored_target_does_not_construct_provider( + refresh_context: tuple[Mock, Mock, Mock], target: str +) -> None: + from superset.commands.semantic_layer import refresh_metadata as module + from superset.commands.semantic_layer.exceptions import ( + SemanticLayerNotFoundError, + SemanticViewNotFoundError, + ) + + view: Mock + provider: Mock + manager: Mock + view, provider, manager = refresh_context + if target == "view": + cast(Mock, module.SemanticViewDAO.find_by_uuid).return_value = None + else: + view.semantic_layer = None + with pytest.raises((SemanticViewNotFoundError, SemanticLayerNotFoundError)): + module.RefreshMetadataCommand(VIEW_UUID).run() + provider.from_configuration.assert_not_called() + cast(Mock, module.guarded_store).assert_not_called() + + [email protected]("principal", ["anonymous", "guest", "inactive"]) +def test_ineligible_principal_has_no_capability_or_provider_work( + refresh_context: tuple[Mock, Mock, Mock], principal: str +) -> None: + from superset.commands.semantic_layer import refresh_metadata as module + from superset.commands.semantic_layer.exceptions import SemanticLayerForbiddenError + + view: Mock + provider: Mock + manager: Mock + view, provider, manager = refresh_context + user: Mock = Mock( + id=1, + is_anonymous=principal == "anonymous", + is_guest_user=principal == "guest", + is_active=principal != "inactive", + ) + cast( + Mock, module.Session + ).return_value.__enter__.return_value.get.return_value = user + with patch.object(g, "user", user): + with pytest.raises(SemanticLayerForbiddenError): + module.RefreshMetadataCommand(VIEW_UUID).run() + assert module.can_refresh_metadata(view) is False + provider.from_configuration.assert_not_called() + cast(Mock, module.guarded_store).assert_not_called() + + +def test_unregistered_provider_has_no_capability_or_construction( + refresh_context: tuple[Mock, Mock, Mock], +) -> None: + from superset.commands.semantic_layer import refresh_metadata as module + + view: Mock + provider: Mock + manager: Mock + view, provider, manager = refresh_context + with patch.dict(module.registry, {}, clear=True): + with pytest.raises(MetadataRefreshError, match="^unsupported$"): + module.RefreshMetadataCommand(VIEW_UUID).run() + assert module.can_refresh_metadata(view) is False + provider.from_configuration.assert_not_called() + cast(Mock, module.guarded_store).assert_not_called() + + [email protected]("namespace", [None, "", 7]) +def test_missing_trusted_namespace_denies_before_provider_construction( + refresh_context: tuple[Mock, Mock, Mock], namespace: object +) -> None: + from flask import current_app + + from superset.commands.semantic_layer import refresh_metadata as module + from superset.semantic_layers.metadata_binding import connection_metadata_scope + + view: Mock + provider: Mock + manager: Mock + view, provider, manager = refresh_context + with ( + patch.object(module, "connection_metadata_scope", connection_metadata_scope), + patch.dict( + current_app.config, {"SEMANTIC_LAYER_METADATA_NAMESPACE": namespace} + ), + ): + with pytest.raises(MetadataRefreshError, match="^configuration$"): + module.RefreshMetadataCommand(VIEW_UUID).run() + assert module.can_refresh_metadata(view) is False + provider.from_configuration.assert_not_called() + cast(Mock, module.guarded_store).assert_not_called() + + +def test_capability_projection_does_not_construct_or_acquire_metadata( + refresh_context: tuple[Mock, Mock, Mock], +) -> None: + from superset.commands.semantic_layer import refresh_metadata as module + + view: Mock + provider: Mock + manager: Mock + view, provider, manager = refresh_context + assert module.can_refresh_metadata(view) is True + provider.supports_metadata_refresh.assert_called_once_with({"token": "test"}) + provider.from_configuration.assert_not_called() + cast(Mock, module.guarded_store).assert_not_called() + + [email protected]( + "command_name,method", + [ + ("InvalidateCatalogCommand", "invalidate_catalog"), + ("InvalidateCompatibilityCommand", "invalidate_compatibility"), + ("InspectCatalogCommand", "inspect_catalog"), + ], +) +def test_separate_controls_never_construct_a_provider( + refresh_context: tuple[Mock, Mock, Mock], + command_name: str, + method: str, +) -> None: + from superset.commands.semantic_layer import refresh_metadata as module + + view: Mock + provider: Mock + manager: Mock + view, provider, manager = refresh_context + getattr(module, command_name)(VIEW_UUID).run() Review Comment: Added in ac2497f6b92619a707677c571f8c6e83d9c003fe: both clear commands reject authority revocation or binding changes between validation and clearing. The tests use the real store with an in-memory backend and assert neither catalog nor compatibility generation changes. Existing production revalidation already passes these cases. ########## superset/common/query_context_factory.py: ########## @@ -114,6 +124,12 @@ def create( # pylint: disable=too-many-arguments custom_cache_timeout=custom_cache_timeout, cache_values=cache_values, ) + if defer_discovery: + security_manager.raise_for_access(query_context=context) Review Comment: I could not reproduce the reported data-access bypass. The early check is additive: ChartDataCommand.validate still authorizes the completed query, including tooltip columns, before execution. ac2497f6b92619a707677c571f8c6e83d9c003fe adds a regression using the real factory and guest payload comparator with refresh enabled and disabled; an unshared tooltip column is rejected in both cases. This is a resolved-chart fixture, not evidence of persisted semantic guest-chart support. ########## superset/semantic_layers/models.py: ########## @@ -257,9 +259,18 @@ def after_delete( security_manager.semantic_layer_after_delete(mapper, connection, target) - @cached_property + @property def implementation( self, + ) -> SemanticLayerABC[Any, SemanticViewABC]: + if metadata_binding.participates(self): + # Read authorization belongs to callers with full request context. + return metadata_binding.layer_implementation(self) + return self._legacy_implementation + + @cached_property + def _legacy_implementation( + self, ) -> SemanticLayerABC[Any, SemanticViewABC]: Review Comment: The cached property is the legacy provider path, whose caching behavior is unchanged from pinned master. Participating providers are scoped by stored configuration; the regression in 74a604376434379748549c4d7f6e3857e4deedd4 changes configuration on the same ORM object within one operation and verifies a new provider, while unchanged configuration reuses it. I preserved the legacy path's existing behavior rather than changing the flag-off contract in this stack. ########## superset/semantic_layers/metadata_binding.py: ########## @@ -0,0 +1,227 @@ +# 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 math +import time +from collections.abc import Iterator +from contextlib import contextmanager +from contextvars import ContextVar, Token +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.orm import Session +from superset_core.semantic_layers.layer import SemanticLayer as LayerABC +from superset_core.semantic_layers.metadata import ( + 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 + + +@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) + + +_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 participates(layer: SemanticLayer) -> bool: + return metadata_refresh_enabled() and registry[ + layer.type + ].supports_metadata_refresh(json.loads(layer.configuration)) Review Comment: Addressed in the store layer, 74a604376434379748549c4d7f6e3857e4deedd4, and merged forward: participates normalizes unknown providers, malformed JSON and non-object configurations to MetadataRefreshError(configuration). Red-first tests cover each case; the API layer maps the typed category to its established safe response. -- 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]
