geido commented on code in PR #44144:
URL: https://github.com/apache/superset/pull/44144#discussion_r4047654879
##########
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:
Good catch. Fixed in a843e1eb6c: the new Pending generation and pointer are
published before delay(), and the handler never rewrites the pointer after
broker I/O. The controlled lease-expiry regression lets a newer generation
publish and complete before the older publisher resumes, then verifies the
pointer cannot roll back.
##########
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(
+ request_cache_key,
+ next_cache_key,
+ cache_scope,
+ )
+ except ScreenshotCacheError:
+ logger.exception(
+ "Screenshot generation publication failed: %s",
+ next_cache_key,
+ )
+ lock_response = self.response(
+ 503,
+ message=gettext("Screenshot cache is unavailable"),
+ )
+ return lock_response
+ lock_response = build_response(202, next_cache_key,
cache_payload)
+ return lock_response
+ except ReleaseDistributedLockFailedException:
+ if lock_response is not None:
+ logger.warning(
+ "Screenshot request completed but its producer lock
could "
+ "not be released: %s",
+ request_cache_key,
+ exc_info=True,
+ )
+ return lock_response
+ logger.exception(
+ "Could not release screenshot producer lock: %s",
+ request_cache_key,
+ )
+ return self.response(
+ 503,
+ message=gettext("Screenshot request is temporarily
unavailable"),
+ )
+ except LockAlreadyHeldException:
+ try:
+ cache_key, current_payload = get_current_generation()
Review Comment:
Fixed in a843e1eb6c. Contenders now poll only the small generation-pointer
entry; the payload is fetched once, only when the pointer changes and the
request is ready to respond.
test_producer_lock_contention_polls_pointer_without_refetching_payload enforces
this.
##########
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 = (
Review Comment:
Fixed in a843e1eb6c. The synchronous contention budget is capped at one
second, and this API reservation uses a two-second producer lease instead of
the global 30-second default. Tests cover both a permanently held lock
returning 503 within the bound and takeover after the short lease expires.
##########
superset/utils/screenshots.py:
##########
@@ -191,13 +201,46 @@ def is_error_cache_ttl_expired(self) -> bool:
def is_computing_stale(self) -> bool:
"""Check if a COMPUTING status is stale (task likely failed or
stuck)."""
+ return self.is_in_progress_stale()
+
+ def is_in_progress_stale(self) -> bool:
+ """Check if a pending or computing request has exceeded its lease."""
computing_ttl = app.config["THUMBNAIL_COMPUTING_CACHE_TTL"]
return (
datetime.now() - datetime.fromisoformat(self.get_timestamp())
).total_seconds() >= computing_ttl
- def should_trigger_task(
+ def is_in_progress(self) -> bool:
+ """Return whether screenshot computation has not reached a terminal
state."""
+
+ return self.status in (StatusValues.PENDING, StatusValues.COMPUTING)
+
+ def is_updated(self) -> bool:
+ """Return whether screenshot computation completed successfully."""
+
+ return self.status == StatusValues.UPDATED
+
+ def should_enqueue_task(
self, force: bool = False, expected_scope: str | None = None
+ ) -> bool:
+ """Return whether an API producer should enqueue a new generation.
+
+ Fresh pending/computing state is already accepted work, so even forced
+ callers observe it instead of producing another generation. A stale
+ in-progress state remains retryable through the existing lease TTL.
+ """
+
+ if expected_scope is not None and self._scope != expected_scope:
+ return True
+ if self.is_in_progress():
+ return self.is_in_progress_stale()
Review Comment:
Not intended. Fixed in a843e1eb6c: forced requests coalesce only during a
one-second reservation grace, then can replace a dead Pending or Computing
generation instead of waiting 360 seconds. Immediate 12-way force contention
still produces one Celery task; tests also prove delayed force retries while
non-force callers retain the normal computing TTL.
--
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]