mikebridge commented on code in PR #44849: URL: https://github.com/apache/superset/pull/44849#discussion_r4213731095
########## 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: Confirmed the existing opt-in Redis test alone did not protect CI from that swap. Local commit [269b3b409a](https://github.com/apache/superset/commit/269b3b409a35de08302b6c73286a9b4b029d94d8) adds a service-free test through the actual inspection factory with distinct data/coordination DBs, prefixes, timestamps and TTLs for both inspection kinds. Deliberately selecting the coordination config fails both cases. This proves cache selection; it does not claim live Redis or Lua coverage. ########## superset/models/slice.py: ########## @@ -404,6 +404,28 @@ def form_data(self) -> dict[str, Any]: update_time_range(form_data) return form_data + def get_query_context_datasource(self) -> Datasource | None: + """Resolve the saved query source without constructing queries or metadata.""" + from superset.daos.datasource import DatasourceDAO + from superset.daos.exceptions import DatasourceNotFound + + if self.query_context: + try: + datasource: utils.DatasourceDict = json.loads(self.query_context)[ + "datasource" + ] + except json.JSONDecodeError: + # Match get_query_context's missing-context fallback. + return self.resolved_datasource + try: + return DatasourceDAO.get_datasource( + datasource_type=utils.DatasourceType(datasource["type"]), + database_id_or_uuid=datasource["id"], + ) + except DatasourceNotFound: + return self.resolved_datasource Review Comment: Confirmed. Store [34f90b5244](https://github.com/apache/superset/commit/34f90b5244a22e8b7e2d7be62e6953314562c9f0), inherited by API [269b3b409a](https://github.com/apache/superset/commit/269b3b409a35de08302b6c73286a9b4b029d94d8), handles malformed/null/missing saved datasource shapes and the DAO's unsupported/incorrect datasource exceptions through the existing type-guarded relationship fallback. Regressions cover both present and absent fallback relationships; unrelated SQL and access errors are not swallowed. ########## 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: Confirmed and pending design. The real deadline backend rejects a running asyncio loop before transport I/O, so entering metadata_operation alone cannot make these async tools work even with a warm cache. I reproduced that boundary without opening a connection. This round explicitly holds connection and revalidation changes; the decision is between a supported synchronous execution boundary preserving admission, authorization, cancellation and Session ownership, or an async provider/store path. No transport fix is included in API candidate [269b3b409a](https://github.com/apache/superset/commit/269b3b409a35de08302b6c73286a9b4b029d94d8). ########## 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: Confirmed the missing operation scope. Store [34f90b5244](https://github.com/apache/superset/commit/34f90b5244a22e8b7e2d7be62e6953314562c9f0), inherited by API [269b3b409a](https://github.com/apache/superset/commit/269b3b409a35de08302b6c73286a9b4b029d94d8), places compatible-dimensions discovery under the same metadata-operation scope as its siblings, with a regression that fails before the wrapper. This fixes the configuration/deadline scope; the separate running-loop transport issue in r4212699155 remains pending design. ########## 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: Confirmed. Store [34f90b5244](https://github.com/apache/superset/commit/34f90b5244a22e8b7e2d7be62e6953314562c9f0) preserves MetadataRefreshError through annotation capture instead of turning it into validation400; API [269b3b409a](https://github.com/apache/superset/commit/269b3b409a35de08302b6c73286a9b4b029d94d8) inherits it and already wraps both chart-data routes in the typed mapper. The new processor-plus-adapter regression fails400 before the fix and yields503/504/502 afterward. It checks annotation access precedes token acquisition; existing GET/POST mapping tests also pass. ########## 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: Confirmed; this shares the store timeout fix in commit [34f90b5244](https://github.com/apache/superset/commit/34f90b5244a22e8b7e2d7be62e6953314562c9f0), inherited by API [269b3b409a](https://github.com/apache/superset/commit/269b3b409a35de08302b6c73286a9b4b029d94d8). Both generation read paths recheck the deadline on a Redis failure, distinguishing spent-budget deadline from remaining-budget unavailable. The regression covers both paths and both timings. ########## 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: Confirmed the disabled flag made the old guard test ineffective. API commit [269b3b409a](https://github.com/apache/superset/commit/269b3b409a35de08302b6c73286a9b4b029d94d8) enables both configuration and feature flag in an app context without an HTTP request, and verifies no fingerprint call or captured result key. Removing the request-context guard makes it fail; the restored result-inspection file passes10 tests. -- 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]
