sadpandajoe commented on code in PR #44835: URL: https://github.com/apache/superset/pull/44835#discussion_r4212199666
########## superset/semantic_layers/metadata_binding.py: ########## @@ -0,0 +1,319 @@ +# 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 +) + +_worker_chart: ContextVar[bool] = ContextVar( + "superset.semantic_metadata.chart", default=False +) + + +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 _worker_chart.get(): + return _worker_operation.get() + 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) + + +@contextmanager +def chart_metadata_operation(*, allow_request: bool = False) -> Iterator[None]: + """Give a worker/export chart a fresh budget; nested chart work shares it.""" + if ( + (has_request_context() and not allow_request) + or _worker_chart.get() + or not metadata_refresh_enabled() + ): + yield + return + # A task may have spent its fallback budget on earlier charts or other work. + # Restore that state after this chart, including on cancellation or failure. + operation_token: Token[MetadataOperation | None] = _worker_operation.set( + MetadataOperation(time.monotonic() + FETCH_DEADLINE_SECONDS) + ) + chart_token: Token[bool] = _worker_chart.set(True) + try: + operation_deadline() + yield + finally: + _worker_chart.reset(chart_token) + _worker_operation.reset(operation_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(require_budget=False) + scope: str = connection_metadata_scope(layer) + if scope in state.stores: + return state.stores[scope] + 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=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( Review Comment: Revalidation opens a fresh `Session` on the metastore engine while the request's `db.session` still holds its connection, and the checkout isn't bounded by the metadata deadline. With a small pool (for example `pool_size=1, max_overflow=0`) or a burst of cold-catalog requests, this waits for a connection the caller itself is holding until the pool timeout, so publication fails and the catalog stays cold. Could this reuse the caller's connection, or check out under the remaining budget? ########## 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 Review Comment: `SoftTimeLimitExceeded` subclasses `Exception`, so this converts it into `MetadataRefreshError("upstream")` when the soft limit fires inside a provider fetch. The Excel export loop re-raises `SoftTimeLimitExceeded` specifically but swallows other exceptions, so the chart is logged as skipped and the export keeps running until the hard limit kills the worker. Should this let `SoftTimeLimitExceeded` propagate? -- 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]
