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


##########
superset/tasks/query_cancel.py:
##########
@@ -163,3 +172,150 @@ def cancel_chart_query(
             exc_info=True,
         )
         return False
+
+
+def _registry_key(user_id: int, client_id: str) -> str:
+    """Cache key for a user's in-flight, cancellable chart query.
+
+    The user id is part of the key rather than a field compared after lookup, 
so
+    a ``client_id`` is only ever resolvable within the namespace of the user 
who
+    registered it. A caller passing somebody else's ``client_id`` gets a miss —
+    it cannot read, cancel, or overwrite another user's entry.
+    """
+    return f"chart-query-cancel:{user_id}:{client_id}"
+
+
+def _registry_ttl() -> int:
+    """How long a cancel handle stays resolvable.
+
+    A synchronous chart query cannot outlive the web request running it, so the
+    webserver timeout is the natural upper bound. Entries are discarded as soon
+    as the query returns; this TTL only bounds the leak when a worker dies
+    mid-query.
+    """
+    return int(current_app.config.get("SUPERSET_WEBSERVER_TIMEOUT", 60))
+
+
+@contextmanager
+def cancellable_chart_query(
+    client_id: "str | None", database: "Database | None"
+) -> Iterator[None]:
+    """Let the requesting user cancel this synchronous chart query by 
``client_id``.
+
+    Captures the engine cancel id off the live cursor (for engines that expose
+    one before execution) and publishes it so a concurrent Stop request — which
+    lands on a different worker while this one is blocked on the query — can 
kill
+    the backend session. Engines without cancel support capture nothing and the
+    query stays non-cancellable, exactly as before.
+
+    A no-op without a ``client_id``, without a database (e.g. the annotation
+    datasource, which queries Superset's own metadata DB), or for an
+    unauthenticated request — an anonymous viewer of a public dashboard has no
+    user id to scope the handle to, and an unscoped handle would be cancellable
+    by any other anonymous visitor.
+    """
+    # Inline (also below, and in _publish_cancel_handle/_discard_cancel_handle/
+    # cancel_chart_query_for_user): the unit tests patch get_user_id and
+    # cache_manager at their defining modules (superset.utils.core,
+    # superset.extensions), not here. A module-level `from ... import X` binds
+    # X once at import time, before any test patch runs, so the patch would
+    # silently never take effect; re-importing on every call picks up the
+    # patched object instead.
+    from superset.utils.core import get_user_id
+
+    user_id = get_user_id()
+    if not client_id or database is None or user_id is None:
+        yield
+        return
+
+    # Rebound as non-optional locals: mypy does not carry the narrowing above
+    # into the nested function below.
+    owner_id: int = user_id
+    query_id: str = client_id
+    target: "Database" = database
+    database_id = target.id
+    captured = False
+
+    def _sink(cursor: Any) -> None:
+        # Republish for every cursor rather than only the first: one query
+        # object can execute several statements in turn (e.g. the grouping-sets
+        # fallback), each on whatever connection the pool hands out, and the
+        # handle must always name the session that is executing right now.
+        nonlocal captured
+        cancel_id = capture_cancel_query_id(target, cursor)
+        if cancel_id is None:
+            return
+        captured = True
+        _publish_cancel_handle(owner_id, query_id, database_id, cancel_id)
+
+    try:
+        with capture_cancel_id(_sink):
+            yield
+    finally:
+        if captured:
+            _discard_cancel_handle(owner_id, query_id)

Review Comment:
   The handle is republished on every cursor via `_sink` (for multi-statement 
query objects like the grouping-sets fallback), but discarded only once here, 
after the whole `get_query_result` call returns. A Stop request that reads the 
cache in the gap between one cursor finishing and the next cursor's republish 
resolves a stale `cancel_query_id`. For Postgres and Redshift, `cancel_query` 
reports success whenever the terminate statement executes without raising, even 
when the PID matches no row — so the stale handle is treated as successfully 
cancelled, `stopped: true` is returned, and this line then deletes the cache 
entry the *next*, still-running statement just published, while that statement 
keeps executing to completion. Could the handle carry something that lets 
discard/cancel confirm it's still the current statement's before acting on it?



##########
superset/charts/data/api.py:
##########
@@ -345,6 +360,63 @@ def data(  # noqa: C901
             expected_rows=expected_rows,
         )
 
+    @expose("/data/stop", methods=("POST",))
+    @protect()
+    @statsd_metrics
+    @event_logger.log_this_with_context(
+        action=lambda self, *args, **kwargs: 
f"{self.__class__.__name__}.stop_data",
+        log_to_statsd=False,
+    )
+    def stop_data(self) -> Response:
+        """
+        Cancel a running chart-data query.
+        ---
+        post:
+          summary: Cancel a running chart-data query
+          description: >-
+            Cancels the warehouse query started by a chart-data request that
+            carried the given `client_id`, for databases whose engine supports
+            query cancellation. The `client_id` is resolved only within the
+            requesting user's own in-flight queries, so it cannot be used to
+            reach another user's query.
+          requestBody:
+            required: true
+            content:
+              application/json:
+                schema:
+                  $ref: '#/components/schemas/ChartDataStopSchema'
+          responses:
+            200:
+              description: Cancellation outcome
+              content:
+                application/json:
+                  schema:
+                    type: object
+                    properties:
+                      result:
+                        type: object
+                        properties:
+                          stopped:
+                            type: boolean
+                            description: >-
+                              Whether the engine reported the query cancelled.
+                              False when there was no such in-flight query for
+                              this user, or the engine could not cancel it.
+            400:
+              $ref: '#/components/responses/400'
+            401:
+              $ref: '#/components/responses/401'
+            500:

Review Comment:
   This endpoint is decorated with `@protect()` and mapped to 
`method_permission_name["stop_data"] = "read"`, so a caller without `can_read` 
on Chart gets a 403 before this handler runs — but the response spec here only 
documents 400/401/500. The sibling `data` endpoint above documents 403 for the 
same permission check. Should 403 be added here too so a client generated from 
this spec models the permission-denied response?



##########
superset/common/query_context_processor.py:
##########
@@ -305,7 +306,15 @@ def get_df_payload_result(
                         )
                     )
 
-                query_result = self.get_query_result(query_obj)
+                # Make the warehouse query cancellable by its owner for the 
span
+                # of its execution, so Explore's Stop button kills the query
+                # rather than just abandoning the response. No-op for engines
+                # without cancel support, and for requests carrying no 
client_id.
+                with cancellable_chart_query(

Review Comment:
   Under `GLOBAL_ASYNC_QUERIES`, this query reaches the Celery worker through 
`serialize_query`/`load_serialized_query`, and neither carries `client_id` 
(only `force_nonce` and `custom_cache_timeout` are threaded through the 
`SerializedQuery` payload). The worker's reconstructed `QueryContext.client_id` 
is therefore always `None`, so `cancellable_chart_query` is a permanent no-op 
here for every async chart query — Stop silently never reaches the warehouse 
query in this mode. Should `client_id` be added to the serialized payload so 
async chart queries get the same cancellation coverage as synchronous ones?



##########
superset/tasks/query_cancel.py:
##########
@@ -163,3 +172,150 @@ def cancel_chart_query(
             exc_info=True,
         )
         return False
+
+
+def _registry_key(user_id: int, client_id: str) -> str:
+    """Cache key for a user's in-flight, cancellable chart query.
+
+    The user id is part of the key rather than a field compared after lookup, 
so
+    a ``client_id`` is only ever resolvable within the namespace of the user 
who
+    registered it. A caller passing somebody else's ``client_id`` gets a miss —
+    it cannot read, cancel, or overwrite another user's entry.
+    """
+    return f"chart-query-cancel:{user_id}:{client_id}"
+
+
+def _registry_ttl() -> int:
+    """How long a cancel handle stays resolvable.
+
+    A synchronous chart query cannot outlive the web request running it, so the
+    webserver timeout is the natural upper bound. Entries are discarded as soon
+    as the query returns; this TTL only bounds the leak when a worker dies
+    mid-query.
+    """
+    return int(current_app.config.get("SUPERSET_WEBSERVER_TIMEOUT", 60))
+
+
+@contextmanager
+def cancellable_chart_query(
+    client_id: "str | None", database: "Database | None"
+) -> Iterator[None]:
+    """Let the requesting user cancel this synchronous chart query by 
``client_id``.
+
+    Captures the engine cancel id off the live cursor (for engines that expose
+    one before execution) and publishes it so a concurrent Stop request — which
+    lands on a different worker while this one is blocked on the query — can 
kill
+    the backend session. Engines without cancel support capture nothing and the
+    query stays non-cancellable, exactly as before.
+
+    A no-op without a ``client_id``, without a database (e.g. the annotation
+    datasource, which queries Superset's own metadata DB), or for an
+    unauthenticated request — an anonymous viewer of a public dashboard has no
+    user id to scope the handle to, and an unscoped handle would be cancellable
+    by any other anonymous visitor.
+    """
+    # Inline (also below, and in _publish_cancel_handle/_discard_cancel_handle/
+    # cancel_chart_query_for_user): the unit tests patch get_user_id and
+    # cache_manager at their defining modules (superset.utils.core,
+    # superset.extensions), not here. A module-level `from ... import X` binds
+    # X once at import time, before any test patch runs, so the patch would
+    # silently never take effect; re-importing on every call picks up the
+    # patched object instead.
+    from superset.utils.core import get_user_id
+
+    user_id = get_user_id()
+    if not client_id or database is None or user_id is None:
+        yield
+        return
+
+    # Rebound as non-optional locals: mypy does not carry the narrowing above
+    # into the nested function below.
+    owner_id: int = user_id
+    query_id: str = client_id
+    target: "Database" = database
+    database_id = target.id
+    captured = False
+
+    def _sink(cursor: Any) -> None:
+        # Republish for every cursor rather than only the first: one query
+        # object can execute several statements in turn (e.g. the grouping-sets
+        # fallback), each on whatever connection the pool hands out, and the
+        # handle must always name the session that is executing right now.
+        nonlocal captured
+        cancel_id = capture_cancel_query_id(target, cursor)
+        if cancel_id is None:
+            return
+        captured = True
+        _publish_cancel_handle(owner_id, query_id, database_id, cancel_id)
+
+    try:
+        with capture_cancel_id(_sink):
+            yield
+    finally:
+        if captured:
+            _discard_cancel_handle(owner_id, query_id)
+
+
+def _publish_cancel_handle(
+    user_id: int, client_id: str, database_id: int, cancel_query_id: str
+) -> None:
+    """Publish a cancel handle for the owning user. Best-effort."""
+    # Inline for test-patchability; see cancellable_chart_query's comment.
+    from superset.extensions import cache_manager
+
+    try:
+        cache_manager.cache.set(

Review Comment:
   Superset's default `CACHE_CONFIG` is `NullCache`, so this 
`cache_manager.cache.set(...)` call is a no-op there — the cancel handle is 
never actually stored, the Stop endpoint's lookup always misses, and the 
warehouse query keeps running after the user clicks Stop, with no error 
surfaced (the endpoint just returns `stopped: false`, which the UI treats as 
normal). Should this registry use a cache guaranteed to persist across 
requests/workers by default, e.g. a dedicated required cache the way 
`FILTER_STATE_CACHE_CONFIG` is configured, rather than relying on the general 
`CACHE_CONFIG` an operator may leave at its default?



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