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 {}