EnxDev commented on code in PR #44336:
URL: https://github.com/apache/superset/pull/44336#discussion_r4059912487
##########
superset/commands/report/execute.py:
##########
@@ -2262,11 +2454,44 @@ def run(self) -> None:
self._execution_id,
report_execution_context,
).get_dashboard_urls()
+ execution_claim = None
+ if self._model.last_state != ReportState.WORKING:
+ execution_claim = claim_execution(
+ db.session,
+ self._model.id,
+ str(self._execution_id),
+ self._scheduled_dttm,
+ is_retry=self._is_retry,
+ expected_owner=self._expected_owner,
+
retries_enabled=feature_flag_manager.is_feature_enabled(
+ "ALERT_REPORTS_RETRY"
+ ),
+ stale_retry_seconds=(
+ app.config.get(
+ "ALERT_REPORTS_RETRY_MAX_DELAY_SECONDS", 3600
+ )
+ + resolve_report_execution_budget_seconds(
Review Comment:
Is the extra execution budget intentional here?
`ReportNotTriggeredErrorState` explicitly treats a `RETRYING` chain as dead
after `ALERT_REPORTS_RETRY_MAX_DELAY_SECONDS`, but this claim runs first and
uses max delay plus the execution budget. With the defaults, an `apply_async`
failure after committing `RETRYING` leaves the schedule unclaimable for two
hours, so the one-hour recovery branch below cannot run. Aligning these
thresholds, or centralizing stale recovery in the claim, would make that
recovery behavior effective.
##########
superset/commands/report/execute.py:
##########
@@ -2159,34 +2333,51 @@ class AsyncExecuteReportScheduleCommand(BaseCommand):
- On Alerts uses related Command AlertCommand and sends configured
notifications
"""
- def __init__(self, task_id: str, model_id: int, scheduled_dttm: datetime):
+ def __init__(
+ self,
+ task_id: str,
+ model_id: int,
+ scheduled_dttm: datetime,
+ *,
+ is_retry: bool = False,
+ expected_owner: str | None = None,
+ ):
self._model_id = model_id
self._model: Optional[ReportSchedule] = None
self._scheduled_dttm = scheduled_dttm
self._execution_id = UUID(task_id)
+ self._is_retry = is_retry
+ self._expected_owner = expected_owner
- def run(self) -> None:
+ def run(self) -> None: # noqa: C901
monotonic_started_at = time.monotonic()
report_execution_context: ReportExecutionContext | None = None
- owns_report_working_state = False
+ owns_working_state = False
try:
self.validate()
if not self._model:
raise ReportScheduleExecuteUnexpectedError()
- # Reports always run under an execution context; alerts join them
- # only when they deliver a rendered screenshot, so a blank/partial
- # capture fails closed instead of being delivered. Ownership and
- # terminal-error persistence remain report-only recovery semantics.
+ if self._is_retry and (
+ not
feature_flag_manager.is_feature_enabled("ALERT_REPORTS_RETRY")
+ or not self._model.retry_on_failure
+ or self._model.last_state != ReportState.RETRYING
+ or normalize_window(self._model.retry_scheduled_dttm)
+ != normalize_window(self._scheduled_dttm)
+ ):
+ logger.info(
+ "report_retry_discarded report_schedule_id=%s
execution_id=%s",
+ self._model_id,
+ self._execution_id,
+ )
+ return
+
+ # All scheduled executions share ownership and retry fencing.
if _should_build_execution_context(self._model):
# An invocation that enters on WORKING is a duplicate or stale
# recovery, not the owner that created the active row. Its
state
# handler may terminalize a stale execution, but the command
# boundary must never infer ownership from a replayed UUID.
- owns_report_working_state = (
- self._model.type == ReportScheduleType.REPORT
- and self._model.last_state != ReportState.WORKING
- )
total_seconds = resolve_report_execution_budget_seconds(
Review Comment:
Could we keep the report deadline semantics report-only, or explicitly
define an alert deadline here? Because every alert now gets an execution
context, this call caps alerts at `ALERT_REPORTS_EXECUTION_BUDGET_SECONDS`.
`_send()` enforces that deadline before delivery, while alert Celery tasks
still use `working_timeout + lag`, and `config.py` says alerts retain their
per-schedule timeout behavior. An alert configured for two hours can therefore
finish its query after one hour but be moved to `ERROR`/retry instead of
notifying. If the shared context is needed for ownership, separating that
concern from the report deadline would preserve the alert timeout contract.
##########
superset/commands/report/execute.py:
##########
@@ -2159,34 +2333,51 @@ class AsyncExecuteReportScheduleCommand(BaseCommand):
- On Alerts uses related Command AlertCommand and sends configured
notifications
"""
- def __init__(self, task_id: str, model_id: int, scheduled_dttm: datetime):
+ def __init__(
+ self,
+ task_id: str,
+ model_id: int,
+ scheduled_dttm: datetime,
+ *,
+ is_retry: bool = False,
+ expected_owner: str | None = None,
+ ):
self._model_id = model_id
self._model: Optional[ReportSchedule] = None
self._scheduled_dttm = scheduled_dttm
self._execution_id = UUID(task_id)
+ self._is_retry = is_retry
+ self._expected_owner = expected_owner
- def run(self) -> None:
+ def run(self) -> None: # noqa: C901
monotonic_started_at = time.monotonic()
report_execution_context: ReportExecutionContext | None = None
- owns_report_working_state = False
+ owns_working_state = False
try:
self.validate()
if not self._model:
raise ReportScheduleExecuteUnexpectedError()
- # Reports always run under an execution context; alerts join them
- # only when they deliver a rendered screenshot, so a blank/partial
- # capture fails closed instead of being delivered. Ownership and
- # terminal-error persistence remain report-only recovery semantics.
+ if self._is_retry and (
+ not
feature_flag_manager.is_feature_enabled("ALERT_REPORTS_RETRY")
+ or not self._model.retry_on_failure
+ or self._model.last_state != ReportState.RETRYING
+ or normalize_window(self._model.retry_scheduled_dttm)
+ != normalize_window(self._scheduled_dttm)
+ ):
+ logger.info(
+ "report_retry_discarded report_schedule_id=%s
execution_id=%s",
+ self._model_id,
+ self._execution_id,
+ )
+ return
Review Comment:
Could we terminalize or reset this retry when the flag or schedule opt-in
was turned off, provided the owner and window still match? A retry queued while
enabled can reach this return after an operator disables `ALERT_REPORTS_RETRY`
or an editor clears `retry_on_failure`, leaving the row in `RETRYING` with its
old owner and counter. Fresh cron runs are then rejected by `claim_execution()`
until the stale threshold expires (two hours with the default settings), so
disabling retries can also pause the schedule. A CAS on the same owner/window
before moving to `ERROR` or another terminal state would keep the stale-message
protection.
--
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]