potiuk commented on code in PR #70653:
URL: https://github.com/apache/airflow/pull/70653#discussion_r4064418347


##########
task-sdk/src/airflow/sdk/bases/resumablejobmixin.py:
##########
@@ -149,42 +149,50 @@ def execute_resumable(self, context: Context) -> Any:
                 if external_id:
                     stats.incr("resumable_job.reconnect_attempt", 
tags=stats_tags)
 
-                    status = self.get_job_status(external_id, context)
-
-                    span.set_attribute("resumable.external_id", 
str(external_id))
-                    span.set_attribute("resumable.prior_status", status)
-
-                    if self.is_job_active(status):
-                        # Job is still running, skip submission and reconnect 
to it.
-                        span.set_attribute("resumable.decision", "reconnect")
-                        stats.incr("resumable_job.reconnect_success", 
tags=stats_tags)
-                        self.log.info(
-                            "Reconnecting to existing job",
-                            external_id_key=self.external_id_key,
-                            external_id=external_id,
-                            status=status,
-                        )
-                        reconnect_to = external_id
-                    elif self.is_job_succeeded(status):
-                        # Job already finished successfully, skip polling and 
return result directly.
-                        span.set_attribute("resumable.decision", 
"already_succeeded")
-                        stats.incr("resumable_job.already_succeeded", 
tags=stats_tags)
-                        self.log.info(
-                            "Job already completed successfully, skipping 
resubmission",
-                            external_id_key=self.external_id_key,
-                            external_id=external_id,
-                        )
-                        already_succeeded_id = external_id
-                    else:
-                        # Job is in a terminal failed state, fall through and 
submit a new job.
-                        span.set_attribute("resumable.decision", 
"terminal_resubmit")
-                        stats.incr("resumable_job.terminal_resubmit", 
tags=stats_tags)
-                        self.log.warning(
-                            "Prior job in terminal state, resubmitting fresh",
-                            external_id_key=self.external_id_key,
-                            external_id=external_id,
-                            status=status,
-                        )
+                    try:
+                        status = self.get_job_status(external_id, context)
+
+                        span.set_attribute("resumable.external_id", 
str(external_id))
+                        span.set_attribute("resumable.prior_status", status)
+
+                        if self.is_job_active(status):
+                            # Job is still running, skip submission and 
reconnect to it.
+                            span.set_attribute("resumable.decision", 
"reconnect")
+                            stats.incr("resumable_job.reconnect_success", 
tags=stats_tags)
+                            self.log.info(
+                                "Reconnecting to existing job",
+                                external_id_key=self.external_id_key,
+                                external_id=external_id,
+                                status=status,
+                            )
+                            reconnect_to = external_id
+                        elif self.is_job_succeeded(status):
+                            # Job already finished successfully, skip polling 
and return result directly.
+                            span.set_attribute("resumable.decision", 
"already_succeeded")
+                            stats.incr("resumable_job.already_succeeded", 
tags=stats_tags)
+                            self.log.info(
+                                "Job already completed successfully, skipping 
resubmission",
+                                external_id_key=self.external_id_key,
+                                external_id=external_id,
+                            )
+                            already_succeeded_id = external_id
+                        else:
+                            # Job is in a terminal failed state, fall through 
and submit a new job.
+                            span.set_attribute("resumable.decision", 
"terminal_resubmit")
+                            stats.incr("resumable_job.terminal_resubmit", 
tags=stats_tags)
+                            self.log.warning(
+                                "Prior job in terminal state, resubmitting 
fresh",
+                                external_id_key=self.external_id_key,
+                                external_id=external_id,
+                                status=status,
+                            )
+                    except Exception:

Review Comment:
   The `try` starts at line 152 and wraps the branch bodies, not just the 
decision — so an outcome counter and this one can both fire for a single 
attempt.
   
   Each branch increments and *then* logs: `:161` → `:162`, `:172` → `:173`, 
`:182` → `:183`. If one of those logging calls raises, the outcome has already 
been counted and this `except` adds `reconnect_failure` on top, giving 
`reconnect_attempt` 1 and outcomes 2. That's the reconciliation your E2E table 
demonstrates, and it's the thing this PR exists to establish. Line 193 has the 
same problem in the trace: it overwrites a `resumable.decision` that was 
already set to `reconnect` / `already_succeeded` / `terminal_resubmit`.
   
   Worth saying why I don't think this is academic: `self.log` is structlog, 
and remote log processors are injected into the global chain — #66633 is 
currently open about exactly that path. A logging call failing because a remote 
sink is unreachable is the same class of event this metric is meant to catch.
   
   Narrowing the `try` to the fetch and classification, then branching on the 
result outside it, makes the exclusivity structural:
   
   ```python
   try:
       status = self.get_job_status(external_id, context)
       if self.is_job_active(status):
           decision = "reconnect"
       elif self.is_job_succeeded(status):
           decision = "already_succeeded"
       else:
           decision = "terminal_resubmit"
   except Exception:
       span.set_attribute("resumable.decision", "reconnect_failure")
       stats.incr("resumable_job.reconnect_failure", tags=stats_tags)
       raise
   ```



##########
task-sdk/tests/task_sdk/bases/test_resumablejobmixin.py:
##########
@@ -296,6 +296,26 @@ def test_terminal_resubmit_fires_when_job_failed(self):
         assert "resumable_job.reconnect_success" not in called_names
         assert "resumable_job.fresh_submit" not in called_names
 
+    @pytest.mark.parametrize("failing_method", ["get_job_status", 
"is_job_active"])

Review Comment:
   Both parameters here raise *before* any outcome counter fires, so they're 
the two positions where double-counting can't happen — the test passes for the 
same reason the bug hides.
   
   Two additions worth having: `"is_job_succeeded"` (the third classification 
call, at `resumablejobmixin.py:169`) is one more token in this list; and a case 
where the *log* call raises after an outcome was counted would pin the point in 
my comment on line 189 — it fails today, and passes once the `try` is narrowed.



-- 
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]

Reply via email to