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]