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.
This commit is contained in:
Josh Patterson committed 2026-10-01 09:06:28 -04:00
1 parent 9732e1c639
commit 8de8ba811a
2 files changed
+94 -25

No files matched your search

+49 -24
View File
@@ -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():
@@ -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):