shepherd44 opened a new pull request, #73906:
URL: https://github.com/apache/airflow/pull/73906

   ## Problem
   
   `EksPodOperator.invoke_defer_method` takes a shortcut when the pod already 
reached a terminal state before the operator got to defer. Instead of calling 
`self.defer(...)`, it invokes `self.trigger_reentry(...)` inline — and 
**discards its return value**.
   
   ```python
   if context and (container_state == ContainerState.TERMINATED or 
container_state == ContainerState.FAILED):
       self.log.info("Skipping deferral as pod is already in a terminal state")
       self.trigger_reentry(          # <- return value dropped
           context=context,
           event={...},
       )
   else:
       self.defer(trigger=trigger, method_name="trigger_reentry")
   ```
   
   `trigger_reentry()` produces the task result, which Airflow pushes to XCom 
as `return_value` when `do_xcom_push=True`. Dropping it means **the task still 
succeeds but `return_value` never reaches XCom**. Only `pod_name` and 
`pod_namespace` are present, because `execute_async` pushes those separately — 
so the failure looks like "XCom exists but `return_value` is missing".
   
   Downstream tasks that read the upstream XCom then fail at template 
rendering, or worse, silently treat the missing value as "no data" and skip.
   
   The condition is easy to hit with short-lived pods: the container finishes 
before the operator reaches the defer call. The distinguishing log line is 
`Skipping deferral as pod is already in a terminal state` with no `Pushing 
xcom` afterwards, versus `Pausing task as DEFERRED` + `Pushing xcom` on the 
normal path.
   
   ## Relationship to #72821
   
   This mirrors #72821, which fixes the same discarded-return bug in the base 
`KubernetesPodOperator`. **That fix alone does not cover this one**, because 
`EksPodOperator` overrides `invoke_defer_method` entirely — the base 
implementation is never reached.
   
   The two are independent; this one stands on its own.
   
   ## Change
   
   - `return self.trigger_reentry(...)` in the shortcut path
   - return annotation `-> None` → `-> Any` (the method now returns the task 
result)
   - the `else` after the new `return` is redundant, so it is dropped (ruff 
RET505)
   
   ## Test
   
   Added 
`test_invoke_defer_method_returns_trigger_reentry_result_when_pod_already_terminal`,
 which forces the terminal-state path and asserts the return value propagates. 
It fails without this change:
   
   ```
   FAILED 
providers/amazon/tests/unit/amazon/aws/operators/test_eks.py::TestEksPodOperator::
     
test_invoke_defer_method_returns_trigger_reentry_result_when_pod_already_terminal
   ```
   
   `providers/amazon/tests/unit/amazon/aws/operators/test_eks.py`: 65 passed. 
`ruff check` and `ruff format --check` clean.
   
   ## Gen-AI Assisted contribution
   
   This change was prepared with the assistance of Claude (Claude Code). Per 
the Gen-AI contribution guidelines: the diff is small and I have reviewed and 
understand every line, the root cause was traced against the actual provider 
source, and the tests were written and run locally. I take responsibility for 
the content of this PR and can explain the reasoning behind it.
   


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