diff --git a/salt/_beacons/postgres_pillar_beacon.py b/salt/_beacons/postgres_pillar_beacon.py index 074e9e8af..53b160679 100644 --- a/salt/_beacons/postgres_pillar_beacon.py +++ b/salt/_beacons/postgres_pillar_beacon.py @@ -131,6 +131,8 @@ 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 66fa67e2c..a09d1cb43 100644 --- a/salt/manager/tools/sbin/so-push-drainer +++ b/salt/manager/tools/sbin/so-push-drainer @@ -19,6 +19,8 @@ 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 @@ -30,6 +32,7 @@ import json import logging import logging.handlers import os +import re import subprocess import sys import time @@ -40,8 +43,19 @@ 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') @@ -113,14 +127,167 @@ 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 False + return None except subprocess.TimeoutExpired: log.error('dispatch timed out after 60s') - return False + return None 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 - log.info('dispatch accepted: %s', (result.stdout or '').strip()) + 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) return True @@ -143,6 +310,9 @@ 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: @@ -208,10 +378,16 @@ 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')) - if not _dispatch(deduped, log): + jid = _dispatch(deduped, log) + if jid is None: 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 new file mode 100644 index 000000000..002011810 --- /dev/null +++ b/salt/manager/tools/sbin/so-push-drainer_test.py @@ -0,0 +1,477 @@ +# 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 ba8676961..2fff2f646 100644 --- a/salt/orch/push_batch.sls +++ b/salt/orch/push_batch.sls @@ -3,6 +3,8 @@ {% 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 }}: @@ -12,8 +14,7 @@ apply_highstate_{{ loop.index }}: - highstate: True - batch: {{ action.get('batch', BATCH) }} - batch_wait: {{ action.get('batch_wait', BATCH_WAIT) }} - - kwarg: - queue: 2 + - queue: True {% else %} refresh_pillar_{{ loop.index }}: salt.function: @@ -29,8 +30,7 @@ apply_{{ action.state | replace('.', '_') }}_{{ loop.index }}: - {{ action.state }} - batch: {{ action.get('batch', BATCH) }} - batch_wait: {{ action.get('batch_wait', BATCH_WAIT) }} - - kwarg: - queue: 2 + - queue: True - require: - salt: refresh_pillar_{{ loop.index }} {% endif %} diff --git a/salt/reactor/push_pillar.sls b/salt/reactor/push_pillar.sls index 8d28ef5cf..30417d6a2 100644 --- a/salt/reactor/push_pillar.sls +++ b/salt/reactor/push_pillar.sls @@ -138,6 +138,7 @@ 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) @@ -150,8 +151,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)', - app, setting_id, + 'picked up at the next scheduled highstate (setting_id=%s audit_id=%s)', + app, setting_id, audit_id, ) return {} @@ -165,12 +166,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)', - app, node_id, setting_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) 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)', app, setting_id) + LOG.info('push_pillar: app intent updated for %s (setting_id=%s audit_id=%s)', app, setting_id, audit_id) return {}