mirror of
https://github.com/Security-Onion-Solutions/securityonion.git
synced 2026-10-01 20:14:45 +02:00
Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
8de8ba811a | ||
|
|
9732e1c639 | ||
|
|
53f9ebcd46 |
No files matched your search
@@ -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
|
||||
|
||||
|
||||
@@ -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 (stderr by default).
|
||||
JID_RE = re.compile(r'salt/run/(\d{20})')
|
||||
|
||||
|
||||
def _make_logger():
|
||||
logger = logging.getLogger('so-push-drainer')
|
||||
@@ -113,14 +126,154 @@ def _dispatch(actions, log):
|
||||
except subprocess.CalledProcessError as exc:
|
||||
log.error('dispatch failed (rc=%s): stdout=%s stderr=%s',
|
||||
exc.returncode, exc.stdout, exc.stderr)
|
||||
return False
|
||||
return None
|
||||
except subprocess.TimeoutExpired:
|
||||
log.error('dispatch timed out after 60s')
|
||||
return False
|
||||
return None
|
||||
except Exception:
|
||||
log.exception('dispatch raised')
|
||||
return None
|
||||
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))
|
||||
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 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 _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 _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']
|
||||
try:
|
||||
if _report_result(record, age, log):
|
||||
_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
|
||||
log.info('dispatch accepted: %s', (result.stdout or '').strip())
|
||||
failures = _orch_failures(ret)
|
||||
if failures:
|
||||
log.error('push failed jid=%s paths=%s; change will be applied at the next scheduled highstate: %s',
|
||||
jid, paths, ' | '.join(failures))
|
||||
else:
|
||||
log.info('push succeeded jid=%s paths=%s', jid, paths)
|
||||
return True
|
||||
|
||||
|
||||
@@ -143,6 +296,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 +364,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:
|
||||
|
||||
@@ -0,0 +1,444 @@
|
||||
# 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')
|
||||
|
||||
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_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_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, '')
|
||||
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_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))
|
||||
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_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):
|
||||
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()
|
||||
@@ -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 %}
|
||||
|
||||
@@ -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 {}
|
||||
Reference in new issue
Block a user