From 53f9ebcd46cbef56b9a05ce4b25800507d6e368e Mon Sep 17 00:00:00 2001 From: Josh Patterson Date: Wed, 30 Sep 2026 15:15:50 -0400 Subject: [PATCH 1/6] FIX: queue auto-applied state runs instead of failing on conflict orch.push_batch passed `kwarg: {queue: 2}` to salt.state, but in Salt 3006 queue is a top-level salt.state argument and salt.state always sets the minion's queue kwarg from it (default False), so the kwarg block was silently dropped and every pushed state ran with queue=False. The drainer dispatches a separate async orchestration each 15s pass, so settings saved more than ~15s apart overlap on the same minion and every run after the first fails immediately with 'The function "state.sls" is running as PID ...'. The change then waits for the next scheduled highstate. Seen on a 3.4.0 standalone: hydra.enabled, telegraf.output, and two soc settings (including soc.config.licenseKey) were saved within 30s. The soc state was dispatched while the telegraf state was still running and was rejected, so the license key was not applied. Use `queue: True`, as orch.deploy_newnode already does. An int is treated as max_queue and still falls through to the conflict error once that many state runs are active. The failure was only visible in the master log, since the drainer dispatches with --async and logged only "dispatch accepted". The drainer now: - logs each dispatched action - parses the orchestration jid from salt-run's stderr (the only place --async reports it) and records it under /opt/so/state/push_dispatched - on later passes looks each jid up with jobs.lookup_jid and logs either "push succeeded" or an ERROR with the failed step, the per-minion failed states or rejection text, and the triggering paths Lookups run outside the pending-intent lock since the reactors share it. The beacon now logs each audit_settings row it emits and the reactor logs the audit row id, so a single change can be traced from audit_settings to its push result. Adds so-push-drainer_test.py; the drainer is now held to the 100% coverage requirement in python-test. Verified on the standalone: a soc push dispatched while a 90s state run was in progress queued behind it (queue=True in the job args), completed, and the drainer logged "push succeeded" for its jid. The new result parsing reports the original soc conflict and the hydra license failure from the job cache. --- salt/_beacons/postgres_pillar_beacon.py | 2 + salt/manager/tools/sbin/so-push-drainer | 145 ++++++- .../tools/sbin/so-push-drainer_test.py | 391 ++++++++++++++++++ salt/orch/push_batch.sls | 8 +- salt/reactor/push_pillar.sls | 11 +- 5 files changed, 542 insertions(+), 15 deletions(-) create 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 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..d80000ccb 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,18 @@ 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_MAX_AGE = 7200 +RESULT_CHECKS_PER_PASS = 5 +TEXT_LIMIT = 500 + +# salt-run --async reports the jid only in a log line on stderr. +JID_RE = re.compile(r'salt/run/(\d{20})') + def _make_logger(): logger = logging.getLogger('so-push-drainer') @@ -113,15 +126,126 @@ 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 False - log.info('dispatch accepted: %s', (result.stdout or '').strip()) - return True + return None + match = JID_RE.search(result.stderr or '') + if not match: + log.warning('dispatch accepted but no jid found, result will not be tracked: stderr=%s', + _trim(result.stderr)) + 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) + text = text.strip() + return text if len(text) <= TEXT_LIMIT else text[:TEXT_LIMIT] + '...' + + +def _unlink(path, log): + try: + os.unlink(path) + except OSError: + log.exception('failed to remove %s', path) + + +def _record_dispatch(jid, actions, paths, log): + record = {'jid': jid, 'dispatched_at': time.time(), 'actions': actions, 'paths': paths} + path = os.path.join(DISPATCHED_DIR, '{}.json'.format(jid)) + 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 OSError: + log.exception('failed to record dispatch %s', jid) + + +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 _orch_failures(ret): + failures = [] + for job in ret.values(): + if not isinstance(job, dict): + continue + job_ret = job.get('return') + data = (job_ret.get('data') or {}) if isinstance(job_ret, dict) else {} + for steps in data.values(): + if not isinstance(steps, dict): + failures.append(_trim(steps)) + continue + for step in steps.values(): + if not isinstance(step, dict) or step.get('result') is not False: + continue + failures.append('{}: {}'.format(step.get('__id__', step.get('name')), step.get('comment'))) + minion_rets = (step.get('changes') or {}).get('ret') or {} + for minion, minion_ret in minion_rets.items(): + text = _minion_failure(minion_ret) + if text: + failures.append('{}: {}'.format(minion, text)) + if job.get('success') is False and not failures: + failures.append('orchestration reported failure: {}'.format(_trim(job.get('return')))) + return failures + + +def _check_dispatched(log, now): + checked = 0 + for path in sorted(glob.glob(os.path.join(DISPATCHED_DIR, '*.json'))): + if checked >= RESULT_CHECKS_PER_PASS: + break + 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) + if age < RESULT_CHECK_DELAY: + continue + checked += 1 + 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) + _unlink(path, log) + continue + 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) + _unlink(path, log) def main(): @@ -143,6 +267,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 +335,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..65c1c042b --- /dev/null +++ b/salt/manager/tools/sbin/so-push-drainer_test.py @@ -0,0 +1,391 @@ +# 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 + +# salt is not installed where these tests run; the drainer only needs salt.client.Caller. +_salt = MagicMock() +sys.modules.setdefault('salt', _salt) +sys.modules.setdefault('salt.client', _salt.client) + +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) +_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') + self.addCleanup(logger.handlers.clear) + 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_unlink_missing_logs(self): + drainer._unlink(os.path.join(self.tmpdir, 'missing'), 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_no_jid(self): + jid, _ = self.run_dispatch(return_value=MagicMock(stdout='', stderr=None)) + self.assertEqual(jid, '') + self.log.warning.assert_called_once() + + 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_oserror(self): + with patch.object(drainer.os, 'makedirs', side_effect=OSError('ro')): + drainer._record_dispatch(JID, [], [], self.log) + self.log.exception.assert_called_once() + + +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_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)) + self.assertTrue(os.path.exists(paths['3_pending'])) + 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_limit(self): + now = time.time() + for i in range(drainer.RESULT_CHECKS_PER_PASS + 2): + self.record('{:02d}'.format(i), 60, now) + with patch.object(drainer, '_lookup_jid', return_value={}) as lookup: + drainer._check_dispatched(self.log, now) + self.assertEqual(lookup.call_count, drainer.RESULT_CHECKS_PER_PASS) + + +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 {} From 9732e1c639358785bcbc28e2d04d72f1da45e547 Mon Sep 17 00:00:00 2001 From: Josh Patterson Date: Wed, 30 Sep 2026 15:43:04 -0400 Subject: [PATCH 2/6] Trim tracebacks in push failure log lines When an orchestration step raises, salt returns the full traceback as the step comment, and the drainer wrote it verbatim, putting ~70 lines into so-push-drainer.log per failure. Collapse comments to one line and, for tracebacks, keep only the lead-in and the raised exception, e.g. "apply_soc_1: An exception occurred in this state: salt.exceptions.AuthenticationError: Authentication error occurred." Seen on a standalone when a pushed highstate restarted salt-master while two queued pushes were waiting: their orchestrations lost the master connection and failed with AuthenticationError, although the minion completed both state runs. --- salt/manager/tools/sbin/so-push-drainer | 8 ++++++-- salt/manager/tools/sbin/so-push-drainer_test.py | 9 +++++++++ 2 files changed, 15 insertions(+), 2 deletions(-) diff --git a/salt/manager/tools/sbin/so-push-drainer b/salt/manager/tools/sbin/so-push-drainer index d80000ccb..369fcc0fe 100644 --- a/salt/manager/tools/sbin/so-push-drainer +++ b/salt/manager/tools/sbin/so-push-drainer @@ -144,7 +144,11 @@ def _dispatch(actions, log): def _trim(value): text = value if isinstance(value, str) else json.dumps(value, default=str) - text = text.strip() + 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] + '...' @@ -207,7 +211,7 @@ def _orch_failures(ret): for step in steps.values(): if not isinstance(step, dict) or step.get('result') is not False: continue - failures.append('{}: {}'.format(step.get('__id__', step.get('name')), step.get('comment'))) + failures.append('{}: {}'.format(step.get('__id__', step.get('name')), _trim(step.get('comment', '')))) minion_rets = (step.get('changes') or {}).get('ret') or {} for minion, minion_ret in minion_rets.items(): text = _minion_failure(minion_ret) diff --git a/salt/manager/tools/sbin/so-push-drainer_test.py b/salt/manager/tools/sbin/so-push-drainer_test.py index 65c1c042b..f758b6107 100644 --- a/salt/manager/tools/sbin/so-push-drainer_test.py +++ b/salt/manager/tools/sbin/so-push-drainer_test.py @@ -166,6 +166,15 @@ class TestHelpers(DrainerTestCase): 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_missing_logs(self): drainer._unlink(os.path.join(self.tmpdir, 'missing'), self.log) self.log.exception.assert_called_once() From 8de8ba811aed3e16fb0464200cc22d37781f36f9 Mon Sep 17 00:00:00 2001 From: Josh Patterson Date: Thu, 1 Oct 2026 09:06:28 -0400 Subject: [PATCH 3/6] Harden so-push-drainer result parsing _orch_failures assumed every level of a jobs.lookup_jid result was a dict. A list or string at the top level, in return.data, in a step's changes, or in changes.ret raised AttributeError. Because result checks run before the drain and a record is only removed after it is evaluated, one such record would have failed every 15s pass and stopped all pushes until it was removed by hand. Guard each shape, and evaluate each record under its own exception handler so an unreadable result is logged and dropped instead of blocking the drainer. Per-step parsing moves to _step_failures. Search stdout as well as stderr for the async jid, in case salt-run logging is routed to stdout. Close the RotatingFileHandler in test_make_logger_adds_handler_once to avoid a ResourceWarning on Python 3.12+. Verified on a 3.4.0 standalone: real failed and successful orchestration results parse as before, a record whose evaluation raises is logged and removed while the next record still reports, and a replicated SOC change to telegraf.output (and its revert) is pushed, rendered and logged as succeeded. --- salt/manager/tools/sbin/so-push-drainer | 73 +++++++++++++------ .../tools/sbin/so-push-drainer_test.py | 46 +++++++++++- 2 files changed, 94 insertions(+), 25 deletions(-) diff --git a/salt/manager/tools/sbin/so-push-drainer b/salt/manager/tools/sbin/so-push-drainer index 369fcc0fe..48c3e5b7e 100644 --- a/salt/manager/tools/sbin/so-push-drainer +++ b/salt/manager/tools/sbin/so-push-drainer @@ -52,7 +52,7 @@ RESULT_MAX_AGE = 7200 RESULT_CHECKS_PER_PASS = 5 TEXT_LIMIT = 500 -# salt-run --async reports the jid only in a log line on stderr. +# salt-run --async reports the jid only in a log line (stderr by default). JID_RE = re.compile(r'salt/run/(\d{20})') @@ -133,7 +133,7 @@ def _dispatch(actions, log): except Exception: log.exception('dispatch raised') return None - match = JID_RE.search(result.stderr or '') + match = JID_RE.search('{}\n{}'.format(result.stderr or '', result.stdout or '')) if not match: log.warning('dispatch accepted but no jid found, result will not be tracked: stderr=%s', _trim(result.stderr)) @@ -197,26 +197,39 @@ def _minion_failure(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') or {}) if isinstance(job_ret, dict) else {} + 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(): - if not isinstance(step, dict) or step.get('result') is not False: - continue - failures.append('{}: {}'.format(step.get('__id__', step.get('name')), _trim(step.get('comment', '')))) - minion_rets = (step.get('changes') or {}).get('ret') or {} - for minion, minion_ret in minion_rets.items(): - text = _minion_failure(minion_ret) - if text: - failures.append('{}: {}'.format(minion, text)) + 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 @@ -236,20 +249,32 @@ def _check_dispatched(log, now): continue checked += 1 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) + try: + if _report_result(record, age, log): _unlink(path, log) - continue - 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) - _unlink(path, 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) + return True def main(): diff --git a/salt/manager/tools/sbin/so-push-drainer_test.py b/salt/manager/tools/sbin/so-push-drainer_test.py index f758b6107..0c415710e 100644 --- a/salt/manager/tools/sbin/so-push-drainer_test.py +++ b/salt/manager/tools/sbin/so-push-drainer_test.py @@ -126,7 +126,13 @@ class TestHelpers(DrainerTestCase): def test_make_logger_adds_handler_once(self): logger = logging.getLogger('so-push-drainer') - self.addCleanup(logger.handlers.clear) + + 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() @@ -194,6 +200,10 @@ class TestDispatch(DrainerTestCase): 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='', stderr=None)) self.assertEqual(jid, '') @@ -260,6 +270,21 @@ class TestResults(DrainerTestCase): 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}} @@ -295,6 +320,25 @@ class TestResults(DrainerTestCase): 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_check_dispatched_limit(self): now = time.time() for i in range(drainer.RESULT_CHECKS_PER_PASS + 2): From ba95b9bbc2941c72f5c93fac78d406b6cdf09713 Mon Sep 17 00:00:00 2001 From: Josh Patterson Date: Fri, 2 Oct 2026 10:35:30 -0400 Subject: [PATCH 4/6] Address review feedback on so-push-drainer result tracking Result checks walked dispatch records oldest-first with a cap of five lookups per pass, counting records whose push was still running. Five long-running pushes therefore used every slot on every 15s pass and newer, finished pushes were not reported until one cleared. Check the least recently checked records first and back off on pushes that are still running (30s for the first two minutes, then age/4 up to 5 minutes), recording checked_at in the dispatch record. Catch any exception when writing a dispatch record so a failed write cannot skip intent cleanup and re-dispatch the same intents every pass. Log both output streams when no jid is found, and stop logging a traceback when a record has already been removed. Scope the test's salt mock to the drainer import. Run from the repo root, 'salt' resolves to this repo's salt/ directory as a namespace package, so setdefault left it in place and test_load_push_cfg failed. Verified on a 3.4.0 managersearch + sensor: a pushed highstate with soc and telegraf pushes dispatched into it all reported success, with 25 result lookups across the three pushes instead of one per record per pass. --- salt/manager/tools/sbin/so-push-drainer | 44 ++++++++----- .../tools/sbin/so-push-drainer_test.py | 65 ++++++++++++++----- 2 files changed, 78 insertions(+), 31 deletions(-) diff --git a/salt/manager/tools/sbin/so-push-drainer b/salt/manager/tools/sbin/so-push-drainer index 48c3e5b7e..a09d1cb43 100644 --- a/salt/manager/tools/sbin/so-push-drainer +++ b/salt/manager/tools/sbin/so-push-drainer @@ -48,6 +48,7 @@ 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 @@ -133,10 +134,11 @@ def _dispatch(actions, log): except Exception: log.exception('dispatch raised') return None - match = JID_RE.search('{}\n{}'.format(result.stderr or '', result.stdout or '')) + 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: stderr=%s', - _trim(result.stderr)) + 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) @@ -155,21 +157,26 @@ def _trim(value): def _unlink(path, log): try: os.unlink(path) + except FileNotFoundError: + pass except OSError: log.exception('failed to remove %s', path) -def _record_dispatch(jid, actions, paths, log): - record = {'jid': jid, 'dispatched_at': time.time(), 'actions': actions, 'paths': paths} - path = os.path.join(DISPATCHED_DIR, '{}.json'.format(jid)) +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 OSError: - log.exception('failed to record dispatch %s', jid) + 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): @@ -235,23 +242,30 @@ def _orch_failures(ret): return failures +def _recheck_delay(age): + return min(RESULT_RECHECK_MAX, max(RESULT_CHECK_DELAY, age / 4)) + + def _check_dispatched(log, now): - checked = 0 - for path in sorted(glob.glob(os.path.join(DISPATCHED_DIR, '*.json'))): - if checked >= RESULT_CHECKS_PER_PASS: - break + 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) - if age < RESULT_CHECK_DELAY: - continue - checked += 1 + 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) diff --git a/salt/manager/tools/sbin/so-push-drainer_test.py b/salt/manager/tools/sbin/so-push-drainer_test.py index 0c415710e..002011810 100644 --- a/salt/manager/tools/sbin/so-push-drainer_test.py +++ b/salt/manager/tools/sbin/so-push-drainer_test.py @@ -16,17 +16,17 @@ import unittest from importlib.machinery import SourceFileLoader from unittest.mock import MagicMock, patch -# salt is not installed where these tests run; the drainer only needs salt.client.Caller. -_salt = MagicMock() -sys.modules.setdefault('salt', _salt) -sys.modules.setdefault('salt.client', _salt.client) - 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) -_loader.exec_module(drainer) + +# 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' @@ -181,8 +181,10 @@ class TestHelpers(DrainerTestCase): 'salt.exceptions.AuthenticationError: Authentication error occurred.') self.assertEqual(drainer._trim('line one\n line two\n'), 'line one line two') - def test_unlink_missing_logs(self): + 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() @@ -205,9 +207,9 @@ class TestDispatch(DrainerTestCase): self.assertEqual(jid, JID) def test_no_jid(self): - jid, _ = self.run_dispatch(return_value=MagicMock(stdout='', stderr=None)) + jid, _ = self.run_dispatch(return_value=MagicMock(stdout='unexpected output', stderr=None)) self.assertEqual(jid, '') - self.log.warning.assert_called_once() + self.assertIn('output=unexpected output', self.logged('warning')) def test_failures_return_none(self): for exc in (subprocess.CalledProcessError(1, 'salt-run', 'out', 'err'), @@ -224,10 +226,12 @@ class TestDispatch(DrainerTestCase): self.assertEqual(record['paths'], ['audit:soc.config.licenseKey']) self.assertIn('dispatched_at', record) - def test_record_dispatch_oserror(self): + def test_record_dispatch_errors(self): with patch.object(drainer.os, 'makedirs', side_effect=OSError('ro')): drainer._record_dispatch(JID, [], [], self.log) - self.log.exception.assert_called_once() + 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): @@ -311,7 +315,8 @@ class TestResults(DrainerTestCase): drainer._check_dispatched(self.log, now) self.assertTrue(os.path.exists(young)) - self.assertTrue(os.path.exists(paths['3_pending'])) + 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)) @@ -339,13 +344,41 @@ class TestResults(DrainerTestCase): 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_check_dispatched_limit(self): + 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() - for i in range(drainer.RESULT_CHECKS_PER_PASS + 2): - self.record('{:02d}'.format(i), 60, now) + 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(lookup.call_count, drainer.RESULT_CHECKS_PER_PASS) + 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): From 90b3d37be6a7f6f56687d1021b193a5ffcb65abe Mon Sep 17 00:00:00 2001 From: defensivedepth Date: Mon, 5 Oct 2026 08:14:13 -0400 Subject: [PATCH 5/6] Clarify req --- salt/soc/config.sls | 1 + 1 file changed, 1 insertion(+) diff --git a/salt/soc/config.sls b/salt/soc/config.sls index d613145e9..be2887897 100644 --- a/salt/soc/config.sls +++ b/salt/soc/config.sls @@ -118,6 +118,7 @@ crondetectionsbackup: - month: '*' - dayweek: '*' +# sigma-cli only loads *.yml from the pipelines dir socsigmafinalpipeline: file.managed: - name: /opt/so/conf/soc/sigma_pipelines/sigma_final_pipeline.yml From 0a628bb7e78473c452136ec47f403263261116c6 Mon Sep 17 00:00:00 2001 From: Josh Patterson Date: Mon, 5 Oct 2026 17:43:08 -0400 Subject: [PATCH 6/6] 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 {}