EnxDev commented on code in PR #44144:
URL: https://github.com/apache/superset/pull/44144#discussion_r4045739950
##########
superset/dashboards/api.py:
##########
@@ -1965,44 +1984,249 @@ 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_generation() -> tuple[
+ str | None, ScreenshotCachePayload | None
+ ]:
+ cache_key = screenshot_obj.get_current_api_generation_cache_key(
+ request_cache_key,
+ cache_scope,
+ )
+ 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)
+
+ try:
+ observed_cache_key, _ = get_current_generation()
+ except ScreenshotCacheError:
+ logger.exception("Screenshot cache read failed: %s",
request_cache_key)
+ return self.response(
+ 503,
+ message=gettext("Screenshot cache is unavailable"),
+ )
+
+ producer_lock_ttl = get_default_lock_ttl()
+ # A contender must be able to outwait a live producer's lease. In
+ # particular, Celery broker publication can legitimately exceed one
+ # second; timing out sooner would reject a caller just before the
+ # producer publishes the generation it should join.
+ lock_deadline = (
+ time.monotonic() + producer_lock_ttl +
SCREENSHOT_API_LOCK_RETRY_SECONDS
+ )
+ while True:
+ lock_response: WerkzeugResponse | None = None
+ try:
+ with DistributedLock(
+ namespace=SCREENSHOT_API_LOCK_NAMESPACE,
+ request_cache_key=request_cache_key,
+ ttl_seconds=producer_lock_ttl,
+ ):
+ try:
+ cache_key, cached_payload = get_current_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 cache_payload.should_enqueue_task(
+ force,
+ expected_scope=cache_scope,
+ )
+ ):
+ 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,
+ )
+ except ScreenshotCacheError:
+ logger.exception(
+ "Screenshot task preparation failed: %s",
+ next_cache_key,
+ )
+ lock_response = self.response(
+ 503,
+ message=gettext("Screenshot cache is unavailable"),
+ )
+ return lock_response
+
+ try:
+ cache_dashboard_screenshot.delay(
+ username=get_current_user(),
+ guest_token=(
+ g.user.guest_token
+ if get_current_user() and isinstance(g.user,
GuestUser)
+ else None
+ ),
+ dashboard_id=dashboard.id,
+ dashboard_url=dashboard_url,
+ thumb_size=thumb_size,
+ window_size=window_size,
+ cache_key=next_cache_key,
+ # The API has already selected a fresh generation.
+ # Duplicate deliveries should never force a
completed
+ # result to recompute.
+ force=False,
+ )
+ except Exception: # pylint: disable=broad-except
+ screenshot_obj.mark_cache_error_if_incomplete(
+ next_cache_key,
+ cache_scope,
+ )
+ raise
+ try:
+ # Publish only after Celery accepts the task.
Otherwise a
+ # process exit between these operations strands a fresh
+ # Pending generation that no worker can complete.
+ screenshot_obj.set_current_api_generation_cache_key(
Review Comment:
Could we protect this publication against the producer's lease expiring
during `delay()`? The default lease is 30 seconds, but expiration doesn't stop
the request body. If A is still publishing at t=31, B can acquire the expired
lock, enqueue its own generation, and publish it. When A resumes, this
unconditional SET points the request back to A's older generation. The lock's
ownership-safe release doesn't protect this cache write.
I reproduced that ordering with the actual handler body and cache helpers
under a controlled clock/expiring-lock double: two tasks were enqueued, B's
payload reached `Updated`, then A replaced the pointer and the next non-force
poll returned `Pending` for A. If A's worker fails, polling can report `Error`
despite B's completed artifact. Could we make publication preserve a newer
generation after lease loss, through ownership-aware coordination or
reacquiring and rechecking the generation, and add an expiry-during-publish
regression?
--
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]