EnxDev commented on code in PR #44144:
URL: https://github.com/apache/superset/pull/44144#discussion_r4059911294
##########
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:
Moving this before `delay()` handles the slow-broker case, but could we also
protect it against the lease expiring during `store_cache_payload()` just
above? That cache call can outlast the two-second lease, and this pointer write
is still unconditional.
I reproduced this on `a843e1eb6c` with the actual handler body and cache
helpers, using a controlled clock and an expiring-lock double: A pauses in its
Pending cache SET; at t=2.1s, B acquires the expired lease, publishes its
generation, and completes as `Updated`; then A resumes and replaces B's pointer
with its own `Pending` generation. The next non-force POST returns A even
though B's artifact is ready. Ownership-checked lock release doesn't prevent
this write.
Could we add the slow-cache-write variant to the lease-expiry regression and
make pointer publication conditional on the expected generation/ownership
atomically, so an expired producer can't overwrite its successor?
--
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]