geido commented on code in PR #44144:
URL: https://github.com/apache/superset/pull/44144#discussion_r4062635325
##########
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:
Good catch — this holds. Fixed in c4d477abce: pointer publication is now an
atomic compare-and-set against the generation observed before the Pending
write. A stale producer that loses the CAS adopts the readable winner and never
enqueues its orphan; the pointer lives in the metadata cache so custom/S3
thumbnail backends do not need CAS support. I added the exact slow-cache-write
regression plus concurrent CAS, expiry, and custom-backend coverage.
--
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]