This is an automated email from the ASF dual-hosted git repository.

ashb pushed a commit to branch local-executor-bookkeeping
in repository https://gitbox.apache.org/repos/asf/airflow.git

commit 2b2b21b2082d23fb02c623db21ce353f50b1aac2
Author: Ash Berlin-Taylor <[email protected]>
AuthorDate: Sat Oct 3 12:36:01 2026 +0100

    fixup! Improve LocalExecutor bookkeeping: correctly add tasks to the 
`running` list
---
 airflow-core/src/airflow/executors/local_executor.py | 3 ++-
 airflow-core/tests/unit/jobs/test_scheduler_job.py   | 9 +++++----
 2 files changed, 7 insertions(+), 5 deletions(-)

diff --git a/airflow-core/src/airflow/executors/local_executor.py 
b/airflow-core/src/airflow/executors/local_executor.py
index a279073131f..9e1f1f8bbf8 100644
--- a/airflow-core/src/airflow/executors/local_executor.py
+++ b/airflow-core/src/airflow/executors/local_executor.py
@@ -190,7 +190,8 @@ class LocalExecutor(BaseExecutor):
         for pid, proc in self.workers.items():
             if not proc.is_alive():
                 self._read_results()
-                # A worker killed between dequeue and START cannot identify 
its workload.
+                # A worker killed between dequeue and START has no entry here; 
the scheduler's
+                # stuck-in-queued handling releases that workload through 
revoke_task.
                 if (key := self._worker_tasks.pop(pid, None)) is not None:
                     self._finish_dispatch(key, state_class_for_key(key).FAILED)
                 to_remove.add(pid)
diff --git a/airflow-core/tests/unit/jobs/test_scheduler_job.py 
b/airflow-core/tests/unit/jobs/test_scheduler_job.py
index 08d9588b8bc..17fa1c693b6 100644
--- a/airflow-core/tests/unit/jobs/test_scheduler_job.py
+++ b/airflow-core/tests/unit/jobs/test_scheduler_job.py
@@ -4330,12 +4330,13 @@ class TestSchedulerJob:
         ti.queued_dttm = timezone.utcnow() - timedelta(minutes=15)
         session.commit()
 
+        class NoRevokeExecutor(BaseExecutor):
+            pass
+
         assert "revoke_task" in BaseExecutor.__dict__
-        # this is just verifying that LocalExecutor is good enough for this 
test
-        # in that it does not implement revoke_task
-        assert "revoke_task" not in LocalExecutor.__dict__
+        assert "revoke_task" not in NoRevokeExecutor.__dict__
         scheduler_job = Job()
-        job_runner = SchedulerJobRunner(job=scheduler_job, num_runs=0, 
executors=[LocalExecutor()])
+        job_runner = SchedulerJobRunner(job=scheduler_job, num_runs=0, 
executors=[NoRevokeExecutor()])
         job_runner._task_queued_timeout = 300
         job_runner._handle_tasks_stuck_in_queued()
 

Reply via email to