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:
   @/tmp/reply_body.txt



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