EnxDev commented on code in PR #44144:
URL: https://github.com/apache/superset/pull/44144#discussion_r4063751834


##########
superset/dashboards/api.py:
##########
@@ -1965,44 +1992,260 @@ def cache_dashboard_screenshot(self, pk: int, 
**kwargs: Any) -> WerkzeugResponse
 
         dashboard_url = get_url_path("Superset.dashboard_permalink", 
key=permalink_key)
         screenshot_obj = DashboardScreenshot(dashboard_url, dashboard.digest)
-        cache_key = screenshot_obj.get_cache_key(window_size, thumb_size, 
permalink_key)
-        image_url = get_url_path(
-            "DashboardRestApi.screenshot", pk=dashboard.id, digest=cache_key
-        )
-        cache_payload = (
-            screenshot_obj.get_from_cache_key(cache_key) or 
ScreenshotCachePayload()
+        cache_scope = f"dashboard:{dashboard.id}"
+        request_cache_key = screenshot_obj.get_api_request_cache_key(
+            window_size,
+            thumb_size,
+            permalink_key,
+            cache_scope,
         )
 
-        def build_response(status_code: int) -> WerkzeugResponse:
+        def build_response(
+            status_code: int,
+            cache_key: str,
+            cache_payload: ScreenshotCachePayload,
+        ) -> WerkzeugResponse:
             return self.response(
                 status_code,
                 cache_key=cache_key,
                 dashboard_url=dashboard_url,
-                image_url=image_url,
+                image_url=get_url_path(
+                    "DashboardRestApi.screenshot",
+                    pk=dashboard.id,
+                    digest=cache_key,
+                ),
+                task_timeout_seconds=(
+                    2 * current_app.config["THUMBNAIL_COMPUTING_CACHE_TTL"]
+                ),
                 task_updated_at=cache_payload.get_timestamp(),
                 task_status=cache_payload.get_status(),
             )
 
-        if cache_payload.should_trigger_task(
-            force, expected_scope=f"dashboard:{dashboard.id}"
-        ):
-            logger.info("Triggering screenshot ASYNC")
-            cache_dashboard_screenshot.delay(
-                username=get_current_user(),
-                guest_token=(
-                    g.user.guest_token
-                    if get_current_user() and isinstance(g.user, GuestUser)
+        def get_current_cache_key() -> str | None:
+            return screenshot_obj.get_current_api_generation_cache_key(
+                request_cache_key,
+                cache_scope,
+            )
+
+        def get_generation(
+            cache_key: str | None = None,
+        ) -> tuple[str | None, ScreenshotCachePayload | None]:
+            cache_key = cache_key or get_current_cache_key()
+            return (
+                cache_key,
+                (
+                    screenshot_obj.get_from_cache_key(
+                        cache_key,
+                        raise_on_error=True,
+                    )
+                    if cache_key
                     else None
                 ),
-                dashboard_id=dashboard.id,
-                dashboard_url=dashboard_url,
-                thumb_size=thumb_size,
-                window_size=window_size,
-                cache_key=cache_key,
-                force=force,
             )
-            return build_response(202)
-        return build_response(200)
+
+        def should_enqueue(cache_payload: ScreenshotCachePayload) -> bool:
+            return cache_payload.should_enqueue_task(
+                force,
+                expected_scope=cache_scope,
+                force_retry_after_seconds=SCREENSHOT_API_FORCE_RETRY_SECONDS,
+            )
+
+        try:
+            observed_cache_key, observed_payload = get_generation()
+        except ScreenshotCacheError:
+            logger.exception("Screenshot cache read failed: %s", 
request_cache_key)
+            return self.response(
+                503,
+                message=gettext("Screenshot cache is unavailable"),
+            )
+
+        if (
+            observed_cache_key
+            and observed_payload
+            and not should_enqueue(observed_payload)
+        ):
+            return build_response(200, observed_cache_key, observed_payload)
+
+        # Publish Pending before broker I/O so contenders can join it without
+        # waiting for the producer lock's full lease. Keep broker publication
+        # inside the lock to coalesce simultaneous forced requests, but never
+        # mutate the pointer after broker I/O because the lease may have 
expired.
+        lock_deadline = time.monotonic() + SCREENSHOT_API_LOCK_WAIT_SECONDS
+        while True:
+            lock_response: WerkzeugResponse | None = None
+            try:
+                with DistributedLock(
+                    namespace=SCREENSHOT_API_LOCK_NAMESPACE,
+                    request_cache_key=request_cache_key,
+                    ttl_seconds=SCREENSHOT_API_LOCK_TTL_SECONDS,
+                ):
+                    try:
+                        cache_key, cached_payload = get_generation()
+                    except ScreenshotCacheError:
+                        logger.exception(
+                            "Screenshot cache read failed: %s",
+                            request_cache_key,
+                        )
+                        lock_response = self.response(
+                            503,
+                            message=gettext("Screenshot cache is unavailable"),
+                        )
+                        return lock_response
+
+                    cache_payload = cached_payload or ScreenshotCachePayload(
+                        scope=cache_scope
+                    )
+                    if cached_payload is not None:
+                        if cache_key is None:
+                            logger.error(
+                                "Screenshot generation payload has no cache 
key: %s",
+                                request_cache_key,
+                            )
+                            lock_response = self.response(
+                                503,
+                                message=gettext("Screenshot cache is 
unavailable"),
+                            )
+                            return lock_response
+                        if cache_key != observed_cache_key or not 
should_enqueue(
+                            cache_payload
+                        ):
+                            lock_response = build_response(
+                                200, cache_key, cache_payload
+                            )
+                            return lock_response
+
+                    logger.info("Triggering screenshot ASYNC")
+                    next_cache_key = 
screenshot_obj.get_next_api_generation_cache_key(
+                        request_cache_key,
+                        cache_key,
+                    )
+                    cache_payload = ScreenshotCachePayload(scope=cache_scope)
+                    cache_payload.pending()
+                    try:
+                        screenshot_obj.store_cache_payload(
+                            next_cache_key,
+                            cache_payload,
+                        )
+                        # Make the reservation visible before broker I/O. A
+                        # failed publish is converted to terminal Error, and an
+                        # explicit force can replace an abandoned reservation
+                        # after the producer lease.
+                        screenshot_obj.set_current_api_generation_cache_key(
+                            request_cache_key,
+                            next_cache_key,
+                            cache_scope,
+                        )

Review Comment:
   Thanks, this addresses the race I raised. I replayed the slow Pending write 
on `c4d477abce` using the actual handler body, a controlled lease clock, and 
the real metadata-cache CAS against isolated SQLite. After B publishes and 
completes, A returns B's `Updated` generation with HTTP 200; the pointer stays 
on B and only one task is enqueued.
   
   All 224 focused backend tests and 7 download-hook tests pass. I also ran 
eight checks against the real metadata cache, covering concurrent 
creation/replacement and expiry with both JSON and pickle codecs. I found no 
further actionable issues in this pass. I haven't rerun live browser exports or 
the full dashboard API integration suite locally.



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