sadpandajoe commented on code in PR #44849:
URL: https://github.com/apache/superset/pull/44849#discussion_r4213840188


##########
superset/semantic_layers/metadata_errors.py:
##########
@@ -0,0 +1,158 @@
+# 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
+
+import logging
+from collections.abc import Callable
+from functools import wraps
+from typing import Any
+
+from flask_appbuilder.api import BaseApi
+from flask_babel import gettext as t
+from sqlalchemy.exc import SQLAlchemyError
+from superset_core.semantic_layers.metadata import (
+    MetadataRefreshError,
+    MetadataRefreshErrorCategory,
+)
+
+from superset.semantic_layers.metadata_binding import metadata_refresh_enabled
+from superset.superset_typing import FlaskResponse
+
+logger: logging.Logger = logging.getLogger(__name__)
+
+
+def metadata_error_response(
+    api: BaseApi | type[BaseApi], error: MetadataRefreshError
+) -> FlaskResponse:
+    """Return only stable categories and actionable, localized safe 
messages."""
+    errors: dict[str, tuple[int, str]] = {
+        "unsupported": (
+            422,
+            str(t("This semantic layer does not support metadata sync.")),
+        ),
+        "configuration": (
+            422,
+            str(
+                t("Complete the semantic layer configuration before syncing 
metadata.")
+            ),
+        ),
+        "in_progress": (
+            409,
+            str(t("A metadata sync is already in progress. Try again 
shortly.")),
+        ),
+        "configuration_changed": (
+            409,
+            str(
+                t(
+                    "The semantic view or connection changed. "
+                    "Reopen the editor and try again."
+                )
+            ),
+        ),
+        "upstream": (
+            502,
+            str(t("The semantic layer could not return its catalog. Try again 
later.")),
+        ),
+        "invalid_payload": (
+            502,
+            str(t("The semantic layer returned an invalid or oversized 
catalog.")),
+        ),
+        "deadline": (504, str(t("Metadata sync timed out. Try again later."))),
+        "unavailable": (
+            503,
+            str(t("Semantic metadata is unavailable. Try again later.")),
+        ),
+        "indeterminate": (
+            503,
+            str(
+                t(
+                    "Metadata sync could not be confirmed. "
+                    "Reload fields before trying again."
+                )
+            ),
+        ),
+    }
+    status: int
+    message: str
+    status, message = errors[error.category]
+    return api.response(status, error=error.category, message=message)

Review Comment:
   Putting the machine category in `error` takes over the field that existing 
clients display. `parseErrorJson` only copies `message` into `error` when 
`error` is unset, so `getClientErrorObject` returns `error: "unavailable"` and 
`SemanticViewEditModal` (and `StatefulChart`) show a toast that reads just 
"unavailable" instead of the translated message. Should the category move to a 
separate key such as `error_type` (or `errors[0]`) so `error` stays the 
human-readable text?



##########
superset/semantic_layers/metadata_binding.py:
##########
@@ -0,0 +1,383 @@
+# 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, NoReturn, TYPE_CHECKING
+
+from flask import current_app, has_app_context, has_request_context, request
+from sqlalchemy import func, select
+from sqlalchemy.engine import Connection, NestedTransaction, RowMapping
+from sqlalchemy.exc import SQLAlchemyError
+from sqlalchemy.sql.elements import ColumnElement
+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 _revalidation_unavailable(reason: str) -> NoReturn:
+    """Explain a metadata DB constraint without exposing configuration or 
SQL."""
+    logger.warning("Metadata database revalidation unavailable: %s", reason)
+    raise MetadataRefreshError("unavailable")
+
+
+@contextmanager
+def _revalidation_savepoint(connection: Connection) -> Iterator[None]:
+    """Contain statement failures without flushing or ending the caller's 
work."""
+    savepoint: NestedTransaction = connection.begin_nested()
+    try:
+        yield
+        # End only this SAVEPOINT; @transaction would end the caller's work.
+        savepoint.commit()  # pylint: disable=consider-using-transaction
+    except BaseException:
+        try:
+            savepoint.rollback()  # pylint: disable=consider-using-transaction
+        except SQLAlchemyError:
+            # A lost connection may prevent recovery. Preserve the original
+            # failure, especially worker cancellation, rather than masking it.
+            logger.warning(
+                "Metadata revalidation savepoint recovery failed", 
exc_info=True
+            )
+        raise
+
+
+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 _revalidate_layer(layer: SemanticLayer, scope: str) -> None:
+    """Require a fresh committed scope without borrowing another connection."""
+    # 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")
+    if db.session.new or db.session.dirty or db.session.deleted:

Review Comment:
   Revalidation refuses to publish whenever the caller's session has any 
pending ORM object or an assigned transaction ID, and `unavailable` is raised 
for writes that have nothing to do with the semantic layer. For example, with 
`STORE_CACHE_KEYS_IN_METADATA_DB` on, a result-cache write in 
`set_and_log_cache` does `db.session.add(CacheKey(...))` without committing, so 
a later semantic chart with a cold catalog in the same request (a multi-chart 
dashboard export) hits this branch and is dropped. The async task path looks 
similar: `_execute_task_body` commits the status transition and then assigns 
`task.status` again, which autoflushes into the next transaction before the 
revalidation `SELECT`. Should this check be scoped to writes that can affect 
the layer (for example by revalidating on an independent connection or in a 
fresh session) so ordinary cache or task bookkeeping doesn't turn a healthy 
refresh into a 503?



-- 
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