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]