sadpandajoe commented on code in PR #44849: URL: https://github.com/apache/superset/pull/44849#discussion_r4172126409
########## superset/coordination/deadline_backend.py: ########## @@ -0,0 +1,343 @@ +# 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 os +import time +from _thread import LockType +from contextlib import AsyncExitStack +from typing import Any, TYPE_CHECKING + +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 + +if TYPE_CHECKING: + from threading import local + + from gevent.threadpool import ThreadPool + +_METADATA_POOL_SIZE: int = 4 +_metadata_pool_state: local | None = None + + +def _metadata_threadpool() -> ThreadPool: + """Reuse a lazy metadata-only pool in its owning native thread and hub.""" + from gevent import get_hub + from gevent.hub import Hub + from gevent.monkey import get_original + from gevent.threadpool import ThreadPool + + global _metadata_pool_state # pylint: disable=global-statement + if _metadata_pool_state is None: + # Patched threading.local is greenlet-local; the pool belongs to a hub. + _metadata_pool_state = get_original("threading", "local")() + hub: Hub = get_hub() + pool: ThreadPool | None = getattr(_metadata_pool_state, "pool", None) + if pool is None or pool.hub is not hub or pool.pid != os.getpid(): + pool = ThreadPool(_METADATA_POOL_SIZE, hub=hub, idle_task_timeout=30) + _metadata_pool_state.pool = pool + return pool + + +_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 _CommandCancellation: + """Transfer caller cancellation to its private native-thread asyncio task.""" + + def __init__(self, lock: LockType) -> None: + self._lock: LockType = lock + self._cancelled: bool = False + self._loop: asyncio.AbstractEventLoop | None = None + self._task: asyncio.Task[Any] | None = None + + def bind(self, loop: asyncio.AbstractEventLoop, task: asyncio.Task[Any]) -> None: + """A cancelled queued command must not begin Redis work.""" + with self._lock: + self._loop, self._task = loop, task + if self._cancelled: + task.cancel() + + def cancel(self) -> None: + """Schedule cancellation without calling task methods across threads.""" + with self._lock: + self._cancelled = True + if self._loop is not None and self._task is not None: + self._loop.call_soon_threadsafe(self._task.cancel) + + def detach(self) -> None: + """Prevent racing cancellation from targeting a closed loop.""" + with self._lock: + self._loop, self._task = None, None + + +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 + + def _socket_timeout(self, key: str, remaining: float) -> float: + """Preserve shorter node timeouts inside the shared operation deadline.""" + configured: object = self._config.get(key) + if ( + isinstance(configured, (int, float)) + and not isinstance(configured, bool) + and 0 < configured < math.inf + ): + return min(configured, remaining) + return remaining + + 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": self._socket_timeout( + "CACHE_REDIS_SOCKET_TIMEOUT", remaining + ), + "socket_connect_timeout": self._socket_timeout( + "CACHE_REDIS_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": options["socket_timeout"], + "socket_connect_timeout": options["socket_connect_timeout"], + "retry": Retry(NoBackoff(), 0), + "protocol": 2, + }, + **options, + ) + sentinel_client: Redis + for sentinel_client in sentinel.sentinels: + stack.push_async_callback( + getattr(sentinel_client, "aclose", None) + or sentinel_client.close + ) + client = sentinel.master_for( + self._config.get("CACHE_REDIS_SENTINEL_MASTER", "mymaster") + ) + else: + client = Redis( + host=self._config.get("CACHE_REDIS_HOST", "localhost"), + port=self._config.get("CACHE_REDIS_PORT", 6379), + **options, + ) + await stack.enter_async_context(client) + return await client.execute_command(*args) + + def execute(self, *args: str | int) -> Any: + """Keep request greenlets from sharing asyncio's native-thread loop state.""" + self._remaining() + try: + from gevent import getcurrent, Greenlet, Timeout + from gevent.event import AsyncResult + from gevent.monkey import get_original + except ImportError: + return self._execute_sync(*args) + if not isinstance(getcurrent(), Greenlet): + return self._execute_sync(*args) + cancellation: _CommandCancellation = _CommandCancellation( + get_original("_thread", "allocate_lock")() + ) + try: + # The metadata-only native pool leaves the hub DNS pool available. + # Its queue wait consumes the original deadline too. + with Timeout( + self._remaining(), RedisTimeoutError("Metadata deadline expired") + ): + result: AsyncResult = _metadata_threadpool().spawn( + self._execute_sync, *args, cancellation=cancellation + ) + return result.get() + finally: + cancellation.cancel() + + def _execute_sync( + self, + *args: str | int, + cancellation: _CommandCancellation | None = None, + ) -> Any: + """Run a private loop and cancel its socket work within the call budget.""" + self._remaining() + try: + asyncio.get_running_loop() + except RuntimeError: + pass + else: + raise RedisError("Metadata backend requires a synchronous caller") + loop: asyncio.AbstractEventLoop = asyncio.new_event_loop() + task: asyncio.Task[Any] = loop.create_task(self._command(*args)) + if cancellation is not None: + cancellation.bind(loop, task) + try: + try: + return loop.run_until_complete(task) + except asyncio.CancelledError: + raise RedisTimeoutError("Metadata command cancelled") from None + except TimeoutError: + raise RedisTimeoutError("Metadata deadline expired") from None + finally: + # asyncio.run waits for the default executor on shutdown, including + # uncancellable system DNS resolution. Closing our private loop does + # not wait for that resolver. The cancelled command cannot connect or + # publish when the resolver eventually finishes. + if cancellation is not None: + cancellation.detach() + loop.close() Review Comment: Closing this per-command loop does not stop its running default-executor DNS lookup. When name resolution stalls, repeated timed-out requests can leave one resolver thread each after the metadata-pool slot is released; could resolver admission stay bounded until those underlying lookups actually finish? ########## superset/tasks/async_queries.py: ########## @@ -178,30 +181,63 @@ def _inject_contribution_totals( cached dataframe and inject the sums into this query's contribution post-processing before it runs — the same result the synchronous ``ensure_totals_available`` produces, but reading the cache the prerequisite - populated instead of re-running the totals query. ``contribution_totals`` is + populated. Participating semantic queries resolve the totals context under + the dependent operation so an obsolete denominator is recomputed. + ``contribution_totals`` is stripped from the cache key, so this affects only the result, not the key. """ from superset.common.query_context_processor import is_summable from superset.common.utils.query_cache_manager import QueryCacheManager - cache = QueryCacheManager.get(key=totals_cache_key, region=CacheRegion.DATA) - if not cache.is_loaded or cache.df is None: - # The depends_on prerequisite guarantees the totals task succeeded and wrote - # this cache entry, so a miss is unexpected (e.g. it was evicted between the - # totals task finishing and this task reading). Fail loudly rather than - # caching a silently un-normalized result the client would then re-request: - # this task's single query cannot reproduce the synchronous path's - # ensure_totals_available (it has no totals query to run). - raise SupersetException( - f"Contribution totals not found in cache under {totals_cache_key}" + df: pd.DataFrame + if totals_context is not None: + # The ordinary key captures this operation's catalog. A matching totals + # entry is reused; a refreshed catalog recomputes before normalization. + totals_context.force = ( + totals_context.force + and _query_task_cache_key(totals_context, 0) != totals_cache_key ) - df = cache.df - totals = {col: df[col].sum() for col in df.columns if is_summable(df[col])} + totals_context.is_async_execution = True + with _capture_query_cancellation(totals_context): + df = totals_context.get_df_payload_result( + totals_context.queries[0] + ).payload["df"] + if df is None: Review Comment: When a refresh makes the totals query fail (for example, it removes a metric not used by the dependent), `get_df_payload_result()` returns a failed payload with an empty DataFrame, not `None`, so this proceeds with `{}` and can cache percentages normalized against the displayed rows rather than the shared total. Could we check the acquisition status/error before consuming the dataframe? ########## superset/coordination/deadline_backend.py: ########## @@ -0,0 +1,343 @@ +# 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 os +import time +from _thread import LockType +from contextlib import AsyncExitStack +from typing import Any, TYPE_CHECKING + +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 + +if TYPE_CHECKING: + from threading import local + + from gevent.threadpool import ThreadPool + +_METADATA_POOL_SIZE: int = 4 +_metadata_pool_state: local | None = None + + +def _metadata_threadpool() -> ThreadPool: + """Reuse a lazy metadata-only pool in its owning native thread and hub.""" + from gevent import get_hub + from gevent.hub import Hub + from gevent.monkey import get_original + from gevent.threadpool import ThreadPool + + global _metadata_pool_state # pylint: disable=global-statement + if _metadata_pool_state is None: + # Patched threading.local is greenlet-local; the pool belongs to a hub. + _metadata_pool_state = get_original("threading", "local")() + hub: Hub = get_hub() + pool: ThreadPool | None = getattr(_metadata_pool_state, "pool", None) + if pool is None or pool.hub is not hub or pool.pid != os.getpid(): + pool = ThreadPool(_METADATA_POOL_SIZE, hub=hub, idle_task_timeout=30) + _metadata_pool_state.pool = pool + return pool + + +_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 _CommandCancellation: + """Transfer caller cancellation to its private native-thread asyncio task.""" + + def __init__(self, lock: LockType) -> None: + self._lock: LockType = lock + self._cancelled: bool = False + self._loop: asyncio.AbstractEventLoop | None = None + self._task: asyncio.Task[Any] | None = None + + def bind(self, loop: asyncio.AbstractEventLoop, task: asyncio.Task[Any]) -> None: + """A cancelled queued command must not begin Redis work.""" + with self._lock: + self._loop, self._task = loop, task + if self._cancelled: + task.cancel() + + def cancel(self) -> None: + """Schedule cancellation without calling task methods across threads.""" + with self._lock: + self._cancelled = True + if self._loop is not None and self._task is not None: + self._loop.call_soon_threadsafe(self._task.cancel) + + def detach(self) -> None: + """Prevent racing cancellation from targeting a closed loop.""" + with self._lock: + self._loop, self._task = None, None + + +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 + + def _socket_timeout(self, key: str, remaining: float) -> float: + """Preserve shorter node timeouts inside the shared operation deadline.""" + configured: object = self._config.get(key) + if ( + isinstance(configured, (int, float)) + and not isinstance(configured, bool) + and 0 < configured < math.inf + ): + return min(configured, remaining) + return remaining + + 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": self._socket_timeout( + "CACHE_REDIS_SOCKET_TIMEOUT", remaining + ), + "socket_connect_timeout": self._socket_timeout( + "CACHE_REDIS_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": options["socket_timeout"], + "socket_connect_timeout": options["socket_connect_timeout"], + "retry": Retry(NoBackoff(), 0), + "protocol": 2, + }, + **options, + ) + sentinel_client: Redis + for sentinel_client in sentinel.sentinels: + stack.push_async_callback( + getattr(sentinel_client, "aclose", None) + or sentinel_client.close + ) + client = sentinel.master_for( + self._config.get("CACHE_REDIS_SENTINEL_MASTER", "mymaster") + ) + else: + client = Redis( + host=self._config.get("CACHE_REDIS_HOST", "localhost"), + port=self._config.get("CACHE_REDIS_PORT", 6379), + **options, + ) + await stack.enter_async_context(client) + return await client.execute_command(*args) Review Comment: With supported redis-py 5.0.0, `Sentinel.master_for()` supplies an external pool and the client context manager does not disconnect it, leaving each command’s master socket cleanup to garbage collection as its event loop closes. Could we explicitly close this privately owned pool and cover the supported-version boundary? ########## superset/semantic_layers/models.py: ########## @@ -355,8 +366,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: Enabled semantic discovery still returns 500 through the live legacy metadata routes: dashboard add-chart calls `/fetch_datasource_metadata`, and `/datasource/get/semantic_view/<id>/` also reads `datasource.data` without the typed-error mapping. Could we translate unavailable/deadline failures at those boundaries (or migrate their consumers) and cover both HTTP paths? -- 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]
