From 5ae24ee841e2dad77f72cdfdc0f3991e6cde1eba Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Guilherme=20Da=20Silva=20Gon=C3=A7alves?= Date: Sun, 12 Jul 2026 19:44:57 +0200 Subject: [PATCH 01/10] Fence KubernetesExecutor launches with executor token --- .../execution_api/datamodels/taskinstance.py | 3 + .../execution_api/routes/task_instances.py | 15 ++++ .../src/airflow/executors/workloads/task.py | 2 +- .../versions/head/test_task_instances.py | 29 +++++++ .../executors/kubernetes_executor.py | 82 ++++++++++++++++++- .../executors/kubernetes_executor_utils.py | 3 + .../cncf/kubernetes/pod_generator.py | 3 + .../executors/test_kubernetes_executor.py | 43 ++++++++++ task-sdk/src/airflow/sdk/api/client.py | 12 ++- .../airflow/sdk/api/datamodels/_generated.py | 2 + .../airflow/sdk/execution_time/supervisor.py | 7 +- task-sdk/tests/task_sdk/api/test_client.py | 3 +- .../execution_time/test_supervisor.py | 3 + 13 files changed, 201 insertions(+), 6 deletions(-) diff --git a/airflow-core/src/airflow/api_fastapi/execution_api/datamodels/taskinstance.py b/airflow-core/src/airflow/api_fastapi/execution_api/datamodels/taskinstance.py index ad051b3e6d340..fb9ff13262943 100644 --- a/airflow-core/src/airflow/api_fastapi/execution_api/datamodels/taskinstance.py +++ b/airflow-core/src/airflow/api_fastapi/execution_api/datamodels/taskinstance.py @@ -64,6 +64,8 @@ class TIEnterRunningPayload(StrictBaseModel): """Process Identifier on `hostname`""" start_date: UtcDateTime """When the task started executing""" + external_executor_id: str | None = None + """Executor launch token assigned when the task was queued""" # Create an enum to give a nice name in the generated datamodels @@ -290,6 +292,7 @@ class TaskInstance(BaseModel): # hand-built instances (tests, dry runs) valid; the executor workload # always sends the real value. queue: str = "default" + external_executor_id: str | None = None class AssetReferenceAssetEventDagRun(StrictBaseModel): diff --git a/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py b/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py index 34c3dc35406f8..213cd50368f53 100644 --- a/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py +++ b/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py @@ -163,6 +163,7 @@ def ti_run( TI.hostname, TI.unixname, TI.pid, + TI.external_executor_id, # This selects the raw JSON value, bypassing the deserialization -- we want that to happen on the # client column("next_kwargs", JSON), @@ -190,6 +191,7 @@ def ti_run( # We exclude_unset to avoid updating fields that are not set in the payload data = ti_run_payload.model_dump(exclude_unset=True) + payload_external_executor_id = data.pop("external_executor_id", None) # don't update start date when resuming from deferral if ti.next_kwargs: @@ -208,6 +210,19 @@ def ti_run( ti_run_payload.pid, ): log.info("Duplicate start request received", hostname=ti_run_payload.hostname) + elif ti.external_executor_id is not None and ti.external_executor_id != payload_external_executor_id: + log.warning( + "Cannot start Task Instance with stale executor launch token", + expected_external_executor_id=ti.external_executor_id, + provided_external_executor_id=payload_external_executor_id, + ) + raise HTTPException( + status_code=status.HTTP_409_CONFLICT, + detail={ + "reason": "stale_executor_launch", + "message": "TI executor launch token does not match the current queued task instance", + }, + ) elif previous_state not in (TaskInstanceState.QUEUED, TaskInstanceState.RESTARTING): log.warning( "Cannot start Task Instance in invalid state", diff --git a/airflow-core/src/airflow/executors/workloads/task.py b/airflow-core/src/airflow/executors/workloads/task.py index 3099fe1d77485..7271a6e479b80 100644 --- a/airflow-core/src/airflow/executors/workloads/task.py +++ b/airflow-core/src/airflow/executors/workloads/task.py @@ -45,7 +45,7 @@ class TaskInstanceDTO(TaskInstance): pool_slots: int priority_weight: int - external_executor_id: str | None = Field(default=None, exclude=True) + external_executor_id: str | None = None executor_config: dict | None = Field(default=None, exclude=True) # TODO: Task-SDK: Can we replace TaskInstanceKey with just the uuid across the codebase? diff --git a/airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py b/airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py index 8a152bebe0d3f..ed2b9bc4561de 100644 --- a/airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py +++ b/airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py @@ -323,6 +323,35 @@ def test_ti_run_state_to_running( ) assert response.status_code == 409 + def test_ti_run_rejects_stale_external_executor_id(self, client, session, create_task_instance): + """A worker can only start the launch token currently assigned to the queued TI.""" + ti = create_task_instance( + task_id="test_ti_run_rejects_stale_external_executor_id", + state=State.QUEUED, + dagrun_state=DagRunState.RUNNING, + session=session, + dag_id=str(uuid4()), + ) + ti.external_executor_id = "current-launch-token" + session.commit() + + response = client.patch( + f"/execution/task-instances/{ti.id}/run", + json={ + "state": "running", + "hostname": "random-hostname", + "unixname": "random-unixname", + "pid": 100, + "start_date": "2024-09-30T12:00:00Z", + "external_executor_id": "stale-launch-token", + }, + ) + + assert response.status_code == 409 + assert response.json()["detail"]["reason"] == "stale_executor_launch" + session.refresh(ti) + assert ti.state == State.QUEUED + def test_ti_run_returns_execution_token( self, client, exec_app, session, create_task_instance, time_machine ): diff --git a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py index 8aaa430b92108..0b702a09d584c 100644 --- a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py +++ b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py @@ -36,7 +36,7 @@ from datetime import datetime, timedelta from itertools import chain from queue import Empty, Queue -from typing import TYPE_CHECKING, Any +from typing import TYPE_CHECKING, Any, ClassVar from deprecated import deprecated from kubernetes.dynamic import DynamicClient @@ -100,6 +100,7 @@ class KubernetesExecutor(BaseExecutor): RUNNING_POD_LOG_LINES = 100 supports_ad_hoc_ti_run: bool = True supports_multi_team: bool = True + pre_assigns_external_executor_id: ClassVar[bool] = True if TYPE_CHECKING and AIRFLOW_V_3_0_PLUS: # In the v3 path, we store workloads, not commands as strings. @@ -394,6 +395,82 @@ def _process_workloads(self, workloads: Sequence[workloads.All]) -> None: self.execute_async(key=key, command=command, queue=queue, executor_config=executor_config) self.running.add(key) + @provide_session + def _should_create_pod_for_job(self, task: KubernetesJob, *, session: Session = NEW_SESSION) -> bool: + """Check whether a queued Kubernetes job still owns the current task launch token.""" + from airflow.executors.workloads import ExecuteTask + from airflow.models.taskinstance import TaskInstance + + if not task.command or not isinstance(task.command[0], ExecuteTask): + return True + + workload = task.command[0] + workload_ti = workload.ti + launch_token = workload_ti.external_executor_id + + try: + scheduler_job_id = int(self.scheduler_job_id) if self.scheduler_job_id is not None else None + except ValueError: + self.log.debug( + "Skipping stale Kubernetes workload check because scheduler_job_id %r is not numeric", + self.scheduler_job_id, + ) + return True + + ti = session.execute( + select( + TaskInstance.id, + TaskInstance.state, + TaskInstance.try_number, + TaskInstance.queued_by_job_id, + TaskInstance.external_executor_id, + ).where(TaskInstance.id == workload_ti.id) + ).one_or_none() + if ti is None: + self.log.info( + "Dropping stale Kubernetes workload for %s because task instance id %s no longer exists", + task.key, + workload_ti.id, + ) + return False + + _, state, try_number, queued_by_job_id, external_executor_id = ti + if ( + state == TaskInstanceState.QUEUED + and try_number == workload_ti.try_number + and queued_by_job_id == scheduler_job_id + and external_executor_id == launch_token + ): + return True + + self.log.info( + "Dropping stale Kubernetes workload for %s because current task instance launch does not " + "match queued workload: ti_id=%s state=%s try_number=%s queued_by_job_id=%s " + "external_executor_id=%s workload_try_number=%s workload_external_executor_id=%s " + "scheduler_job_id=%s", + task.key, + workload_ti.id, + state, + try_number, + queued_by_job_id, + external_executor_id, + workload_ti.try_number, + launch_token, + scheduler_job_id, + ) + return False + + def _discard_stale_pod_creation_task(self, task: KubernetesJob) -> None: + """Remove executor bookkeeping for a stale job that will not create a pod.""" + self.running.discard(task.key) + if self.event_buffer.get(task.key) == (TaskInstanceState.QUEUED, self.scheduler_job_id): + self.event_buffer.pop(task.key, None) + self.task_publish_retries.pop(task.key, None) + Stats.incr( + "kubernetes_executor.stale_workload_dropped", + tags=prune_dict({"team_name": self.team_name}), + ) + def sync(self) -> None: """Synchronize task state.""" if TYPE_CHECKING: @@ -471,6 +548,9 @@ def sync(self) -> None: try: key = task.key + if not self._should_create_pod_for_job(task): + self._discard_stale_pod_creation_task(task) + continue self.kube_scheduler.run_next(task) self.task_publish_retries.pop(key, None) except PodReconciliationError as e: diff --git a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor_utils.py b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor_utils.py index 2c4bd306e52d7..0862b1f2799da 100644 --- a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor_utils.py +++ b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor_utils.py @@ -575,11 +575,13 @@ def run_next(self, next_job: KubernetesJob) -> None: kube_image = next_job.kube_image or self.kube_config.kube_image dag_id, task_id, run_id, try_number, map_index = key + external_executor_id = None if len(command) == 1: from airflow.executors.workloads import ExecuteTask if isinstance(command[0], ExecuteTask): workload = command[0] + external_executor_id = workload.ti.external_executor_id command = workload_to_command_args(workload) else: raise ValueError( @@ -606,6 +608,7 @@ def run_next(self, next_job: KubernetesJob) -> None: map_index=map_index, date=None, run_id=run_id, + external_executor_id=external_executor_id, args=list(command), pod_override_object=kube_executor_config, base_worker_pod=base_worker_pod, diff --git a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/pod_generator.py b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/pod_generator.py index ca499d363fc77..7f4458c99bd54 100644 --- a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/pod_generator.py +++ b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/pod_generator.py @@ -373,6 +373,7 @@ def construct_pod( scheduler_job_id: str, run_id: str | None = None, map_index: int = -1, + external_executor_id: str | None = None, *, with_mutation_hook: bool = False, ) -> k8s.V1Pod: @@ -410,6 +411,8 @@ def construct_pod( annotations[get_logical_date_key()] = date.isoformat() if run_id: annotations["run_id"] = run_id + if external_executor_id: + annotations["external_executor_id"] = external_executor_id main_container = k8s.V1Container( name="base", diff --git a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/executors/test_kubernetes_executor.py b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/executors/test_kubernetes_executor.py index ee541bd83fdbe..d221926a8a9ae 100644 --- a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/executors/test_kubernetes_executor.py +++ b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/executors/test_kubernetes_executor.py @@ -969,6 +969,49 @@ def test_run_next_exception_requeue( finally: kubernetes_executor.end() + @pytest.mark.db_test + @pytest.mark.skipif(not AIRFLOW_V_3_0_PLUS, reason="workloads are used on Airflow 3+") + @mock.patch("airflow.providers.cncf.kubernetes.executors.kubernetes_executor_utils.KubernetesJobWatcher") + @mock.patch("airflow.providers.cncf.kubernetes.kube_client.get_kube_client") + def test_sync_drops_workload_with_stale_external_executor_id( + self, + mock_get_kube_client, + mock_kubernetes_job_watcher, + create_task_instance, + session, + ): + """A delayed Kubernetes workload must not create a pod after its launch token changes.""" + from airflow.executors.workloads import ExecuteTask + + executor = self.kubernetes_executor + executor.start() + try: + ti = create_task_instance(state=TaskInstanceState.QUEUED) + ti.queued_by_job_id = executor.job_id + ti.external_executor_id = "current-launch-token" + session.merge(ti) + session.commit() + + workload = ExecuteTask.make(ti) + executor.queue_workload(workload, session=session) + executor._process_workloads([workload]) + + ti.external_executor_id = "new-launch-token" + session.merge(ti) + session.commit() + + assert executor.kube_scheduler is not None + executor.kube_scheduler.run_next = mock.Mock() + executor.sync() + + executor.kube_scheduler.run_next.assert_not_called() + assert executor.task_queue is not None + assert executor.task_queue.empty() + assert ti.key not in executor.running + assert ti.key not in executor.event_buffer + finally: + executor.end() + @pytest.mark.skipif( AirflowKubernetesScheduler is None, reason="kubernetes python package is not installed" ) diff --git a/task-sdk/src/airflow/sdk/api/client.py b/task-sdk/src/airflow/sdk/api/client.py index 7e681c1114956..092bdc1e8e0a9 100644 --- a/task-sdk/src/airflow/sdk/api/client.py +++ b/task-sdk/src/airflow/sdk/api/client.py @@ -247,9 +247,17 @@ class TaskInstanceOperations: def __init__(self, client: Client): self.client = client - def start(self, id: uuid.UUID, pid: int, when: datetime) -> TIRunContext: + def start( + self, id: uuid.UUID, pid: int, when: datetime, external_executor_id: str | None = None + ) -> TIRunContext: """Tell the API server that this TI has started running.""" - body = TIEnterRunningPayload(pid=pid, hostname=get_hostname(), unixname=getuser(), start_date=when) + body = TIEnterRunningPayload( + pid=pid, + hostname=get_hostname(), + unixname=getuser(), + start_date=when, + external_executor_id=external_executor_id, + ) try: resp = self.client.patch(f"task-instances/{id}/run", content=body.model_dump_json()) diff --git a/task-sdk/src/airflow/sdk/api/datamodels/_generated.py b/task-sdk/src/airflow/sdk/api/datamodels/_generated.py index 8e5bfc1d076e5..0594feb1e8150 100644 --- a/task-sdk/src/airflow/sdk/api/datamodels/_generated.py +++ b/task-sdk/src/airflow/sdk/api/datamodels/_generated.py @@ -292,6 +292,7 @@ class TIEnterRunningPayload(BaseModel): unixname: Annotated[str, Field(title="Unixname")] pid: Annotated[int, Field(title="Pid")] start_date: Annotated[AwareDatetime, Field(title="Start Date")] + external_executor_id: Annotated[str | None, Field(title="External Executor Id")] = None class TIHeartbeatInfo(BaseModel): @@ -561,6 +562,7 @@ class TaskInstance(BaseModel): hostname: Annotated[str | None, Field(title="Hostname")] = None context_carrier: Annotated[dict[str, Any] | None, Field(title="Context Carrier")] = None queue: Annotated[str | None, Field(title="Queue")] = "default" + external_executor_id: Annotated[str | None, Field(title="External Executor Id")] = None class BundleInfo(BaseModel): diff --git a/task-sdk/src/airflow/sdk/execution_time/supervisor.py b/task-sdk/src/airflow/sdk/execution_time/supervisor.py index 87311f02da7a1..16c3e8ba2156e 100644 --- a/task-sdk/src/airflow/sdk/execution_time/supervisor.py +++ b/task-sdk/src/airflow/sdk/execution_time/supervisor.py @@ -1385,7 +1385,12 @@ def _on_child_started( # We've forked, but the task won't start doing anything until we send it the StartupDetails # message. But before we do that, we need to tell the server it's started (so it has the chance to # tell us "no, stop!" for any reason) - ti_context = self.client.task_instances.start(ti.id, self.pid, datetime.now(tz=timezone.utc)) + ti_context = self.client.task_instances.start( + ti.id, + self.pid, + datetime.now(tz=timezone.utc), + getattr(ti, "external_executor_id", None), + ) self._should_retry = ti_context.should_retry self._last_successful_heartbeat = time.monotonic() except Exception: diff --git a/task-sdk/tests/task_sdk/api/test_client.py b/task-sdk/tests/task_sdk/api/test_client.py index 0f9b8130e2654..5d2c05b30f85c 100644 --- a/task-sdk/tests/task_sdk/api/test_client.py +++ b/task-sdk/tests/task_sdk/api/test_client.py @@ -344,6 +344,7 @@ def handle_request(request: httpx.Request) -> httpx.Response: assert actual_body["pid"] == 100 assert actual_body["start_date"] == start_date assert actual_body["state"] == "running" + assert actual_body["external_executor_id"] == "launch-token" return httpx.Response( status_code=200, json=ti_context.model_dump(mode="json"), @@ -351,7 +352,7 @@ def handle_request(request: httpx.Request) -> httpx.Response: return httpx.Response(status_code=400, json={"detail": "Bad Request"}) client = make_client(transport=httpx.MockTransport(handle_request)) - resp = client.task_instances.start(ti_id, 100, start_date) + resp = client.task_instances.start(ti_id, 100, start_date, external_executor_id="launch-token") assert resp == ti_context assert call_count == 3 diff --git a/task-sdk/tests/task_sdk/execution_time/test_supervisor.py b/task-sdk/tests/task_sdk/execution_time/test_supervisor.py index 7ba463567a17a..7aa993975fdd4 100644 --- a/task-sdk/tests/task_sdk/execution_time/test_supervisor.py +++ b/task-sdk/tests/task_sdk/execution_time/test_supervisor.py @@ -957,6 +957,8 @@ def test_start_raises_task_already_running_and_kills_subprocess(self): def handle_request(request: httpx.Request) -> httpx.Response: if request.url.path == f"/task-instances/{ti_id}/run": + body = json.loads(request.read()) + assert body["external_executor_id"] == "launch-token" return httpx.Response( 409, json={ @@ -985,6 +987,7 @@ def subprocess_main(): try_number=1, dag_version_id=uuid7(), queue="default", + external_executor_id="launch-token", ), client=make_client(transport=httpx.MockTransport(handle_request)), target=subprocess_main, From f7d87b7c7bdd2c1c1f85e358d8bc2ac348abe0d1 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Guilherme=20Da=20Silva=20Gon=C3=A7alves?= Date: Mon, 13 Jul 2026 11:44:52 +0200 Subject: [PATCH 02/10] Make launch-token fence upgrade-safe and cover matching-token paths Only reject a task-instance /run when the worker actually presents an external_executor_id. Older Task SDK workers (mid rolling-upgrade, before the field existed) omit it entirely; fencing them would 409 every launch for pre-assigning executors (KubernetesExecutor, CeleryExecutor) until all workers are upgraded. Absent tokens now fall back to state-based validation. Also document that dropping a stale Kubernetes workload intentionally leaves the task instance in place (owned by the newer launch / reclaimed by queued-task timeout or adoption) rather than failing it. Tests: - API: matching token accepted; missing token accepted for old clients. - KubernetesExecutor: matching token creates the pod (positive counterpart to the stale-drop test). Co-Authored-By: Claude Opus 4.8 --- .../execution_api/routes/task_instances.py | 12 +++- .../versions/head/test_task_instances.py | 60 +++++++++++++++++++ .../executors/kubernetes_executor.py | 10 +++- .../executors/test_kubernetes_executor.py | 40 +++++++++++++ 4 files changed, 120 insertions(+), 2 deletions(-) diff --git a/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py b/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py index 213cd50368f53..2dade78453db1 100644 --- a/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py +++ b/airflow-core/src/airflow/api_fastapi/execution_api/routes/task_instances.py @@ -191,6 +191,12 @@ def ti_run( # We exclude_unset to avoid updating fields that are not set in the payload data = ti_run_payload.model_dump(exclude_unset=True) + # Only fence on the launch token when the worker actually presented one. Older Task SDK + # clients (e.g. mid rolling-upgrade, before this field existed) omit it entirely; treating + # that as a stale launch would 409 every task start for pre-assigning executors + # (KubernetesExecutor, CeleryExecutor) until every worker is upgraded. When it is absent we + # fall back to state-based validation, matching the pre-token behavior. + payload_has_external_executor_id = "external_executor_id" in data payload_external_executor_id = data.pop("external_executor_id", None) # don't update start date when resuming from deferral @@ -210,7 +216,11 @@ def ti_run( ti_run_payload.pid, ): log.info("Duplicate start request received", hostname=ti_run_payload.hostname) - elif ti.external_executor_id is not None and ti.external_executor_id != payload_external_executor_id: + elif ( + payload_has_external_executor_id + and ti.external_executor_id is not None + and ti.external_executor_id != payload_external_executor_id + ): log.warning( "Cannot start Task Instance with stale executor launch token", expected_external_executor_id=ti.external_executor_id, diff --git a/airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py b/airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py index ed2b9bc4561de..01df9c285f394 100644 --- a/airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py +++ b/airflow-core/tests/unit/api_fastapi/execution_api/versions/head/test_task_instances.py @@ -352,6 +352,66 @@ def test_ti_run_rejects_stale_external_executor_id(self, client, session, create session.refresh(ti) assert ti.state == State.QUEUED + def test_ti_run_accepts_matching_external_executor_id(self, client, session, create_task_instance): + """A worker presenting the launch token currently assigned to the TI can start it.""" + ti = create_task_instance( + task_id="test_ti_run_accepts_matching_external_executor_id", + state=State.QUEUED, + dagrun_state=DagRunState.RUNNING, + session=session, + dag_id=str(uuid4()), + ) + ti.external_executor_id = "current-launch-token" + session.commit() + + response = client.patch( + f"/execution/task-instances/{ti.id}/run", + json={ + "state": "running", + "hostname": "random-hostname", + "unixname": "random-unixname", + "pid": 100, + "start_date": "2024-09-30T12:00:00Z", + "external_executor_id": "current-launch-token", + }, + ) + + assert response.status_code == 200 + session.refresh(ti) + assert ti.state == State.RUNNING + # The launch token is preserved; matching it must not clear or overwrite the field. + assert ti.external_executor_id == "current-launch-token" + + def test_ti_run_allows_missing_external_executor_id_for_old_clients( + self, client, session, create_task_instance + ): + """An older worker that omits the launch token must not be fenced out (rolling upgrade).""" + ti = create_task_instance( + task_id="test_ti_run_allows_missing_external_executor_id_for_old_clients", + state=State.QUEUED, + dagrun_state=DagRunState.RUNNING, + session=session, + dag_id=str(uuid4()), + ) + ti.external_executor_id = "current-launch-token" + session.commit() + + # Payload deliberately omits ``external_executor_id`` entirely, as a pre-token Task SDK would. + response = client.patch( + f"/execution/task-instances/{ti.id}/run", + json={ + "state": "running", + "hostname": "random-hostname", + "unixname": "random-unixname", + "pid": 100, + "start_date": "2024-09-30T12:00:00Z", + }, + ) + + assert response.status_code == 200 + session.refresh(ti) + assert ti.state == State.RUNNING + def test_ti_run_returns_execution_token( self, client, exec_app, session, create_task_instance, time_machine ): diff --git a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py index 0b702a09d584c..84a6ff5e1f602 100644 --- a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py +++ b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py @@ -461,7 +461,15 @@ def _should_create_pod_for_job(self, task: KubernetesJob, *, session: Session = return False def _discard_stale_pod_creation_task(self, task: KubernetesJob) -> None: - """Remove executor bookkeeping for a stale job that will not create a pod.""" + """ + Remove executor bookkeeping for a stale job that will not create a pod. + + We intentionally do not fail the task instance here. Dropping the pod creation leaves the + row in its current DB state (typically still QUEUED under the *newer* launch that replaced + this workload). The newer launch owns the task from now on: either its own workload creates + the pod, or the scheduler's queued-task timeout / task-instance adoption reclaims the row. + Failing it here would clobber that newer launch. + """ self.running.discard(task.key) if self.event_buffer.get(task.key) == (TaskInstanceState.QUEUED, self.scheduler_job_id): self.event_buffer.pop(task.key, None) diff --git a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/executors/test_kubernetes_executor.py b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/executors/test_kubernetes_executor.py index d221926a8a9ae..203d9a8137029 100644 --- a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/executors/test_kubernetes_executor.py +++ b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/executors/test_kubernetes_executor.py @@ -1012,6 +1012,46 @@ def test_sync_drops_workload_with_stale_external_executor_id( finally: executor.end() + @pytest.mark.db_test + @pytest.mark.skipif(not AIRFLOW_V_3_0_PLUS, reason="workloads are used on Airflow 3+") + @mock.patch("airflow.providers.cncf.kubernetes.executors.kubernetes_executor_utils.KubernetesJobWatcher") + @mock.patch("airflow.providers.cncf.kubernetes.kube_client.get_kube_client") + def test_sync_creates_pod_when_external_executor_id_matches( + self, + mock_get_kube_client, + mock_kubernetes_job_watcher, + create_task_instance, + session, + ): + """A Kubernetes workload whose launch token still matches the DB must create a pod.""" + from airflow.executors.workloads import ExecuteTask + + executor = self.kubernetes_executor + executor.start() + try: + ti = create_task_instance(state=TaskInstanceState.QUEUED) + ti.queued_by_job_id = executor.job_id + ti.external_executor_id = "current-launch-token" + session.merge(ti) + session.commit() + + workload = ExecuteTask.make(ti) + executor.queue_workload(workload, session=session) + executor._process_workloads([workload]) + + # The DB token is left untouched, so the queued workload still owns the launch. + assert executor.kube_scheduler is not None + executor.kube_scheduler.run_next = mock.Mock() + executor.sync() + + executor.kube_scheduler.run_next.assert_called_once() + created_job = executor.kube_scheduler.run_next.call_args.args[0] + assert created_job.key == ti.key + assert executor.task_queue is not None + assert executor.task_queue.empty() + finally: + executor.end() + @pytest.mark.skipif( AirflowKubernetesScheduler is None, reason="kubernetes python package is not installed" ) From 23abee03dad6d0ef1e9b0709999131f4c622391a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Guilherme=20Da=20Silva=20Gon=C3=A7alves?= Date: Mon, 13 Jul 2026 14:12:43 +0200 Subject: [PATCH 03/10] Preserve KubernetesExecutor launch token in queued events --- .../kubernetes/executors/kubernetes_executor.py | 14 ++++++++++++-- .../executors/test_kubernetes_executor.py | 1 + 2 files changed, 13 insertions(+), 2 deletions(-) diff --git a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py index 84a6ff5e1f602..8665a52f75825 100644 --- a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py +++ b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py @@ -364,9 +364,18 @@ def execute_async( queue, ) - self.event_buffer[key] = (TaskInstanceState.QUEUED, self.scheduler_job_id) job = KubernetesJob(key, command, kube_executor_config, pod_template_file, coordinator_kube_image) self.pod_launch_attempts[key] = _PodLaunchAttempt(job=job) + event_info = self.scheduler_job_id + try: + from airflow.executors.workloads import ExecuteTask + + if len(command) == 1 and isinstance(command[0], ExecuteTask): + event_info = command[0].ti.external_executor_id or self.scheduler_job_id + except (ImportError, TypeError): + pass + + self.event_buffer[key] = (TaskInstanceState.QUEUED, event_info) self.task_queue.put(job) def queue_workload(self, workload: workloads.All, session: Session | None) -> None: @@ -471,7 +480,8 @@ def _discard_stale_pod_creation_task(self, task: KubernetesJob) -> None: Failing it here would clobber that newer launch. """ self.running.discard(task.key) - if self.event_buffer.get(task.key) == (TaskInstanceState.QUEUED, self.scheduler_job_id): + queued_event = self.event_buffer.get(task.key) + if queued_event is not None and queued_event[0] == TaskInstanceState.QUEUED: self.event_buffer.pop(task.key, None) self.task_publish_retries.pop(task.key, None) Stats.incr( diff --git a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/executors/test_kubernetes_executor.py b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/executors/test_kubernetes_executor.py index 203d9a8137029..7bb70542a7af1 100644 --- a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/executors/test_kubernetes_executor.py +++ b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/executors/test_kubernetes_executor.py @@ -1038,6 +1038,7 @@ def test_sync_creates_pod_when_external_executor_id_matches( workload = ExecuteTask.make(ti) executor.queue_workload(workload, session=session) executor._process_workloads([workload]) + assert executor.event_buffer[ti.key] == (TaskInstanceState.QUEUED, "current-launch-token") # The DB token is left untouched, so the queued workload still owns the launch. assert executor.kube_scheduler is not None From f674c5f0a769b88cbbe50a411f4cdb24719c2749 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Guilherme=20Da=20Silva=20Gon=C3=A7alves?= Date: Mon, 20 Jul 2026 12:59:53 +0200 Subject: [PATCH 04/10] Fix launch-token CI failures --- .../kubernetes/executors/kubernetes_executor.py | 17 ++++++++++++----- .../edge3/worker_api/v2-edge-generated.yaml | 5 +++++ .../observability/metrics/metrics_template.yaml | 6 ++++++ .../sdk/execution_time/schema/schema.json | 12 ++++++++++++ .../task_sdk/execution_time/test_supervisor.py | 2 +- 5 files changed, 36 insertions(+), 6 deletions(-) diff --git a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py index 8665a52f75825..4b6b296c60d27 100644 --- a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py +++ b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py @@ -369,11 +369,14 @@ def execute_async( event_info = self.scheduler_job_id try: from airflow.executors.workloads import ExecuteTask - - if len(command) == 1 and isinstance(command[0], ExecuteTask): - event_info = command[0].ti.external_executor_id or self.scheduler_job_id - except (ImportError, TypeError): + except ImportError: pass + else: + try: + if len(command) == 1 and isinstance(command[0], ExecuteTask): + event_info = command[0].ti.external_executor_id or self.scheduler_job_id + except TypeError: + pass self.event_buffer[key] = (TaskInstanceState.QUEUED, event_info) self.task_queue.put(job) @@ -407,7 +410,11 @@ def _process_workloads(self, workloads: Sequence[workloads.All]) -> None: @provide_session def _should_create_pod_for_job(self, task: KubernetesJob, *, session: Session = NEW_SESSION) -> bool: """Check whether a queued Kubernetes job still owns the current task launch token.""" - from airflow.executors.workloads import ExecuteTask + try: + from airflow.executors.workloads import ExecuteTask + except ImportError: + return True + from airflow.models.taskinstance import TaskInstance if not task.command or not isinstance(task.command[0], ExecuteTask): diff --git a/providers/edge3/src/airflow/providers/edge3/worker_api/v2-edge-generated.yaml b/providers/edge3/src/airflow/providers/edge3/worker_api/v2-edge-generated.yaml index 4bd6a68af4326..814dfb38e870a 100644 --- a/providers/edge3/src/airflow/providers/edge3/worker_api/v2-edge-generated.yaml +++ b/providers/edge3/src/airflow/providers/edge3/worker_api/v2-edge-generated.yaml @@ -1318,6 +1318,11 @@ components: type: string title: Queue default: default + external_executor_id: + anyOf: + - type: string + - type: 'null' + title: External Executor Id pool_slots: type: integer title: Pool Slots diff --git a/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml b/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml index a561db2d7aef6..a60f741bf4e12 100644 --- a/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml +++ b/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml @@ -626,6 +626,12 @@ metrics: legacy_name: "-" name_variables: ["status"] + - name: "kubernetes_executor.stale_workload_dropped" + description: "Number of stale queued KubernetesExecutor workloads dropped before worker pod creation." + type: "counter" + legacy_name: "-" + name_variables: [] + # ========== # Timers # ========== diff --git a/task-sdk/src/airflow/sdk/execution_time/schema/schema.json b/task-sdk/src/airflow/sdk/execution_time/schema/schema.json index 0ec8fe4e49a1c..04f1fe6251dc3 100644 --- a/task-sdk/src/airflow/sdk/execution_time/schema/schema.json +++ b/task-sdk/src/airflow/sdk/execution_time/schema/schema.json @@ -4998,6 +4998,18 @@ "default": "default", "title": "Queue", "type": "string" + }, + "external_executor_id": { + "anyOf": [ + { + "type": "string" + }, + { + "type": "null" + } + ], + "default": null, + "title": "External Executor Id" } }, "required": [ diff --git a/task-sdk/tests/task_sdk/execution_time/test_supervisor.py b/task-sdk/tests/task_sdk/execution_time/test_supervisor.py index 7aa993975fdd4..6530944ac5027 100644 --- a/task-sdk/tests/task_sdk/execution_time/test_supervisor.py +++ b/task-sdk/tests/task_sdk/execution_time/test_supervisor.py @@ -801,7 +801,7 @@ def mock_monotonic(): assert exit_code == 0, captured_logs # Validate calls to the client - mock_client.task_instances.start.assert_called_once_with(ti.id, mocker.ANY, mocker.ANY) + mock_client.task_instances.start.assert_called_once_with(ti.id, mocker.ANY, mocker.ANY, None) mock_client.task_instances.heartbeat.assert_called_once_with(ti.id, pid=mocker.ANY) mock_client.task_instances.defer.assert_called_once_with( ti.id, From 5cd26c897d7d401a72f3a86e8338e49dff76ae64 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Guilherme=20Da=20Silva=20Gon=C3=A7alves?= Date: Mon, 20 Jul 2026 16:09:55 +0200 Subject: [PATCH 05/10] Add execution API version for launch token --- .../api_fastapi/execution_api/versions/__init__.py | 2 ++ .../execution_api/versions/v2026_06_30.py | 12 ++++++++++++ 2 files changed, 14 insertions(+) diff --git a/airflow-core/src/airflow/api_fastapi/execution_api/versions/__init__.py b/airflow-core/src/airflow/api_fastapi/execution_api/versions/__init__.py index dc7035d31e3c9..5a620533674d4 100644 --- a/airflow-core/src/airflow/api_fastapi/execution_api/versions/__init__.py +++ b/airflow-core/src/airflow/api_fastapi/execution_api/versions/__init__.py @@ -47,6 +47,7 @@ AddPartitionDateField, AddRetryPolicyFields, AddTaskAndAssetStateStoreEndpoints, + AddTaskInstanceExternalExecutorIdField, AddTaskInstanceQueueField, AddTeamNameField, AddVariableKeysEndpoint, @@ -59,6 +60,7 @@ AddVariableKeysEndpoint, AddConnectionTestEndpoint, AddAwaitingInputStatePayload, + AddTaskInstanceExternalExecutorIdField, AddTaskInstanceQueueField, AddRetryPolicyFields, AddTeamNameField, diff --git a/airflow-core/src/airflow/api_fastapi/execution_api/versions/v2026_06_30.py b/airflow-core/src/airflow/api_fastapi/execution_api/versions/v2026_06_30.py index e89e2ed04cc5d..8762c887d8611 100644 --- a/airflow-core/src/airflow/api_fastapi/execution_api/versions/v2026_06_30.py +++ b/airflow-core/src/airflow/api_fastapi/execution_api/versions/v2026_06_30.py @@ -29,6 +29,7 @@ DagRun, TaskInstance, TIAwaitingInputStatePayload, + TIEnterRunningPayload, TIRetryStatePayload, TIRunContext, ) @@ -61,6 +62,17 @@ class AddTaskInstanceQueueField(VersionChange): instructions_to_migrate_to_previous_version = (schema(TaskInstance).field("queue").didnt_exist,) +class AddTaskInstanceExternalExecutorIdField(VersionChange): + """Add the `external_executor_id` field to task instance launch payloads.""" + + description = __doc__ + + instructions_to_migrate_to_previous_version = ( + schema(TaskInstance).field("external_executor_id").didnt_exist, + schema(TIEnterRunningPayload).field("external_executor_id").didnt_exist, + ) + + class AddAwaitingInputStatePayload(VersionChange): """Add the awaiting_input task instance state transition payload (Human-in-the-loop, no trigger).""" From fe3a8387cf502db2c54d016d29988d54c8c7b80e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Guilherme=20Da=20Silva=20Gon=C3=A7alves?= Date: Mon, 20 Jul 2026 16:51:31 +0200 Subject: [PATCH 06/10] Add supervisor schema version marker --- .../schema/versions/v2026_06_16.py | 25 +++++++++++++++++++ ts-sdk/src/generated/supervisor.ts | 2 ++ 2 files changed, 27 insertions(+) create mode 100644 task-sdk/src/airflow/sdk/execution_time/schema/versions/v2026_06_16.py diff --git a/task-sdk/src/airflow/sdk/execution_time/schema/versions/v2026_06_16.py b/task-sdk/src/airflow/sdk/execution_time/schema/versions/v2026_06_16.py new file mode 100644 index 0000000000000..c62fd6551f537 --- /dev/null +++ b/task-sdk/src/airflow/sdk/execution_time/schema/versions/v2026_06_16.py @@ -0,0 +1,25 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +""" +Supervisor schema changes for the initial 2026-06-16 schema version. + +Cadwyn does not allow VersionChange entries on the first supervisor schema +version because there is no previous version to migrate to. This file is kept +with the schema-changing PR so the supervisor schema version hook records that +the current head-version schema was intentionally updated. +""" diff --git a/ts-sdk/src/generated/supervisor.ts b/ts-sdk/src/generated/supervisor.ts index ab2632831ab95..68afa76289434 100644 --- a/ts-sdk/src/generated/supervisor.ts +++ b/ts-sdk/src/generated/supervisor.ts @@ -187,6 +187,7 @@ export type ContextCarrier = { [k: string]: unknown; } | null; export type Queue = string; +export type ExternalExecutorId = string | null; export type IsFailureCallback = boolean | null; export type Type12 = "DagCallbackRequest"; export type File = string; @@ -949,6 +950,7 @@ export interface TaskInstance { hostname?: Hostname; context_carrier?: ContextCarrier; queue?: Queue; + external_executor_id?: ExternalExecutorId; } /** * Request for DAG File Parsing. From 3f73e8cf96b0bba2adeb58c0c282038639aa43db Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Guilherme=20Da=20Silva=20Gon=C3=A7alves?= Date: Mon, 20 Jul 2026 17:29:54 +0200 Subject: [PATCH 07/10] Update workload serialization test for launch token --- airflow-core/tests/unit/executors/test_workloads.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/airflow-core/tests/unit/executors/test_workloads.py b/airflow-core/tests/unit/executors/test_workloads.py index 37fbcd96ce950..6d41a472dd2a2 100644 --- a/airflow-core/tests/unit/executors/test_workloads.py +++ b/airflow-core/tests/unit/executors/test_workloads.py @@ -159,7 +159,7 @@ def test_workload_ti_round_trips_through_sdk_generated_model(): ) dumped = ti.model_dump(mode="json") - assert "external_executor_id" not in dumped + assert dumped["external_executor_id"] == "celery-id" assert "executor_config" not in dumped # Executor-side scheduling fields stay on the workload wire (older workers # deserialize the workload with a model that requires them) but are not @@ -169,6 +169,7 @@ def test_workload_ti_round_trips_through_sdk_generated_model(): received = GeneratedTaskInstance.model_validate(dumped) assert received.queue == "jdk-17" + assert received.external_executor_id == "celery-id" assert received.map_index == 3 assert not hasattr(received, "pool_slots") From acc6fdb81e20621e870b2afeff17970142dce0c5 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Guilherme=20Da=20Silva=20Gon=C3=A7alves?= Date: Tue, 21 Jul 2026 00:09:17 +0200 Subject: [PATCH 08/10] Handle launch tokens on older Airflow workloads --- .../cncf/kubernetes/executors/kubernetes_executor.py | 6 +++--- .../cncf/kubernetes/executors/test_kubernetes_executor.py | 8 ++++++-- 2 files changed, 9 insertions(+), 5 deletions(-) diff --git a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py index 4b6b296c60d27..9319aabe5f2ff 100644 --- a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py +++ b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py @@ -374,7 +374,7 @@ def execute_async( else: try: if len(command) == 1 and isinstance(command[0], ExecuteTask): - event_info = command[0].ti.external_executor_id or self.scheduler_job_id + event_info = getattr(command[0].ti, "external_executor_id", None) or self.scheduler_job_id except TypeError: pass @@ -422,7 +422,7 @@ def _should_create_pod_for_job(self, task: KubernetesJob, *, session: Session = workload = task.command[0] workload_ti = workload.ti - launch_token = workload_ti.external_executor_id + launch_token = getattr(workload_ti, "external_executor_id", None) try: scheduler_job_id = int(self.scheduler_job_id) if self.scheduler_job_id is not None else None @@ -455,7 +455,7 @@ def _should_create_pod_for_job(self, task: KubernetesJob, *, session: Session = state == TaskInstanceState.QUEUED and try_number == workload_ti.try_number and queued_by_job_id == scheduler_job_id - and external_executor_id == launch_token + and (launch_token is None or external_executor_id == launch_token) ): return True diff --git a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/executors/test_kubernetes_executor.py b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/executors/test_kubernetes_executor.py index 7bb70542a7af1..fd61437fb8ef3 100644 --- a/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/executors/test_kubernetes_executor.py +++ b/providers/cncf/kubernetes/tests/unit/cncf/kubernetes/executors/test_kubernetes_executor.py @@ -970,7 +970,9 @@ def test_run_next_exception_requeue( kubernetes_executor.end() @pytest.mark.db_test - @pytest.mark.skipif(not AIRFLOW_V_3_0_PLUS, reason="workloads are used on Airflow 3+") + @pytest.mark.skipif( + not AIRFLOW_V_3_2_PLUS, reason="workloads include external_executor_id on Airflow 3.2+" + ) @mock.patch("airflow.providers.cncf.kubernetes.executors.kubernetes_executor_utils.KubernetesJobWatcher") @mock.patch("airflow.providers.cncf.kubernetes.kube_client.get_kube_client") def test_sync_drops_workload_with_stale_external_executor_id( @@ -1013,7 +1015,9 @@ def test_sync_drops_workload_with_stale_external_executor_id( executor.end() @pytest.mark.db_test - @pytest.mark.skipif(not AIRFLOW_V_3_0_PLUS, reason="workloads are used on Airflow 3+") + @pytest.mark.skipif( + not AIRFLOW_V_3_2_PLUS, reason="workloads include external_executor_id on Airflow 3.2+" + ) @mock.patch("airflow.providers.cncf.kubernetes.executors.kubernetes_executor_utils.KubernetesJobWatcher") @mock.patch("airflow.providers.cncf.kubernetes.kube_client.get_kube_client") def test_sync_creates_pod_when_external_executor_id_matches( From 32c0ffe21b5288e9c83a54ad8edb6b3958037412 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Guilherme=20Da=20Silva=20Gon=C3=A7alves?= Date: Tue, 21 Jul 2026 09:51:20 +0200 Subject: [PATCH 09/10] Avoid DB access for legacy Kubernetes jobs --- .../cncf/kubernetes/executors/kubernetes_executor.py | 9 ++++++--- 1 file changed, 6 insertions(+), 3 deletions(-) diff --git a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py index 9319aabe5f2ff..f63e4962b40ba 100644 --- a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py +++ b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py @@ -59,7 +59,7 @@ from airflow.providers.common.compat.sdk import Stats, conf from airflow.utils.helpers import prune_dict from airflow.utils.log.logging_mixin import remove_escape_codes -from airflow.utils.session import NEW_SESSION, provide_session +from airflow.utils.session import NEW_SESSION, create_session, provide_session from airflow.utils.state import TaskInstanceState if TYPE_CHECKING: @@ -407,8 +407,7 @@ def _process_workloads(self, workloads: Sequence[workloads.All]) -> None: self.execute_async(key=key, command=command, queue=queue, executor_config=executor_config) self.running.add(key) - @provide_session - def _should_create_pod_for_job(self, task: KubernetesJob, *, session: Session = NEW_SESSION) -> bool: + def _should_create_pod_for_job(self, task: KubernetesJob, *, session: Session | None = None) -> bool: """Check whether a queued Kubernetes job still owns the current task launch token.""" try: from airflow.executors.workloads import ExecuteTask @@ -424,6 +423,10 @@ def _should_create_pod_for_job(self, task: KubernetesJob, *, session: Session = workload_ti = workload.ti launch_token = getattr(workload_ti, "external_executor_id", None) + if session is None: + with create_session() as session: + return self._should_create_pod_for_job(task, session=session) + try: scheduler_job_id = int(self.scheduler_job_id) if self.scheduler_job_id is not None else None except ValueError: From 43d4e84d1f287abbd9ed5d01e9901691eaf3c83d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Guilherme=20Da=20Silva=20Gon=C3=A7alves?= Date: Fri, 24 Jul 2026 19:14:56 +0200 Subject: [PATCH 10/10] Remove stale workload metric from launch-token PR --- .../cncf/kubernetes/executors/kubernetes_executor.py | 4 ---- .../observability/metrics/metrics_template.yaml | 6 ------ 2 files changed, 10 deletions(-) diff --git a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py index f63e4962b40ba..a8d93729ea804 100644 --- a/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py +++ b/providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py @@ -494,10 +494,6 @@ def _discard_stale_pod_creation_task(self, task: KubernetesJob) -> None: if queued_event is not None and queued_event[0] == TaskInstanceState.QUEUED: self.event_buffer.pop(task.key, None) self.task_publish_retries.pop(task.key, None) - Stats.incr( - "kubernetes_executor.stale_workload_dropped", - tags=prune_dict({"team_name": self.team_name}), - ) def sync(self) -> None: """Synchronize task state.""" diff --git a/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml b/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml index a60f741bf4e12..a561db2d7aef6 100644 --- a/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml +++ b/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml @@ -626,12 +626,6 @@ metrics: legacy_name: "-" name_variables: ["status"] - - name: "kubernetes_executor.stale_workload_dropped" - description: "Number of stale queued KubernetesExecutor workloads dropped before worker pod creation." - type: "counter" - legacy_name: "-" - name_variables: [] - # ========== # Timers # ==========