Compare commits

..
Author SHA1 Message Date
Josh Patterson 9732e1c639 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.
2026-09-30 16:18:37 -04:00
Josh Patterson 53f9ebcd46 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.
2026-09-30 16:18:37 -04:00
coreyogburn 47d74f1ae1 Merge pull request #16270 from Security-Onion-Solutions/cogburn/automation
New Automation Fields
2026-09-30 11:10:51 -06:00
Corey Ogburn 855716846a New Automation Fields 2026-09-29 16:50:04 -06:00
Mike Reeves 8e35d70595 Merge pull request #16266 from Security-Onion-Solutions/mreeves/soai-context-1m
Raise SOAI Sonnet default small context limit to 1M
2026-09-29 12:19:20 -04:00
Jason Ertel e4625cfcae Merge pull request #16267 from Security-Onion-Solutions/jertel/wip
resolve startup errors
2026-09-29 12:07:26 -04:00
Jason Ertel eb803dce0e resolve startup errors 2026-09-29 12:02:48 -04:00
Mike Reeves a06f08217a Raise SOAI Sonnet default small context limit to 1M
Context is now flat-priced, so match contextLimitSmall to contextLimitLarge.
With equal limits the SOC assistant hides the increase-context toggle.
2026-09-29 11:42:29 -04:00
Josh Patterson 29d27cf255 Merge pull request #16261 from Security-Onion-Solutions/fix/telegraf-drop-docker-socket
FIX: remove the docker socket from so-telegraf
2026-09-29 10:46:54 -04:00
Jason Ertel 235a60e587 Merge pull request #16264 from Security-Onion-Solutions/jertel/wip
Metric alarms and more NTF annotations
2026-09-29 08:31:28 -04:00
Jason Ertel a8f7c46b0d Merge branch '3/dev' into jertel/wip 2026-09-28 13:47:19 -04:00
Jason Ertel 26d895ccb7 alarms and ntf 2026-09-28 13:47:16 -04:00
Josh Patterson b43efc458f Merge pull request #16263 from Security-Onion-Solutions/fix/service-account-nologin
FIX: use /sbin/nologin for service accounts
2026-09-28 13:05:26 -04:00
Josh Patterson 21222ff119 FIX: use /sbin/nologin for service accounts
These accounts existed only for container UID mapping and filesystem
ownership, but user.present omitted shell:, so Salt fell through to the
platform useradd default and every one of them got /bin/bash. Pin them to
/sbin/nologin so none can be used as an interactive login or `su -` target.

socore keeps /bin/bash: `su socore -c '/usr/sbin/so-repo-sync'` in soup and
so-kernel-upgrade execs the account's passwd shell, and operator docs tell
users to su to socore. soqemussh keeps /bin/bash as an SSH login account.

elastic-agent, elastic-agent-pr and kafka are included alongside the accounts
named in the issue, being the same class with the same unset shell, so the
default is uniform.

Cron is unaffected: cronie runs jobs via the crontab SHELL (default /bin/sh),
not the passwd shell. suricata is the only account changed here that owns a
crontab, and somon has shipped as nologin with a working cron job already.
The zeek `runuser -l zeek` calls all run inside so-zeek via docker.run/exec,
so they resolve the shell from the image, not the host.

Verified on a 3.4.0 managersearch + sensor grid: highstate converges with the
shell as the only change and no failures, is idempotent on a second run, all
containers stay up, SOC still issues a Kratos login flow, and the suricata
surilogcompress cron job runs post-change ((suricata) CMD/CMDEND in
/var/log/cron) while `su - suricata` is now refused.

Closes #16256
2026-09-25 09:23:09 -04:00
Mike Reeves 88fa7e7fb4 Merge pull request #16262 from Security-Onion-Solutions/TOoSmOotH-patch-4
Add openai_embeddings to the YAML configuration
2026-09-24 15:29:17 -04:00
Mike Reeves efe0581892 Add openai_embeddings to the YAML configuration 2026-09-24 15:27:51 -04:00
19 changed files with 676 additions and 17 deletions

No files matched your search

+2
View File
@@ -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
+1
View File
@@ -177,6 +177,7 @@ if [[ $EXCLUDE_FALSE_POSITIVE_ERRORS == 'Y' ]]; then
EXCLUDED_ERRORS="$EXCLUDED_ERRORS|Unexpected authorization header" # expected WARN log lines indicating invalid auth header
EXCLUDED_ERRORS="$EXCLUDED_ERRORS|Missing ory_kratos_session cookie" # expected WARN log lines indicating invalid auth header
EXCLUDED_ERRORS="$EXCLUDED_ERRORS|Static assets preprocessor only supports GET and HEAD requests" # expected WARN log lines indicating invalid auth header
EXCLUDED_ERRORS="$EXCLUDED_ERRORS|respondError" # respondError is a function name, output via http middleware as standard request logging
fi
if [[ $EXCLUDE_KNOWN_ERRORS == 'Y' ]]; then
+1
View File
@@ -21,6 +21,7 @@ elastalert:
- gid: 933
- home: /opt/so/conf/elastalert
- createhome: False
- shell: /sbin/nologin
elastalogdir:
file.directory:
@@ -19,6 +19,7 @@ elastic-agent-pr:
- gid: 948
- home: /opt/so/conf/elastic-fleet-pr
- createhome: False
- shell: /sbin/nologin
{% else %}
+1
View File
@@ -20,6 +20,7 @@ elastic-agent:
- gid: 949
- home: /opt/so/conf/elastic-agent
- createhome: False
- shell: /sbin/nologin
elasticagentconfdir:
file.directory:
+1
View File
@@ -26,6 +26,7 @@ elastic-fleet:
- gid: 947
- home: /opt/so/conf/elastic-fleet
- createhome: False
- shell: /sbin/nologin
elasticfleet_sbin:
file.recurse:
+1
View File
@@ -32,6 +32,7 @@ elasticsearch:
- gid: 930
- home: /opt/so/conf/elasticsearch
- createhome: False
- shell: /sbin/nologin
elasticsearch_sbin:
file.recurse:
+1
View File
@@ -21,6 +21,7 @@ kafka_user:
- gid: 960
- home: /opt/so/conf/kafka
- createhome: False
- shell: /sbin/nologin
kafka_home_dir:
file.absent:
+1
View File
@@ -22,6 +22,7 @@ kibana:
- gid: 932
- home: /opt/so/conf/kibana
- createhome: False
- shell: /sbin/nologin
# Drop the correct nginx config based on role
+1
View File
@@ -27,6 +27,7 @@ kratos:
- uid: 928
- gid: 928
- home: /opt/so/conf/kratos
- shell: /sbin/nologin
kratosdir:
file.directory:
+1
View File
@@ -35,6 +35,7 @@ logstash:
- uid: 931
- gid: 931
- home: /opt/so/conf/logstash
- shell: /sbin/nologin
logstash_sbin:
file.recurse:
+143 -6
View File
@@ -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,130 @@ 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)
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 _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')), _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))
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 +271,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 +339,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,400 @@
# 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_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_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()
+4 -4
View File
@@ -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 %}
+6 -5
View File
@@ -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 {}
+10 -2
View File
@@ -1496,7 +1496,7 @@ soc:
verifyCert: false
notification:
dismissedPruneDays: 30
enabled: false
enabled: true
playbook:
autoUpdateEnabled: true
playbookImportFrequencySeconds: 86400
@@ -1561,6 +1561,14 @@ soc:
reconcilePersona: ""
toolUseTurnAttempts: 12
toolUseTurnDelayMs: 175
agentSessionMaxTurns: 20
agentStreamFlushIntervalMs: 1000
agentStreamIdleTimeoutSeconds: 300
automationSettings:
tickIntervalSeconds: 60
maxConcurrentItems: 4
maxQueuedItems: 0
alertTriageEpoch: "2026-09-24T00:00:00Z"
tools:
filterEventFields:
- "@timestamp"
@@ -2799,7 +2807,7 @@ soc:
- id: sonnet
displayName: Claude Sonnet
origin: USA
contextLimitSmall: 200000
contextLimitSmall: 1000000
contextLimitLarge: 1000000
lowBalanceColorAlert: 500000
enabled: true
+99
View File
@@ -160,6 +160,7 @@ soc:
description: Schedules that are shared across the Security Onion product. Modify via one of the SOC Schedules view.
readonlyUi: True
global: True
advanced: True
forcedType: string
syntax: json
storage: db
@@ -495,6 +496,7 @@ soc:
description: JSON list of notifications. Modify via the SOC Notifications view.
readonlyUi: True
global: True
advanced: True
forcedType: string
syntax: json
storage: db
@@ -503,6 +505,10 @@ soc:
description: The number of days to retain dismissed notifications. When a notification is dismissed, it will be pruned after this many days. Only one user need dismiss a notification for it to be pruned.
forcedType: int
global: True
maxListLimit:
description: Maximum number of notifications to display.
forcedType: int
global: True
enabled:
description: Enables or disables the SOC notification module.
forcedType: bool
@@ -533,6 +539,48 @@ soc:
global: True
sensitive: True
advanced: True
postgresmetrics:
host:
description: Hostname or IP address of the PostgreSQL server used by Telegraf. Defaults to the manager hostname.
global: True
advanced: True
port:
description: Port of the PostgreSQL server used by Telegraf.
global: True
advanced: True
sslMode:
description: "Use encrypted connections to the PostgreSQL server used by Telegraf. Must be one of the following values: disable, allow, prefer, require, verify-ca, verify-full."
global: True
advanced: True
database:
description: Database to authenticate to on the PostgreSQL server.
global: True
advanced: True
user:
description: Username to authenticate to the PostgreSQL server used by Telegraf.
global: True
advanced: True
password:
description: Password used to authenticate to the PostgreSQL server used by Telegraf.
global: True
sensitive: True
advanced: True
cacheExpirationMs:
description: The interval (in milliseconds) to wait before querying the DB for updated metrics.
global: True
advanced: True
maxMetricAgeSeconds:
description: The maximum age (in seconds) of metrics to display in the SOC Grid Metrics view. Metrics older than this value will not be displayed.
global: True
advanced: True
alarms:
description: JSON list of metric alarms. Modify via the SOC Grid Alarms view.
readonlyUi: True
advanced: True
global: True
forcedType: string
syntax: json
storage: db
salt:
longRelayTimeoutMs:
description: Duration (in milliseconds) to wait for a response from the Salt API when executing tasks known for being long running before giving up and showing an error on the SOC UI.
@@ -791,6 +839,7 @@ soc:
- gemini
- openai_responses
- openai_chat
- openai_embeddings
- field: apiUrl
label: API URL
required: False
@@ -812,6 +861,16 @@ soc:
description: Indicates if the Assistant Module should operate in agentic mode or not. If true, agents can work together to solve tasks.
global: True
forcedType: bool
automations:
template:
description: Scheduled automations for the Onion AI assistant, managed from the Agent Studio. Each automation is stored under its own generated ID so its history and rollback are independent of every other automation.
global: True
advanced: False
readonlyUi: True
duplicates: True
forcedType: string
syntax: json
helpLink: onion-ai
agents:
description: Agent definitions for the Onion AI assistant, managed from the Agent Studio. An entry naming a system agent overrides only the fields an admin may change; everything else comes from the built-in definition.
global: True
@@ -844,6 +903,9 @@ soc:
- field: persona
label: Persona
multiline: True
- field: maxConcurrentInstances
label: Max Concurrent Instances
forcedType: int
skills:
description: Skill definitions for the Onion AI assistant, managed from the Agent Studio. An entry naming a system skill overrides only its enabled state and persona addendum; its tool set comes from the built-in definition.
global: True
@@ -947,6 +1009,43 @@ soc:
description: The number of times to retry extracting memories from a session if errors occur.
global: True
advanced: True
agentSessionMaxTurns:
description: Maximum number of model turns a headless agent session, such as one started by an automation, may take before it is stopped. Turns taken by delegated sub-agents count toward this limit. A session that reaches the limit is recorded as failed.
global: True
advanced: True
forcedType: int
agentStreamFlushIntervalMs:
description: Milliseconds between writes of a streaming headless agent turn to the database. Lower values show progress sooner in the Agent Studio at the cost of more frequent Elasticsearch updates.
global: True
advanced: True
forcedType: int
agentStreamIdleTimeoutSeconds:
description: Seconds a streaming headless agent turn may go without receiving any output before it is abandoned and the session is recorded as failed. Set to 0 to disable the timeout.
global: True
advanced: True
forcedType: int
automationSettings:
tickIntervalSeconds:
description: How often, in seconds, the automation scheduler checks for automations that are due to run. Must be greater than 0.
global: True
advanced: True
forcedType: int
maxConcurrentItems:
description: Maximum number of automation work items that may run at the same time. Additional work items wait in the queue until a running item finishes. User chat sessions count toward this limit but are never held back by it. Set to 0 to disable the limit.
global: True
advanced: True
forcedType: int
maxQueuedItems:
description: Maximum number of automation work items that may wait to start. Once the queue is full, no new work items are created until the backlog drains. Set to 0 to disable the limit.
global: True
advanced: True
forcedType: int
alertTriageEpoch:
description: The earliest alert time the Alert Triage automation will consider. Alerts before this time are never triaged, which keeps a first run on an existing deployment from working through old history. Must be in UTC format (2026-09-24T00:00:00Z).
regex: '^(\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}(\.\d+)?Z)?$'
regexFailureMessage: Expecting date in RFC3339 format (2026-09-24T00:00:00Z)
global: True
advanced: True
tools:
filterEventFields:
description: A whitelist of fields to return when OnionAI uses the query_events tool. All other fields are removed. One field per line.
+1
View File
@@ -64,6 +64,7 @@ suricata:
- gid: 940
- home: /nsm/suricata
- createhome: False
- shell: /sbin/nologin
socoregroupwithsuricata:
group.present:
+1
View File
@@ -23,6 +23,7 @@ zeek:
- gid: 937
- home: /opt/so/conf/zeek
- createhome: False
- shell: /sbin/nologin
# Create some directories
zeekpolicydir: