Repository navigation
Pass the scheduler's session to set_state when processing executor events - #73810
namanjain24-sudo wants to merge 2 commits into
Conversation
b9d0f3f to
84fca1f
Compare
| 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 |
There was a problem hiding this 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).
There was a problem hiding this comment.
Good catch — addressed in cf2bac6: 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
…ents process_executor_events calls ti.set_state() without a session when the Dag of a finished task instance cannot be loaded. set_state is @provide_session and settings.Session is scoped, so create_session() returned the scheduler's own session and committed and closed it on exit, in the middle of the batch. A second call site with the same bug (a cleared task instance reported terminated) was superseded by apache#73554, which replaced that whole code path with TaskInstance.complete_restart(), already passing the scheduler's session correctly.
After the attempt-UUID executor event rework, the event buffer is keyed by TaskInstanceUuid rather than the legacy TaskInstanceKey. Keying this regression test's event by ti1.key meant the event never matched a captured task identity and was discarded before reaching set_state, so the test passed even without the session fix it exists to guard. Key it the same way the other tests in this file do, and assert the state actually flips to FAILED before the rollback so a regression is caught again.
84fca1f to
cf2bac6
Compare
Opened fresh from #73213, which was closed under the new open-PR-limit policy and then failed to reopen (both
gh pr reopenand the UI button errored with a generic "Could not open the pull request") — following @potiuk's guidance on Slack to push/open fresh from the branch when reopening doesn't work.SchedulerJobRunner.process_executor_eventscallsti.set_state()withoutsessionin two places: when a cleared (RESTARTING) task instance is reported as successfully terminated, and when the Dag of a finished task instance cannot be loaded.TaskInstance.set_stateis@provide_sessionandsettings.Sessionis scoped, socreate_session()returns the scheduler's own session and commits and closes it on exit. This is the same mechanism as #67850 and #71968.process_executor_eventsruns insidewith create_session()in_run_scheduler_loop, not underprohibit_commit, so nothing raises. Instead, in the middle of an executor-event batch:FOR UPDATE SKIP LOCKEDrow locks taken on the batch's task instances;close()detaches the task instances loaded for the batch, so changes made to them afterwards are never written. In a local run with oneRESTARTING→SUCCESSevent and fourQUEUEDevents carrying an external executor id, none of the fourexternal_executor_idvalues reached the database without this change, and all four did with it (SQLite and PostgreSQL 16).This passes the scheduler's session at both call sites, as
_enqueue_task_instances_with_queued_statealready does for its ownti.set_state()call.Tests:
test_process_executor_events_sets_state_in_callers_transactioncovers both paths: the state change has to roll back with the caller's transaction. On main both cases fail (the task instance is already committed asNone/failed); with this change they pass.test_scheduler_job.pyand prek, including mypy for airflow-core, pass locally on SQLite and PostgreSQL 16.Not covered here. In the same loop,
executor.send_callback()reachesDatabaseCallbackSink.send, which is also@provide_sessionand gets no session. So when the task instance has anon_failure_callback/on_retry_callback(reproduced locally with this change applied) or email configured, the scheduler's transaction is still committed at that call._maybe_requeue_stuck_tiand_purge_task_instances_without_heartbeatscallsend_callback()the same way. Fixing that means either adding asessionargument toBaseExecutor.send_callbackandBaseCallbackSink.send, which executors that overridesend_callbackwould then have to accept, or handing the session to the sink directly, asTriggeralready does withDatabaseCallbackSink().send(callback=request, session=session). I kept this PR to theset_statecalls and am happy to follow up with whichever approach you prefer.Was generative AI tooling used to co-author this PR?
Generated-by: a Gen-AI coding assistant, following the guidelines. I reviewed the change and ran the checks above.