From 0a628bb7e78473c452136ec47f403263261116c6 Mon Sep 17 00:00:00 2001 From: Josh Patterson Date: Mon, 5 Oct 2026 17:43:08 -0400 Subject: [PATCH] Revert "Fix/auto apply state queue" --- salt/_beacons/postgres_pillar_beacon.py | 2 - salt/manager/tools/sbin/so-push-drainer | 184 +------ .../tools/sbin/so-push-drainer_test.py | 477 ------------------ salt/orch/push_batch.sls | 8 +- salt/reactor/push_pillar.sls | 11 +- 5 files changed, 13 insertions(+), 669 deletions(-) delete mode 100644 salt/manager/tools/sbin/so-push-drainer_test.py diff --git a/salt/_beacons/postgres_pillar_beacon.py b/salt/_beacons/postgres_pillar_beacon.py index 53b160679..074e9e8af 100644 --- a/salt/_beacons/postgres_pillar_beacon.py +++ b/salt/_beacons/postgres_pillar_beacon.py @@ -131,8 +131,6 @@ def beacon(config): # noqa: C901 'setting_id': setting_id, 'node_id': node_id, }) - log.info('postgres_pillar_beacon: audit_settings id=%d setting_id=%s node_id=%s', - row_id, setting_id, node_id) if row_id > max_id: max_id = row_id diff --git a/salt/manager/tools/sbin/so-push-drainer b/salt/manager/tools/sbin/so-push-drainer index a09d1cb43..66fa67e2c 100644 --- a/salt/manager/tools/sbin/so-push-drainer +++ b/salt/manager/tools/sbin/so-push-drainer @@ -19,8 +19,6 @@ is older than debounce_seconds, this script: * dispatches a single `salt-run state.orchestrate orch.push_batch --async` with the deduped actions list passed as pillar kwargs * deletes the contributed intent files on successful dispatch - * records the orchestration jid under /opt/so/state/push_dispatched and, on - later passes, looks up its result and logs success or per-minion failures Reactor sls files (push_files, push_pillar) write intents but never dispatch directly @@ -32,7 +30,6 @@ import json import logging import logging.handlers import os -import re import subprocess import sys import time @@ -43,19 +40,8 @@ PENDING_DIR = '/opt/so/state/push_pending' LOCK_FILE = os.path.join(PENDING_DIR, '.lock') LOG_FILE = '/opt/so/log/salt/so-push-drainer.log' -DISPATCHED_DIR = '/opt/so/state/push_dispatched' - HIGHSTATE_SENTINEL = '__highstate__' -RESULT_CHECK_DELAY = 30 -RESULT_RECHECK_MAX = 300 -RESULT_MAX_AGE = 7200 -RESULT_CHECKS_PER_PASS = 5 -TEXT_LIMIT = 500 - -# salt-run --async reports the jid only in a log line (stderr by default). -JID_RE = re.compile(r'salt/run/(\d{20})') - def _make_logger(): logger = logging.getLogger('so-push-drainer') @@ -127,167 +113,14 @@ def _dispatch(actions, log): except subprocess.CalledProcessError as exc: log.error('dispatch failed (rc=%s): stdout=%s stderr=%s', exc.returncode, exc.stdout, exc.stderr) - return None + return False except subprocess.TimeoutExpired: log.error('dispatch timed out after 60s') - return None + return False except Exception: log.exception('dispatch raised') - return None - output = '{}\n{}'.format(result.stderr or '', result.stdout or '') - match = JID_RE.search(output) - if not match: - log.warning('dispatch accepted but no jid found, result will not be tracked: output=%s', - _trim(output)) - return '' - log.info('dispatch accepted: jid=%s', match.group(1)) - return match.group(1) - - -def _trim(value): - text = value if isinstance(value, str) else json.dumps(value, default=str) - lines = [line.strip() for line in text.splitlines() if line.strip()] - if 'Traceback (most recent call last):' in text: - # Keep the lead-in and the raised exception; the frames are noise in a log line. - lines = [text.split('Traceback (most recent call last):', 1)[0].strip(), lines[-1]] - text = ' '.join(line for line in lines if line) - return text if len(text) <= TEXT_LIMIT else text[:TEXT_LIMIT] + '...' - - -def _unlink(path, log): - try: - os.unlink(path) - except FileNotFoundError: - pass - except OSError: - log.exception('failed to remove %s', path) - - -def _write_record(path, record, log): - try: - os.makedirs(DISPATCHED_DIR, exist_ok=True) - tmp_path = path + '.tmp' - with open(tmp_path, 'w') as f: - json.dump(record, f) - os.rename(tmp_path, path) - except Exception: - log.exception('failed to record dispatch %s', record.get('jid')) - - -def _record_dispatch(jid, actions, paths, log): - record = {'jid': jid, 'dispatched_at': time.time(), 'actions': actions, 'paths': paths} - _write_record(os.path.join(DISPATCHED_DIR, '{}.json'.format(jid)), record, log) - - -def _lookup_jid(jid, log): - """Returns the job cache entry for jid, {} while it is still running, or None on error.""" - cmd = ['salt-run', 'jobs.lookup_jid', jid, '--out=json'] - try: - result = subprocess.run(cmd, check=True, capture_output=True, text=True, timeout=60) - return json.loads(result.stdout or '{}') - except (subprocess.CalledProcessError, subprocess.TimeoutExpired, ValueError) as exc: - log.warning('lookup of jid %s failed: %s', jid, exc) - return None - - -def _minion_failure(minion_ret): - if isinstance(minion_ret, dict): - return '; '.join( - '{}: {}'.format(state.get('__id__', state_key), _trim(state.get('comment', ''))) - for state_key, state in minion_ret.items() - if isinstance(state, dict) and state.get('result') is False - ) - # A state run rejected before it starts (e.g. another state run is in - # progress) returns a list of error strings instead of state results. - if isinstance(minion_ret, (list, str)): - return _trim(minion_ret) - return '' - - -def _step_failures(step): - if not isinstance(step, dict) or step.get('result') is not False: - return [] - failures = ['{}: {}'.format(step.get('__id__', step.get('name')), _trim(step.get('comment', '')))] - changes = step.get('changes') - minion_rets = changes.get('ret') if isinstance(changes, dict) else None - if isinstance(minion_rets, dict): - for minion, minion_ret in minion_rets.items(): - text = _minion_failure(minion_ret) - if text: - failures.append('{}: {}'.format(minion, text)) - return failures - - -def _orch_failures(ret): - if not isinstance(ret, dict): - return [_trim(ret)] - failures = [] - for job in ret.values(): - if not isinstance(job, dict): - continue - job_ret = job.get('return') - data = job_ret.get('data') if isinstance(job_ret, dict) else {} - if not isinstance(data, dict): - if data: - failures.append(_trim(data)) - data = {} - for steps in data.values(): - if not isinstance(steps, dict): - failures.append(_trim(steps)) - continue - for step in steps.values(): - failures.extend(_step_failures(step)) - if job.get('success') is False and not failures: - failures.append('orchestration reported failure: {}'.format(_trim(job.get('return')))) - return failures - - -def _recheck_delay(age): - return min(RESULT_RECHECK_MAX, max(RESULT_CHECK_DELAY, age / 4)) - - -def _check_dispatched(log, now): - due = [] - for path in glob.glob(os.path.join(DISPATCHED_DIR, '*.json')): - record = _read_intent(path, log) - if not isinstance(record, dict) or not record.get('jid'): - _unlink(path, log) - continue - age = now - record.get('dispatched_at', 0) - last_check = record.get('checked_at', record.get('dispatched_at', 0)) - if now - last_check >= _recheck_delay(age): - due.append((last_check, path, record, age)) - # Least recently checked first, so pushes that are still running can't starve finished ones. - for _, path, record, age in sorted(due, key=lambda item: item[:2])[:RESULT_CHECKS_PER_PASS]: - jid = record['jid'] - try: - if _report_result(record, age, log): - _unlink(path, log) - else: - record['checked_at'] = now - _write_record(path, record, log) - except Exception: - # Drop the record so one unreadable result can't fail every pass ahead of the drain. - log.exception('cannot evaluate result for jid=%s; no longer tracking', jid) - _unlink(path, log) - - -def _report_result(record, age, log): - """Logs the outcome of a dispatched push. Returns True once the record is finished with.""" - jid = record['jid'] - paths = record.get('paths', []) - ret = _lookup_jid(jid, log) - if not ret: - if age > RESULT_MAX_AGE: - log.warning('no result for jid=%s after %ds, no longer tracking; paths=%s', jid, age, paths) - return True return False - failures = _orch_failures(ret) - if failures: - log.error('push failed jid=%s paths=%s; change will be applied at the next scheduled highstate: %s', - jid, paths, ' | '.join(failures)) - else: - log.info('push succeeded jid=%s paths=%s', jid, paths) + log.info('dispatch accepted: %s', (result.stdout or '').strip()) return True @@ -310,9 +143,6 @@ def main(): debounce_seconds = int(push.get('debounce_seconds', 30)) - # Outside the lock: lookups are slow and the reactors take the same lock. - _check_dispatched(log, time.time()) - os.makedirs(PENDING_DIR, exist_ok=True) lock_fd = os.open(LOCK_FILE, os.O_CREAT | os.O_RDWR, 0o644) try: @@ -378,16 +208,10 @@ def main(): len(ready), len(deduped), len(combined_actions), debounce_duration, all_paths[:20], ) - for action in deduped: - log.info('action: %s tgt=%s', 'highstate' if action.get('highstate') else action.get('state'), - action.get('tgt')) - jid = _dispatch(deduped, log) - if jid is None: + if not _dispatch(deduped, log): log.warning('dispatch failed; leaving intent files in place for retry') return 1 - if jid: - _record_dispatch(jid, deduped, all_paths[:20], log) for path, _ in ready: try: diff --git a/salt/manager/tools/sbin/so-push-drainer_test.py b/salt/manager/tools/sbin/so-push-drainer_test.py deleted file mode 100644 index 002011810..000000000 --- a/salt/manager/tools/sbin/so-push-drainer_test.py +++ /dev/null @@ -1,477 +0,0 @@ -# Copyright Security Onion Solutions LLC and/or licensed to Security Onion Solutions LLC under one -# or more contributor license agreements. Licensed under the Elastic License 2.0 as shown at -# https://securityonion.net/license; you may not use this file except in compliance with the -# Elastic License 2.0. - -import importlib.util -import json -import logging -import os -import shutil -import subprocess -import sys -import tempfile -import time -import unittest -from importlib.machinery import SourceFileLoader -from unittest.mock import MagicMock, patch - -HERE = os.path.dirname(os.path.abspath(__file__)) -SCRIPT = os.path.join(HERE, 'so-push-drainer') -_loader = SourceFileLoader('so_push_drainer', SCRIPT) -_spec = importlib.util.spec_from_loader('so_push_drainer', _loader) -drainer = importlib.util.module_from_spec(_spec) - -# salt is not installed where these tests run; the drainer only needs salt.client.Caller. -# Mocked only while the drainer loads: run from the repo root, 'salt' is this repo's salt/ directory. -_salt = MagicMock() -with patch.dict(sys.modules, {'salt': _salt, 'salt.client': _salt.client}): - _loader.exec_module(drainer) - -MASTER = 'manager.localdomain_master' -JID = '20260930171554259426' -ASYNC_STDERR = ('[WARNING ] Running in asynchronous mode. Results of this execution may be collected ' - 'by attaching to the master event bus or by examining the master job cache, if ' - 'configured. This execution is running under tag salt/run/{}\n'.format(JID)) -CONFLICT = ('The function "state.sls" is running as PID 372218 and was started at ' - '2026, Sep 30 17:15:40.466233 with jid 20260930171540466233') - - -def _orch_ret(steps, success=True): - return {MASTER: { - 'fun': 'runner.state.orchestrate', - 'jid': JID, - 'return': {'data': {MASTER: steps}, 'outputter': 'highstate', 'retcode': 0 if success else 1}, - 'success': success, - }} - - -REFRESH_STEP = { - 'salt_|-refresh_pillar_1_|-saltutil.refresh_pillar_|-function': { - '__id__': 'refresh_pillar_1', 'result': True, - 'changes': {'ret': {'manager_standalone': True}}, - 'comment': 'Function ran successfully.', - }, -} - -CONFLICT_RET = _orch_ret(dict(REFRESH_STEP, **{ - 'salt_|-apply_soc_1_|-apply_soc_1_|-state': { - '__id__': 'apply_soc_1', 'result': False, - 'changes': {'out': 'highstate', 'ret': {'manager_standalone': [CONFLICT]}}, - 'comment': 'Run failed on minions: manager_standalone', - }, -}), success=False) - -STATE_FAIL_RET = _orch_ret({ - 'salt_|-apply_hydra_1_|-apply_hydra_1_|-state': { - '__id__': 'apply_hydra_1', 'result': False, - 'changes': {'out': 'highstate', 'ret': {'manager_standalone': { - 'test_|-no_license_|-no_license_|-fail_without_changes': { - '__id__': 'hydra.enabled_no_license_detected', 'result': False, - 'comment': 'This is a feature supported only for customers with a valid license.', - }, - 'file_|-hydra_conf_|-/opt/so/conf/hydra_|-managed': {'result': True, 'comment': 'ok'}, - }}}, - 'comment': 'Run failed on minions: manager_standalone', - }, -}, success=False) - -SUCCESS_RET = _orch_ret(dict(REFRESH_STEP, **{ - 'salt_|-apply_telegraf_1_|-apply_telegraf_1_|-state': { - '__id__': 'apply_telegraf_1', 'result': True, - 'changes': {'out': 'highstate', 'ret': {'manager_standalone': { - 'file_|-tgrafconf_|-/opt/so/conf/telegraf/etc/telegraf.conf_|-managed': {'result': True}, - }}}, - 'comment': 'States ran successfully.', - }, -})) - - -class DrainerTestCase(unittest.TestCase): - - def setUp(self): - self.tmpdir = tempfile.mkdtemp() - self.pending = os.path.join(self.tmpdir, 'push_pending') - self.dispatched = os.path.join(self.tmpdir, 'push_dispatched') - os.makedirs(self.pending) - for name, value in ( - ('PENDING_DIR', self.pending), - ('LOCK_FILE', os.path.join(self.pending, '.lock')), - ('DISPATCHED_DIR', self.dispatched), - ('LOG_FILE', os.path.join(self.tmpdir, 'log', 'so-push-drainer.log')), - ): - patcher = patch.object(drainer, name, value) - patcher.start() - self.addCleanup(patcher.stop) - self.log = MagicMock() - - def tearDown(self): - shutil.rmtree(self.tmpdir, ignore_errors=True) - - def write_json(self, directory, name, data): - os.makedirs(directory, exist_ok=True) - path = os.path.join(directory, name) - with open(path, 'w') as f: - if isinstance(data, str): - f.write(data) - else: - json.dump(data, f) - return path - - def logged(self, level): - return ' '.join(c.args[0] % c.args[1:] for c in getattr(self.log, level).call_args_list) - - -class TestHelpers(DrainerTestCase): - - def test_make_logger_adds_handler_once(self): - logger = logging.getLogger('so-push-drainer') - - def close_handlers(): - for handler in logger.handlers: - handler.close() - logger.handlers.clear() - - self.addCleanup(close_handlers) - logger.handlers.clear() - self.assertIs(drainer._make_logger(), logger) - drainer._make_logger() - self.assertEqual(len(logger.handlers), 1) - self.assertTrue(os.path.isdir(os.path.dirname(drainer.LOG_FILE))) - - def test_load_push_cfg(self): - with patch.object(drainer.salt.client, 'Caller') as caller: - caller.return_value.cmd.return_value = {'enabled': False} - self.assertEqual(drainer._load_push_cfg(), {'enabled': False}) - caller.return_value.cmd.return_value = 'garbage' - self.assertEqual(drainer._load_push_cfg(), {}) - - def test_read_intent(self): - good = self.write_json(self.pending, 'good.json', {'a': 1}) - bad = self.write_json(self.pending, 'bad.json', '{nope') - self.assertEqual(drainer._read_intent(good, self.log), {'a': 1}) - self.assertIsNone(drainer._read_intent(bad, self.log)) - with patch('builtins.open', side_effect=RuntimeError('boom')): - self.assertIsNone(drainer._read_intent(good, self.log)) - self.log.exception.assert_called_once() - - def test_dedupe_actions(self): - actions = [ - 'not a dict', - {'state': 'soc'}, - {'state': 'soc', 'tgt': '*'}, - {'state': 'soc', 'tgt': '*', 'tgt_type': 'compound'}, - {'highstate': True, 'tgt': '*'}, - {'state': 'soc', 'tgt': 'node1', 'tgt_type': 'glob'}, - ] - self.assertEqual(drainer._dedupe_actions(actions), [actions[2], actions[4], actions[5]]) - - def test_trim(self): - self.assertEqual(drainer._trim(' text \n'), 'text') - self.assertEqual(drainer._trim(['a']), '["a"]') - self.assertEqual(drainer._trim(None), 'null') - self.assertEqual(drainer._trim('x' * 600), 'x' * drainer.TEXT_LIMIT + '...') - - def test_trim_traceback(self): - comment = ('An exception occurred in this state: Traceback (most recent call last):\n' - ' File "salt/client/__init__.py", line 1934, in pub\n' - ' raise AuthenticationError(err_msg)\n' - 'salt.exceptions.AuthenticationError: Authentication error occurred.\n') - self.assertEqual(drainer._trim(comment), 'An exception occurred in this state: ' - 'salt.exceptions.AuthenticationError: Authentication error occurred.') - self.assertEqual(drainer._trim('line one\n line two\n'), 'line one line two') - - def test_unlink(self): - drainer._unlink(os.path.join(self.tmpdir, 'missing'), self.log) - self.log.exception.assert_not_called() - drainer._unlink(self.tmpdir, self.log) - self.log.exception.assert_called_once() - - -class TestDispatch(DrainerTestCase): - - def run_dispatch(self, **kwargs): - with patch.object(drainer.subprocess, 'run', **kwargs) as run: - jid = drainer._dispatch([{'state': 'soc', 'tgt': '*'}], self.log) - return jid, run - - def test_jid_parsed_from_stderr(self): - jid, run = self.run_dispatch(return_value=MagicMock(stdout='', stderr=ASYNC_STDERR)) - self.assertEqual(jid, JID) - cmd = run.call_args[0][0] - self.assertEqual(cmd[:3], ['salt-run', 'state.orchestrate', 'orch.push_batch']) - self.assertIn('--async', cmd) - - def test_jid_parsed_from_stdout(self): - jid, _ = self.run_dispatch(return_value=MagicMock(stdout=ASYNC_STDERR, stderr=None)) - self.assertEqual(jid, JID) - - def test_no_jid(self): - jid, _ = self.run_dispatch(return_value=MagicMock(stdout='unexpected output', stderr=None)) - self.assertEqual(jid, '') - self.assertIn('output=unexpected output', self.logged('warning')) - - def test_failures_return_none(self): - for exc in (subprocess.CalledProcessError(1, 'salt-run', 'out', 'err'), - subprocess.TimeoutExpired('salt-run', 60), - RuntimeError('boom')): - jid, _ = self.run_dispatch(side_effect=exc) - self.assertIsNone(jid) - - def test_record_dispatch(self): - drainer._record_dispatch(JID, [{'state': 'soc'}], ['audit:soc.config.licenseKey'], self.log) - with open(os.path.join(self.dispatched, JID + '.json')) as f: - record = json.load(f) - self.assertEqual(record['jid'], JID) - self.assertEqual(record['paths'], ['audit:soc.config.licenseKey']) - self.assertIn('dispatched_at', record) - - def test_record_dispatch_errors(self): - with patch.object(drainer.os, 'makedirs', side_effect=OSError('ro')): - drainer._record_dispatch(JID, [], [], self.log) - drainer._record_dispatch(JID, [object()], [], self.log) - self.assertEqual(self.log.exception.call_count, 2) - self.assertFalse(os.path.exists(os.path.join(self.dispatched, JID + '.json'))) - - -class TestResults(DrainerTestCase): - - def test_lookup_jid(self): - with patch.object(drainer.subprocess, 'run') as run: - run.return_value = MagicMock(stdout=json.dumps(SUCCESS_RET)) - self.assertEqual(drainer._lookup_jid(JID, self.log), SUCCESS_RET) - self.assertEqual(run.call_args[0][0], ['salt-run', 'jobs.lookup_jid', JID, '--out=json']) - run.return_value = MagicMock(stdout='') - self.assertEqual(drainer._lookup_jid(JID, self.log), {}) - run.return_value = MagicMock(stdout='not json') - self.assertIsNone(drainer._lookup_jid(JID, self.log)) - run.side_effect = subprocess.TimeoutExpired('salt-run', 60) - self.assertIsNone(drainer._lookup_jid(JID, self.log)) - - def test_minion_failure_shapes(self): - self.assertEqual(drainer._minion_failure([CONFLICT]), json.dumps([CONFLICT])) - self.assertEqual(drainer._minion_failure('Rendering SLS failed'), 'Rendering SLS failed') - self.assertEqual(drainer._minion_failure(True), '') - self.assertEqual(drainer._minion_failure({'a': {'result': True}}), '') - - def test_orch_failures_conflict(self): - failures = drainer._orch_failures(CONFLICT_RET) - self.assertEqual(failures[0], 'apply_soc_1: Run failed on minions: manager_standalone') - self.assertIn('manager_standalone', failures[1]) - self.assertIn('is running as PID 372218', failures[1]) - self.assertEqual(len(failures), 2) - - def test_orch_failures_failed_state(self): - failures = drainer._orch_failures(STATE_FAIL_RET) - self.assertEqual(len(failures), 2) - self.assertIn('hydra.enabled_no_license_detected: This is a feature', failures[1]) - self.assertNotIn('hydra_conf', failures[1]) - - def test_orch_failures_success(self): - self.assertEqual(drainer._orch_failures(SUCCESS_RET), []) - - def test_orch_failures_render_error(self): - ret = {MASTER: {'return': {'data': {MASTER: ['Rendering SLS failed']}}, 'success': False}} - self.assertEqual(drainer._orch_failures(ret), ['["Rendering SLS failed"]']) - - def test_orch_failures_not_a_dict(self): - self.assertEqual(drainer._orch_failures(['No minions matched']), ['["No minions matched"]']) - self.assertEqual(drainer._orch_failures('Runner error'), ['Runner error']) - - def test_orch_failures_data_not_a_dict(self): - ret = {MASTER: {'return': {'data': ["Rendering SLS 'orch.push_batch' failed"]}, 'success': False}} - self.assertEqual(drainer._orch_failures(ret), ['["Rendering SLS \'orch.push_batch\' failed"]']) - - def test_orch_failures_odd_changes(self): - for changes in ('Run failed', {'ret': ['manager_standalone']}): - ret = _orch_ret({'salt_|-apply_soc_1_|-apply_soc_1_|-state': { - '__id__': 'apply_soc_1', 'result': False, 'changes': changes, 'comment': 'Run failed on minions', - }}, success=False) - self.assertEqual(drainer._orch_failures(ret), ['apply_soc_1: Run failed on minions']) - - def test_orch_failures_unparsed(self): - self.assertEqual(drainer._orch_failures({MASTER: 'odd'}), []) - ret = {MASTER: {'return': 'Exception occurred', 'success': False}} - self.assertEqual(drainer._orch_failures(ret), ['orchestration reported failure: Exception occurred']) - - def record(self, jid, age, now): - return self.write_json(self.dispatched, jid + '.json', { - 'jid': jid, 'dispatched_at': now - age, 'actions': [], 'paths': ['audit:' + jid], - }) - - def test_check_dispatched(self): - now = time.time() - results = { - '1_failed': CONFLICT_RET, - '2_ok': SUCCESS_RET, - '3_pending': {}, - '4_expired': None, - } - young = self.record('0_young', 5, now) - paths = {jid: self.record(jid, 60, now) for jid in results} - paths['4_expired'] = self.record('4_expired', drainer.RESULT_MAX_AGE + 1, now) - bad = self.write_json(self.dispatched, '5_bad.json', '{nope') - with patch.object(drainer, '_lookup_jid', side_effect=lambda jid, log: results[jid]): - drainer._check_dispatched(self.log, now) - - self.assertTrue(os.path.exists(young)) - with open(paths['3_pending']) as f: - self.assertEqual(json.load(f)['checked_at'], now) - for jid in ('1_failed', '2_ok', '4_expired'): - self.assertFalse(os.path.exists(paths[jid]), jid) - self.assertFalse(os.path.exists(bad)) - self.assertIn('push failed jid=1_failed', self.logged('error')) - self.assertIn('is running as PID 372218', self.logged('error')) - self.assertIn('push succeeded jid=2_ok', self.logged('info')) - self.assertIn('no result for jid=4_expired', self.logged('warning')) - - def test_check_dispatched_survives_bad_result(self): - now = time.time() - bad = self.record('1_bad', 60, now) - good = self.record('2_ok', 60, now) - - def orch_failures(ret): - if ret == 'boom': - raise ValueError('unexpected shape') - return [] - - with patch.object(drainer, '_lookup_jid', side_effect=lambda jid, log: 'boom' if jid == '1_bad' else SUCCESS_RET), \ - patch.object(drainer, '_orch_failures', side_effect=orch_failures): - drainer._check_dispatched(self.log, now) - self.assertFalse(os.path.exists(bad)) - self.assertFalse(os.path.exists(good)) - self.log.exception.assert_called_once() - self.assertIn('jid=1_bad', self.log.exception.call_args[0][0] % self.log.exception.call_args[0][1:]) - self.assertIn('push succeeded jid=2_ok', self.logged('info')) - - def test_recheck_delay(self): - self.assertEqual(drainer._recheck_delay(10), drainer.RESULT_CHECK_DELAY) - self.assertEqual(drainer._recheck_delay(400), 100) - self.assertEqual(drainer._recheck_delay(drainer.RESULT_MAX_AGE), drainer.RESULT_RECHECK_MAX) - - def test_check_dispatched_limit_rotates(self): - now = time.time() - limit = drainer.RESULT_CHECKS_PER_PASS - jids = ['{:02d}'.format(i) for i in range(limit + 2)] - for jid in jids: - self.record(jid, 60, now) - with patch.object(drainer, '_lookup_jid', return_value={}) as lookup: - drainer._check_dispatched(self.log, now) - self.assertEqual([c.args[0] for c in lookup.call_args_list], jids[:limit]) - lookup.reset_mock() - drainer._check_dispatched(self.log, now + 40) - self.assertEqual([c.args[0] for c in lookup.call_args_list], jids[limit:] + jids[:limit - 2]) - - def test_check_dispatched_not_blocked_by_running(self): - now = time.time() - for i in range(drainer.RESULT_CHECKS_PER_PASS): - self.record('1_running{}'.format(i), 600, now) - done = self.record('2_done', 60, now) - - def lookup(jid, log): - return SUCCESS_RET if jid == '2_done' else {} - - with patch.object(drainer, '_lookup_jid', side_effect=lookup) as lookup_jid: - drainer._check_dispatched(self.log, now) - self.assertTrue(os.path.exists(done)) - lookup_jid.reset_mock() - drainer._check_dispatched(self.log, now + 15) - self.assertEqual([c.args[0] for c in lookup_jid.call_args_list], ['2_done']) - self.assertFalse(os.path.exists(done)) - self.assertIn('push succeeded jid=2_done', self.logged('info')) - - -class TestMain(DrainerTestCase): - - def setUp(self): - super().setUp() - self.cfg = {'enabled': True, 'debounce_seconds': 30} - for name, kwargs in ( - ('_make_logger', {'return_value': self.log}), - ('_load_push_cfg', {'side_effect': lambda: self.cfg}), - ('_check_dispatched', {}), - ): - patcher = patch.object(drainer, name, **kwargs) - setattr(self, name, patcher.start()) - self.addCleanup(patcher.stop) - - def intent(self, name, age=60, actions=None, paths=None): - now = time.time() - return self.write_json(self.pending, name, { - 'first_touch': now - age - 5, 'last_touch': now - age, - 'actions': [{'state': 'soc', 'tgt': '*'}] if actions is None else actions, - 'paths': paths or ['audit:soc.config.licenseKey'], - }) - - def test_no_pending_dir(self): - shutil.rmtree(self.pending) - self.assertEqual(drainer.main(), 0) - self._load_push_cfg.assert_not_called() - - def test_cfg_error(self): - self._load_push_cfg.side_effect = RuntimeError('no salt') - self.assertEqual(drainer.main(), 1) - - def test_disabled(self): - self.cfg['enabled'] = False - self.assertEqual(drainer.main(), 0) - self._check_dispatched.assert_not_called() - - def test_no_intents_still_checks_results(self): - self.assertEqual(drainer.main(), 0) - self._check_dispatched.assert_called_once() - - def test_debounce_and_broken(self): - young = self.intent('young.json', age=1) - broken = self.write_json(self.pending, 'broken.json', '{nope') - with patch.object(drainer, '_dispatch') as dispatch: - self.assertEqual(drainer.main(), 0) - dispatch.assert_not_called() - self.assertTrue(os.path.exists(young)) - self.assertFalse(os.path.exists(broken)) - - def test_broken_unlink_error_ignored(self): - self.write_json(self.pending, 'broken.json', '{nope') - with patch.object(drainer.os, 'unlink', side_effect=OSError('busy')): - self.assertEqual(drainer.main(), 0) - - def test_no_usable_actions(self): - path = self.intent('empty.json', actions=[{'state': 'soc'}]) - self.assertEqual(drainer.main(), 0) - self.assertFalse(os.path.exists(path)) - self.intent('empty.json', actions=[{'state': 'soc'}]) - with patch.object(drainer.os, 'unlink', side_effect=OSError('busy')): - self.assertEqual(drainer.main(), 0) - - def test_dispatch_failure_keeps_intents(self): - path = self.intent('pillar_soc.json') - with patch.object(drainer, '_dispatch', return_value=None): - self.assertEqual(drainer.main(), 1) - self.assertTrue(os.path.exists(path)) - - def test_dispatch_records_jid(self): - soc = self.intent('pillar_soc.json') - hs = self.intent('pillar_global.json', actions=[{'highstate': True, 'tgt': '*'}], paths=['audit:global.x']) - with patch.object(drainer, '_dispatch', return_value=JID) as dispatch, \ - patch.object(drainer, '_record_dispatch') as record: - self.assertEqual(drainer.main(), 0) - self.assertEqual(len(dispatch.call_args[0][0]), 2) - record.assert_called_once() - self.assertEqual(record.call_args[0][0], JID) - self.assertEqual(sorted(record.call_args[0][2]), ['audit:global.x', 'audit:soc.config.licenseKey']) - self.assertFalse(os.path.exists(soc)) - self.assertFalse(os.path.exists(hs)) - self.assertIn('action: highstate tgt=*', self.logged('info')) - - def test_dispatch_without_jid_not_recorded(self): - self.intent('pillar_soc.json') - with patch.object(drainer, '_dispatch', return_value=''), \ - patch.object(drainer, '_record_dispatch') as record, \ - patch.object(drainer.os, 'unlink', side_effect=OSError('busy')): - self.assertEqual(drainer.main(), 0) - record.assert_not_called() - self.log.exception.assert_called_once() - - -if __name__ == '__main__': - unittest.main() diff --git a/salt/orch/push_batch.sls b/salt/orch/push_batch.sls index 2fff2f646..ba8676961 100644 --- a/salt/orch/push_batch.sls +++ b/salt/orch/push_batch.sls @@ -3,8 +3,6 @@ {% set BATCH = AUTOAPPLY.batch %} {% set BATCH_WAIT = AUTOAPPLY.batch_wait %} -{# queue must be a top-level salt.state arg (kwarg is ignored); an int is max_queue and still fails on conflict #} - {% for action in actions %} {% if action.get('highstate') %} apply_highstate_{{ loop.index }}: @@ -14,7 +12,8 @@ apply_highstate_{{ loop.index }}: - highstate: True - batch: {{ action.get('batch', BATCH) }} - batch_wait: {{ action.get('batch_wait', BATCH_WAIT) }} - - queue: True + - kwarg: + queue: 2 {% else %} refresh_pillar_{{ loop.index }}: salt.function: @@ -30,7 +29,8 @@ apply_{{ action.state | replace('.', '_') }}_{{ loop.index }}: - {{ action.state }} - batch: {{ action.get('batch', BATCH) }} - batch_wait: {{ action.get('batch_wait', BATCH_WAIT) }} - - queue: True + - kwarg: + queue: 2 - require: - salt: refresh_pillar_{{ loop.index }} {% endif %} diff --git a/salt/reactor/push_pillar.sls b/salt/reactor/push_pillar.sls index 30417d6a2..8d28ef5cf 100644 --- a/salt/reactor/push_pillar.sls +++ b/salt/reactor/push_pillar.sls @@ -138,7 +138,6 @@ def run(): # top level so the reactor is robust to either shape. event = data.get('data', data) # noqa: F821 -- data provided by reactor setting_id = event.get('setting_id', '') - audit_id = event.get('id') node_id = (event.get('node_id') or '').strip() app = _app_from_setting(setting_id) @@ -151,8 +150,8 @@ def run(): if not entry: LOG.warning( 'push_pillar: app "%s" is not in pillar_push_map.yaml; change will be ' - 'picked up at the next scheduled highstate (setting_id=%s audit_id=%s)', - app, setting_id, audit_id, + 'picked up at the next scheduled highstate (setting_id=%s)', + app, setting_id, ) return {} @@ -166,12 +165,12 @@ def run(): 'node_{}_{}'.format(node_id, app), actions, 'audit:{}@{}'.format(setting_id, node_id), ) - LOG.info('push_pillar: per-node intent updated for %s on %s (setting_id=%s audit_id=%s)', - app, node_id, setting_id, audit_id) + LOG.info('push_pillar: per-node intent updated for %s on %s (setting_id=%s)', + app, node_id, setting_id) return {} # Branch B: grid-wide app change -> use the map entry's actions as-is. actions = list(entry) # copy to avoid mutating the cache _write_intent('pillar_{}'.format(app), actions, 'audit:{}'.format(setting_id)) - LOG.info('push_pillar: app intent updated for %s (setting_id=%s audit_id=%s)', app, setting_id, audit_id) + LOG.info('push_pillar: app intent updated for %s (setting_id=%s)', app, setting_id) return {}