Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 14 additions & 8 deletions airflow-core/tests/unit/jobs/test_triggerer_job.py
Original file line number Diff line number Diff line change
Expand Up @@ -1183,18 +1183,24 @@ def test_trigger_log(mock_monotonic, trigger, watcher_count, trigger_count, sess
Checks that the triggerer will log watcher and trigger in separate lines.
"""
create_trigger_in_db(session, trigger)
trigger_line = f"{trigger_count} triggers currently running"
watcher_line = f"{watcher_count} watchers currently running"

trigger_runner_supervisor = TriggerRunnerSupervisor.start(job=Job(id=123456), capacity=10)
trigger_runner_supervisor.load_triggers()

for _ in range(30):
trigger_runner_supervisor._service_subprocess(0.1)
try:
trigger_runner_supervisor.load_triggers()

stdout = capsys.readouterr().out
assert f"{trigger_count} triggers currently running" in stdout
assert f"{watcher_count} watchers currently running" in stdout
stdout = ""
for _ in range(300):
trigger_runner_supervisor._service_subprocess(0.1)
stdout += capsys.readouterr().out
if trigger_line in stdout and watcher_line in stdout:
break
finally:
trigger_runner_supervisor.kill(force=False)

trigger_runner_supervisor.kill(force=False)
assert trigger_line in stdout
assert watcher_line in stdout


def test_trigger_logger_close():
Expand Down