sadpandajoe commented on code in PR #44835:
URL: https://github.com/apache/superset/pull/44835#discussion_r4209191363
##########
superset/datasource/api.py:
##########
@@ -203,6 +212,11 @@ def _column_values_response(
result=[],
suggestions_status="unavailable_versioned_view",
)
+ metadata_token: str | None = cast(
+ SemanticView, datasource
+ ).metadata_cache_token
Review Comment:
With metadata refresh on and a participating view, this property can raise
`MetadataRefreshError` (shared store down, discovery budget exhausted) before
the cache lookup, and `get_column_values` is not wrapped by
`metadata_api_errors`. `@safe` then returns a generic 500 instead of the mapped
503/504 that `compatible` and runtime-schema return. The `/views` and
`/structure` routes in `semantic_layers/api.py` have the same gap: their broad
`except Exception` turns the same typed failure into a 400 "check the layer
configuration" or a 422. Should these routes use the shared mapper and re-raise
typed errors the way runtime-schema now does?
##########
superset/common/query_context_processor.py:
##########
@@ -307,8 +309,19 @@ def get_df_payload_result(
)
)
+ self._capture_annotation_metadata(query_obj)
query_result = self.get_query_result(query_obj)
annotation_data = self.get_annotation_data(query_obj)
+ if query_obj.annotation_layers:
+ from superset.semantic_layers.metadata_binding import (
+ metadata_refresh_enabled,
+ )
+
+ if metadata_refresh_enabled():
+ # Discovery on a miss can capture a newer annotation
+ # snapshot than the lookup peek. Store only under the
+ # identity actually used by that annotation query.
+ cache_key, cacheable = self._query_cache_key(query_obj)
Review Comment:
This re-key re-runs `datasource.get_extra_cache_keys()` after the query has
executed, not just the annotation part. For a SQL dataset that re-renders the
parent's Jinja, so a `cache_key_wrapper(...)` value such as
`latest_partition()` can differ after a long-running query. The rows fetched
under partition P0 would then be stored under P1's key, and later P1 requests
would get stale data. Could the parent's key material captured at lookup be
kept, with only the annotation identity recomputed?
##########
superset/models/slice.py:
##########
@@ -404,6 +404,24 @@ 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
+
+ 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
+ return DatasourceDAO.get_datasource(
Review Comment:
This now resolves the annotation chart's saved datasource while the cache
key is being built, before the `try` in `get_df_payload_result`. If that
datasource was deleted, `DatasourceNotFound` propagates out of the whole
request, even with metadata refresh off. Before this change the lookup happened
during annotation execution, where a `SupersetException` becomes a failed-query
payload for that annotation. Should a missing saved datasource fall back to
`resolved_datasource`, or be translated into the same validation error?
##########
superset/semantic_layers/metadata_binding.py:
##########
@@ -0,0 +1,308 @@
+# 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 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")
Review Comment:
`_operation()` raises `MetadataRefreshError("configuration")` whenever no
operation is active, and the docs say other synchronous callers must enter
`metadata_operation()`. The in-tree MCP semantic-layer tools
(`mcp_service/semantic_layer/tool/get_table.py`, `get_compatible_metrics.py`,
`list_metrics.py`) read `view.implementation` outside a request hook and
outside `metadata_operation()`, and catch only `JSONDecodeError`. With the flag
on and a participating layer, would those tools now fail with an unhandled
error? Should they enter `metadata_operation()` or map the typed error?
##########
tests/unit_tests/semantic_layers/metadata_redis_test.py:
##########
@@ -0,0 +1,428 @@
+# 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.
+
+"""Opt-in actual Redis acceptance; never flush a database or delete shared
keys."""
+
+from __future__ import annotations
+
+import multiprocessing
+import os
+import time
+from collections.abc import Iterator
+from concurrent.futures import Future
+from multiprocessing.context import SpawnContext
+from multiprocessing.process import BaseProcess
+from multiprocessing.queues import Queue
+from multiprocessing.synchronize import Event
+from typing import Any, cast
+from urllib.parse import ParseResult, urlparse
+from uuid import uuid4
+
+import pytest
+from flask import Flask
+from superset_core.semantic_layers.metadata import CatalogSnapshot,
MetadataRefreshError
+
+from superset.coordination.deadline_backend import DeadlineRedisBackend
+from superset.semantic_layers.cache_inspection import CacheEntryInfo
+from superset.semantic_layers.metadata import ScopedMetadataStore
+
+
[email protected]
+def redis_config() -> dict[str, Any]:
+ url: str | None = os.environ.get("SEMANTIC_METADATA_TEST_REDIS_URL")
+ if url is None:
+ pytest.skip(
+ "Set SEMANTIC_METADATA_TEST_REDIS_URL for isolated Redis
acceptance"
+ )
+ assert url is not None
+ parsed: ParseResult = urlparse(url)
+ return {
+ "CACHE_TYPE": "RedisCache",
+ "CACHE_REDIS_HOST": parsed.hostname,
+ "CACHE_REDIS_PORT": parsed.port or 6379,
+ "CACHE_REDIS_DB": int(parsed.path.lstrip("/") or "0"),
+ }
+
+
[email protected]
+def redis_scope(redis_config: dict[str, Any]) -> Iterator[str]:
+ scope: str = "sc121047-test-" + uuid4().hex
+ yield scope
+ backend: DeadlineRedisBackend = DeadlineRedisBackend(
+ redis_config, deadline=time.monotonic() + 5
+ )
+ backend.delete(
+ *(
+ f"semantic-metadata:{{{scope}}}:{suffix}"
+ for suffix in ("snapshot", "lease", "compatibility")
+ )
+ )
+
+
+def _reader(
+ config: dict[str, Any], scope: str, commands: Queue[bool], results:
Queue[str]
+) -> None:
+ store_deadline: float = time.monotonic() + 30
+ store: ScopedMetadataStore = ScopedMetadataStore(
+ DeadlineRedisBackend(config, deadline=time.monotonic() + 30),
+ scope,
+ deadline=store_deadline,
+ )
+
+ def forbidden(deadline: float) -> str:
+ raise AssertionError("warm reader fetched upstream")
+
+ while commands.get(timeout=10):
+ results.put(store.read(forbidden, deadline=store_deadline).cache_token)
+
+
+def _paused_writer(
+ config: dict[str, Any],
+ scope: str,
+ started: Event,
+ release: Event,
+ results: Queue[str],
+) -> None:
+ deadline: float = time.monotonic() + 30
+ store_deadline: float = deadline
+ store: ScopedMetadataStore = ScopedMetadataStore(
+ DeadlineRedisBackend(config, deadline=deadline), scope,
deadline=store_deadline
+ )
+
+ def fetch(budget: float) -> str:
+ started.set()
+ assert release.wait(10)
+ return '["old"]'
+
+ try:
+ store.refresh(fetch, deadline=store_deadline)
+ results.put("published")
+ except MetadataRefreshError as error:
+ results.put(error.category)
+
+
+def test_two_processes_alternate_100_reads_of_one_observation(
+ redis_config: dict[str, Any], redis_scope: str
+) -> None:
+ deadline: float = time.monotonic() + 30
+ backend: DeadlineRedisBackend = DeadlineRedisBackend(
+ redis_config, deadline=deadline
+ )
+ acquisition_deadline_1: float = deadline
+ snapshot: CatalogSnapshot = ScopedMetadataStore(
+ backend, redis_scope, deadline=acquisition_deadline_1
+ ).read(lambda budget: '["orders"]', deadline=acquisition_deadline_1)
+ context: SpawnContext = cast(SpawnContext,
multiprocessing.get_context("spawn"))
+ commands: list[Queue[bool]] = [context.Queue(), context.Queue()]
+ results: Queue[str] = context.Queue()
+ processes: list[BaseProcess] = [
+ context.Process(
+ target=_reader, args=(redis_config, redis_scope, queue, results)
+ )
+ for queue in commands
+ ]
+ process: BaseProcess
+ index: int
+ queue: Queue[bool]
+ try:
+ for process in processes:
+ process.start()
+ for index in range(100):
+ commands[index % 2].put(True)
+ assert results.get(timeout=10) == snapshot.cache_token
+ finally:
+ for queue in commands:
+ queue.put(False)
+ for process in processes:
+ process.join(10)
+ if process.is_alive():
+ process.terminate()
+ process.join(2)
+ assert process.exitcode == 0
+
+
+def test_real_invalidation_fences_a_paused_process(
+ redis_config: dict[str, Any], redis_scope: str
+) -> None:
+ context: SpawnContext = cast(SpawnContext,
multiprocessing.get_context("spawn"))
+ started: Event = context.Event()
+ release: Event = context.Event()
+ results: Queue[str] = context.Queue()
+ process: BaseProcess = context.Process(
+ target=_paused_writer,
+ args=(redis_config, redis_scope, started, release, results),
+ )
+ process.start()
+ deadline: float = time.monotonic() + 30
+ store_deadline: float = deadline
+ store: ScopedMetadataStore = ScopedMetadataStore(
+ DeadlineRedisBackend(redis_config, deadline=deadline),
+ redis_scope,
+ deadline=store_deadline,
+ )
+ try:
+ assert started.wait(10)
+ store.invalidate_catalog()
Review Comment:
This does not exercise owner fencing against a replacement writer. The
successor finishes publishing before the paused writer resumes, so no lease
exists and the old writer is rejected by the lease-existence check. The test
would still pass if the Lua owner comparison were removed. Could it have writer
A invalidated, writer B acquire a lease and pause, then A resume, asserting A
gets `configuration_changed` and B's lease survives A's publish and cleanup so
B can still publish?
##########
superset/initialization/__init__.py:
##########
@@ -163,10 +163,23 @@ class AppContextTask(task_base): # type: ignore
# nested context here would silently hand the task a second,
# blind session unable to see the caller's uncommitted work.
def __call__(self, *args: Any, **kwargs: Any) -> Any:
- if has_app_context():
- return task_base.__call__(self, *args, **kwargs)
- with superset_app.app_context():
- return task_base.__call__(self, *args, **kwargs)
+ # Avoid circular import through superset.app during
initialization.
+ from superset.semantic_layers.metadata_binding import
metadata_operation
+
+ with (
+ contextlib.nullcontext()
+ if has_app_context()
+ else superset_app.app_context()
+ ):
+ with (
+ metadata_operation()
Review Comment:
The per-chart scope covers the Celery task and async chart queries, but the
inline export path still runs every chart under the one request deadline:
`chart_metadata_operation()` returns immediately when `has_request_context()`
is true, and `_export_xlsx_inline` calls `build_workbook` inside the request.
Once cumulative chart time passes 30 seconds, each later chart whose view is
not yet captured raises a `deadline` error, which `build_workbook` records as a
general error. The download is a 200 workbook silently missing those charts.
Should the inline path give each chart its own acquisition scope, or should the
row budget also limit how long the inline export can run?
--
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]