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


##########
superset/coordination/deadline_backend.py:
##########
@@ -0,0 +1,228 @@
+# 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.
+
+"""Private, cancellation-bounded Redis operations for synchronous metadata 
callers.
+
+Each command owns its async client and loop. Cancelling the command disconnects
+its socket; shared coordinator pools and their retry/timeout policy are 
untouched.
+"""
+
+from __future__ import annotations
+
+import asyncio
+import math
+import time
+from contextlib import AsyncExitStack
+from typing import Any
+
+from redis.asyncio import Redis
+from redis.asyncio.retry import Retry
+from redis.asyncio.sentinel import Sentinel
+from redis.backoff import NoBackoff
+from redis.exceptions import RedisError, TimeoutError as RedisTimeoutError
+from superset_core.semantic_layers.metadata import (
+    MetadataRefreshError,
+    remaining_budget,
+)
+
+from superset.coordination.cache_backend import _COMPARE_AND_DELETE_LUA
+
+_COMPARE_AND_PUBLISH_LUA: str = """
+if redis.call('get', KEYS[1]) ~= ARGV[1] then
+    return 0
+end
+redis.call('psetex', KEYS[2], ARGV[3], ARGV[2])
+redis.call('del', KEYS[1])
+return 1
+"""
+
+_GET_WITH_TTL_LUA: str = """
+return {redis.call('get', KEYS[1]), redis.call('pttl', KEYS[1])}
+"""
+
+_GET_OR_CREATE_LUA: str = """
+local value = redis.call('get', KEYS[1])
+if value then
+    return value
+end
+redis.call('set', KEYS[1], ARGV[1], 'EX', ARGV[2])
+return ARGV[1]
+"""
+
+
+class DeadlineRedisBackend:
+    """Use the coordinator configuration without sharing mutable 
connections."""
+
+    def __init__(self, config: dict[str, Any], *, deadline: float) -> None:
+        if not math.isfinite(deadline) or config.get("CACHE_TYPE") not in {
+            "RedisCache",
+            "RedisSentinelCache",
+        }:
+            raise ValueError("Unsupported metadata coordination configuration")
+        self._config: dict[str, Any] = dict(config)
+        self._deadline: float = deadline
+
+    def with_deadline(self, deadline: float) -> DeadlineRedisBackend:
+        """Create a private call budget without extending the operation 
ceiling."""
+        if not math.isfinite(deadline):
+            raise ValueError("Metadata deadline must be finite")
+        return DeadlineRedisBackend(
+            self._config, deadline=min(self._deadline, deadline)
+        )
+
+    def _remaining(self) -> float:
+        """Keep the transport's Redis error boundary while sharing SDK 
validation."""
+        try:
+            return remaining_budget(self._deadline, now=time.monotonic())
+        except MetadataRefreshError:
+            raise RedisTimeoutError("Metadata deadline invalid or expired") 
from None
+
+    async def _command(self, *args: str | int) -> Any:
+        remaining: float = self._remaining()
+        options: dict[str, Any] = {
+            "db": self._config.get("CACHE_REDIS_DB", 0),
+            "username": self._config.get("CACHE_REDIS_USER"),
+            "password": self._config.get("CACHE_REDIS_PASSWORD"),
+            "socket_timeout": remaining,
+            "socket_connect_timeout": remaining,
+            "retry": Retry(NoBackoff(), 0),
+            "protocol": 2,
+        }
+        if self._config.get("CACHE_REDIS_SSL", False):
+            options.update(
+                {
+                    "ssl": True,
+                    "ssl_certfile": 
self._config.get("CACHE_REDIS_SSL_CERTFILE"),
+                    "ssl_keyfile": self._config.get("CACHE_REDIS_SSL_KEYFILE"),
+                    "ssl_ca_certs": 
self._config.get("CACHE_REDIS_SSL_CA_CERTS"),
+                    "ssl_cert_reqs": self._config.get(
+                        "CACHE_REDIS_SSL_CERT_REQS", "required"
+                    ),
+                }
+            )
+        # One cancellation deadline covers DNS, Sentinel discovery, 
authentication,
+        # response parsing (including trickled responses) and connection 
cleanup.
+        stack: AsyncExitStack
+        async with asyncio.timeout(remaining), AsyncExitStack() as stack:
+            client: Redis
+            if self._config["CACHE_TYPE"] == "RedisSentinelCache":
+                sentinel: Sentinel = Sentinel(
+                    self._config.get("CACHE_REDIS_SENTINELS", [("127.0.0.1", 
26379)]),
+                    sentinel_kwargs={
+                        "password": 
self._config.get("CACHE_REDIS_SENTINEL_PASSWORD"),
+                        "socket_timeout": remaining,
+                        "socket_connect_timeout": remaining,
+                        "retry": Retry(NoBackoff(), 0),
+                        "protocol": 2,
+                    },
+                    **options,
+                )
+                sentinel_client: Redis
+                for sentinel_client in sentinel.sentinels:
+                    stack.push_async_callback(sentinel_client.aclose)

Review Comment:
   Thanks — fixed in b6b40bd3e8 with a follow-up in 3e4e330dc7 (both now 
pushed): prefer `aclose()` with a `close()` fallback. Both client shapes pass; 
the close-only case failed first with the reported AttributeError.



##########
superset/semantic_layers/metadata_binding.py:
##########
@@ -0,0 +1,234 @@
+# 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:
+    """Classify stored configuration without leaking parser or registry 
errors."""
+    if not metadata_refresh_enabled():
+        return False
+    try:
+        configuration: Any = json.loads(layer.configuration)
+        if not isinstance(configuration, dict):
+            raise MetadataRefreshError("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": 
json.loads(layer.configuration)},
+        sort_keys=True,
+    )
+    return metadata_scope(secret, namespace, str(layer.uuid), configuration)
+
+
+def connection_store(layer: SemanticLayer) -> ScopedMetadataStore:
+    """Resolve only stored, server-owned scope and recheck it before 
publication."""
+    state: MetadataOperation = _operation()
+    scope: str = connection_metadata_scope(layer)
+    if scope in state.stores:
+        return state.stores[scope]
+    config: Any = current_app.config.get("DISTRIBUTED_COORDINATION_CONFIG")
+    if not isinstance(config, dict):
+        raise MetadataRefreshError("unavailable")
+    try:
+        backend: DeadlineRedisBackend = DeadlineRedisBackend(
+            config, deadline=state.deadline
+        )
+    except ValueError:
+        raise MetadataRefreshError("configuration") from None
+
+    def revalidate() -> None:
+        # Avoid app-init regression: binding loads before encrypted model 
fields.
+        from superset.semantic_layers.models import SemanticLayer
+
+        if not participates(layer):
+            raise MetadataRefreshError("configuration_changed")
+        session: Session
+        with Session(
+            bind=db.session.get_bind(mapper=SemanticLayer), autoflush=False
+        ) as session:

Review Comment:
   Fixed in b6b40bd3e8 with a follow-up in 3e4e330dc7 (both now pushed): 
fresh-read SQLAlchemy errors become `unavailable` with the cause suppressed; 
configuration changes retain `configuration_changed`. The failed acquisition 
publishes no snapshot; operator logs retain a warning with exception 
information.



##########
superset/semantic_layers/models.py:
##########
@@ -700,7 +718,18 @@ def data_for_slices(self, slices: list[Any]) -> 
ExplorableData:
         return self.data
 
     def get_extra_cache_keys(self, query_obj: QueryObjectDict) -> 
list[Hashable]:
-        return []
+        token: str | None = self.metadata_cache_token
+        return [token] if token is not None else []

Review Comment:
   Fixed in b6b40bd3e8 with a follow-up in 3e4e330dc7 (both now pushed): 
metadata identity resolves through `chart.resolved_datasource`, while RLS stays 
on table-only `chart.datasource`. Real Slice/semantic_view relationships 
reproduce the stale warm-parent cache for line and table annotations, then pass 
with refreshed data; SQL/flag-off keys remain unchanged. The first version of 
this fix read the table-only `chart.datasource` and its test patched that 
attribute, so it did not take effect for a real semantic-view chart; an 
independent review caught that before push and 3e4e330dc7 corrects it.



##########
superset/semantic_layers/metadata_binding.py:
##########
@@ -0,0 +1,234 @@
+# 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:
+    """Classify stored configuration without leaking parser or registry 
errors."""
+    if not metadata_refresh_enabled():
+        return False
+    try:
+        configuration: Any = json.loads(layer.configuration)
+        if not isinstance(configuration, dict):
+            raise MetadataRefreshError("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": 
json.loads(layer.configuration)},
+        sort_keys=True,
+    )
+    return metadata_scope(secret, namespace, str(layer.uuid), configuration)
+
+
+def connection_store(layer: SemanticLayer) -> ScopedMetadataStore:
+    """Resolve only stored, server-owned scope and recheck it before 
publication."""
+    state: MetadataOperation = _operation()
+    scope: str = connection_metadata_scope(layer)
+    if scope in state.stores:
+        return state.stores[scope]
+    config: Any = current_app.config.get("DISTRIBUTED_COORDINATION_CONFIG")
+    if not isinstance(config, dict):
+        raise MetadataRefreshError("unavailable")
+    try:
+        backend: DeadlineRedisBackend = DeadlineRedisBackend(
+            config, deadline=state.deadline
+        )
+    except ValueError:
+        raise MetadataRefreshError("configuration") from None
+
+    def revalidate() -> None:
+        # Avoid app-init regression: binding loads before encrypted model 
fields.
+        from superset.semantic_layers.models import SemanticLayer
+
+        if not participates(layer):
+            raise MetadataRefreshError("configuration_changed")
+        session: Session
+        with Session(
+            bind=db.session.get_bind(mapper=SemanticLayer), autoflush=False
+        ) as session:
+            fresh: SemanticLayer | None = session.get(SemanticLayer, 
layer.uuid)
+            if fresh is None or connection_metadata_scope(fresh) != scope:
+                raise MetadataRefreshError("configuration_changed")
+
+    store: ScopedMetadataStore = ScopedMetadataStore(
+        backend, scope, deadline=state.deadline, before_publish=revalidate
+    )
+    state.stores[scope] = store
+    return store
+
+
+def layer_implementation(layer: SemanticLayer) -> LayerABC[Any, ViewABC]:
+    """Construct once per operation, binding the store before any discovery."""
+    state: MetadataOperation = _operation(require_budget=False)
+    scope: str = connection_metadata_scope(layer)
+    if scope not in state.layers:
+        operation_deadline()
+        store: ScopedMetadataStore = connection_store(layer)
+        implementation: LayerABC[Any, ViewABC] = registry[
+            layer.type
+        ].from_configuration(json.loads(layer.configuration))
+        if implementation.metadata_refresh is None:
+            raise MetadataRefreshError("configuration")
+        implementation.metadata_refresh.bind(store, deadline=state.deadline)
+        state.layers[scope] = implementation
+    return state.layers[scope]
+
+
+def view_implementation(view: SemanticView) -> ViewABC:
+    """A view captures one observation for this operation, never across 
requests."""
+    state: MetadataOperation = _operation(require_budget=False)
+    scope: str = connection_metadata_scope(view.semantic_layer)
+    key: tuple[str, str, str] = (
+        scope,
+        str(view.uuid),
+        json.dumps([view.name, json.loads(view.configuration)], 
sort_keys=True),
+    )
+    if key not in state.views:
+        operation_deadline()
+        implementation: ViewABC = layer_implementation(
+            view.semantic_layer
+        ).get_semantic_view(view.name, json.loads(view.configuration))

Review Comment:
   Fixed in b6b40bd3e8 with a follow-up in 3e4e330dc7 (both now pushed): 
participating dependent tasks carry the totals query and resolve it under their 
captured catalog. Matching cached totals are reused; a changed identity 
recomputes before normalization/caching. Red-first scheduled-task and real 
result-cache tests cover T0/T1. Deploy workers before web nodes, or 
participating tasks can fail until both are upgraded; resubmit older queued 
tasks missing totals after upgrade.



##########
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):

Review Comment:
   Addressed in b6b40bd3e8 with a follow-up in 3e4e330dc7 (both now pushed): 
parsed configuration is cached per operation by stored JSON text, with 
defensive copies for providers. Tests cover repeated reads, changed text, new 
operations and provider mutation; flag-off construction retains its existing 
cache.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to