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

shahar1 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/main by this push:
     new 7800fc44c07 Fix scheduler crash on restart with queued or running Edge 
tasks (#73832)
7800fc44c07 is described below

commit 7800fc44c07a2fc2457a4aee41cf12db9579538c
Author: PoAn Yang <[email protected]>
AuthorDate: Sat Oct 3 01:18:09 2026 +0800

    Fix scheduler crash on restart with queued or running Edge tasks (#73832)
    
    * Fix scheduler crash on restart with queued or running Edge tasks
    
    Signed-off-by: PoAn Yang <[email protected]>
    
    * Add the Edge scheduler-crash fix to the edge3 5.0.0 changelog
    
    ---------
    
    Signed-off-by: PoAn Yang <[email protected]>
    Co-authored-by: Shahar Epstein <[email protected]>
---
 providers/edge3/docs/changelog.rst                 |  1 +
 .../providers/edge3/executors/edge_executor.py     | 19 +++---
 .../unit/edge3/executors/test_edge_executor.py     | 71 ++++++++++++++++++++++
 3 files changed, 82 insertions(+), 9 deletions(-)

diff --git a/providers/edge3/docs/changelog.rst 
b/providers/edge3/docs/changelog.rst
index 1741821c983..65a12ff1dde 100644
--- a/providers/edge3/docs/changelog.rst
+++ b/providers/edge3/docs/changelog.rst
@@ -49,6 +49,7 @@ Features
 Bug Fixes
 ~~~~~~~~~
 
+* ``Fix scheduler crash on restart with queued or running Edge tasks (#73832)``
 * ``Fix airflow edge list-workers always showing null concurrency (#72959)``
 
 Misc
diff --git 
a/providers/edge3/src/airflow/providers/edge3/executors/edge_executor.py 
b/providers/edge3/src/airflow/providers/edge3/executors/edge_executor.py
index 6990b0097b7..cbc5e5e9811 100644
--- a/providers/edge3/src/airflow/providers/edge3/executors/edge_executor.py
+++ b/providers/edge3/src/airflow/providers/edge3/executors/edge_executor.py
@@ -43,7 +43,7 @@ from airflow.providers.edge3.models.types import (
 from airflow.providers.edge3.version_compat import AIRFLOW_V_3_4_PLUS
 from airflow.utils.db import DBLocks, create_global_lock
 from airflow.utils.helpers import prune_dict
-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 AIRFLOW_V_3_4_PLUS:
@@ -422,10 +422,7 @@ class EdgeExecutor(BaseExecutor):
         )
         self.log.info("Revoked task instance %s from EdgeExecutor", ti.key)
 
-    @provide_session
-    def try_adopt_task_instances(
-        self, tis: Sequence[TaskInstance], *, session: Session = NEW_SESSION
-    ) -> Sequence[TaskInstance]:
+    def try_adopt_task_instances(self, tis: Sequence[TaskInstance]) -> 
Sequence[TaskInstance]:
         """
         Adopt the task instances whose job is still in flight in the edge_job 
table.
 
@@ -435,10 +432,14 @@ class EdgeExecutor(BaseExecutor):
 
         :return: any TaskInstances that were unable to be adopted
         """
-        tracked_keys = self._get_tracked_job_keys(
-            session,
-            states=(TaskInstanceState.QUEUED, TaskInstanceState.RESTARTING, 
TaskInstanceState.RUNNING),
-        )
+        # The scheduler calls this without passing a session, while its own 
scoped session still holds
+        # ``tis``. A scoped session here would be that same session, and 
create_session() would commit
+        # and close it on exit, which detaches ``tis`` before the scheduler 
reads them again.
+        with create_session(scoped=False) as session:
+            tracked_keys = self._get_tracked_job_keys(
+                session,
+                states=(TaskInstanceState.QUEUED, 
TaskInstanceState.RESTARTING, TaskInstanceState.RUNNING),
+            )
         self.running.update(ti.key for ti in tis if ti.key in tracked_keys)
         return [ti for ti in tis if ti.key not in tracked_keys]
 
diff --git a/providers/edge3/tests/unit/edge3/executors/test_edge_executor.py 
b/providers/edge3/tests/unit/edge3/executors/test_edge_executor.py
index 463c515cc3f..bb636de8382 100644
--- a/providers/edge3/tests/unit/edge3/executors/test_edge_executor.py
+++ b/providers/edge3/tests/unit/edge3/executors/test_edge_executor.py
@@ -29,6 +29,8 @@ import time_machine
 from sqlalchemy import delete, select
 
 from airflow.executors.workloads import BundleInfo, ExecuteTask
+from airflow.jobs.job import Job
+from airflow.jobs.scheduler_job_runner import SchedulerJobRunner
 from airflow.models.taskinstance import TaskInstance
 from airflow.providers.common.compat.sdk import Stats, TaskInstanceKey, conf, 
timezone
 from airflow.providers.edge3.executors.edge_executor import EdgeExecutor
@@ -38,6 +40,7 @@ from airflow.providers.edge3.models.types import 
EXECUTE_CALLBACK_TAG
 from airflow.utils.session import create_session
 from airflow.utils.state import TaskInstanceState
 
+from tests_common.test_utils.compat import EmptyOperator
 from tests_common.test_utils.config import conf_vars
 from tests_common.test_utils.version_compat import AIRFLOW_V_3_2_PLUS, 
AIRFLOW_V_3_3_PLUS
 
@@ -407,6 +410,74 @@ class TestEdgeExecutor:
         assert key not in executor.running
         assert key not in executor.queued_tasks
 
+    def test_try_adopt_task_instances_keeps_caller_session_open(self):
+        key = TaskInstanceKey(
+            dag_id="test_dag", run_id="test_run", task_id="test_task", 
map_index=-1, try_number=1
+        )
+        with create_session() as session:
+            session.add(
+                EdgeJobModel(
+                    dag_id="test_dag",
+                    task_id="test_task",
+                    run_id="test_run",
+                    map_index=-1,
+                    try_number=1,
+                    state=TaskInstanceState.QUEUED,
+                    queue="default",
+                    command="mock",
+                    concurrency_slots=1,
+                )
+            )
+            session.commit()
+        executor = EdgeExecutor()
+
+        # The scheduler calls try_adopt_task_instances() without a session and 
keeps using the objects
+        # it loaded into its own scoped session.
+        with create_session() as session:
+            job = session.scalar(select(EdgeJobModel))
+            executor.try_adopt_task_instances([mock.Mock(spec=TaskInstance, 
key=key)])
+
+            assert job in session
+        assert executor.running == {key}
+
+    
@mock.patch("airflow.executors.executor_loader.ExecutorLoader.init_executors", 
autospec=True)
+    def test_scheduler_restart_adopts_queued_edge_task(self, 
mock_init_executors, dag_maker, session):
+        with dag_maker("test_dag", session=session):
+            EmptyOperator(task_id="test_task")
+        dag_run = dag_maker.create_dagrun()
+        previous_scheduler_job = Job()
+        restarted_scheduler_job = Job()
+        session.add_all([previous_scheduler_job, restarted_scheduler_job])
+        session.flush()
+        ti = dag_run.get_task_instance("test_task", session=session)
+        ti.state = TaskInstanceState.QUEUED
+        ti.queued_by_job_id = previous_scheduler_job.id
+        session.add(
+            EdgeJobModel(
+                dag_id=ti.dag_id,
+                task_id=ti.task_id,
+                run_id=ti.run_id,
+                map_index=ti.map_index,
+                try_number=ti.try_number,
+                state=TaskInstanceState.QUEUED,
+                queue="default",
+                command="mock",
+                concurrency_slots=1,
+            )
+        )
+        session.commit()
+        # A restarted scheduler loads the task instance from the database, not 
from this session.
+        session.expunge_all()
+        executor = EdgeExecutor()
+        mock_init_executors.return_value = [executor]
+
+        SchedulerJobRunner(job=restarted_scheduler_job, 
num_runs=0).adopt_or_reset_orphaned_tasks()
+
+        ti.refresh_from_db(session=session)
+        assert ti.state == TaskInstanceState.QUEUED
+        assert ti.queued_by_job_id == restarted_scheduler_job.id
+        assert executor.running == {ti.key}
+
 
 class TestEdgeExecutorMultiTeam:
     """Tests for multi-team (AIP-67) support in EdgeExecutor."""

Reply via email to