From b3158a0e240f43ddd91c5b48fc20f6cfc6be8924 Mon Sep 17 00:00:00 2001 From: JH0917 Date: Sun, 21 Jun 2026 16:07:09 +0900 Subject: [PATCH 1/3] Detect supervisor-subprocess trigger count mismatch in Triggerer --- .../src/airflow/jobs/triggerer_job_runner.py | 20 ++++++++ .../tests/unit/jobs/test_triggerer_job.py | 50 +++++++++++++++++++ 2 files changed, 70 insertions(+) diff --git a/airflow-core/src/airflow/jobs/triggerer_job_runner.py b/airflow-core/src/airflow/jobs/triggerer_job_runner.py index 8f8e930b07fce..8042e6bd2d373 100644 --- a/airflow-core/src/airflow/jobs/triggerer_job_runner.py +++ b/airflow-core/src/airflow/jobs/triggerer_job_runner.py @@ -332,6 +332,7 @@ class TriggerStateChanges(BaseModel): # Format of list[str] is the exc traceback format failures: list[tuple[int, list[str] | None]] | None = None finished: list[int] | None = None + num_running: int = 0 class TriggerStateSync(BaseModel): type: Literal["TriggerStateSync"] = "TriggerStateSync" @@ -603,6 +604,8 @@ def _handle_request(self, msg: ToTriggerSupervisor, log: FilteringBoundLogger, r # handle leaks for every failed upload. factory.close() + self.check_for_unhandled_triggers(msg.num_running) + # Drain the persist confirmations accumulated since the last sync. events_persisted: list[int] = [] while self.persisted_event_seqs: @@ -783,6 +786,21 @@ def clean_unused(self) -> None: """Remove triggers that are no longer needed.""" Trigger.clean_unused() + def check_for_unhandled_triggers(self, num_running: int) -> None: + """ + Shut down if the subprocess trigger count disagrees with the supervisor. + + Only valid between finished-removal and to_create-addition in ``_handle_request``. + """ + expected = len(self.running_triggers) + if expected != num_running: + log.error( + "Trigger count mismatch: expected %d, subprocess reports %d. Shutting down.", + expected, + num_running, + ) + self.stop = True + def handle_failed_triggers(self): """ Handle "failed" triggers. - ones that errored or exited before they sent an event. @@ -1513,6 +1531,7 @@ def process_trigger_events(self, finished_ids: list[int]) -> messages.TriggerSta events=events_to_send if events_to_send else None, finished=finished_ids if finished_ids else None, failures=failures_to_send if failures_to_send else None, + num_running=len(self.triggers), ) def sanitize_trigger_events(self, msg: messages.TriggerStateChanges) -> messages.TriggerStateChanges: @@ -1541,6 +1560,7 @@ def sanitize_trigger_events(self, msg: messages.TriggerStateChanges) -> messages events=events_to_send if events_to_send else None, finished=msg.finished, failures=msg.failures, + num_running=msg.num_running, ) async def sync_state_to_supervisor(self, finished_ids: list[int]) -> None: diff --git a/airflow-core/tests/unit/jobs/test_triggerer_job.py b/airflow-core/tests/unit/jobs/test_triggerer_job.py index 648ecf602fab0..5f03e7d0ab316 100644 --- a/airflow-core/tests/unit/jobs/test_triggerer_job.py +++ b/airflow-core/tests/unit/jobs/test_triggerer_job.py @@ -3253,3 +3253,53 @@ async def _drive(): trigger_id, _event, seq = events[0] assert trigger_id == 1 assert seq is None + + +class TestCheckForUnhandledTriggers: + """Tests for supervisor-vs-subprocess trigger count mismatch detection.""" + + def test_mismatch_shuts_down(self, jobless_supervisor): + jobless_supervisor.running_triggers = {1, 2, 3} + + jobless_supervisor.check_for_unhandled_triggers(num_running=0) + assert jobless_supervisor.stop is True + + def test_matching_counts_no_shutdown(self, jobless_supervisor): + jobless_supervisor.running_triggers = {1, 2, 3} + + jobless_supervisor.check_for_unhandled_triggers(num_running=3) + assert jobless_supervisor.stop is False + + def test_no_triggers_no_shutdown(self, jobless_supervisor): + jobless_supervisor.running_triggers = set() + + jobless_supervisor.check_for_unhandled_triggers(num_running=0) + assert jobless_supervisor.stop is False + + def test_subprocess_more_than_expected_shuts_down(self, jobless_supervisor): + jobless_supervisor.running_triggers = {1} + + jobless_supervisor.check_for_unhandled_triggers(num_running=3) + assert jobless_supervisor.stop is True + + def test_handle_request_checks_before_adding_to_create(self, jobless_supervisor, mocker): + """The check fires after finished processing but before to_create is added to running_triggers.""" + mocker.patch.object(TriggerRunnerSupervisor, "send_msg", autospec=True) + jobless_supervisor.running_triggers = {1, 2} + jobless_supervisor.creating_triggers.append( + mocker.MagicMock(id=3), + ) + + jobless_supervisor._handle_request( + messages.TriggerStateChanges( + events=None, + failures=None, + finished=[1], + num_running=1, + ), + log=MagicMock(spec=FilteringBoundLogger), + req_id=1, + ) + + assert jobless_supervisor.stop is False + assert jobless_supervisor.running_triggers == {2, 3} From 630fc2b0df29b434273c2c34156092df0244e21e Mon Sep 17 00:00:00 2001 From: JH0917 Date: Sun, 12 Jul 2026 13:48:59 +0900 Subject: [PATCH 2/3] Fix false-positive path in trigger count mismatch check Remove the incorrect difference_update on msg.failures (which would break for serialization failures where the trigger is still running). Instead, include creation-failure IDs in finished_ids so the supervisor naturally removes them from running_triggers via the existing finished processing path. Co-Authored-By: Claude Opus 4.6 --- .../src/airflow/jobs/triggerer_job_runner.py | 2 + .../tests/unit/jobs/test_triggerer_job.py | 42 +++++++++++++++++++ 2 files changed, 44 insertions(+) diff --git a/airflow-core/src/airflow/jobs/triggerer_job_runner.py b/airflow-core/src/airflow/jobs/triggerer_job_runner.py index 8042e6bd2d373..3a616e4cf82f7 100644 --- a/airflow-core/src/airflow/jobs/triggerer_job_runner.py +++ b/airflow-core/src/airflow/jobs/triggerer_job_runner.py @@ -1526,6 +1526,8 @@ def process_trigger_events(self, finished_ids: list[int]) -> messages.TriggerSta trigger_id, exc = self.failed_triggers.popleft() tb = format_exception(type(exc), exc, exc.__traceback__) if exc else None failures_to_send.append((trigger_id, tb)) + if trigger_id not in self.triggers: + finished_ids.append(trigger_id) return messages.TriggerStateChanges( events=events_to_send if events_to_send else None, diff --git a/airflow-core/tests/unit/jobs/test_triggerer_job.py b/airflow-core/tests/unit/jobs/test_triggerer_job.py index 5f03e7d0ab316..bbb8eb5e745d8 100644 --- a/airflow-core/tests/unit/jobs/test_triggerer_job.py +++ b/airflow-core/tests/unit/jobs/test_triggerer_job.py @@ -3303,3 +3303,45 @@ def test_handle_request_checks_before_adding_to_create(self, jobless_supervisor, assert jobless_supervisor.stop is False assert jobless_supervisor.running_triggers == {2, 3} + + def test_creation_failure_reported_in_finished(self, jobless_supervisor, mocker): + """A creation failure appears in both failures and finished, so running_triggers stays in sync.""" + mocker.patch.object(TriggerRunnerSupervisor, "send_msg", autospec=True) + jobless_supervisor.running_triggers = {1, 2, 3} + + jobless_supervisor._handle_request( + messages.TriggerStateChanges( + events=None, + failures=[(3, ["Traceback..."])], + finished=[3], + num_running=2, + ), + log=MagicMock(spec=FilteringBoundLogger), + req_id=1, + ) + + assert jobless_supervisor.stop is False + assert 3 not in jobless_supervisor.running_triggers + + +class TestCreationFailureInFinished: + """Tests that creation failures (not in self.triggers) are included in finished_ids.""" + + def test_creation_failure_included_in_finished(self): + runner = TriggerRunner() + runner.failed_triggers.append((42, ValueError("bad classpath"))) + + msg = runner.process_trigger_events(finished_ids=[10]) + + assert msg.finished == [10, 42] + assert msg.num_running == 0 + + def test_serialization_failure_not_in_finished(self): + runner = TriggerRunner() + runner.triggers = {42: {"task": MagicMock(), "is_watcher": False, "name": "t", "events": 0}} + runner.failed_triggers.append((42, ValueError("not serializable"))) + + msg = runner.process_trigger_events(finished_ids=[]) + + assert msg.finished is None + assert msg.num_running == 1 From d0f5f2dd37bffa966cfd53f15e37863455010266 Mon Sep 17 00:00:00 2001 From: JH0917 Date: Sun, 2 Aug 2026 15:56:41 +0900 Subject: [PATCH 3/3] Recover lost Triggerer triggers instead of shutting down --- .../src/airflow/jobs/triggerer_job_runner.py | 36 +++++----- .../tests/unit/jobs/test_triggerer_job.py | 69 ++++++++++++------- .../metrics/metrics_template.yaml | 7 ++ 3 files changed, 72 insertions(+), 40 deletions(-) diff --git a/airflow-core/src/airflow/jobs/triggerer_job_runner.py b/airflow-core/src/airflow/jobs/triggerer_job_runner.py index 3a616e4cf82f7..26d557c3c3acd 100644 --- a/airflow-core/src/airflow/jobs/triggerer_job_runner.py +++ b/airflow-core/src/airflow/jobs/triggerer_job_runner.py @@ -332,7 +332,8 @@ class TriggerStateChanges(BaseModel): # Format of list[str] is the exc traceback format failures: list[tuple[int, list[str] | None]] | None = None finished: list[int] | None = None - num_running: int = 0 + # Ids the runner has a live coroutine for + running_ids: set[int] = set() class TriggerStateSync(BaseModel): type: Literal["TriggerStateSync"] = "TriggerStateSync" @@ -604,7 +605,7 @@ def _handle_request(self, msg: ToTriggerSupervisor, log: FilteringBoundLogger, r # handle leaks for every failed upload. factory.close() - self.check_for_unhandled_triggers(msg.num_running) + self.check_for_unhandled_triggers(msg.running_ids) # Drain the persist confirmations accumulated since the last sync. events_persisted: list[int] = [] @@ -786,20 +787,20 @@ def clean_unused(self) -> None: """Remove triggers that are no longer needed.""" Trigger.clean_unused() - def check_for_unhandled_triggers(self, num_running: int) -> None: + def check_for_unhandled_triggers(self, running_ids: set[int]) -> None: """ - Shut down if the subprocess trigger count disagrees with the supervisor. + Re-create triggers we track as running that the runner has no coroutine for. - Only valid between finished-removal and to_create-addition in ``_handle_request``. + Only valid between finished-removal and to_create-addition in ``_handle_request``, where the + two sides agree. Dropping the leftovers lets :meth:`update_triggers` rebuild them. """ - expected = len(self.running_triggers) - if expected != num_running: - log.error( - "Trigger count mismatch: expected %d, subprocess reports %d. Shutting down.", - expected, - num_running, - ) - self.stop = True + unhandled = self.running_triggers - running_ids + if not unhandled: + return + log.error("Triggers have no coroutine in the runner; re-creating", trigger_ids=sorted(unhandled)) + self.running_triggers -= unhandled + self.cancelling_triggers -= unhandled + stats.incr("triggers.state_mismatch", len(unhandled), tags=prune_dict({"team_name": self.team_name})) def handle_failed_triggers(self): """ @@ -1518,6 +1519,7 @@ def process_trigger_events(self, finished_ids: list[int]) -> messages.TriggerSta # Copy out of our dequeues in threadsafe manner to sync state with parent events_to_send: list[TriggerEventEntry] = [] failures_to_send: list[tuple[int, list[str] | None]] = [] + finished_to_send = list(finished_ids) while self.events: events_to_send.append(self.events.popleft()) @@ -1527,13 +1529,13 @@ def process_trigger_events(self, finished_ids: list[int]) -> messages.TriggerSta tb = format_exception(type(exc), exc, exc.__traceback__) if exc else None failures_to_send.append((trigger_id, tb)) if trigger_id not in self.triggers: - finished_ids.append(trigger_id) + finished_to_send.append(trigger_id) return messages.TriggerStateChanges( events=events_to_send if events_to_send else None, - finished=finished_ids if finished_ids else None, + finished=finished_to_send if finished_to_send else None, failures=failures_to_send if failures_to_send else None, - num_running=len(self.triggers), + running_ids=set(self.triggers), ) def sanitize_trigger_events(self, msg: messages.TriggerStateChanges) -> messages.TriggerStateChanges: @@ -1562,7 +1564,7 @@ def sanitize_trigger_events(self, msg: messages.TriggerStateChanges) -> messages events=events_to_send if events_to_send else None, finished=msg.finished, failures=msg.failures, - num_running=msg.num_running, + running_ids=msg.running_ids, ) async def sync_state_to_supervisor(self, finished_ids: list[int]) -> None: diff --git a/airflow-core/tests/unit/jobs/test_triggerer_job.py b/airflow-core/tests/unit/jobs/test_triggerer_job.py index bbb8eb5e745d8..7b7e0394c0c57 100644 --- a/airflow-core/tests/unit/jobs/test_triggerer_job.py +++ b/airflow-core/tests/unit/jobs/test_triggerer_job.py @@ -3256,31 +3256,45 @@ async def _drive(): class TestCheckForUnhandledTriggers: - """Tests for supervisor-vs-subprocess trigger count mismatch detection.""" + """Tests for recovery of triggers the runner has no coroutine for.""" - def test_mismatch_shuts_down(self, jobless_supervisor): - jobless_supervisor.running_triggers = {1, 2, 3} - - jobless_supervisor.check_for_unhandled_triggers(num_running=0) - assert jobless_supervisor.stop is True - - def test_matching_counts_no_shutdown(self, jobless_supervisor): - jobless_supervisor.running_triggers = {1, 2, 3} - - jobless_supervisor.check_for_unhandled_triggers(num_running=3) - assert jobless_supervisor.stop is False + @pytest.mark.parametrize( + ("running_triggers", "running_ids", "expected_after"), + [ + pytest.param({1, 2, 3}, {1, 2, 3}, {1, 2, 3}, id="in-sync"), + pytest.param(set(), set(), set(), id="no-triggers"), + pytest.param({1, 2, 3}, {1, 2}, {1, 2}, id="one-unhandled"), + pytest.param({1, 2, 3}, set(), set(), id="all-unhandled"), + ], + ) + def test_unhandled_triggers_are_dropped_not_shut_down( + self, jobless_supervisor, mocker, running_triggers, running_ids, expected_after + ): + incr = mocker.patch("airflow.jobs.triggerer_job_runner.stats.incr", autospec=True) + jobless_supervisor.running_triggers = set(running_triggers) + jobless_supervisor.cancelling_triggers = set(running_triggers) - def test_no_triggers_no_shutdown(self, jobless_supervisor): - jobless_supervisor.running_triggers = set() + jobless_supervisor.check_for_unhandled_triggers(running_ids) - jobless_supervisor.check_for_unhandled_triggers(num_running=0) + unhandled = running_triggers - running_ids assert jobless_supervisor.stop is False + assert jobless_supervisor.running_triggers == expected_after + assert jobless_supervisor.cancelling_triggers == expected_after + assert ( + mocker.call("triggers.state_mismatch", len(unhandled), tags={}) in incr.call_args_list + ) is bool(unhandled) + + def test_dropped_trigger_is_recreated_next_loop(self, jobless_supervisor, mocker): + """Dropping an unhandled id is what lets update_triggers rebuild its workload.""" + build = mocker.patch.object( + TriggerRunnerSupervisor, "build_trigger_workloads", autospec=True, return_value=[] + ) + jobless_supervisor.running_triggers = {1, 2} - def test_subprocess_more_than_expected_shuts_down(self, jobless_supervisor): - jobless_supervisor.running_triggers = {1} + jobless_supervisor.check_for_unhandled_triggers({1}) + jobless_supervisor.update_triggers({1, 2}) - jobless_supervisor.check_for_unhandled_triggers(num_running=3) - assert jobless_supervisor.stop is True + assert build.call_args.args[1] == {2} def test_handle_request_checks_before_adding_to_create(self, jobless_supervisor, mocker): """The check fires after finished processing but before to_create is added to running_triggers.""" @@ -3295,7 +3309,7 @@ def test_handle_request_checks_before_adding_to_create(self, jobless_supervisor, events=None, failures=None, finished=[1], - num_running=1, + running_ids={2}, ), log=MagicMock(spec=FilteringBoundLogger), req_id=1, @@ -3314,7 +3328,7 @@ def test_creation_failure_reported_in_finished(self, jobless_supervisor, mocker) events=None, failures=[(3, ["Traceback..."])], finished=[3], - num_running=2, + running_ids={1, 2}, ), log=MagicMock(spec=FilteringBoundLogger), req_id=1, @@ -3334,7 +3348,7 @@ def test_creation_failure_included_in_finished(self): msg = runner.process_trigger_events(finished_ids=[10]) assert msg.finished == [10, 42] - assert msg.num_running == 0 + assert msg.running_ids == set() def test_serialization_failure_not_in_finished(self): runner = TriggerRunner() @@ -3344,4 +3358,13 @@ def test_serialization_failure_not_in_finished(self): msg = runner.process_trigger_events(finished_ids=[]) assert msg.finished is None - assert msg.num_running == 1 + assert msg.running_ids == {42} + + def test_caller_finished_ids_not_mutated(self): + runner = TriggerRunner() + runner.failed_triggers.append((42, ValueError("bad classpath"))) + finished_ids = [10] + + runner.process_trigger_events(finished_ids=finished_ids) + + assert finished_ids == [10] diff --git a/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml b/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml index 536e439f876e2..8714a4f70ad0a 100644 --- a/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml +++ b/shared/observability/src/airflow_shared/observability/metrics/metrics_template.yaml @@ -286,6 +286,13 @@ metrics: legacy_name: "-" name_variables: [] + - name: "triggers.state_mismatch" + description: "Number of triggers the Triggerer tracked as running that the trigger runner + had no coroutine for, and which were re-created" + type: "counter" + legacy_name: "-" + name_variables: [] + - name: "triggers.succeeded" description: "Number of triggers that have fired at least one event" type: "counter"