luc-pimentel commented on code in PR #73810:
URL: https://github.com/apache/airflow/pull/73810#discussion_r4178208937


##########
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:
   After #73916, this `ti1.key` event is dropped before it reaches `set_state`, 
so the test passes even without the fix...
   
   Should we rebase and use `TaskInstanceUuid(ti1.id)` like the other tests do 
now? With that change and adding an assertion that the state is `FAILED` before 
the rollback, the test fails again without the fix (which is the correct 
behavior).



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