/
githubmirror
/
salt
Обзор
Документация
Войти
/
githubmirror
/
salt
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
master
tests/pytests/unit/test_minion.py
2 246 строк
76 KB
Joseph Hall
Fix MinionManager skipping connect_master when file_client is local (#64921)
04 июл 2026, 23:45
Не верифицирован
04 июл 2026, 23:45
cafcd99
Код
Авторство
О чём код?
import asyncio import contextlib import copy import logging import os import pathlib import signal import threading import time import uuid import pytest import tornado import tornado.gen import tornado.ioloop import tornado.testing import salt.minion import salt.modules.test as test_mod import salt.syspaths import salt.utils.crypt import salt.utils.jid import salt.utils.platform import salt.utils.process import salt.utils.state from salt._compat import ipaddress from salt.exceptions import ( SaltClientError, SaltMasterUnresolvableError, SaltReqTimeoutError, SaltSystemExit, ) from tests.support.mock import MagicMock, patch log = logging.getLogger(__name__) @pytest.fixture def connect_master_mock(): class ConnectMasterMock: """ Mock connect master call. The first call will raise an exception stored on the exc attribute. Subsequent calls will return True. """ def __init__(self): self.calls = 0 self.exc = Exception @tornado.gen.coroutine def __call__(self, *args, **kwargs): self.calls += 1 if self.calls == 1: raise self.exc() else: return True return ConnectMasterMock() def test_minion_load_grains_false(minion_opts): """ Minion does not generate grains when load_grains is False """ minion_opts["grains"] = {"foo": "bar"} with patch("salt.loader.grains") as grainsfunc: minion = salt.minion.Minion(minion_opts, load_grains=False) try: assert minion.opts["grains"] == minion_opts["grains"] grainsfunc.assert_not_called() finally: minion.destroy() def test_minion_load_grains_true(minion_opts): """ Minion generates grains when load_grains is True """ with patch("salt.loader.grains") as grainsfunc: minion = salt.minion.Minion(minion_opts, load_grains=True) try: assert minion.opts["grains"] != {} grainsfunc.assert_called() finally: minion.destroy() def test_minion_load_grains_default(minion_opts): """ Minion load_grains defaults to True """ with patch("salt.loader.grains") as grainsfunc: minion = salt.minion.Minion(minion_opts) try: assert minion.opts["grains"] != {} grainsfunc.assert_called() finally: minion.destroy() @pytest.mark.parametrize( "event", [ ( "fire_event", lambda data, tag, cb=None, timeout=60: True, ), ( "fire_event_async", lambda data, tag, cb=None, timeout=60: tornado.gen.maybe_future(True), ), ], ) def test_send_req_fires_completion_event(event, minion_opts): req_id = uuid.uuid4() event_enter = MagicMock() event_enter.send.side_effect = event[1] event_enter.get_event.return_value = {"ret": True} event = MagicMock() event.__enter__.return_value = event_enter with patch("salt.utils.event.get_event", return_value=event), patch( "uuid.uuid4", return_value=req_id ): minion_opts["random_startup_delay"] = 0 minion_opts["return_retry_tries"] = 30 minion_opts["grains"] = {} with patch("salt.loader.grains"): minion = salt.minion.Minion(minion_opts) try: load = {"load": "value"} timeout = 60 # XXX This is buggy because "async" in event[0] will never evaluate # to True and if it *did* evaluate to true the test would fail # because you Mock isn't a co-routine. if "async" in event[0]: rtn = minion._send_req_async(load, timeout).result() else: rtn = minion._send_req_sync(load, timeout) fire_event_called = False # get the for idx, call in enumerate(event.mock_calls, 1): if "fire_event" in call[0]: condition_event_tag = ( len(call.args) > 1 and call.args[1] == f"__master_req_channel_payload/{req_id}/{minion_opts['master']}" ) condition_event_tag_error = ( "{} != {}; Call(number={}): {}".format( idx, call, call.args[1], "__master_req_channel_payload" ) ) condition_timeout = ( len(call.kwargs) == 1 and call.kwargs["timeout"] == timeout ) condition_timeout_error = ( "{} != {}; Call(number={}): {}".format( idx, call, call.kwargs["timeout"], timeout ) ) fire_event_called = True assert condition_event_tag, condition_event_tag_error assert condition_timeout, condition_timeout_error assert fire_event_called assert rtn finally: minion.destroy() async def test_send_req_async_regression_62453(minion_opts): class MockEvent: def __init__(self, *args, **kwargs): pass @tornado.gen.coroutine def fire_event_async(self, *args, **kwargs): return def get_event(self, *args, **kwargs): return def __enter__(self): return self def __exit__(self, *args): return def get_event(*args, **kwargs): return MockEvent() minion_opts["random_startup_delay"] = 0 minion_opts["return_retry_tries"] = 5 minion_opts["grains"] = {} minion_opts["ipc_mode"] = "tcp" with patch("salt.loader.grains"): minion = salt.minion.Minion(minion_opts) load = {"load": "value"} timeout = 1 with patch("salt.utils.event.get_event", get_event): # We are just validating no exception is raised with pytest.raises(SaltReqTimeoutError): rtn = await minion._send_req_async(load, timeout) def test_mine_send_tries(minion_opts): channel_enter = MagicMock() channel_enter.send.side_effect = lambda load, timeout, tries: tries channel = MagicMock() channel.__enter__.return_value = channel_enter minion_opts["return_retry_tries"] = 20 with patch("salt.channel.client.ReqChannel.factory", return_value=channel), patch( "salt.loader.grains" ): minion = salt.minion.Minion(minion_opts) minion.tok = "token" data = {} tag = "tag" rtn = minion._mine_send(tag, data) assert rtn == 20 def test_invalid_master_address(minion_opts): minion_opts.update( { "ipv6": False, "master": float("127.0"), "master_port": "4555", "retry_dns": False, } ) with pytest.raises(SaltSystemExit): salt.minion.resolve_dns(minion_opts) def test_source_int_name_local(minion_opts): """ test when file_client local and source_interface_name is set """ interfaces = { "bond0.1234": { "hwaddr": "01:01:01:d0:d0:d0", "up": True, "inet": [ { "broadcast": "111.1.111.255", "netmask": "111.1.0.0", "label": "bond0", "address": "111.1.0.1", } ], } } minion_opts.update( { "ipv6": False, "master": "127.0.0.1", "master_port": "4555", "file_client": "local", "source_interface_name": "bond0.1234", "source_ret_port": 49017, "source_publish_port": 49018, }, ) with patch("salt.utils.network.interfaces", MagicMock(return_value=interfaces)): assert salt.minion.resolve_dns(minion_opts) == { "master_ip": "127.0.0.1", "source_ip": "111.1.0.1", "source_ret_port": 49017, "source_publish_port": 49018, "master_uri": "tcp://127.0.0.1:4555", } @pytest.mark.slow_test def test_source_int_name_remote(minion_opts): """ test when file_client remote and source_interface_name is set and interface is down """ interfaces = { "bond0.1234": { "hwaddr": "01:01:01:d0:d0:d0", "up": False, "inet": [ { "broadcast": "111.1.111.255", "netmask": "111.1.0.0", "label": "bond0", "address": "111.1.0.1", } ], } } minion_opts.update( { "ipv6": False, "master": "127.0.0.1", "master_port": "4555", "file_client": "remote", "source_interface_name": "bond0.1234", "source_ret_port": 49017, "source_publish_port": 49018, }, ) with patch("salt.utils.network.interfaces", MagicMock(return_value=interfaces)): assert salt.minion.resolve_dns(minion_opts) == { "master_ip": "127.0.0.1", "source_ret_port": 49017, "source_publish_port": 49018, "master_uri": "tcp://127.0.0.1:4555", } @pytest.mark.slow_test def test_source_address(minion_opts): """ test when source_address is set """ interfaces = { "bond0.1234": { "hwaddr": "01:01:01:d0:d0:d0", "up": False, "inet": [ { "broadcast": "111.1.111.255", "netmask": "111.1.0.0", "label": "bond0", "address": "111.1.0.1", } ], } } minion_opts.update( { "ipv6": False, "master": "127.0.0.1", "master_port": "4555", "file_client": "local", "source_interface_name": "", "source_address": "111.1.0.1", "source_ret_port": 49017, "source_publish_port": 49018, }, ) with patch("salt.utils.network.interfaces", MagicMock(return_value=interfaces)): assert salt.minion.resolve_dns(minion_opts) == { "source_publish_port": 49018, "source_ret_port": 49017, "master_uri": "tcp://127.0.0.1:4555", "source_ip": "111.1.0.1", "master_ip": "127.0.0.1", } # Tests for _handle_decoded_payload in the salt.minion.Minion() class: 3 @pytest.mark.slow_test async def test_handle_decoded_payload_jid_match_in_jid_queue(minion_opts, io_loop): """ Tests that the _handle_decoded_payload function returns when a jid is given that is already present in the jid_queue. Note: This test doesn't contain all of the patch decorators above the function like the other tests for _handle_decoded_payload below. This is essential to this test as the call to the function must return None BEFORE any of the processes are spun up because we should be avoiding firing duplicate jobs. """ mock_data = {"fun": "foo.bar", "jid": 123} mock_jid_queue = [123] minion = salt.minion.Minion( minion_opts, jid_queue=copy.copy(mock_jid_queue), io_loop=io_loop, ) try: ret = await minion._handle_decoded_payload(mock_data) assert minion.jid_queue == mock_jid_queue assert ret is None finally: minion.destroy() @pytest.mark.slow_test async def test_handle_decoded_payload_jid_queue_addition(minion_opts, io_loop): """ Tests that the _handle_decoded_payload function adds a jid to the minion's jid_queue when the new jid isn't already present in the jid_queue. """ with patch("salt.minion.Minion.ctx", MagicMock(return_value={})), patch( "salt.utils.process.SignalHandlingProcess.start", MagicMock(return_value=True), ), patch( "salt.utils.process.SignalHandlingProcess.join", MagicMock(return_value=True), ): mock_jid = 11111 mock_data = {"fun": "foo.bar", "jid": mock_jid} mock_jid_queue = [123, 456] minion = salt.minion.Minion( minion_opts, jid_queue=copy.copy(mock_jid_queue), io_loop=io_loop, ) try: # Assert that the minion's jid_queue attribute matches the mock_jid_queue as a baseline # This can help debug any test failures if the _handle_decoded_payload call fails. assert minion.jid_queue == mock_jid_queue # Call the _handle_decoded_payload function and update the mock_jid_queue to include the new # mock_jid. The mock_jid should have been added to the jid_queue since the mock_jid wasn't # previously included. The minion's jid_queue attribute and the mock_jid_queue should be equal. await minion._handle_decoded_payload(mock_data) mock_jid_queue.append(mock_jid) assert minion.jid_queue == mock_jid_queue finally: minion.destroy() @pytest.mark.slow_test async def test_handle_decoded_payload_jid_queue_reduced_minion_jid_queue_hwm( minion_opts, io_loop ): """ Tests that the _handle_decoded_payload function removes a jid from the minion's jid_queue when the minion's jid_queue high water mark (minion_jid_queue_hwm) is hit. """ with patch("salt.minion.Minion.ctx", MagicMock(return_value={})), patch( "salt.utils.process.SignalHandlingProcess.start", MagicMock(return_value=True), ), patch( "salt.utils.process.SignalHandlingProcess.join", MagicMock(return_value=True), ): minion_opts["minion_jid_queue_hwm"] = 2 mock_data = {"fun": "foo.bar", "jid": 789} mock_jid_queue = [123, 456] minion = salt.minion.Minion( minion_opts, jid_queue=copy.copy(mock_jid_queue), io_loop=io_loop, ) try: # Assert that the minion's jid_queue attribute matches the mock_jid_queue as a baseline # This can help debug any test failures if the _handle_decoded_payload call fails. assert minion.jid_queue == mock_jid_queue # Call the _handle_decoded_payload function and check that the queue is smaller by one item # and contains the new jid await minion._handle_decoded_payload(mock_data) assert len(minion.jid_queue) == 2 assert minion.jid_queue == [456, 789] finally: minion.destroy() @pytest.mark.slow_test def test_process_count_max(minion_opts, io_loop): """ Tests that the _handle_decoded_payload function does not spawn more than the configured amount of processes, as per process_count_max. """ start_mock = MagicMock(return_value=True) def mock_proc_side_effect(*args, **kwargs): m = MagicMock(name="MockProcess") m.is_alive.return_value = True m.start = start_mock return m @contextlib.asynccontextmanager async def mock_await_lock(*args, **kwargs): yield fopen_mock = MagicMock() with patch("salt.minion.Minion.ctx", MagicMock(return_value={})), patch( "salt.minion.SignalHandlingProcess", MagicMock(side_effect=mock_proc_side_effect), ), patch( "salt.minion.SignalHandlingProcess.join", MagicMock(return_value=True), ), patch( "os.path.exists", MagicMock(return_value=True) ), patch( "os.makedirs", MagicMock() ), patch( "salt.utils.files.fopen", fopen_mock ), patch( "salt.payload.dump", MagicMock() ), patch( "salt.utils.files.await_lock", side_effect=mock_await_lock ), patch( "salt.loader.grains", MagicMock(return_value={"id": "foo", "os": "Linux"}) ): process_count_max = 10 minion_opts["__role"] = "minion" minion_opts["minion_jid_queue_hwm"] = 100 minion_opts["process_count_max"] = process_count_max # cachedir needed for lock; master pins the per-master subpath minion_opts["cachedir"] = "/tmp/salt_test_cache" minion_opts["master"] = "master-a" minion = salt.minion.Minion(minion_opts, jid_queue=[], io_loop=io_loop) try: # up until process_count_max: processes are started normally for i in range(process_count_max): mock_data = {"fun": "foo.bar", "jid": str(i)} io_loop.run_sync( lambda data=mock_data: minion._handle_decoded_payload(data) ) assert start_mock.call_count == i + 1 assert len(minion.jid_queue) == i + 1 # above process_count_max: Queue logic kicks in mock_data = {"fun": "foo.bar", "jid": str(process_count_max + 1)} # Run execution io_loop.run_sync(lambda: minion._handle_decoded_payload(mock_data)) # Assert NO new process started assert start_mock.call_count == process_count_max # Assert Job was queued (payload dumped) assert salt.payload.dump.called # Assert JID added to active queue (deduplication cache) assert len(minion.jid_queue) == process_count_max + 1 # Assert the queued job file landed under the per-master job_queue dir, # not the legacy shared cachedir/job_queue path. expected_dir = salt.utils.state.job_queue_dir(minion_opts) queue_paths = [c.args[0] for c in fopen_mock.call_args_list] assert any(p.startswith(expected_dir) for p in queue_paths), queue_paths finally: minion.destroy() def test_queue_job_preserves_master_jid_69386(minion_opts): """ Regression test for #69386 (job_queue side). The companion fix in ``salt/modules/state.py:_check_queue`` covers the state-queue write path. This test pins down the contract for the job-queue write path in ``salt.minion.Minion._queue_job``: when the minion shelves a payload to disk because ``process_count_max`` was reached, the master-supplied JID must end up unchanged in both the serialized payload and the queue filename. The job-queue path was never broken by the state-queue refactor that introduced #69386, but it is the most natural place for a future regression to creep back in -- so we assert the invariant explicitly. """ master_jid = "20260601000000123456" payload = { "fun": "state.apply", "arg": ["highstate"], "jid": master_jid, "tgt": "minion-1", "ret": "", "user": "root", } minion_opts["__role"] = "minion" minion_opts["cachedir"] = "/tmp/salt_test_cache_69386" minion_opts["master"] = "master-a" dump_mock = MagicMock() rename_calls = [] def _rename(src, dst): rename_calls.append((src, dst)) with patch("salt.minion.Minion.ctx", MagicMock(return_value={})), patch( "salt.loader.grains", MagicMock(return_value={"id": "foo", "os": "Linux"}) ), patch("os.path.exists", MagicMock(return_value=True)), patch( "os.makedirs", MagicMock() ), patch( "salt.utils.files.fopen", MagicMock() ), patch( "salt.payload.dump", dump_mock ), patch( "salt.utils.atomicfile.atomic_rename", side_effect=_rename ), patch( "salt.utils.jid.gen_jid", side_effect=AssertionError( "_queue_job must never mint a new JID for a master-published " "payload (#69386 regression)" ), ): io_loop = tornado.ioloop.IOLoop() minion = salt.minion.Minion(minion_opts, jid_queue=[], io_loop=io_loop) try: minion._queue_job(payload) # Payload written to disk must carry the master JID, unchanged. assert dump_mock.called dumped_payload = dump_mock.call_args.args[0] assert dumped_payload["jid"] == master_jid assert dumped_payload is payload # _queue_job dumps the dict by reference # Final on-disk filename must embed the master JID and land in # the per-master job_queue dir. expected_dir = salt.utils.state.job_queue_dir(minion_opts) assert rename_calls, "expected an atomic_rename into place" final_path = rename_calls[-1][1] assert final_path.startswith(expected_dir), final_path assert final_path.endswith(f"_{master_jid}.p"), final_path finally: minion.destroy() async def test_process_queue_rechecks_count_per_job(minion_opts): """ Test that job queue processing re-checks process count before each individual job, preventing race conditions where process count changes during batch processing. """ # Create a simple test that just verifies the queue processing method exists from salt.minion import Minion minion = Minion(minion_opts) try: # Just test that the method exists and can be called without crashing await minion._process_process_queue_async_impl() # If we get here without exception, test passes assert True finally: minion.destroy() def test_cleanup_orphaned_queue_files(minion_opts): """ Test that orphaned running_ queue files are cleaned up on minion startup. This prevents stale files from blocking future jobs after minion crashes. """ # Create a simple test that just verifies the method exists and can be called from salt.minion import Minion minion = Minion(minion_opts) try: # Just test that the method exists and doesn't crash when called minion._cleanup_orphaned_queue_files() # If we get here without exception, test passes assert True finally: minion.destroy() @pytest.mark.slow_test async def test_beacons_before_connect(minion_opts): """ Tests that the 'beacons_before_connect' option causes the beacons to be initialized before connect. """ with patch("salt.minion.Minion.ctx", MagicMock(return_value={})), patch( "salt.minion.Minion.sync_connect_master", MagicMock(side_effect=RuntimeError("stop execution")), ), patch( "salt.utils.process.SignalHandlingProcess.start", MagicMock(return_value=True), ), patch( "salt.utils.process.SignalHandlingProcess.join", MagicMock(return_value=True), ): minion_opts["beacons_before_connect"] = True io_loop = tornado.ioloop.IOLoop() minion = salt.minion.Minion(minion_opts, io_loop=io_loop) try: try: await minion.tune_in(start=True) except RuntimeError: pass # Make sure beacons are initialized but the sheduler is not assert "beacons" in minion.periodic_callbacks assert "schedule" not in minion.periodic_callbacks finally: minion.destroy() @pytest.mark.slow_test async def test_scheduler_before_connect(minion_opts): """ Tests that the 'scheduler_before_connect' option causes the scheduler to be initialized before connect. """ with patch("salt.minion.Minion.ctx", MagicMock(return_value={})), patch( "salt.minion.Minion.sync_connect_master", MagicMock(side_effect=RuntimeError("stop execution")), ), patch( "salt.utils.process.SignalHandlingProcess.start", MagicMock(return_value=True), ), patch( "salt.utils.process.SignalHandlingProcess.join", MagicMock(return_value=True), ): minion_opts["scheduler_before_connect"] = True io_loop = tornado.ioloop.IOLoop() minion = salt.minion.Minion(minion_opts, io_loop=io_loop) try: try: await minion.tune_in(start=True) except RuntimeError: pass # Make sure the scheduler is initialized but the beacons are not assert "schedule" in minion.periodic_callbacks assert "beacons" not in minion.periodic_callbacks finally: minion.destroy() def test_minion_module_refresh(minion_opts): """ Tests that the 'module_refresh' just return in case there is no 'schedule' because destroy method was already called. """ with patch("salt.minion.Minion.ctx", MagicMock(return_value={})), patch( "salt.utils.process.SignalHandlingProcess.start", MagicMock(return_value=True), ), patch( "salt.utils.process.SignalHandlingProcess.join", MagicMock(return_value=True), ): try: minion = salt.minion.Minion( minion_opts, io_loop=tornado.ioloop.IOLoop(), ) minion.schedule = salt.utils.schedule.Schedule( minion_opts, {}, returners={} ) assert hasattr(minion, "schedule") minion.destroy() assert not hasattr(minion, "schedule") assert not minion.module_refresh() finally: minion.destroy() def test_minion_module_refresh_beacons_refresh(minion_opts): """ Tests that 'module_refresh' calls beacons_refresh and that the minion object has a beacons attribute with beacons. """ with patch("salt.minion.Minion.ctx", MagicMock(return_value={})), patch( "salt.utils.process.SignalHandlingProcess.start", MagicMock(return_value=True), ), patch( "salt.utils.process.SignalHandlingProcess.join", MagicMock(return_value=True), ): try: minion = salt.minion.Minion( minion_opts, io_loop=tornado.ioloop.IOLoop(), ) minion.schedule = salt.utils.schedule.Schedule( minion_opts, {}, returners={} ) assert not hasattr(minion, "beacons") minion.module_refresh() assert hasattr(minion, "beacons") assert hasattr(minion.beacons, "beacons") assert "service.beacon" in minion.beacons.beacons minion.destroy() finally: if minion is not None: minion.destroy() def test_beacons_refresh_preserves_interval_map(minion_opts): """ Tests that 'beacons_refresh' preserves the interval_map so that beacon intervals are not reset during module refresh. """ with patch("salt.minion.Minion.ctx", MagicMock(return_value={})), patch( "salt.utils.process.SignalHandlingProcess.start", MagicMock(return_value=True), ), patch( "salt.utils.process.SignalHandlingProcess.join", MagicMock(return_value=True), ): minion = None try: minion = salt.minion.Minion( minion_opts, io_loop=tornado.ioloop.IOLoop.current(), ) minion.schedule = salt.utils.schedule.Schedule( minion_opts, {}, returners={} ) minion.module_refresh() assert hasattr(minion, "beacons") assert hasattr(minion.beacons, "interval_map") test_interval_map = {"status": 50, "diskusage": 30} minion.beacons.interval_map = test_interval_map.copy() old_beacons = minion.beacons minion.beacons_refresh() assert minion.beacons is not old_beacons assert minion.beacons.interval_map == test_interval_map assert minion.beacons.interval_map["status"] == 50 assert minion.beacons.interval_map["diskusage"] == 30 finally: if minion is not None: minion.destroy() def test_beacons_refresh_closes_old_beacons(minion_opts): """ Tests that 'beacons_refresh' calls close_beacons() on the old Beacon instance before replacing it, preventing inotify fd leaks. See: https://github.com/saltstack/salt/issues/66449 See: https://github.com/saltstack/salt/issues/58907 """ with patch("salt.minion.Minion.ctx", MagicMock(return_value={})), patch( "salt.utils.process.SignalHandlingProcess.start", MagicMock(return_value=True), ), patch( "salt.utils.process.SignalHandlingProcess.join", MagicMock(return_value=True), ): minion = None try: minion = salt.minion.Minion( minion_opts, io_loop=tornado.ioloop.IOLoop.current(), ) minion.schedule = salt.utils.schedule.Schedule( minion_opts, {}, returners={} ) minion.module_refresh() assert hasattr(minion, "beacons") old_beacons = minion.beacons with patch.object(old_beacons, "close_beacons") as close_mock: minion.beacons_refresh() close_mock.assert_called_once() assert minion.beacons is not old_beacons finally: if minion is not None: minion.destroy() @pytest.mark.slow_test async def test_when_ping_interval_is_set_the_callback_should_be_added_to_periodic_callbacks( minion_opts, ): with patch("salt.minion.Minion.ctx", MagicMock(return_value={})), patch( "salt.minion.Minion.sync_connect_master", MagicMock(side_effect=RuntimeError("stop execution")), ), patch( "salt.utils.process.SignalHandlingProcess.start", MagicMock(return_value=True), ), patch( "salt.utils.process.SignalHandlingProcess.join", MagicMock(return_value=True), ): minion_opts["ping_interval"] = 10 io_loop = tornado.ioloop.IOLoop() minion = salt.minion.Minion(minion_opts, io_loop=io_loop) try: try: minion.connected = MagicMock(side_effect=(False, True)) # _fire_master_minion_start is now called as a coroutine via create_task # so it must be an async function async def async_mock(): pass minion._fire_master_minion_start = async_mock minion.tune_in(start=False) except RuntimeError: pass # Make sure the scheduler is initialized but the beacons are not assert "ping" in minion.periodic_callbacks finally: minion.destroy() @pytest.mark.slow_test def test_when_passed_start_event_grains(minion_opts): # provide mock opts an os grain since we'll look for it later. minion_opts["grains"]["os"] = "linux" minion_opts["start_event_grains"] = ["os"] io_loop = tornado.ioloop.IOLoop() minion = salt.minion.Minion(minion_opts, io_loop=io_loop) try: minion.tok = MagicMock() minion._send_req_sync = MagicMock() minion._fire_master( "Minion has started", "minion_start", include_startup_grains=True ) load = minion._send_req_sync.call_args[0][0] assert "grains" in load assert "os" in load["grains"] finally: minion.destroy() @pytest.mark.slow_test def test_when_not_passed_start_event_grains(minion_opts): io_loop = tornado.ioloop.IOLoop() minion = salt.minion.Minion(minion_opts, io_loop=io_loop) try: minion.tok = MagicMock() minion._send_req_sync = MagicMock() minion._fire_master("Minion has started", "minion_start") load = minion._send_req_sync.call_args[0][0] assert "grains" not in load finally: minion.destroy() @pytest.mark.slow_test def test_when_other_events_fired_and_start_event_grains_are_set(minion_opts): minion_opts["start_event_grains"] = ["os"] io_loop = tornado.ioloop.IOLoop() minion = salt.minion.Minion(minion_opts, io_loop=io_loop) try: minion.tok = MagicMock() minion._send_req_sync = MagicMock() minion._fire_master("Custm_event_fired", "custom_event") load = minion._send_req_sync.call_args[0][0] assert "grains" not in load finally: minion.destroy() @pytest.mark.slow_test def test_fire_start_event_minimal_payload(minion_opts): """ A minimal published job with start_event=True should produce a salt/job/<jid>/start/<minion_id> event whose payload contains the core identifying fields and excludes the function arguments. """ minion_opts["id"] = "minion-under-test" io_loop = tornado.ioloop.IOLoop() minion = salt.minion.Minion(minion_opts, io_loop=io_loop) try: minion.tok = MagicMock() minion._fire_master = MagicMock() data = { "jid": "20260429000000000000", "fun": "test.ping", "tgt": "*", "tgt_type": "glob", "user": "root", "arg": ["should-not-leak"], } minion._fire_start_event(data) minion._fire_master.assert_called_once() load, tag = minion._fire_master.call_args[0] assert tag == "salt/job/20260429000000000000/start/{}".format(minion_opts["id"]) assert load["id"] == minion_opts["id"] assert load["jid"] == data["jid"] assert load["fun"] == data["fun"] assert load["tgt"] == data["tgt"] assert load["tgt_type"] == data["tgt_type"] assert load["user"] == data["user"] assert "arg" not in load assert "fun_args" not in load assert "master_id" not in load assert "metadata" not in load finally: minion.destroy() @pytest.mark.slow_test def test_fire_start_event_includes_master_id_and_metadata(minion_opts): """ When the published load carries master_id and metadata, the start event should propagate both. """ io_loop = tornado.ioloop.IOLoop() minion = salt.minion.Minion(minion_opts, io_loop=io_loop) try: minion.tok = MagicMock() minion._fire_master = MagicMock() data = { "jid": "20260429000000000001", "fun": "test.ping", "tgt": "*", "tgt_type": "glob", "user": "root", "master_id": "master-a", "metadata": {"ticket": "INC-1234"}, } minion._fire_start_event(data) load, _tag = minion._fire_master.call_args[0] assert load["master_id"] == "master-a" assert load["metadata"] == {"ticket": "INC-1234"} finally: minion.destroy() @pytest.mark.slow_test def test_fire_start_event_swallows_failures(minion_opts): """ A failure inside _fire_master must not propagate out of _fire_start_event; firing the start event must never abort job execution. """ io_loop = tornado.ioloop.IOLoop() minion = salt.minion.Minion(minion_opts, io_loop=io_loop) try: minion.tok = MagicMock() minion._fire_master = MagicMock(side_effect=RuntimeError("transport down")) data = { "jid": "20260429000000000002", "fun": "test.ping", "tgt": "*", "tgt_type": "glob", "user": "root", } minion._fire_start_event(data) finally: minion.destroy() @pytest.mark.slow_test def test_fire_start_event_empty_metadata_is_propagated(minion_opts): """ An empty metadata dict is still meaningful (caller deliberately sent one) and must be propagated. Distinguishes ``metadata={}`` from a missing key. """ io_loop = tornado.ioloop.IOLoop() minion = salt.minion.Minion(minion_opts, io_loop=io_loop) try: minion.tok = MagicMock() minion._fire_master = MagicMock() data = { "jid": "20260429000000000010", "fun": "test.ping", "tgt": "*", "tgt_type": "glob", "user": "root", "metadata": {}, } minion._fire_start_event(data) load, _tag = minion._fire_master.call_args[0] assert "metadata" in load assert load["metadata"] == {} finally: minion.destroy() @pytest.mark.slow_test def test_fire_start_event_omits_falsy_master_id(minion_opts): """ A master_id of None or empty string must not appear in the start event load (the gate is truthiness, matching how master_id flows in the publish path). """ io_loop = tornado.ioloop.IOLoop() minion = salt.minion.Minion(minion_opts, io_loop=io_loop) try: minion.tok = MagicMock() minion._fire_master = MagicMock() for falsy in (None, ""): minion._fire_master.reset_mock() data = { "jid": "20260429000000000011", "fun": "test.ping", "tgt": "*", "tgt_type": "glob", "user": "root", "master_id": falsy, } minion._fire_start_event(data) load, _tag = minion._fire_master.call_args[0] assert "master_id" not in load finally: minion.destroy() @pytest.mark.slow_test def test_fire_start_event_multi_fun_passes_list(minion_opts): """ For multi-fun jobs, ``data["fun"]`` is a list. The start event should still fire successfully and propagate the list of function names verbatim. """ io_loop = tornado.ioloop.IOLoop() minion = salt.minion.Minion(minion_opts, io_loop=io_loop) try: minion.tok = MagicMock() minion._fire_master = MagicMock() data = { "jid": "20260429000000000012", "fun": ["test.ping", "test.echo"], "arg": [[], ["hello"]], "tgt": "*", "tgt_type": "glob", "user": "root", } minion._fire_start_event(data) minion._fire_master.assert_called_once() load, _tag = minion._fire_master.call_args[0] assert load["fun"] == ["test.ping", "test.echo"] assert "arg" not in load finally: minion.destroy() def _make_thread_return_minion_mock(tmp_path, opts): """ Build a MagicMock that quacks like a Minion well enough for _thread_return to reach (or skip) the start-event gate. Heavy dependencies (executors, returners, the actual function) are short-circuited; only the start-event gating logic is exercised. """ proc_dir = tmp_path / "proc" proc_dir.mkdir(exist_ok=True) minion_instance = MagicMock() minion_instance.proc_dir = str(proc_dir) minion_instance.opts = opts minion_instance.executors = {} minion_instance.module_executors = [] minion_instance.functions = MagicMock() minion_instance.functions.__contains__.return_value = True minion_instance.functions.pack = {"__context__": {"retcode": 0}} minion_instance._execute_job_function = MagicMock(return_value="ok") minion_instance._return_pub = MagicMock() minion_instance._fire_master = MagicMock() minion_instance._fire_start_event = MagicMock() minion_instance._return_retry_timer = MagicMock(return_value=1) minion_instance.connected = False # skip _return_pub minion_instance.returners = MagicMock() minion_instance.function_errors = {} return minion_instance @pytest.mark.slow_test def test_thread_return_fires_start_event_when_requested(tmp_path, minion_opts): """ _thread_return must invoke _fire_start_event exactly once when the published load carries start_event=True. """ minion_opts["multiprocessing"] = False minion_opts["id"] = "minion-thread-return" minion_instance = _make_thread_return_minion_mock(tmp_path, minion_opts) data = { "jid": "20260429000000000020", "fun": "test.ping", "arg": [], "tgt": "*", "tgt_type": "glob", "user": "root", "ret": "", "start_event": True, } salt.minion.Minion._thread_return(minion_instance, minion_opts, data) minion_instance._fire_start_event.assert_called_once_with(data) @pytest.mark.slow_test def test_thread_return_skips_start_event_when_not_requested(tmp_path, minion_opts): """ _thread_return must NOT invoke _fire_start_event when the published load has no start_event key. This is the default behavior for all pre-existing callers and must not regress. """ minion_opts["multiprocessing"] = False minion_opts["id"] = "minion-thread-return" minion_instance = _make_thread_return_minion_mock(tmp_path, minion_opts) data = { "jid": "20260429000000000021", "fun": "test.ping", "arg": [], "tgt": "*", "tgt_type": "glob", "user": "root", "ret": "", } salt.minion.Minion._thread_return(minion_instance, minion_opts, data) minion_instance._fire_start_event.assert_not_called() @pytest.mark.slow_test def test_thread_return_skips_start_event_when_falsy(tmp_path, minion_opts): """ A start_event of False/None/empty string must be treated as opt-out; _fire_start_event should not be invoked. """ minion_opts["multiprocessing"] = False minion_opts["id"] = "minion-thread-return" for falsy in (False, None, "", 0): minion_instance = _make_thread_return_minion_mock(tmp_path, minion_opts) data = { "jid": "20260429000000000022", "fun": "test.ping", "arg": [], "tgt": "*", "tgt_type": "glob", "user": "root", "ret": "", "start_event": falsy, } salt.minion.Minion._thread_return(minion_instance, minion_opts, data) assert ( minion_instance._fire_start_event.call_count == 0 ), f"Unexpectedly fired start event for falsy value {falsy!r}" def _make_thread_multi_return_minion_mock(tmp_path, opts): """ Sibling helper for _thread_multi_return tests. """ proc_dir = tmp_path / "proc" proc_dir.mkdir(exist_ok=True) minion_instance = MagicMock() minion_instance.proc_dir = str(proc_dir) minion_instance.opts = opts minion_instance.executors = {} minion_instance.module_executors = [] minion_instance.functions = MagicMock() minion_instance.functions.__contains__.return_value = True minion_instance.functions.pack = {"__context__": {"retcode": 0}} minion_instance._execute_job_function = MagicMock(return_value="ok") minion_instance._return_pub_multi = MagicMock() minion_instance._fire_master = MagicMock() minion_instance._fire_start_event = MagicMock() minion_instance._return_retry_timer = MagicMock(return_value=1) minion_instance.connected = False minion_instance.returners = MagicMock() minion_instance.function_errors = {} return minion_instance @pytest.mark.slow_test def test_thread_multi_return_fires_start_event_once_per_jid(tmp_path, minion_opts): """ Multi-fun jobs share a single jid; the start event must fire exactly once for the jid regardless of how many sub-functions are in the load. """ minion_opts["multiprocessing"] = False minion_opts["id"] = "minion-multi" minion_instance = _make_thread_multi_return_minion_mock(tmp_path, minion_opts) data = { "jid": "20260429000000000030", "fun": ["test.ping", "test.echo", "test.version"], "arg": [[], ["hello"], []], "tgt": "*", "tgt_type": "glob", "user": "root", "ret": "", "start_event": True, } salt.minion.Minion._thread_multi_return(minion_instance, minion_opts, data) minion_instance._fire_start_event.assert_called_once_with(data) @pytest.mark.slow_test def test_thread_multi_return_skips_start_event_when_not_requested( tmp_path, minion_opts ): """ Multi-fun jobs without start_event must not produce a start event, matching the single-fun gating. """ minion_opts["multiprocessing"] = False minion_opts["id"] = "minion-multi" minion_instance = _make_thread_multi_return_minion_mock(tmp_path, minion_opts) data = { "jid": "20260429000000000031", "fun": ["test.ping", "test.echo"], "arg": [[], []], "tgt": "*", "tgt_type": "glob", "user": "root", "ret": "", } salt.minion.Minion._thread_multi_return(minion_instance, minion_opts, data) minion_instance._fire_start_event.assert_not_called() @pytest.mark.slow_test def test_fire_start_event_missing_jid_does_not_raise(minion_opts): """ A malformed payload (missing jid) must not propagate an exception out of _fire_start_event. Job execution must never be aborted by a failure in the start-event helper. """ io_loop = tornado.ioloop.IOLoop() minion = salt.minion.Minion(minion_opts, io_loop=io_loop) try: minion.tok = MagicMock() minion._fire_master = MagicMock() data = {"fun": "test.ping"} # Should not raise even though "jid" is required and missing. minion._fire_start_event(data) # _fire_master should not have been called because payload # construction failed before reaching it. minion._fire_master.assert_not_called() finally: minion.destroy() @pytest.mark.slow_test def test_return_pub_handles_send_req_timeout(minion_opts): """ Ensure _return_pub catches SaltReqTimeoutError from _send_req_sync and returns an empty string rather than letting the exception propagate. This is the end-to-end contract between the two methods: _send_req_sync must raise SaltReqTimeoutError (not the bare TimeoutError builtin) so that _return_pub's except clause fires correctly. """ io_loop = tornado.ioloop.IOLoop() io_loop.make_current() minion = salt.minion.Minion(minion_opts, io_loop=io_loop) try: minion.proc_dir = salt.minion.get_proc_dir(minion_opts["cachedir"]) minion._send_req_sync = MagicMock(side_effect=SaltReqTimeoutError("timed out")) result = minion._return_pub( { "id": minion_opts["id"], "jid": "20260101000000000001", "return": True, "fun": "test.ping", } ) assert result == "" finally: minion.destroy() @pytest.mark.slow_test def test_minion_retry_dns_count(minion_opts): """ Tests that the resolve_dns will retry dns look ups for a maximum of 3 times before raising a SaltMasterUnresolvableError exception. """ minion_opts.update( { "ipv6": False, "master": "dummy", "master_port": "4555", "retry_dns": 1, "retry_dns_count": 3, }, ) with pytest.raises(SaltMasterUnresolvableError): salt.minion.resolve_dns(minion_opts) def test_resolve_dns_retry_aborts_on_shutdown_request_69466(minion_opts): """ Regression test for #69466. The resolve_dns() retry loop must wake up promptly when a shutdown is requested (e.g. SIGTERM via MinionManager.stop()) instead of blocking the io_loop for the full ``retry_dns`` interval. Without the fix the blocking ``time.sleep(opts["retry_dns"])`` inside resolve_dns starved the io_loop and the shutdown callback never ran until systemd sent SIGKILL. """ # The fix exposes a public module-level abort hook used by # MinionManager.stop(). Its absence is itself a regression. assert hasattr(salt.minion, "request_resolve_dns_abort"), ( "salt.minion is missing request_resolve_dns_abort(); the SIGTERM " "path cannot interrupt the DNS retry loop. See #69466." ) assert hasattr(salt.minion, "_RESOLVE_DNS_ABORT"), ( "salt.minion is missing the _RESOLVE_DNS_ABORT event used to " "wake an in-progress resolve_dns() retry. See #69466." ) minion_opts.update( { "ipv6": False, "master": "dummy", "master_port": "4555", # A retry interval that is much larger than the test deadline. # If the abort path is not honored, this test would block for # the full 90 seconds. "retry_dns": 90, "retry_dns_count": None, }, ) # The resolve_dns abort flag is process-wide; make sure we leave it # clean for other tests. salt.minion._RESOLVE_DNS_ABORT.clear() def trip_abort(): # Give resolve_dns a moment to enter its sleep, then request abort # the same way MinionManager.stop() does on SIGTERM. time.sleep(0.25) salt.minion.request_resolve_dns_abort() aborter = threading.Thread(target=trip_abort, daemon=True) started = time.monotonic() try: aborter.start() with pytest.raises(SaltMasterUnresolvableError): salt.minion.resolve_dns(minion_opts) finally: aborter.join(timeout=5) salt.minion._RESOLVE_DNS_ABORT.clear() elapsed = time.monotonic() - started # The fix should wake well under 5s; the broken code would sleep for # the full retry_dns (90s) per iteration. assert elapsed < 5, ( f"resolve_dns did not honor the shutdown abort flag " f"(elapsed={elapsed:.2f}s); regression of #69466." ) def test_minion_manager_stop_unblocks_resolve_dns_69466(minion_opts): """ Regression test for #69466. ``MinionManager.stop()`` is the entry point invoked from the SIGTERM handler. It must trip the resolve_dns abort flag before scheduling the async shutdown so a minion currently stuck in the DNS retry loop yields the io_loop. Without this, ``stop_async`` is queued but never runs and systemd escalates to SIGKILL after 90 seconds. """ # The abort flag must be cleared at entry; stop() should set it. salt.minion._RESOLVE_DNS_ABORT.clear() assert not salt.minion._RESOLVE_DNS_ABORT.is_set() manager = salt.minion.MinionManager.__new__(salt.minion.MinionManager) manager.io_loop = MagicMock() # Populate the attributes __del__ -> destroy() touches so the # interpreter does not log an AttributeError at GC time. manager.minions = [] manager.event_publisher = None manager.event = None try: manager.stop(signal.SIGTERM, lambda *a, **kw: None) assert salt.minion._RESOLVE_DNS_ABORT.is_set(), ( "MinionManager.stop() did not request a resolve_dns abort; " "a SIGTERM during the DNS retry loop will be ignored. See #69466." ) # MinionManager.stop() schedules stop_async via # ``io_loop.create_task`` (the 3007.x refactor replaced the earlier # ``add_callback`` form). Either call is acceptable evidence that # the async shutdown got queued. assert ( manager.io_loop.create_task.call_count + manager.io_loop.add_callback.call_count == 1 ) finally: salt.minion._RESOLVE_DNS_ABORT.clear() @pytest.mark.slow_test def test_gen_modules_executors(minion_opts): """ Ensure gen_modules is called with the correct arguments #54429 """ io_loop = tornado.ioloop.IOLoop() minion = salt.minion.Minion(minion_opts, io_loop=io_loop) class MockPillarCompiler: def compile_pillar(self): return {} try: with patch("salt.pillar.get_pillar", return_value=MockPillarCompiler()): with patch("salt.loader.executors", mock=MagicMock()) as execmock: minion.gen_modules() execmock.assert_called_once_with( minion.opts, functions=minion.functions, proxy=minion.proxy, context={} ) finally: minion.destroy() def test_minion_manage_schedule(minion_opts): """ Tests that the manage_schedule will call the add function, adding schedule data into opts. """ with patch("salt.minion.Minion.ctx", MagicMock(return_value={})), patch( "salt.minion.Minion.sync_connect_master", MagicMock(side_effect=RuntimeError("stop execution")), ), patch( "salt.utils.process.SignalHandlingProcess.start", MagicMock(return_value=True), ), patch( "salt.utils.process.SignalHandlingProcess.join", MagicMock(return_value=True), ): io_loop = tornado.ioloop.IOLoop() with patch("salt.utils.schedule.clean_proc_dir", MagicMock(return_value=None)): try: mock_functions = {"test.ping": None} minion = salt.minion.Minion(minion_opts, io_loop=io_loop) minion.schedule = salt.utils.schedule.Schedule( minion_opts, mock_functions, returners={}, new_instance=True, ) minion.opts["foo"] = "bar" schedule_data = { "test_job": { "function": "test.ping", "return_job": False, "jid_include": True, "maxrunning": 2, "seconds": 10, } } data = { "name": "test-item", "schedule": schedule_data, "func": "add", "persist": False, } tag = "manage_schedule" minion.manage_schedule(tag, data) assert "test_job" in minion.opts["schedule"] finally: del minion.schedule minion.destroy() del minion def test_minion_manage_beacons(minion_opts): """ Tests that the manage_beacons will call the add function, adding beacon data into opts. """ with patch("salt.minion.Minion.ctx", MagicMock(return_value={})), patch( "salt.minion.Minion.sync_connect_master", MagicMock(side_effect=RuntimeError("stop execution")), ), patch( "salt.utils.process.SignalHandlingProcess.start", MagicMock(return_value=True), ), patch( "salt.utils.process.SignalHandlingProcess.join", MagicMock(return_value=True), ): minion = None try: minion_opts["beacons"] = {} # io_loop must be a real Tornado IOLoop because our code calls # salt.utils.asynchronous.aioloop() on it io_loop = tornado.ioloop.IOLoop() mock_functions = {"test.ping": None} minion = salt.minion.Minion(minion_opts, io_loop=io_loop) minion.beacons = salt.beacons.Beacon(minion_opts, mock_functions) bdata = [{"salt-master": "stopped"}, {"apache2": "stopped"}] data = {"name": "ps", "beacon_data": bdata, "func": "add"} tag = "manage_beacons" log.debug("==== minion.opts %s ====", minion.opts) minion.manage_beacons(tag, data) assert "ps" in minion.opts["beacons"] assert minion.opts["beacons"]["ps"] == bdata finally: if minion is not None: minion.destroy() def test_prep_ip_port(): _ip = ipaddress.ip_address opts = {"master": "10.10.0.3", "master_uri_format": "ip_only"} ret = salt.minion.prep_ip_port(opts) assert ret == {"master": _ip("10.10.0.3")} opts = { "master": "10.10.0.3", "master_port": 1234, "master_uri_format": "default", } ret = salt.minion.prep_ip_port(opts) assert ret == {"master": "10.10.0.3"} opts = {"master": "10.10.0.3:1234", "master_uri_format": "default"} ret = salt.minion.prep_ip_port(opts) assert ret == {"master": "10.10.0.3", "master_port": 1234} opts = {"master": "host name", "master_uri_format": "default"} pytest.raises(SaltClientError, salt.minion.prep_ip_port, opts) opts = {"master": "10.10.0.3:abcd", "master_uri_format": "default"} pytest.raises(SaltClientError, salt.minion.prep_ip_port, opts) opts = {"master": "10.10.0.3::1234", "master_uri_format": "default"} pytest.raises(SaltClientError, salt.minion.prep_ip_port, opts) @pytest.mark.skip_on_windows(reason="Skippin, no Salt master running on Windows.") async def test_master_type_failover(minion_opts): """ Tests master_type "failover" to not fall back to 127.0.0.1 address when master does not resolve in DNS """ minion_opts.update( { "master_type": "failover", "master": ["master1", "master2"], "__role": "", "retry_dns": 0, } ) class MockPubChannel: def connect(self): raise SaltClientError("MockedChannel") def close(self): return def mock_resolve_dns(opts, fallback=False): assert not fallback if opts["master"] == "master1": raise SaltClientError("Cannot resolve {}".format(opts["master"])) return { "master_ip": "192.168.2.1", "master_uri": "tcp://192.168.2.1:4505", } def mock_channel_factory(opts, **kwargs): assert opts["master"] == "master2" return MockPubChannel() with patch("salt.minion.resolve_dns", mock_resolve_dns), patch( "salt.channel.client.AsyncPubChannel.factory", mock_channel_factory ), patch("salt.loader.grains", MagicMock(return_value=[])): with pytest.raises(SaltClientError): minion = salt.minion.Minion(minion_opts) await minion.connect_master() async def test_master_type_failover_no_masters(minion_opts): """ Tests master_type "failover" to not fall back to 127.0.0.1 address when no master can be resolved """ minion_opts.update( { "master_type": "failover", "master": ["master1", "master2"], "__role": "", "retry_dns": 0, } ) def mock_resolve_dns(opts, fallback=False): assert not fallback raise SaltClientError("Cannot resolve {}".format(opts["master"])) with patch("salt.minion.resolve_dns", mock_resolve_dns), patch( "salt.loader.grains", MagicMock(return_value=[]) ): with pytest.raises(SaltClientError): minion = salt.minion.Minion(minion_opts) # Mock the io_loop so calls to stop/close won't happen. minion.io_loop = MagicMock() await minion.connect_master() def test_eval_master_single_master_closes_pub_channel_on_failure_68901(minion_opts): """ Regression test for #68901: every AsyncPubChannel constructed by Minion.eval_master in the single-master sign-in path must be close()-d when the connection attempt fails, regardless of which exception type pub_channel.connect() raised. Failing to do so leaks the channel's underlying socket file descriptor on each retry, which over time exhausts the minion's fd limit. """ minion_opts.update( { "master": "127.0.0.1", "master_type": "str", "transport": "zeromq", "__role": "", "retry_dns": 0, "acceptance_wait_time": 0, "acceptance_wait_time_max": 0, "master_tries": 1, } ) created = [] class MockPubChannel: def __init__(self): self.closed = 0 created.append(self) @tornado.gen.coroutine def connect(self): # Non-SaltClientError on purpose: prior to the fix, this leaks # the channel because the single-master path only closes # pub_channel inside an `except SaltClientError` clause. raise OSError("simulated transport failure") def close(self): self.closed += 1 def mock_channel_factory(opts, **kwargs): return MockPubChannel() def mock_resolve_dns(opts, fallback=True): return {"master_ip": "127.0.0.1", "master_uri": "tcp://127.0.0.1:4506"} io_loop = tornado.ioloop.IOLoop() try: with patch("salt.minion.resolve_dns", mock_resolve_dns), patch( "salt.channel.client.AsyncPubChannel.factory", mock_channel_factory ), patch("salt.loader.grains", MagicMock(return_value={})): minion = salt.minion.Minion(minion_opts, io_loop=io_loop, load_grains=False) with pytest.raises(OSError): io_loop.run_sync(lambda: minion.eval_master(minion_opts, timeout=1)) finally: io_loop.close(all_fds=True) assert len(created) == 1, "exactly one pub channel should have been created" assert ( created[0].closed == 1 ), "pub channel was not closed on connection failure (#68901 leak)" def test_config_cache_path_overrides(): cachedir = os.path.abspath("/path/to/master/cache") opts = {"cachedir": cachedir, "conf_file": None} mminion = salt.minion.MasterMinion(opts) assert mminion.opts["cachedir"] == cachedir def test_minion_grains_refresh_pre_exec_false(minion_opts): """ Minion does not refresh grains when grains_refresh_pre_exec is False """ minion_opts["multiprocessing"] = False minion_opts["grains_refresh_pre_exec"] = False mock_data = {"fun": "foo.bar", "jid": 123} with patch("salt.loader.grains") as grainsfunc, patch( "salt.minion.Minion._target", MagicMock(return_value=True) ): loop = tornado.ioloop.IOLoop() minion = salt.minion.Minion( minion_opts, jid_queue=None, io_loop=loop, load_grains=False, ) try: loop.run_sync(lambda: minion._handle_decoded_payload(mock_data)) grainsfunc.assert_not_called() finally: minion.destroy() loop.close(all_fds=True) def test_minion_grains_refresh_pre_exec_true(minion_opts): """ Minion refreshes grains when grains_refresh_pre_exec is True """ minion_opts["multiprocessing"] = False minion_opts["grains_refresh_pre_exec"] = True mock_data = {"fun": "foo.bar", "jid": 123} with patch("salt.loader.grains") as grainsfunc, patch( "salt.minion.Minion._target", MagicMock(return_value=True) ): loop = tornado.ioloop.IOLoop() minion = salt.minion.Minion( minion_opts, jid_queue=None, io_loop=loop, load_grains=False, ) try: loop.run_sync(lambda: minion._handle_decoded_payload(mock_data)) grainsfunc.assert_called() finally: minion.destroy() loop.close(all_fds=True) @pytest.mark.skip_on_darwin( reason="Skip on MacOS, where this does not raise an exception." ) def test_valid_ipv4_master_address_ipv6_enabled(minion_opts): """ Tests that the lookups fail back to ipv4 when ipv6 fails. """ interfaces = { "bond0.1234": { "hwaddr": "01:01:01:d0:d0:d0", "up": False, "inet": [ { "broadcast": "111.1.111.255", "netmask": "111.1.0.0", "label": "bond0", "address": "111.1.0.1", } ], } } minion_opts.update( { "ipv6": True, "master": "127.0.0.1", "master_port": "4555", "retry_dns": False, "source_address": "111.1.0.1", "source_interface_name": "bond0.1234", "source_ret_port": 49017, "source_publish_port": 49018, }, ) with patch("salt.utils.network.interfaces", MagicMock(return_value=interfaces)): expected = { "source_publish_port": 49018, "master_uri": "tcp://127.0.0.1:4555", "source_ret_port": 49017, "master_ip": "127.0.0.1", } assert salt.minion.resolve_dns(minion_opts) == expected async def test_master_type_disable(minion_opts): """ Tests master_type "disable" to not even attempt connecting to a master. """ minion_opts.update( { "master_type": "disable", "master": None, "__role": "", "pub_ret": False, "file_client": "local", } ) minion = salt.minion.Minion(minion_opts) try: try: minion_man = salt.minion.MinionManager(minion_opts) await minion_man._connect_minion(minion) except RuntimeError: pytest.fail("_connect_minion(minion) threw an error, This was not expected") # Make sure beacons and sheduler are initialized assert "beacons" in minion.periodic_callbacks assert "schedule" in minion.periodic_callbacks assert minion.connected is False finally: # Mock the io_loop so calls to stop/close won't happen. minion.io_loop = MagicMock() minion.destroy() async def test_masterless_minion_does_not_connect_to_master(minion_opts): """ When file_client is local and use_master_when_local is False (the default), MinionManager._connect_minion must not call connect_master. Regression test for https://github.com/saltstack/salt/issues/57866 and https://github.com/saltstack/salt/issues/64952. """ minion_opts.update( { "file_client": "local", "use_master_when_local": False, "__role": "", "pub_ret": False, } ) minion = salt.minion.Minion(minion_opts) try: connect_master_called = False async def mock_connect_master(failed=False): nonlocal connect_master_called connect_master_called = True minion.connect_master = mock_connect_master minion_man = salt.minion.MinionManager(minion_opts) await minion_man._connect_minion(minion) assert not connect_master_called, ( "connect_master should not be called when file_client=local and" " use_master_when_local=False" ) finally: minion.io_loop = MagicMock() minion.destroy() async def test_masterless_minion_connects_when_use_master_when_local(minion_opts): """ When file_client is local and use_master_when_local is True, MinionManager._connect_minion must still call connect_master. """ minion_opts.update( { "file_client": "local", "use_master_when_local": True, "__role": "", "pub_ret": False, } ) minion = salt.minion.Minion(minion_opts) try: connect_master_called = False async def mock_connect_master(failed=False): nonlocal connect_master_called connect_master_called = True minion.connect_master = mock_connect_master minion_man = salt.minion.MinionManager(minion_opts) await minion_man._connect_minion(minion) assert connect_master_called, ( "connect_master should be called when file_client=local and" " use_master_when_local=True" ) finally: minion.io_loop = MagicMock() minion.destroy() async def test_syndic_async_req_channel(syndic_opts): syndic_opts["_minion_conf_file"] = "" syndic_opts["master_uri"] = "tcp://127.0.0.1:4506" syndic = salt.minion.Syndic(syndic_opts) syndic.pub_channel = MagicMock() syndic.tune_in_no_block() assert isinstance(syndic.async_req_channel, salt.channel.client.AsyncReqChannel) @pytest.mark.slow_test def test_load_args_and_kwargs(minion_opts): """ Ensure load_args_and_kwargs performs correctly """ _args = [{"max": 40, "__kwarg__": True}] ret = salt.minion.load_args_and_kwargs(test_mod.rand_sleep, _args) assert ret == ([], {"max": 40}) assert all([True if "__kwarg__" in item else False for item in _args]) # Test invalid arguments _args = [{"max_sleep": 40, "__kwarg__": True}] with pytest.raises(salt.exceptions.SaltInvocationError): ret = salt.minion.load_args_and_kwargs(test_mod.rand_sleep, _args) async def test_connect_master_salt_client_error(minion_opts, connect_master_mock): """ Ensure minion's destroy is called on a salt client error while connecting to master. """ minion_opts["acceptance_wait_time"] = 0 mm = salt.minion.MinionManager(minion_opts) minion = salt.minion.Minion(minion_opts) connect_master_mock.exc = SaltClientError minion.connect_master = connect_master_mock minion.destroy = MagicMock() await mm._connect_minion(minion) minion.destroy.assert_called_once() # The first call raised an error which caused minion.destroy to get called, # the second call is a success. assert minion.connect_master.calls == 2 async def test_connect_master_unresolveable_error(minion_opts, connect_master_mock): """ Ensure minion's destroy is called on an unresolvable while connecting to master. """ mm = salt.minion.MinionManager(minion_opts) minion = salt.minion.Minion(minion_opts) connect_master_mock.exc = SaltMasterUnresolvableError minion.connect_master = connect_master_mock minion.destroy = MagicMock() await mm._connect_minion(minion) minion.destroy.assert_called_once() # Unresolvable errors break out of the loop. assert minion.connect_master.calls == 1 async def test_connect_master_general_exception_error(minion_opts, connect_master_mock): """ Ensure minion's destroy is called on an un-handled exception while connecting to master. """ mm = salt.minion.MinionManager(minion_opts) minion = salt.minion.Minion(minion_opts) connect_master_mock.exc = SaltClientError minion.connect_master = connect_master_mock minion.destroy = MagicMock() await mm._connect_minion(minion) minion.destroy.assert_called_once() # The first call raised an error which caused minion.destroy to get called, # the second call is a success. assert minion.connect_master.calls == 2 async def test_minion_manager_async_stop(io_loop, minion_opts, tmp_path): """ Ensure MinionManager's stop method works correctly and calls the stop_async method """ # Setup sock_dir with short path minion_opts["sock_dir"] = str(tmp_path / "sock") os.makedirs(minion_opts["sock_dir"]) # Create a MinionManager instance with a mock minion mm = salt.minion.MinionManager(minion_opts) minion = MagicMock(name="minion") minion.destroy = MagicMock() parent_signal_handler = MagicMock(name="parent_signal_handler") mm.minions.append(minion) # Set up event publisher and event mm._bind() assert mm.event_publisher is not None assert mm.event is not None # Check io_loop is running # mm.io_loop is now an asyncio.AbstractEventLoop (not Tornado IOLoop) assert mm.io_loop.is_running() # Wait for the ipc socket to be created, meaning the publish server is listening. while not list(pathlib.Path(minion_opts["sock_dir"]).glob("*")): await tornado.gen.sleep(0.3) # Set up values for event to send load = {"key": "value"} ret = {} # Connect to minion event bus with salt.utils.event.get_event("minion", opts=minion_opts, listen=True) as event: # call stop to start stopping the minion # mm.stop(signal.SIGTERM, parent_signal_handler) mm.stop(signal.SIGTERM, parent_signal_handler) # Fire an event and ensure we can still read it back while the minion # is stopping assert await event.fire_event_async(load, "test_event", timeout=1) is not False start = time.monotonic() while time.monotonic() - start < 5: ret = event.get_event(tag="test_event", wait=1) if ret: break await tornado.gen.sleep(0.3) assert "key" in ret assert ret["key"] == "value" # Sleep to allow stop_async to complete await tornado.gen.sleep(5) # Ensure stop_async has been called (destroy per minion) minion.destroy.assert_called_once() parent_signal_handler.assert_called_once_with(signal.SIGTERM, None) assert mm.event_publisher is None assert mm.event is None def test_minion_io_loop_is_asyncio_loop(minion_opts): """ Test that Minion io_loop is converted to asyncio.AbstractEventLoop. This verifies the salt.utils.asynchronous.aioloop() conversion. """ minion = salt.minion.Minion(minion_opts, load_grains=False) try: # Verify io_loop is an asyncio loop, not a Tornado IOLoop assert isinstance(minion.io_loop, asyncio.AbstractEventLoop) # Ensure it has asyncio methods assert hasattr(minion.io_loop, "create_task") assert hasattr(minion.io_loop, "call_soon") # Ensure it doesn't have Tornado-specific methods assert not hasattr(minion.io_loop, "spawn_callback") finally: minion.destroy() def test_minion_io_loop_with_provided_loop(minion_opts): """ Test that Minion io_loop conversion works when a loop is provided. """ # Create a Tornado IOLoop tornado_loop = tornado.ioloop.IOLoop() try: minion = salt.minion.Minion( minion_opts, io_loop=tornado_loop, load_grains=False ) try: # Should still be converted to asyncio loop assert isinstance(minion.io_loop, asyncio.AbstractEventLoop) # Should be the same underlying loop assert minion.io_loop is tornado_loop.asyncio_loop finally: minion.destroy() finally: tornado_loop.close() def test_minion_manager_io_loop_is_asyncio_loop(minion_opts): """ Test that MinionManager io_loop is converted to asyncio.AbstractEventLoop. """ with patch("salt.utils.process.SignalHandlingProcess.start"): with patch("salt.utils.verify.valid_id"): mm = salt.minion.MinionManager(minion_opts) try: # Verify io_loop is an asyncio loop assert isinstance(mm.io_loop, asyncio.AbstractEventLoop) # Ensure it has asyncio methods assert hasattr(mm.io_loop, "create_task") assert hasattr(mm.io_loop, "call_soon") # Ensure it doesn't have Tornado-specific methods assert not hasattr(mm.io_loop, "spawn_callback") finally: mm.destroy() def test_syndic_manager_io_loop_is_asyncio_loop(minion_opts): """ Test that SyndicManager io_loop is converted to asyncio.AbstractEventLoop. """ minion_opts["order_masters"] = True sm = salt.minion.SyndicManager(minion_opts) try: # Verify io_loop is an asyncio loop assert isinstance(sm.io_loop, asyncio.AbstractEventLoop) # Ensure it has asyncio methods assert hasattr(sm.io_loop, "create_task") assert hasattr(sm.io_loop, "call_soon") # Ensure it doesn't have Tornado-specific methods assert not hasattr(sm.io_loop, "spawn_callback") finally: sm.destroy() def _run_eval_master(opts): """ Drive MinionBase.eval_master far enough to hit the single-master branch (where the random_master warning lives) without touching the network: DNS resolution is stubbed and the pub channel connects immediately. """ io_loop = tornado.ioloop.IOLoop() minion = salt.minion.MinionBase(opts) mock_channel = MagicMock() mock_channel.connect.return_value = tornado.gen.maybe_future(None) mock_channel.auth.gen_token.return_value = b"token" try: with patch( "salt.channel.client.AsyncPubChannel.factory", return_value=mock_channel ), patch("salt.minion.resolve_dns", return_value={}), patch( "salt.minion.prep_ip_port", return_value={} ): io_loop.run_sync(lambda: minion.eval_master(opts)) finally: io_loop.close() def _single_master_opts(minion_opts): minion_opts["master"] = "salt-master-1" minion_opts["master_type"] = "str" minion_opts["random_master"] = True minion_opts["transport"] = "zeromq" minion_opts["acceptance_wait_time"] = 0 minion_opts["master_tries"] = 1 return minion_opts def test_eval_master_random_master_warning_suppressed_for_multimaster( minion_opts, caplog ): """ In multi-master mode the MinionManager spawns one Minion per master, each bound to a single master but inheriting random_master (multimaster=True). Those children must NOT emit the "random_master ... only one master ... Ignoring" warning. Regression test for the spurious per-master warning. """ opts = _single_master_opts(minion_opts) opts["multimaster"] = True with caplog.at_level(logging.WARNING): _run_eval_master(opts) assert ( "random_master is True but there is only one master specified" not in caplog.text ) def test_eval_master_random_master_warning_for_real_single_master(minion_opts, caplog): """ A genuinely single-master minion (not a multimaster child) with random_master set still gets warned -- random_master really is a no-op there. Guards against over-suppressing the warning. """ opts = _single_master_opts(minion_opts) opts.pop("multimaster", None) with caplog.at_level(logging.WARNING): _run_eval_master(opts) assert "random_master is True but there is only one master specified" in caplog.text