kaxil commented on code in PR #73917:
URL: https://github.com/apache/airflow/pull/73917#discussion_r4157599952


##########
providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py:
##########
@@ -195,6 +202,21 @@ def _list_pods(self, query_kwargs):
 
         return pods
 
+    def _task_instance_id_from_pod(self, pod, kube_client=None) -> UUID | None:
+        if TASK_INSTANCE_ID_ANNOTATION in (pod.metadata.annotations or {}):

Review Comment:
   Thanks for adding the full-pod read. This fast path never fires against a 
real cluster, though: `_list_pods` returns DynamicClient items, so 
`pod.metadata.annotations` is a `ResourceField`, and `in` on it returns False 
even when the annotation is there (I checked on kubernetes 35.0.0 and 36.0.3; 
`.get()` works). So every adoption and revoke candidate gets a 
`read_namespaced_pod`, including pods that are already annotated. The tests 
pass because they build `k8s.V1Pod` with dict annotations. Using 
`.get(TASK_INSTANCE_ID_ANNOTATION) is not None`, as `task_instance_id_from_pod` 
does, plus a test pod built from `ResourceInstance(...).items` would cover this.



##########
providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py:
##########
@@ -1146,21 +1202,31 @@ def adopt_launched_task(
             self.log.error("attempting to adopt taskinstance which was not 
specified by database: %s", ti_key)
             return
 
-        new_worker_id_label = 
self._make_safe_label_value(self.scheduler_job_id)
+        key: TaskInstanceUuid | TaskInstanceKey = ti_key
+        metadata = {"labels": {"airflow-worker": 
self._make_safe_label_value(self.scheduler_job_id)}}
+        if self.supports_task_instance_uuid:
+            task_id = self._task_instance_id_from_pod(pod, kube_client)

Review Comment:
   I see 
`test_metadata_pod_get_handles_missing_pod_and_propagates_other_errors` pins 
this, but a 429 or 5xx on this GET now aborts the whole adoption pass. It sits 
outside the `ApiException` guard the patch below has, and 
`adopt_or_reset_orphaned_tasks` only catches `OperationalError`, so it escapes 
the scheduler's startup call. A failed patch only refuses that one pod. Could 
the GET do the same? The revoke loop has the same problem, and the 
`BaseExecutor.revoke_task` docstring says it should not raise.



##########
providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py:
##########
@@ -1146,21 +1202,31 @@ def adopt_launched_task(
             self.log.error("attempting to adopt taskinstance which was not 
specified by database: %s", ti_key)
             return
 
-        new_worker_id_label = 
self._make_safe_label_value(self.scheduler_job_id)
+        key: TaskInstanceUuid | TaskInstanceKey = ti_key
+        metadata = {"labels": {"airflow-worker": 
self._make_safe_label_value(self.scheduler_job_id)}}
+        if self.supports_task_instance_uuid:
+            task_id = self._task_instance_id_from_pod(pod, kube_client)
+            if task_id is None or task_id != tis_to_flush_by_key[ti_key].id:
+                self.log.warning(
+                    "Cannot adopt pod %s without a matching task instance 
UUID", pod.metadata.name
+                )
+                return
+            key = self.get_task_key(tis_to_flush_by_key[ti_key])
+            metadata["annotations"] = {TASK_INSTANCE_ID_ANNOTATION: str(key)}
         from kubernetes.client.rest import ApiException

Review Comment:
   Nit: `ApiException` is imported at module level now (line 44), so this 
import inside the function body, and the ones at 594, 680 and 1349, can go.



##########
providers/celery/tests/unit/celery/executors/test_celery_executor.py:
##########
@@ -439,8 +444,9 @@ def test_cleanup_stuck_queued_tasks(
             executor.job_id = 1
             if hasattr(executor, "_register_task"):
                 executor._register_task(ti)
-            executor.running = {ti.key}
-            executor.workloads = {ti.key: AsyncResult("231")}
+            key = executor.get_task_key(ti) if 
executor.supports_task_instance_uuid else ti.key
+            executor.running = {key}
+            executor.workloads = {key: AsyncResult("231")}
             assert executor.has_task(ti)
             with pytest.warns(AirflowProviderDeprecationWarning, 
match="cleanup_stuck_queued_tasks"):
                 executor.cleanup_stuck_queued_tasks(tis=tis)

Review Comment:
   The `mock_fail.assert_called()` below still passes if 
`cleanup_stuck_queued_tasks` goes back to `self.fail(ti.key)`. Could it be 
`assert_called_once_with(key)`, so the UUID change on the `fail` line is 
actually tested?



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