namanjain24-sudo commented on code in PR #73810:
URL: https://github.com/apache/airflow/pull/73810#discussion_r4181551863


##########
airflow-core/tests/unit/jobs/test_scheduler_job.py:
##########
@@ -759,6 +759,39 @@ def test_retired_attempt_events_do_not_modify_replacement(
         assert ti.try_number == 2
         assert ti.external_executor_id == ("current_worker" if include_current 
else "replacement")
 
+    def test_process_executor_events_sets_state_in_callers_transaction(self, 
dag_maker):
+        """
+        Setting a task instance's state must not commit the caller's 
transaction.
+
+        ``settings.Session`` is scoped, so ``ti.set_state()`` without 
``session`` resolved to the
+        scheduler's own session and ``create_session()`` committed and closed 
it on exit. That released
+        the scheduler's row locks mid-batch and detached the task instances 
still to be processed, so
+        changes made to them afterwards were never written.
+
+        This covers the Dag-not-found path: the Dag can't be loaded, so the 
task instance is marked
+        with the executor's reported state directly, without a session it 
would otherwise resolve to
+        the scheduler's own scoped session and commit early.
+        """
+        session = settings.Session()
+        with dag_maker(dag_id="test_executor_events_callers_transaction", 
fileloc="/test_path1/"):
+            task1 = EmptyOperator(task_id="test_task", retries=2)
+        ti1 = dag_maker.create_dagrun().get_task_instance(task1.task_id)
+        ti1.state = TaskInstanceState.QUEUED
+        session.merge(ti1)
+        session.commit()
+
+        executor = MockExecutor(do_update=False)
+        job_runner = SchedulerJobRunner(Job(), executors=[executor])
+        job_runner.scheduler_dag_bag = mock.MagicMock()
+        job_runner.scheduler_dag_bag.get_dag_for_run.side_effect = 
Exception("failed")
+        executor.event_buffer[ti1.key] = State.FAILED, None

Review Comment:
   Good catch — addressed in cf2bac6a7a: the test now keys the executor event 
by `TaskInstanceUuid(ti1.id)` (matching the other tests in this file) and 
asserts `ti1.state == TaskInstanceState.FAILED` right after 
`_process_executor_events` and before the `session.rollback()`, so it fails 
again without the session fix.
   
   ---
   Drafted-by: Claude Code (Sonnet 5); reviewed by @namanjain24-sudo before 
posting



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