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]

Reply via email to