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