Compare commits

..
Author SHA1 Message Date
Josh Patterson ba95b9bbc2 Address review feedback on so-push-drainer result tracking
Result checks walked dispatch records oldest-first with a cap of five
lookups per pass, counting records whose push was still running. Five
long-running pushes therefore used every slot on every 15s pass and newer,
finished pushes were not reported until one cleared. Check the least
recently checked records first and back off on pushes that are still
running (30s for the first two minutes, then age/4 up to 5 minutes),
recording checked_at in the dispatch record.

Catch any exception when writing a dispatch record so a failed write
cannot skip intent cleanup and re-dispatch the same intents every pass.
Log both output streams when no jid is found, and stop logging a traceback
when a record has already been removed.

Scope the test's salt mock to the drainer import. Run from the repo root,
'salt' resolves to this repo's salt/ directory as a namespace package, so
setdefault left it in place and test_load_push_cfg failed.

Verified on a 3.4.0 managersearch + sensor: a pushed highstate with soc
and telegraf pushes dispatched into it all reported success, with 25
result lookups across the three pushes instead of one per record per pass.
2026-10-02 10:35:30 -04:00
Josh Patterson 8de8ba811a 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.
2026-10-01 09:06:28 -04:00
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
34 changed files with 737 additions and 1613 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 -1
View File
@@ -241,7 +241,7 @@ if [[ $EXCLUDE_KNOWN_ERRORS == 'Y' ]]; then
EXCLUDED_ERRORS="$EXCLUDED_ERRORS|marked for removal" # docker container getting recycled
EXCLUDED_ERRORS="$EXCLUDED_ERRORS|tcp 127.0.0.1:6791: bind: address already in use" # so-elastic-fleet agent restarting. Seen starting w/ 8.18.8 https://github.com/elastic/kibana/issues/201459
EXCLUDED_ERRORS="$EXCLUDED_ERRORS|TransformTask\] \[logs-.*user so_kibana lacks the required permissions" # Known issue with integrations starting transform jobs that are explicitly not allowed to start as a system user
EXCLUDED_ERRORS="$EXCLUDED_ERRORS|manifest unknown" # so-dockerregistry logs a tag lookup miss during image copy; not tied to one docker version
EXCLUDED_ERRORS="$EXCLUDED_ERRORS|manifest unknown" # appears in so-dockerregistry log for so-tcpreplay following docker upgrade to 29.2.1-1
EXCLUDED_ERRORS="$EXCLUDED_ERRORS|Could not index event to Elasticsearch.*\"version\" => \"9.0.8\"" # Expected during Elastic upgrade temporarily, as policies referencing older pipelines are updated
fi
+4 -4
View File
@@ -18,10 +18,10 @@ dockergroup:
dockerheldpackages:
pkg.installed:
- pkgs:
- containerd.io: 2.3.6-1.el9
- docker-ce: 3:29.8.1-1.el9
- docker-ce-cli: 1:29.8.1-1.el9
- docker-ce-rootless-extras: 29.8.1-1.el9
- containerd.io: 2.2.1-1.el9
- docker-ce: 3:29.2.1-1.el9
- docker-ce-cli: 1:29.2.1-1.el9
- docker-ce-rootless-extras: 29.2.1-1.el9
- hold: True
- update_holds: True
+1 -1
View File
@@ -10,7 +10,7 @@ elastalert:
buffer_time:
minutes: 10
old_query_limit:
minutes: 1440
minutes: 5
es_port: 9200
es_conn_timeout: 55
max_query_size: 5000
@@ -6,14 +6,9 @@
# Elastic License 2.0.
from datetime import datetime
from time import gmtime, strftime
import hashlib
import ipaddress
import re
import requests,json
from elastalert.alerts import Alerter, DateTimeEncoder
from elastalert.util import EAException, elastalert_logger
from elastalert.alerts import Alerter
import urllib3
urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning)
@@ -24,152 +19,17 @@ class SecurityOnionESAlerter(Alerter):
"""
required_options = set(['detection_title', 'sigma_level'])
optional_fields = ['sigma_category', 'sigma_product', 'sigma_service', 'sigma_correlation']
count_labels = {
'event_count': '%count% events',
'value_count': '%count% distinct values',
'event_type_count': '%count% correlated rules matched',
'value_sum': 'total %count%',
'value_avg': 'average %count%',
'value_percentile': 'percentile %count%',
'value_median': 'median %count%',
}
placeholder = re.compile(r'%([^%\s]+)%')
# group-by fields copied into ECS related.*
related_users = {'user.name', 'winlog.event_data.TargetUserName', 'winlog.event_data.SubjectUserName'}
related_hosts = {'host.name', 'host.hostname', 'winlog.computer_name'}
@staticmethod
def lookup(doc, dotted):
""" Resolve a dotted path; ES|QL columns arrive nested. """
node = doc
for part in dotted.split('.'):
if not isinstance(node, dict) or part not in node:
return None
node = node[part]
return node
def query_keys(self):
""" compound_query_key holds the list; query_key is flattened to a string. """
return self.rule.get('compound_query_key') or ([self.rule['query_key']] if self.rule.get('query_key') else [])
@staticmethod
def is_correlation(match):
""" Only correlation rows carry window_start, from the ES|QL stats. """
return 'window_start' in match
def alert_id(self, match):
""" Stable id: window end + group values for correlations, source _id otherwise. """
if self.is_correlation(match):
# ungrouped rows have a hashed _id that changes with the count
values = ''.join(f"|{self.lookup(match, k)}" for k in self.query_keys())
key = f"{self.rule['detection_public_id']}|{self.to_dt(match['@timestamp']).isoformat()}{values}"
else:
key = f"{self.rule['detection_public_id']}|{match.get('_id')}"
return hashlib.sha256(key.encode('utf-8')).hexdigest()
def group(self, match):
""" Group-by values joined as ElastAlert joins them for the realert silence key. """
return ', '.join(str(self.lookup(match, k)) for k in self.query_keys())
def related_bucket(self, key):
if key == 'ip' or key.endswith('.ip'):
return 'ip'
if key in self.related_users or key.endswith('.user.name'):
return 'user'
if key in self.related_hosts:
return 'hosts'
return None
@staticmethod
def valid_ip(value):
try:
ipaddress.ip_address(value)
return True
except ValueError:
return False
def related(self, match):
""" ECS related.* from group-by values; invalid IPs are skipped, as they fail the ip mapping. """
related = {}
for key in self.query_keys():
bucket = self.related_bucket(key)
if not bucket:
continue
value = self.lookup(match, key)
for v in value if isinstance(value, list) else [value]:
if v is None or (bucket == 'ip' and not self.valid_ip(str(v))):
continue
related.setdefault(bucket, {})[str(v)] = None
return {bucket: list(values) for bucket, values in related.items()}
def event_data(self, match):
""" The match without the compound query_key field, which ES would map by its last part (e.g. .ip). """
if not self.rule.get('compound_query_key'):
return match
return {k: v for k, v in match.items() if k != self.rule['query_key']}
@staticmethod
def format_value(value):
if isinstance(value, list):
shown = ', '.join(str(v) for v in value[:3])
return shown if len(value) <= 3 else f"{shown} and {len(value) - 3} more"
return str(value)
@staticmethod
def format_count(value):
if isinstance(value, float) and not value.is_integer():
return f"{value:,.2f}"
if isinstance(value, (int, float)):
return f"{int(value):,}"
return str(value)
@staticmethod
def format_duration(seconds):
for unit, size in (('hour', 3600), ('minute', 60)):
if seconds >= 2 * size:
return f"{seconds // size} {unit}s"
return f"{seconds} second{'' if seconds == 1 else 's'}"
@staticmethod
def to_dt(value):
# ES|QL gives ISO strings; ElastAlert parses @timestamp, except on a retried alert.
return value if isinstance(value, datetime) else datetime.fromisoformat(value)
def summary(self, match):
""" One-line correlation summary; None for single-event rules. """
if not self.is_correlation(match):
return None
start = self.to_dt(match['window_start'])
end = self.to_dt(match['@timestamp'])
name = next((f for f in self.count_labels if f in match), None)
values = {
'count': self.format_count(match.get(name)),
'start': start.strftime('%Y-%m-%d %H:%M:%S UTC'),
'end': end.strftime('%Y-%m-%d %H:%M:%S UTC'),
'duration': self.format_duration(int((end - start).total_seconds())),
}
template = self.rule.get('summary_template')
if not template:
groups = ', '.join(f"{k} %{k}%" for k in self.query_keys())
template = f"{self.count_labels.get(name, '%count%')}{' for ' + groups if groups else ''} in %duration%"
def fill(m):
if m[1] in values:
return values[m[1]]
value = self.lookup(match, m[1])
# unknown placeholders stay visible so typos show
return m[0] if value is None else self.format_value(value)
return self.placeholder.sub(fill, template)
optional_fields = ['sigma_category', 'sigma_product', 'sigma_service']
def alert(self, matches):
for match in matches:
timestamp = strftime("%Y-%m-%d"'T'"%H:%M:%S"'.000Z', gmtime())
headers = {"Content-Type": "application/json"}
creds = None
if 'es_username' in self.rule and 'es_password' in self.rule:
creds = (self.rule['es_username'], self.rule['es_password'])
# Start building the rule dict
rule_info = {
"name": self.rule['detection_title'],
@@ -182,74 +42,22 @@ class SecurityOnionESAlerter(Alerter):
if field in self.rule:
rule_info[rule_key] = self.rule[field]
event_info = {
"kind": "alert",
"severity": self.rule['event.severity'],
"module": self.rule['event.module'],
"dataset": self.rule['event.dataset'],
"severity_label": self.rule['sigma_level']
}
reason = self.summary(match)
if reason:
event_info["reason"] = reason
# Construct the payload with the conditional rule_info
payload = {
"tags": ["alert"],
"tags": "alert",
"rule": rule_info,
"event": event_info,
"event": {
"severity": self.rule['event.severity'],
"module": self.rule['event.module'],
"dataset": self.rule['event.dataset'],
"severity_label": self.rule['sigma_level']
},
"sigma_level": self.rule['sigma_level'],
"event_data": self.event_data(match),
"event_data": match,
"@timestamp": timestamp
}
keys = self.query_keys()
if keys and self.is_correlation(match):
payload["labels"] = {
"correlation_group_by": ', '.join(keys),
"correlation_group": self.group(match),
}
related = self.related(match)
if related:
payload["related"] = related
alert_id = self.alert_id(match)
# _create returns 409 on a repeat id; EAException makes ElastAlert retry
url = (f"https://{self.rule['es_host']}:{self.rule['es_port']}"
f"/logs-detections.alerts-so/_create/{alert_id}")
response = self.put_alert(url, payload)
if response.status_code == 400:
# mapping rejections come from event_data; retry with it as unindexed text
rejection = response.text[:500]
payload = self.without_event_data(payload, rejection)
response = self.put_alert(url, payload)
if response.status_code == 400:
elastalert_logger.error("Dropping alert %s for rule %s, rejected by Elasticsearch even without its event data: %s; first rejection: %s",
alert_id, self.rule['detection_public_id'], response.text[:500], rejection)
continue
elastalert_logger.warning("Stored alert %s for rule %s with its event data as text, rejected by Elasticsearch: %s",
alert_id, self.rule['detection_public_id'], rejection)
if response.status_code != 409 and not response.ok:
raise EAException(f"Unable to write alert: {response.status_code} {response.text[:500]}")
def put_alert(self, url, payload):
creds = None
if 'es_username' in self.rule and 'es_password' in self.rule:
creds = (self.rule['es_username'], self.rule['es_password'])
try:
return requests.put(url, data=json.dumps(payload, cls=DateTimeEncoder),
headers={"Content-Type": "application/json"}, verify=False, auth=creds)
except requests.RequestException as e:
raise EAException(f"Unable to write alert: {e}")
@staticmethod
def without_event_data(payload, rejection):
""" event_data moved to event.original; the tag keeps Fleet's final pipeline from removing it. """
fallback = {k: v for k, v in payload.items() if k != 'event_data'}
fallback['event'] = dict(payload['event'], original=json.dumps(payload['event_data'], cls=DateTimeEncoder))
fallback['error'] = {'message': f"event_data rejected by Elasticsearch: {rejection}"}
fallback['tags'] = payload['tags'] + ['preserve_original_event']
return fallback
url = f"https://{self.rule['es_host']}:{self.rule['es_port']}/logs-detections.alerts-so/_doc/"
requests.post(url, data=json.dumps(payload), headers=headers, verify=False, auth=creds)
def get_info(self):
return {'type': 'SecurityOnionESAlerter'}
@@ -1,211 +0,0 @@
# Copyright Security Onion Solutions LLC and/or licensed to Security Onion Solutions LLC under one
# or more contributor license agreements. Licensed under the Elastic License 2.0 as shown at
# https://securityonion.net/license; you may not use this file except in compliance with the
# Elastic License 2.0.
import copy
from datetime import datetime, timezone
import importlib.util
import json
import logging
import os
import sys
import types
import unittest
from unittest.mock import MagicMock, patch
# Real ElastAlert when installed (so-elastalert); otherwise stand-ins for what the alerter imports.
try:
import elastalert.alerts # noqa: F401
except ImportError:
class Alerter:
def __init__(self, rule):
self.rule = rule
class DateTimeEncoder(json.JSONEncoder):
def default(self, obj):
return obj.isoformat() if hasattr(obj, 'isoformat') else json.JSONEncoder.default(self, obj)
class EAException(Exception):
pass
alerts = types.ModuleType('elastalert.alerts')
alerts.Alerter = Alerter
alerts.DateTimeEncoder = DateTimeEncoder
util = types.ModuleType('elastalert.util')
util.EAException = EAException
util.elastalert_logger = logging.getLogger('elastalert')
package = types.ModuleType('elastalert')
package.alerts = alerts
package.util = util
sys.modules.update({'elastalert': package, 'elastalert.alerts': alerts, 'elastalert.util': util})
spec = importlib.util.spec_from_file_location('securityonion_es', os.path.join(os.path.dirname(__file__), 'securityonion-es.py'))
es = importlib.util.module_from_spec(spec)
spec.loader.exec_module(es)
BASE_RULE = {
'name': 'Many Failed Network Logons To One Host From One Source -- 35a42db6-6629-45af-b8aa-e1fa33c28ef5',
'detection_title': 'Many Failed Network Logons To One Host From One Source',
'detection_public_id': '35a42db6-6629-45af-b8aa-e1fa33c28ef5',
'sigma_level': 'medium',
'sigma_correlation': 'event_count',
'event.severity': 3,
'event.module': 'sigma',
'event.dataset': 'sigma.alert',
'es_host': 'manager',
'es_port': 9200,
'summary_template': '%count% failed network logons to %host.name% from %source.ip% in %duration%',
}
def correlation_match():
return {
'event_count': 3561,
'window_start': '2026-09-30T18:05:10+00:00',
'@timestamp': '2026-09-30T18:07:53+00:00',
'host': {'name': 'sa-delta-02-jb'},
'source': {'ip': '192.168.198.149'},
'_id': '6d1c',
'num_hits': 1,
'num_matches': 1,
}
class TestSecurityOnionESAlerter(unittest.TestCase):
def send(self, rule, match):
""" Run alert() and return the payload it wrote and the URL it wrote to. """
alerter = es.SecurityOnionESAlerter(rule)
response = MagicMock(status_code=201, ok=True)
with patch.object(es.requests, 'put', return_value=response) as put:
alerter.alert([match])
self.assertEqual(put.call_count, 1)
return json.loads(put.call_args.kwargs['data']), put.call_args.args[0]
def test_compound_query_key_left_out_of_event_data(self):
rule = dict(BASE_RULE, compound_query_key=['host.name', 'source.ip'], query_key='host.name,source.ip')
match = correlation_match()
match['host.name,source.ip'] = 'sa-delta-02-jb, 192.168.198.149'
original = copy.deepcopy(match)
payload, _ = self.send(rule, match)
self.assertNotIn('host.name,source.ip', payload['event_data'])
self.assertEqual(payload['event_data']['host'], {'name': 'sa-delta-02-jb'})
self.assertEqual(payload['event_data']['source'], {'ip': '192.168.198.149'})
self.assertEqual(payload['labels'], {'correlation_group_by': 'host.name, source.ip', 'correlation_group': 'sa-delta-02-jb, 192.168.198.149'})
self.assertEqual(payload['related'], {'hosts': ['sa-delta-02-jb'], 'ip': ['192.168.198.149']})
self.assertEqual(payload['event']['kind'], 'alert')
self.assertEqual(payload['event']['reason'], '3,561 failed network logons to sa-delta-02-jb from 192.168.198.149 in 2 minutes')
# ElastAlert reuses the match
self.assertEqual(match, original)
def test_single_query_key(self):
rule = dict(BASE_RULE, query_key='user.name', summary_template='%count% failed SOC logins for %user.name%')
match = {'event_count': 3, 'window_start': '2026-09-30T16:49:52+00:00', '@timestamp': '2026-09-30T16:50:00+00:00',
'user': {'name': 'josh@local.invalid'}}
payload, _ = self.send(rule, match)
self.assertEqual(payload['labels'], {'correlation_group_by': 'user.name', 'correlation_group': 'josh@local.invalid'})
self.assertEqual(payload['related'], {'user': ['josh@local.invalid']})
self.assertEqual(payload['event']['reason'], '3 failed SOC logins for josh@local.invalid')
self.assertEqual(payload['event_data'], match)
def test_related_buckets(self):
rule = dict(BASE_RULE, compound_query_key=['winlog.event_data.TargetUserName', 'source.ip', 'dns.highest_registered_domain'],
query_key='winlog.event_data.TargetUserName,source.ip,dns.highest_registered_domain')
match = correlation_match()
match.update({'winlog': {'event_data': {'TargetUserName': ['admmig', 'svc', 'admmig']}},
'source': {'ip': 'not-an-ip'}, 'dns': {'highest_registered_domain': 'example.com'}})
payload, _ = self.send(rule, match)
# deduped; invalid IP skipped; domain stays in the group only
self.assertEqual(payload['related'], {'user': ['admmig', 'svc']})
self.assertEqual(payload['labels']['correlation_group'], "['admmig', 'svc', 'admmig'], not-an-ip, example.com")
payload, _ = self.send(dict(BASE_RULE, query_key='dns.highest_registered_domain'), match)
self.assertNotIn('related', payload)
def test_plain_rule_has_no_correlation_fields(self):
# a query_key alone, as from an override, does not make a correlation
rule = dict(BASE_RULE, query_key='user.name')
# ElastAlert parses @timestamp for EQL hits
match = {'@timestamp': datetime(2026, 9, 30, 16, 50, tzinfo=timezone.utc), '_id': 'abc',
'process': {'name': 'whoami.exe'}, 'user': {'name': 'josh'}}
payload, url = self.send(rule, match)
alerter = es.SecurityOnionESAlerter(rule)
self.assertNotIn('labels', payload)
self.assertNotIn('related', payload)
self.assertNotIn('reason', payload['event'])
self.assertEqual(payload['event']['kind'], 'alert')
self.assertEqual(payload['event_data'], dict(match, **{'@timestamp': '2026-09-30T16:50:00+00:00'}))
self.assertTrue(url.endswith('/' + alerter.alert_id(match)))
self.assertNotEqual(alerter.alert_id(match), alerter.alert_id(dict(match, _id='abd')))
def test_ungrouped_correlation_id_ignores_row_hash(self):
rule = dict(BASE_RULE, summary_template=None)
first = {k: v for k, v in correlation_match().items() if k not in ('host', 'source')}
# ES|QL hashes the row into _id, so a later count changes it
later = dict(first, event_count=3600, _id='9f2a')
payload, url = self.send(rule, first)
alerter = es.SecurityOnionESAlerter(rule)
self.assertEqual(alerter.alert_id(first), alerter.alert_id(later))
self.assertNotEqual(alerter.alert_id(first), alerter.alert_id(dict(first, **{'@timestamp': '2026-09-30T18:09:00+00:00'})))
self.assertTrue(url.endswith('/' + alerter.alert_id(first)))
self.assertEqual(payload['event']['reason'], '3,561 events in 2 minutes')
self.assertNotIn('labels', payload)
def test_grouped_correlation_id_unchanged(self):
""" Ids of alerts already written must not change. """
rule = dict(BASE_RULE, compound_query_key=['host.name', 'source.ip'], query_key='host.name,source.ip')
key = f"{BASE_RULE['detection_public_id']}|2026-09-30T18:07:53+00:00|sa-delta-02-jb|192.168.198.149"
self.assertEqual(es.SecurityOnionESAlerter(rule).alert_id(correlation_match()), es.hashlib.sha256(key.encode()).hexdigest())
def send_responses(self, rule, match, responses):
""" Run alert() against a sequence of write responses; return the payloads written. """
alerter = es.SecurityOnionESAlerter(rule)
with patch.object(es.requests, 'put', side_effect=responses) as put:
alerter.alert([match])
return [json.loads(c.kwargs['data']) for c in put.call_args_list]
def test_rejected_event_data_is_kept_as_text(self):
rule = dict(BASE_RULE, query_key='user.name', summary_template='%count% failed SOC logins for %user.name%')
match = {'event_count': 3, 'window_start': '2026-09-30T16:49:52+00:00', '@timestamp': '2026-09-30T16:50:00+00:00',
'user': {'name': 'josh@local.invalid'}}
rejected = MagicMock(status_code=400, ok=False, text='{"error":{"type":"document_parsing_exception"}}')
first, second = self.send_responses(rule, match, [rejected, MagicMock(status_code=201, ok=True)])
self.assertIn('event_data', first)
self.assertNotIn('event_data', second)
self.assertEqual(json.loads(second['event']['original']), match)
self.assertEqual(second['tags'], ['alert', 'preserve_original_event'])
self.assertEqual(first['tags'], ['alert'])
self.assertTrue(second['error']['message'].startswith('event_data rejected by Elasticsearch: {"error"'))
# everything else carries over
self.assertEqual({k: v for k, v in second['event'].items() if k != 'original'}, first['event'])
changed = ('event_data', 'event', 'error', 'tags')
self.assertEqual({k: v for k, v in second.items() if k not in changed}, {k: v for k, v in first.items() if k not in changed})
def test_rejected_twice_is_dropped_without_retry(self):
rejected = MagicMock(status_code=400, ok=False, text='{"error":{"type":"document_parsing_exception"}}')
match = {'@timestamp': '2026-09-30T16:50:00+00:00', '_id': 'abc'}
# no EAException, so no retry
payloads = self.send_responses(dict(BASE_RULE), match, [rejected, rejected])
self.assertEqual(len(payloads), 2)
if __name__ == '__main__':
unittest.main()
+1 -1
View File
@@ -120,7 +120,7 @@ elastalert:
helpLink: elastalert
old_query_limit:
minutes:
description: How long ElastAlert can be down, in minutes, and still resume each rule where it stopped. After a longer outage, rules restart from now and skip the gap.
description: Amount of time in minutes between queries to start at the most recently run query.
global: True
helpLink: elastalert
es_conn_timeout:
@@ -29,7 +29,7 @@
"\\.gz$"
],
"include_files": [],
"processors": "- dissect:\n tokenizer: \"/nsm/import/%{import.id}/evtx/%{import.file}\"\n field: \"log.file.path\"\n target_prefix: \"\"\n- decode_json_fields:\n fields: [\"message\"]\n target: \"\"\n- add_fields:\n target: event\n fields:\n dataset: windows.forwarded\n module: windows\n imported: true\n- add_fields:\n target: \"@metadata\"\n fields:\n pipeline: import.evtx\n- if:\n equals:\n winlog.channel: 'Security'\n then: \n - add_fields:\n target: event\n fields:\n dataset: system.security\n module: system\n- if:\n equals:\n winlog.channel: 'Windows PowerShell'\n then: \n - add_fields:\n target: event\n fields:\n dataset: windows.powershell\n module: windows\n- if:\n equals:\n winlog.channel: 'Microsoft-Windows-Sysmon/Operational'\n then: \n - add_fields:\n target: event\n fields:\n dataset: windows.sysmon_operational\n module: windows\n imported: true\n- if:\n equals:\n winlog.channel: 'Application'\n then: \n - add_fields:\n target: event\n fields:\n dataset: system.application\n module: system\n- if:\n equals:\n winlog.channel: 'System'\n then: \n - add_fields:\n target: event\n fields:\n dataset: system.system\n module: system\n \n- if:\n equals:\n winlog.channel: 'Microsoft-Windows-PowerShell/Operational'\n then: \n - add_fields:\n target: event\n fields:\n dataset: windows.powershell_operational\n module: windows\n- add_fields:\n target: data_stream\n fields:\n type: logs\n dataset: import",
"processors": "- dissect:\n tokenizer: \"/nsm/import/%{import.id}/evtx/%{import.file}\"\n field: \"log.file.path\"\n target_prefix: \"\"\n- decode_json_fields:\n fields: [\"message\"]\n target: \"\"\n- drop_fields:\n fields: [\"host\"]\n ignore_missing: true\n- add_fields:\n target: data_stream\n fields:\n type: logs\n dataset: system.security\n- add_fields:\n target: event\n fields:\n dataset: system.security\n module: system\n imported: true\n- add_fields:\n target: \"@metadata\"\n fields:\n pipeline: logs-system.security-2.22.3\n- if:\n equals:\n winlog.channel: 'Microsoft-Windows-Sysmon/Operational'\n then: \n - add_fields:\n target: data_stream\n fields:\n dataset: windows.sysmon_operational\n - add_fields:\n target: event\n fields:\n dataset: windows.sysmon_operational\n module: windows\n imported: true\n - add_fields:\n target: \"@metadata\"\n fields:\n pipeline: logs-windows.sysmon_operational-3.9.0\n- if:\n equals:\n winlog.channel: 'Application'\n then: \n - add_fields:\n target: data_stream\n fields:\n dataset: system.application\n - add_fields:\n target: event\n fields:\n dataset: system.application\n - add_fields:\n target: \"@metadata\"\n fields:\n pipeline: logs-system.application-2.22.3\n- if:\n equals:\n winlog.channel: 'System'\n then: \n - add_fields:\n target: data_stream\n fields:\n dataset: system.system\n - add_fields:\n target: event\n fields:\n dataset: system.system\n - add_fields:\n target: \"@metadata\"\n fields:\n pipeline: logs-system.system-2.22.3\n \n- if:\n equals:\n winlog.channel: 'Microsoft-Windows-PowerShell/Operational'\n then: \n - add_fields:\n target: data_stream\n fields:\n dataset: windows.powershell_operational\n - add_fields:\n target: event\n fields:\n dataset: windows.powershell_operational\n module: windows\n - add_fields:\n target: \"@metadata\"\n fields:\n pipeline: logs-windows.powershell_operational-3.9.0\n- add_fields:\n target: data_stream\n fields:\n dataset: import",
"tags": [
"import"
],
@@ -30,14 +30,17 @@
'azure_metrics.monitor': 'azure.monitor',
'azure_metrics.storage_account': 'azure.storage_account',
'azure_openai.metrics': 'azure.open_ai',
'beat.state': 'beats.stack_monitoring.state',
'beat.stats': 'beats.stack_monitoring.stats',
'enterprisesearch.health': 'enterprisesearch.stack_monitoring.health',
'enterprisesearch.stats': 'enterprisesearch.stack_monitoring.stats',
'kibana.cluster_actions': 'kibana.stack_monitoring.cluster_actions',
'kibana.cluster_rules': 'kibana.stack_monitoring.cluster_rules',
'kibana.node_actions': 'kibana.stack_monitoring.node_actions',
'kibana.node_rules': 'kibana.stack_monitoring.node_rules',
'kibana.stats': 'kibana.stack_monitoring.stats',
'kibana.status': 'kibana.stack_monitoring.status',
'logstash.node': 'logstash.stack_monitoring.node',
'logstash.node_cel': 'logstash.node',
'logstash.node_cel': 'logstash.stack_monitoring.node',
'logstash.node_stats': 'logstash.stack_monitoring.node_stats',
'synthetics.browser': 'synthetics-browser',
'synthetics.browser_network': 'synthetics-browser.network',
-6
View File
@@ -1160,7 +1160,6 @@ elasticsearch:
- so-fleet_agent_id_verification-1
- so-logs-mappings
- so-logs-settings
- detections-alerts-mappings
data_stream:
allow_custom_routing: false
hidden: false
@@ -3310,7 +3309,6 @@ elasticsearch:
composed_of:
- event-mappings
- logs-system.security@package
- so-fleet_system.security_caseless-1
- logs-system.security@custom
- so-fleet_integrations.ip_mappings-1
- so-fleet_globals-1
@@ -4177,7 +4175,6 @@ elasticsearch:
index_template:
composed_of:
- logs-windows.forwarded@package
- so-fleet_process_caseless-1
- logs-windows.forwarded@custom
- so-fleet_integrations.ip_mappings-1
- so-fleet_globals-1
@@ -4227,7 +4224,6 @@ elasticsearch:
index_template:
composed_of:
- logs-windows.powershell@package
- so-fleet_process_caseless-1
- logs-windows.powershell@custom
- so-fleet_integrations.ip_mappings-1
- so-fleet_globals-1
@@ -4277,7 +4273,6 @@ elasticsearch:
index_template:
composed_of:
- logs-windows.powershell_operational@package
- so-fleet_process_caseless-1
- logs-windows.powershell_operational@custom
- so-fleet_integrations.ip_mappings-1
- so-fleet_globals-1
@@ -4327,7 +4322,6 @@ elasticsearch:
index_template:
composed_of:
- logs-windows.sysmon_operational@package
- so-fleet_process_caseless-1
- logs-windows.sysmon_operational@custom
- so-fleet_integrations.ip_mappings-1
- so-fleet_globals-1
@@ -99,7 +99,7 @@
},
{
"set": {
"if": "ctx.tags != null && ctx.tags.contains('import') && ctx._index != null && ctx._index.startsWith('logs-import-')",
"if": "ctx.tags != null && ctx.tags.contains('import')",
"override": true,
"field": "data_stream.dataset",
"value": "import"
@@ -107,7 +107,7 @@
},
{
"set": {
"if": "ctx.tags != null && ctx.tags.contains('import') && ctx._index != null && ctx._index.startsWith('logs-import-')",
"if": "ctx.tags != null && ctx.tags.contains('import')",
"override": true,
"field": "data_stream.namespace",
"value": "so"
@@ -1,31 +0,0 @@
{
"description" : "import.evtx: normalize imported EVTX and reroute to logs-<dataset>-import",
"processors" : [
{ "script": {
"description": "Host from the event, not the importing node",
"lang": "painless",
"source": "Map host = ['os': ['type': 'windows', 'family': 'windows', 'platform': 'windows']]; def cn = ctx.winlog?.computer_name; if (cn != null && cn.toString().length() > 0) { String name = cn.toString(); int dot = name.indexOf('.'); if (dot > 0) { name = name.substring(0, dot); } host.put('hostname', name); host.put('name', name.toLowerCase()); } ctx.host = host;"
} },
{ "script": {
"description": "String event IDs, as Winlogbeat sends",
"lang": "painless",
"source": "if (ctx.winlog?.event_id != null) { ctx.winlog.event_id = ctx.winlog.event_id.toString(); } if (ctx.event?.code != null) { ctx.event.code = ctx.event.code.toString(); }"
} },
{ "script": {
"description": "Unnamed <Data> to param1..N, as Winlogbeat",
"lang": "painless",
"if": "ctx.winlog?.event_data?.Data instanceof Map && ctx.winlog.event_data.Data['#text'] != null",
"source": "def t = ctx.winlog.event_data.Data['#text']; List vals = t instanceof List ? t : [t]; for (int i = 0; i < vals.size(); i++) { ctx.winlog.event_data['param' + (i + 1)] = vals.get(i); } ctx.winlog.event_data.remove('Data');"
} },
{ "script": {
"description": "String values and LF line endings, as Winlogbeat",
"lang": "painless",
"if": "ctx.winlog?.event_data instanceof Map || ctx.winlog?.user_data instanceof Map",
"source": "String lf = String.valueOf((char) 10); String crlf = String.valueOf((char) 13) + lf; for (def key : ['event_data', 'user_data']) { def m = ctx.winlog[key]; if (!(m instanceof Map)) { continue; } for (def e : m.entrySet()) { def v = e.getValue(); if (v instanceof String) { e.setValue(v.replace(crlf, lf)); } else if (v instanceof Number || v instanceof Boolean) { e.setValue(v.toString()); } } }"
} },
{ "set": { "description": "event.kind, as Winlogbeat", "field": "event.kind", "value": "event", "override": false } },
{ "set": { "field": "data_stream.dataset", "copy_from": "event.dataset", "override": true, "ignore_empty_value": true } },
{ "set": { "field": "data_stream.namespace", "value": "import", "override": true } },
{ "reroute": { "dataset": "{{data_stream.dataset}}", "namespace": "{{data_stream.namespace}}" } }
]
}
@@ -1,34 +0,0 @@
{
"version": 1,
"_meta": {
"managed_by": "securityonion",
"managed": true
},
"description": "Custom pipeline for the System integration's auth data stream.",
"processors": [
{
"trim": {
"description": "Grok leaves a leading space on 'invalid user' names (elastic/integrations#12174) and, before 2.23.2, sudo padding",
"field": "user.name",
"ignore_missing": true,
"ignore_failure": true
}
},
{
"trim": {
"description": "Appended from the untrimmed user.name",
"field": "related.user",
"ignore_missing": true,
"ignore_failure": true
}
},
{
"script": {
"description": "Dedupe after trimming",
"if": "ctx.related?.user instanceof List",
"source": "ctx.related.user = new ArrayList(new LinkedHashSet(ctx.related.user));",
"ignore_failure": true
}
}
]
}
@@ -1,123 +0,0 @@
{
"_meta": {
"managed_by": "security_onion",
"managed": true,
"description": "Adds .caseless for Lucene queries. Restates each field's package type and .text."
},
"template": {
"mappings": {
"properties": {
"process": {
"properties": {
"executable": {
"type": "keyword",
"ignore_above": 1024,
"fields": {
"caseless": {
"type": "keyword",
"ignore_above": 1024,
"normalizer": "lowercase"
},
"text": {
"type": "match_only_text"
}
}
},
"name": {
"type": "keyword",
"ignore_above": 1024,
"fields": {
"caseless": {
"type": "keyword",
"ignore_above": 1024,
"normalizer": "lowercase"
},
"text": {
"type": "match_only_text"
}
}
},
"command_line": {
"type": "wildcard",
"ignore_above": 1024,
"fields": {
"caseless": {
"type": "keyword",
"ignore_above": 1024,
"normalizer": "lowercase"
},
"text": {
"type": "match_only_text"
}
}
},
"parent": {
"properties": {
"executable": {
"type": "keyword",
"ignore_above": 1024,
"fields": {
"caseless": {
"type": "keyword",
"ignore_above": 1024,
"normalizer": "lowercase"
},
"text": {
"type": "match_only_text"
}
}
},
"name": {
"type": "keyword",
"ignore_above": 1024,
"fields": {
"caseless": {
"type": "keyword",
"ignore_above": 1024,
"normalizer": "lowercase"
},
"text": {
"type": "match_only_text"
}
}
},
"command_line": {
"type": "wildcard",
"ignore_above": 1024,
"fields": {
"caseless": {
"type": "keyword",
"ignore_above": 1024,
"normalizer": "lowercase"
},
"text": {
"type": "match_only_text"
}
}
}
}
}
}
},
"file": {
"properties": {
"path": {
"type": "keyword",
"ignore_above": 1024,
"fields": {
"caseless": {
"type": "keyword",
"ignore_above": 1024,
"normalizer": "lowercase"
},
"text": {
"type": "match_only_text"
}
}
}
}
}
}
}
}
}
@@ -1,80 +0,0 @@
{
"_meta": {
"managed_by": "security_onion",
"managed": true,
"description": "Adds .caseless for Lucene queries. Keeps each field's existing keyword type."
},
"template": {
"mappings": {
"properties": {
"process": {
"properties": {
"command_line": {
"type": "keyword",
"ignore_above": 1024,
"fields": {
"caseless": {
"type": "keyword",
"ignore_above": 1024,
"normalizer": "lowercase"
}
}
},
"parent": {
"properties": {
"executable": {
"type": "keyword",
"ignore_above": 1024,
"fields": {
"caseless": {
"type": "keyword",
"ignore_above": 1024,
"normalizer": "lowercase"
}
}
},
"name": {
"type": "keyword",
"ignore_above": 1024,
"fields": {
"caseless": {
"type": "keyword",
"ignore_above": 1024,
"normalizer": "lowercase"
}
}
},
"command_line": {
"type": "keyword",
"ignore_above": 1024,
"fields": {
"caseless": {
"type": "keyword",
"ignore_above": 1024,
"normalizer": "lowercase"
}
}
}
}
}
}
},
"file": {
"properties": {
"path": {
"type": "keyword",
"ignore_above": 1024,
"fields": {
"caseless": {
"type": "keyword",
"ignore_above": 1024,
"normalizer": "lowercase"
}
}
}
}
}
}
}
}
}
@@ -50,18 +50,6 @@
"ignore_above": 1024,
"type": "keyword"
},
"ruleType": {
"ignore_above": 1024,
"type": "keyword"
},
"correlationType": {
"ignore_above": 1024,
"type": "keyword"
},
"correlationTimespan": {
"ignore_above": 1024,
"type": "keyword"
},
"content": {
"type": "text"
},
@@ -1,149 +0,0 @@
{
"template": {
"mappings": {
"properties": {
"tags": {
"ignore_above": 1024,
"type": "keyword"
},
"sigma_level": {
"ignore_above": 1024,
"type": "keyword"
},
"rule": {
"properties": {
"name": {
"ignore_above": 1024,
"type": "keyword"
},
"uuid": {
"ignore_above": 1024,
"type": "keyword"
},
"category": {
"ignore_above": 1024,
"type": "keyword"
},
"product": {
"ignore_above": 1024,
"type": "keyword"
},
"service": {
"ignore_above": 1024,
"type": "keyword"
},
"correlation": {
"ignore_above": 1024,
"type": "keyword"
}
}
},
"event": {
"properties": {
"severity": {
"type": "long"
},
"severity_label": {
"ignore_above": 1024,
"type": "keyword"
},
"module": {
"ignore_above": 1024,
"type": "keyword"
},
"dataset": {
"ignore_above": 1024,
"type": "keyword"
},
"kind": {
"ignore_above": 1024,
"type": "keyword"
},
"reason": {
"type": "match_only_text",
"fields": {
"keyword": {
"ignore_above": 1024,
"type": "keyword"
}
}
},
"original": {
"type": "keyword",
"index": false,
"doc_values": false
}
}
},
"event_data": {
"properties": {
"@timestamp": {
"type": "date"
},
"window_start": {
"type": "date"
},
"event_count": {
"type": "long"
},
"value_count": {
"type": "long"
},
"event_type_count": {
"type": "long"
},
"value_sum": {
"type": "double"
},
"value_avg": {
"type": "double"
},
"value_percentile": {
"type": "double"
},
"value_median": {
"type": "double"
}
}
},
"labels": {
"properties": {
"correlation_group_by": {
"ignore_above": 1024,
"type": "keyword"
},
"correlation_group": {
"ignore_above": 1024,
"type": "keyword"
}
}
},
"related": {
"properties": {
"ip": {
"type": "ip"
},
"user": {
"ignore_above": 1024,
"type": "keyword"
},
"hosts": {
"ignore_above": 1024,
"type": "keyword"
}
}
},
"error": {
"properties": {
"message": {
"type": "match_only_text"
}
}
}
}
}
},
"_meta": {
"description": "Fields written by the ElastAlert SecurityOnionESAlerter to logs-detections.alerts-so"
}
}
+180 -4
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,19 @@ PENDING_DIR = '/opt/so/state/push_pending'
LOCK_FILE = os.path.join(PENDING_DIR, '.lock')
LOG_FILE = '/opt/so/log/salt/so-push-drainer.log'
DISPATCHED_DIR = '/opt/so/state/push_dispatched'
HIGHSTATE_SENTINEL = '__highstate__'
RESULT_CHECK_DELAY = 30
RESULT_RECHECK_MAX = 300
RESULT_MAX_AGE = 7200
RESULT_CHECKS_PER_PASS = 5
TEXT_LIMIT = 500
# salt-run --async reports the jid only in a log line (stderr by default).
JID_RE = re.compile(r'salt/run/(\d{20})')
def _make_logger():
logger = logging.getLogger('so-push-drainer')
@@ -113,14 +127,167 @@ 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
output = '{}\n{}'.format(result.stderr or '', result.stdout or '')
match = JID_RE.search(output)
if not match:
log.warning('dispatch accepted but no jid found, result will not be tracked: output=%s',
_trim(output))
return ''
log.info('dispatch accepted: jid=%s', match.group(1))
return match.group(1)
def _trim(value):
text = value if isinstance(value, str) else json.dumps(value, default=str)
lines = [line.strip() for line in text.splitlines() if line.strip()]
if 'Traceback (most recent call last):' in text:
# Keep the lead-in and the raised exception; the frames are noise in a log line.
lines = [text.split('Traceback (most recent call last):', 1)[0].strip(), lines[-1]]
text = ' '.join(line for line in lines if line)
return text if len(text) <= TEXT_LIMIT else text[:TEXT_LIMIT] + '...'
def _unlink(path, log):
try:
os.unlink(path)
except FileNotFoundError:
pass
except OSError:
log.exception('failed to remove %s', path)
def _write_record(path, record, log):
try:
os.makedirs(DISPATCHED_DIR, exist_ok=True)
tmp_path = path + '.tmp'
with open(tmp_path, 'w') as f:
json.dump(record, f)
os.rename(tmp_path, path)
except Exception:
log.exception('failed to record dispatch %s', record.get('jid'))
def _record_dispatch(jid, actions, paths, log):
record = {'jid': jid, 'dispatched_at': time.time(), 'actions': actions, 'paths': paths}
_write_record(os.path.join(DISPATCHED_DIR, '{}.json'.format(jid)), record, log)
def _lookup_jid(jid, log):
"""Returns the job cache entry for jid, {} while it is still running, or None on error."""
cmd = ['salt-run', 'jobs.lookup_jid', jid, '--out=json']
try:
result = subprocess.run(cmd, check=True, capture_output=True, text=True, timeout=60)
return json.loads(result.stdout or '{}')
except (subprocess.CalledProcessError, subprocess.TimeoutExpired, ValueError) as exc:
log.warning('lookup of jid %s failed: %s', jid, exc)
return None
def _minion_failure(minion_ret):
if isinstance(minion_ret, dict):
return '; '.join(
'{}: {}'.format(state.get('__id__', state_key), _trim(state.get('comment', '')))
for state_key, state in minion_ret.items()
if isinstance(state, dict) and state.get('result') is False
)
# A state run rejected before it starts (e.g. another state run is in
# progress) returns a list of error strings instead of state results.
if isinstance(minion_ret, (list, str)):
return _trim(minion_ret)
return ''
def _step_failures(step):
if not isinstance(step, dict) or step.get('result') is not False:
return []
failures = ['{}: {}'.format(step.get('__id__', step.get('name')), _trim(step.get('comment', '')))]
changes = step.get('changes')
minion_rets = changes.get('ret') if isinstance(changes, dict) else None
if isinstance(minion_rets, dict):
for minion, minion_ret in minion_rets.items():
text = _minion_failure(minion_ret)
if text:
failures.append('{}: {}'.format(minion, text))
return failures
def _orch_failures(ret):
if not isinstance(ret, dict):
return [_trim(ret)]
failures = []
for job in ret.values():
if not isinstance(job, dict):
continue
job_ret = job.get('return')
data = job_ret.get('data') if isinstance(job_ret, dict) else {}
if not isinstance(data, dict):
if data:
failures.append(_trim(data))
data = {}
for steps in data.values():
if not isinstance(steps, dict):
failures.append(_trim(steps))
continue
for step in steps.values():
failures.extend(_step_failures(step))
if job.get('success') is False and not failures:
failures.append('orchestration reported failure: {}'.format(_trim(job.get('return'))))
return failures
def _recheck_delay(age):
return min(RESULT_RECHECK_MAX, max(RESULT_CHECK_DELAY, age / 4))
def _check_dispatched(log, now):
due = []
for path in glob.glob(os.path.join(DISPATCHED_DIR, '*.json')):
record = _read_intent(path, log)
if not isinstance(record, dict) or not record.get('jid'):
_unlink(path, log)
continue
age = now - record.get('dispatched_at', 0)
last_check = record.get('checked_at', record.get('dispatched_at', 0))
if now - last_check >= _recheck_delay(age):
due.append((last_check, path, record, age))
# Least recently checked first, so pushes that are still running can't starve finished ones.
for _, path, record, age in sorted(due, key=lambda item: item[:2])[:RESULT_CHECKS_PER_PASS]:
jid = record['jid']
try:
if _report_result(record, age, log):
_unlink(path, log)
else:
record['checked_at'] = now
_write_record(path, record, log)
except Exception:
# Drop the record so one unreadable result can't fail every pass ahead of the drain.
log.exception('cannot evaluate result for jid=%s; no longer tracking', jid)
_unlink(path, log)
def _report_result(record, age, log):
"""Logs the outcome of a dispatched push. Returns True once the record is finished with."""
jid = record['jid']
paths = record.get('paths', [])
ret = _lookup_jid(jid, log)
if not ret:
if age > RESULT_MAX_AGE:
log.warning('no result for jid=%s after %ds, no longer tracking; paths=%s', jid, age, paths)
return True
return False
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 +310,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 +378,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,477 @@
# Copyright Security Onion Solutions LLC and/or licensed to Security Onion Solutions LLC under one
# or more contributor license agreements. Licensed under the Elastic License 2.0 as shown at
# https://securityonion.net/license; you may not use this file except in compliance with the
# Elastic License 2.0.
import importlib.util
import json
import logging
import os
import shutil
import subprocess
import sys
import tempfile
import time
import unittest
from importlib.machinery import SourceFileLoader
from unittest.mock import MagicMock, patch
HERE = os.path.dirname(os.path.abspath(__file__))
SCRIPT = os.path.join(HERE, 'so-push-drainer')
_loader = SourceFileLoader('so_push_drainer', SCRIPT)
_spec = importlib.util.spec_from_loader('so_push_drainer', _loader)
drainer = importlib.util.module_from_spec(_spec)
# salt is not installed where these tests run; the drainer only needs salt.client.Caller.
# Mocked only while the drainer loads: run from the repo root, 'salt' is this repo's salt/ directory.
_salt = MagicMock()
with patch.dict(sys.modules, {'salt': _salt, 'salt.client': _salt.client}):
_loader.exec_module(drainer)
MASTER = 'manager.localdomain_master'
JID = '20260930171554259426'
ASYNC_STDERR = ('[WARNING ] Running in asynchronous mode. Results of this execution may be collected '
'by attaching to the master event bus or by examining the master job cache, if '
'configured. This execution is running under tag salt/run/{}\n'.format(JID))
CONFLICT = ('The function "state.sls" is running as PID 372218 and was started at '
'2026, Sep 30 17:15:40.466233 with jid 20260930171540466233')
def _orch_ret(steps, success=True):
return {MASTER: {
'fun': 'runner.state.orchestrate',
'jid': JID,
'return': {'data': {MASTER: steps}, 'outputter': 'highstate', 'retcode': 0 if success else 1},
'success': success,
}}
REFRESH_STEP = {
'salt_|-refresh_pillar_1_|-saltutil.refresh_pillar_|-function': {
'__id__': 'refresh_pillar_1', 'result': True,
'changes': {'ret': {'manager_standalone': True}},
'comment': 'Function ran successfully.',
},
}
CONFLICT_RET = _orch_ret(dict(REFRESH_STEP, **{
'salt_|-apply_soc_1_|-apply_soc_1_|-state': {
'__id__': 'apply_soc_1', 'result': False,
'changes': {'out': 'highstate', 'ret': {'manager_standalone': [CONFLICT]}},
'comment': 'Run failed on minions: manager_standalone',
},
}), success=False)
STATE_FAIL_RET = _orch_ret({
'salt_|-apply_hydra_1_|-apply_hydra_1_|-state': {
'__id__': 'apply_hydra_1', 'result': False,
'changes': {'out': 'highstate', 'ret': {'manager_standalone': {
'test_|-no_license_|-no_license_|-fail_without_changes': {
'__id__': 'hydra.enabled_no_license_detected', 'result': False,
'comment': 'This is a feature supported only for customers with a valid license.',
},
'file_|-hydra_conf_|-/opt/so/conf/hydra_|-managed': {'result': True, 'comment': 'ok'},
}}},
'comment': 'Run failed on minions: manager_standalone',
},
}, success=False)
SUCCESS_RET = _orch_ret(dict(REFRESH_STEP, **{
'salt_|-apply_telegraf_1_|-apply_telegraf_1_|-state': {
'__id__': 'apply_telegraf_1', 'result': True,
'changes': {'out': 'highstate', 'ret': {'manager_standalone': {
'file_|-tgrafconf_|-/opt/so/conf/telegraf/etc/telegraf.conf_|-managed': {'result': True},
}}},
'comment': 'States ran successfully.',
},
}))
class DrainerTestCase(unittest.TestCase):
def setUp(self):
self.tmpdir = tempfile.mkdtemp()
self.pending = os.path.join(self.tmpdir, 'push_pending')
self.dispatched = os.path.join(self.tmpdir, 'push_dispatched')
os.makedirs(self.pending)
for name, value in (
('PENDING_DIR', self.pending),
('LOCK_FILE', os.path.join(self.pending, '.lock')),
('DISPATCHED_DIR', self.dispatched),
('LOG_FILE', os.path.join(self.tmpdir, 'log', 'so-push-drainer.log')),
):
patcher = patch.object(drainer, name, value)
patcher.start()
self.addCleanup(patcher.stop)
self.log = MagicMock()
def tearDown(self):
shutil.rmtree(self.tmpdir, ignore_errors=True)
def write_json(self, directory, name, data):
os.makedirs(directory, exist_ok=True)
path = os.path.join(directory, name)
with open(path, 'w') as f:
if isinstance(data, str):
f.write(data)
else:
json.dump(data, f)
return path
def logged(self, level):
return ' '.join(c.args[0] % c.args[1:] for c in getattr(self.log, level).call_args_list)
class TestHelpers(DrainerTestCase):
def test_make_logger_adds_handler_once(self):
logger = logging.getLogger('so-push-drainer')
def close_handlers():
for handler in logger.handlers:
handler.close()
logger.handlers.clear()
self.addCleanup(close_handlers)
logger.handlers.clear()
self.assertIs(drainer._make_logger(), logger)
drainer._make_logger()
self.assertEqual(len(logger.handlers), 1)
self.assertTrue(os.path.isdir(os.path.dirname(drainer.LOG_FILE)))
def test_load_push_cfg(self):
with patch.object(drainer.salt.client, 'Caller') as caller:
caller.return_value.cmd.return_value = {'enabled': False}
self.assertEqual(drainer._load_push_cfg(), {'enabled': False})
caller.return_value.cmd.return_value = 'garbage'
self.assertEqual(drainer._load_push_cfg(), {})
def test_read_intent(self):
good = self.write_json(self.pending, 'good.json', {'a': 1})
bad = self.write_json(self.pending, 'bad.json', '{nope')
self.assertEqual(drainer._read_intent(good, self.log), {'a': 1})
self.assertIsNone(drainer._read_intent(bad, self.log))
with patch('builtins.open', side_effect=RuntimeError('boom')):
self.assertIsNone(drainer._read_intent(good, self.log))
self.log.exception.assert_called_once()
def test_dedupe_actions(self):
actions = [
'not a dict',
{'state': 'soc'},
{'state': 'soc', 'tgt': '*'},
{'state': 'soc', 'tgt': '*', 'tgt_type': 'compound'},
{'highstate': True, 'tgt': '*'},
{'state': 'soc', 'tgt': 'node1', 'tgt_type': 'glob'},
]
self.assertEqual(drainer._dedupe_actions(actions), [actions[2], actions[4], actions[5]])
def test_trim(self):
self.assertEqual(drainer._trim(' text \n'), 'text')
self.assertEqual(drainer._trim(['a']), '["a"]')
self.assertEqual(drainer._trim(None), 'null')
self.assertEqual(drainer._trim('x' * 600), 'x' * drainer.TEXT_LIMIT + '...')
def test_trim_traceback(self):
comment = ('An exception occurred in this state: Traceback (most recent call last):\n'
' File "salt/client/__init__.py", line 1934, in pub\n'
' raise AuthenticationError(err_msg)\n'
'salt.exceptions.AuthenticationError: Authentication error occurred.\n')
self.assertEqual(drainer._trim(comment), 'An exception occurred in this state: '
'salt.exceptions.AuthenticationError: Authentication error occurred.')
self.assertEqual(drainer._trim('line one\n line two\n'), 'line one line two')
def test_unlink(self):
drainer._unlink(os.path.join(self.tmpdir, 'missing'), self.log)
self.log.exception.assert_not_called()
drainer._unlink(self.tmpdir, self.log)
self.log.exception.assert_called_once()
class TestDispatch(DrainerTestCase):
def run_dispatch(self, **kwargs):
with patch.object(drainer.subprocess, 'run', **kwargs) as run:
jid = drainer._dispatch([{'state': 'soc', 'tgt': '*'}], self.log)
return jid, run
def test_jid_parsed_from_stderr(self):
jid, run = self.run_dispatch(return_value=MagicMock(stdout='', stderr=ASYNC_STDERR))
self.assertEqual(jid, JID)
cmd = run.call_args[0][0]
self.assertEqual(cmd[:3], ['salt-run', 'state.orchestrate', 'orch.push_batch'])
self.assertIn('--async', cmd)
def test_jid_parsed_from_stdout(self):
jid, _ = self.run_dispatch(return_value=MagicMock(stdout=ASYNC_STDERR, stderr=None))
self.assertEqual(jid, JID)
def test_no_jid(self):
jid, _ = self.run_dispatch(return_value=MagicMock(stdout='unexpected output', stderr=None))
self.assertEqual(jid, '')
self.assertIn('output=unexpected output', self.logged('warning'))
def test_failures_return_none(self):
for exc in (subprocess.CalledProcessError(1, 'salt-run', 'out', 'err'),
subprocess.TimeoutExpired('salt-run', 60),
RuntimeError('boom')):
jid, _ = self.run_dispatch(side_effect=exc)
self.assertIsNone(jid)
def test_record_dispatch(self):
drainer._record_dispatch(JID, [{'state': 'soc'}], ['audit:soc.config.licenseKey'], self.log)
with open(os.path.join(self.dispatched, JID + '.json')) as f:
record = json.load(f)
self.assertEqual(record['jid'], JID)
self.assertEqual(record['paths'], ['audit:soc.config.licenseKey'])
self.assertIn('dispatched_at', record)
def test_record_dispatch_errors(self):
with patch.object(drainer.os, 'makedirs', side_effect=OSError('ro')):
drainer._record_dispatch(JID, [], [], self.log)
drainer._record_dispatch(JID, [object()], [], self.log)
self.assertEqual(self.log.exception.call_count, 2)
self.assertFalse(os.path.exists(os.path.join(self.dispatched, JID + '.json')))
class TestResults(DrainerTestCase):
def test_lookup_jid(self):
with patch.object(drainer.subprocess, 'run') as run:
run.return_value = MagicMock(stdout=json.dumps(SUCCESS_RET))
self.assertEqual(drainer._lookup_jid(JID, self.log), SUCCESS_RET)
self.assertEqual(run.call_args[0][0], ['salt-run', 'jobs.lookup_jid', JID, '--out=json'])
run.return_value = MagicMock(stdout='')
self.assertEqual(drainer._lookup_jid(JID, self.log), {})
run.return_value = MagicMock(stdout='not json')
self.assertIsNone(drainer._lookup_jid(JID, self.log))
run.side_effect = subprocess.TimeoutExpired('salt-run', 60)
self.assertIsNone(drainer._lookup_jid(JID, self.log))
def test_minion_failure_shapes(self):
self.assertEqual(drainer._minion_failure([CONFLICT]), json.dumps([CONFLICT]))
self.assertEqual(drainer._minion_failure('Rendering SLS failed'), 'Rendering SLS failed')
self.assertEqual(drainer._minion_failure(True), '')
self.assertEqual(drainer._minion_failure({'a': {'result': True}}), '')
def test_orch_failures_conflict(self):
failures = drainer._orch_failures(CONFLICT_RET)
self.assertEqual(failures[0], 'apply_soc_1: Run failed on minions: manager_standalone')
self.assertIn('manager_standalone', failures[1])
self.assertIn('is running as PID 372218', failures[1])
self.assertEqual(len(failures), 2)
def test_orch_failures_failed_state(self):
failures = drainer._orch_failures(STATE_FAIL_RET)
self.assertEqual(len(failures), 2)
self.assertIn('hydra.enabled_no_license_detected: This is a feature', failures[1])
self.assertNotIn('hydra_conf', failures[1])
def test_orch_failures_success(self):
self.assertEqual(drainer._orch_failures(SUCCESS_RET), [])
def test_orch_failures_render_error(self):
ret = {MASTER: {'return': {'data': {MASTER: ['Rendering SLS failed']}}, 'success': False}}
self.assertEqual(drainer._orch_failures(ret), ['["Rendering SLS failed"]'])
def test_orch_failures_not_a_dict(self):
self.assertEqual(drainer._orch_failures(['No minions matched']), ['["No minions matched"]'])
self.assertEqual(drainer._orch_failures('Runner error'), ['Runner error'])
def test_orch_failures_data_not_a_dict(self):
ret = {MASTER: {'return': {'data': ["Rendering SLS 'orch.push_batch' failed"]}, 'success': False}}
self.assertEqual(drainer._orch_failures(ret), ['["Rendering SLS \'orch.push_batch\' failed"]'])
def test_orch_failures_odd_changes(self):
for changes in ('Run failed', {'ret': ['manager_standalone']}):
ret = _orch_ret({'salt_|-apply_soc_1_|-apply_soc_1_|-state': {
'__id__': 'apply_soc_1', 'result': False, 'changes': changes, 'comment': 'Run failed on minions',
}}, success=False)
self.assertEqual(drainer._orch_failures(ret), ['apply_soc_1: Run failed on minions'])
def test_orch_failures_unparsed(self):
self.assertEqual(drainer._orch_failures({MASTER: 'odd'}), [])
ret = {MASTER: {'return': 'Exception occurred', 'success': False}}
self.assertEqual(drainer._orch_failures(ret), ['orchestration reported failure: Exception occurred'])
def record(self, jid, age, now):
return self.write_json(self.dispatched, jid + '.json', {
'jid': jid, 'dispatched_at': now - age, 'actions': [], 'paths': ['audit:' + jid],
})
def test_check_dispatched(self):
now = time.time()
results = {
'1_failed': CONFLICT_RET,
'2_ok': SUCCESS_RET,
'3_pending': {},
'4_expired': None,
}
young = self.record('0_young', 5, now)
paths = {jid: self.record(jid, 60, now) for jid in results}
paths['4_expired'] = self.record('4_expired', drainer.RESULT_MAX_AGE + 1, now)
bad = self.write_json(self.dispatched, '5_bad.json', '{nope')
with patch.object(drainer, '_lookup_jid', side_effect=lambda jid, log: results[jid]):
drainer._check_dispatched(self.log, now)
self.assertTrue(os.path.exists(young))
with open(paths['3_pending']) as f:
self.assertEqual(json.load(f)['checked_at'], now)
for jid in ('1_failed', '2_ok', '4_expired'):
self.assertFalse(os.path.exists(paths[jid]), jid)
self.assertFalse(os.path.exists(bad))
self.assertIn('push failed jid=1_failed', self.logged('error'))
self.assertIn('is running as PID 372218', self.logged('error'))
self.assertIn('push succeeded jid=2_ok', self.logged('info'))
self.assertIn('no result for jid=4_expired', self.logged('warning'))
def test_check_dispatched_survives_bad_result(self):
now = time.time()
bad = self.record('1_bad', 60, now)
good = self.record('2_ok', 60, now)
def orch_failures(ret):
if ret == 'boom':
raise ValueError('unexpected shape')
return []
with patch.object(drainer, '_lookup_jid', side_effect=lambda jid, log: 'boom' if jid == '1_bad' else SUCCESS_RET), \
patch.object(drainer, '_orch_failures', side_effect=orch_failures):
drainer._check_dispatched(self.log, now)
self.assertFalse(os.path.exists(bad))
self.assertFalse(os.path.exists(good))
self.log.exception.assert_called_once()
self.assertIn('jid=1_bad', self.log.exception.call_args[0][0] % self.log.exception.call_args[0][1:])
self.assertIn('push succeeded jid=2_ok', self.logged('info'))
def test_recheck_delay(self):
self.assertEqual(drainer._recheck_delay(10), drainer.RESULT_CHECK_DELAY)
self.assertEqual(drainer._recheck_delay(400), 100)
self.assertEqual(drainer._recheck_delay(drainer.RESULT_MAX_AGE), drainer.RESULT_RECHECK_MAX)
def test_check_dispatched_limit_rotates(self):
now = time.time()
limit = drainer.RESULT_CHECKS_PER_PASS
jids = ['{:02d}'.format(i) for i in range(limit + 2)]
for jid in jids:
self.record(jid, 60, now)
with patch.object(drainer, '_lookup_jid', return_value={}) as lookup:
drainer._check_dispatched(self.log, now)
self.assertEqual([c.args[0] for c in lookup.call_args_list], jids[:limit])
lookup.reset_mock()
drainer._check_dispatched(self.log, now + 40)
self.assertEqual([c.args[0] for c in lookup.call_args_list], jids[limit:] + jids[:limit - 2])
def test_check_dispatched_not_blocked_by_running(self):
now = time.time()
for i in range(drainer.RESULT_CHECKS_PER_PASS):
self.record('1_running{}'.format(i), 600, now)
done = self.record('2_done', 60, now)
def lookup(jid, log):
return SUCCESS_RET if jid == '2_done' else {}
with patch.object(drainer, '_lookup_jid', side_effect=lookup) as lookup_jid:
drainer._check_dispatched(self.log, now)
self.assertTrue(os.path.exists(done))
lookup_jid.reset_mock()
drainer._check_dispatched(self.log, now + 15)
self.assertEqual([c.args[0] for c in lookup_jid.call_args_list], ['2_done'])
self.assertFalse(os.path.exists(done))
self.assertIn('push succeeded jid=2_done', self.logged('info'))
class TestMain(DrainerTestCase):
def setUp(self):
super().setUp()
self.cfg = {'enabled': True, 'debounce_seconds': 30}
for name, kwargs in (
('_make_logger', {'return_value': self.log}),
('_load_push_cfg', {'side_effect': lambda: self.cfg}),
('_check_dispatched', {}),
):
patcher = patch.object(drainer, name, **kwargs)
setattr(self, name, patcher.start())
self.addCleanup(patcher.stop)
def intent(self, name, age=60, actions=None, paths=None):
now = time.time()
return self.write_json(self.pending, name, {
'first_touch': now - age - 5, 'last_touch': now - age,
'actions': [{'state': 'soc', 'tgt': '*'}] if actions is None else actions,
'paths': paths or ['audit:soc.config.licenseKey'],
})
def test_no_pending_dir(self):
shutil.rmtree(self.pending)
self.assertEqual(drainer.main(), 0)
self._load_push_cfg.assert_not_called()
def test_cfg_error(self):
self._load_push_cfg.side_effect = RuntimeError('no salt')
self.assertEqual(drainer.main(), 1)
def test_disabled(self):
self.cfg['enabled'] = False
self.assertEqual(drainer.main(), 0)
self._check_dispatched.assert_not_called()
def test_no_intents_still_checks_results(self):
self.assertEqual(drainer.main(), 0)
self._check_dispatched.assert_called_once()
def test_debounce_and_broken(self):
young = self.intent('young.json', age=1)
broken = self.write_json(self.pending, 'broken.json', '{nope')
with patch.object(drainer, '_dispatch') as dispatch:
self.assertEqual(drainer.main(), 0)
dispatch.assert_not_called()
self.assertTrue(os.path.exists(young))
self.assertFalse(os.path.exists(broken))
def test_broken_unlink_error_ignored(self):
self.write_json(self.pending, 'broken.json', '{nope')
with patch.object(drainer.os, 'unlink', side_effect=OSError('busy')):
self.assertEqual(drainer.main(), 0)
def test_no_usable_actions(self):
path = self.intent('empty.json', actions=[{'state': 'soc'}])
self.assertEqual(drainer.main(), 0)
self.assertFalse(os.path.exists(path))
self.intent('empty.json', actions=[{'state': 'soc'}])
with patch.object(drainer.os, 'unlink', side_effect=OSError('busy')):
self.assertEqual(drainer.main(), 0)
def test_dispatch_failure_keeps_intents(self):
path = self.intent('pillar_soc.json')
with patch.object(drainer, '_dispatch', return_value=None):
self.assertEqual(drainer.main(), 1)
self.assertTrue(os.path.exists(path))
def test_dispatch_records_jid(self):
soc = self.intent('pillar_soc.json')
hs = self.intent('pillar_global.json', actions=[{'highstate': True, 'tgt': '*'}], paths=['audit:global.x'])
with patch.object(drainer, '_dispatch', return_value=JID) as dispatch, \
patch.object(drainer, '_record_dispatch') as record:
self.assertEqual(drainer.main(), 0)
self.assertEqual(len(dispatch.call_args[0][0]), 2)
record.assert_called_once()
self.assertEqual(record.call_args[0][0], JID)
self.assertEqual(sorted(record.call_args[0][2]), ['audit:global.x', 'audit:soc.config.licenseKey'])
self.assertFalse(os.path.exists(soc))
self.assertFalse(os.path.exists(hs))
self.assertIn('action: highstate tgt=*', self.logged('info'))
def test_dispatch_without_jid_not_recorded(self):
self.intent('pillar_soc.json')
with patch.object(drainer, '_dispatch', return_value=''), \
patch.object(drainer, '_record_dispatch') as record, \
patch.object(drainer.os, 'unlink', side_effect=OSError('busy')):
self.assertEqual(drainer.main(), 0)
record.assert_not_called()
self.log.exception.assert_called_once()
if __name__ == '__main__':
unittest.main()
+1 -2
View File
@@ -42,8 +42,7 @@ def loadYaml(filename):
try:
with open(filename, "r") as file:
content = file.read()
loaded = yaml.safe_load(content)
return loaded if loaded is not None else {}
return yaml.safe_load(content)
except FileNotFoundError:
print(f"File not found: {filename}", file=sys.stderr)
sys.exit(1)
-97
View File
@@ -95,20 +95,6 @@ class TestRemove(unittest.TestCase):
expected = "key1:\n child1: 123\n child2:\n deep2: ab\nkey2: false\n"
self.assertEqual(actual, expected)
def test_remove_empty_file(self):
filename = "/tmp/so-yaml_test-remove-empty.yaml"
file = open(filename, "w")
file.close()
code = soyaml.remove([filename, "key1"])
self.assertEqual(code, 0)
file = open(filename, "r")
actual = file.read()
file.close()
self.assertEqual(actual, "{}\n")
def test_remove_missing_args(self):
with patch('sys.exit', new=MagicMock()) as sysmock:
with patch('sys.stderr', new=StringIO()) as mock_stderr:
@@ -308,36 +294,6 @@ class TestRemove(unittest.TestCase):
expected = "key1:\n child1: 123\n child2:\n deep1: 45\n deep2: d\nkey2: false\nkey3:\n- e\n- f\n- g\n"
self.assertEqual(actual, expected)
def test_add_empty_file(self):
filename = "/tmp/so-yaml_test-add-empty.yaml"
file = open(filename, "w")
file.close()
code = soyaml.add([filename, "telegraf.output", "BOTH"])
self.assertEqual(code, 0)
file = open(filename, "r")
actual = file.read()
file.close()
expected = "telegraf:\n output: BOTH\n"
self.assertEqual(actual, expected)
def test_add_empty_file_simple(self):
filename = "/tmp/so-yaml_test-add-empty-simple.yaml"
file = open(filename, "w")
file.close()
code = soyaml.add([filename, "telegraf", "BOTH"])
self.assertEqual(code, 0)
file = open(filename, "r")
actual = file.read()
file.close()
expected = "telegraf: BOTH\n"
self.assertEqual(actual, expected)
def test_replace_missing_arg(self):
with patch('sys.exit', new=MagicMock()) as sysmock:
with patch('sys.stderr', new=StringIO()) as mock_stderr:
@@ -390,21 +346,6 @@ class TestRemove(unittest.TestCase):
expected = "key1:\n child1: 123\n child2:\n deep1: 46\nkey2: false\nkey3:\n- e\n- f\n- g\n"
self.assertEqual(actual, expected)
def test_replace_empty_file(self):
filename = "/tmp/so-yaml_test-replace-empty.yaml"
file = open(filename, "w")
file.close()
code = soyaml.replace([filename, "telegraf.output", "BOTH"])
self.assertEqual(code, 0)
file = open(filename, "r")
actual = file.read()
file.close()
expected = "telegraf:\n output: BOTH\n"
self.assertEqual(actual, expected)
def test_convert(self):
self.assertEqual(soyaml.convertType("foo"), "foo")
self.assertEqual(soyaml.convertType("foo.bar"), "foo.bar")
@@ -565,18 +506,6 @@ class TestRemove(unittest.TestCase):
self.assertEqual(result, 2)
self.assertEqual("", mock_stdout.getvalue())
def test_get_empty_file(self):
with patch('sys.stdout', new=StringIO()) as mock_stdout:
with patch('sys.stderr', new=StringIO()) as mock_stderr:
filename = "/tmp/so-yaml_test-get-empty.yaml"
file = open(filename, "w")
file.close()
result = soyaml.get([filename, "telegraf.output"])
self.assertEqual(result, 2)
self.assertEqual("", mock_stdout.getvalue())
self.assertIn("Key 'telegraf.output' not found by so-yaml.py", mock_stderr.getvalue())
def test_get_usage(self):
with patch('sys.exit', new=MagicMock()) as sysmock:
with patch('sys.stderr', new=StringIO()) as mock_stderr:
@@ -1062,29 +991,3 @@ class TestLoadYaml(unittest.TestCase):
soyaml.loadYaml("/tmp/so-yaml_test-unreadable.yaml")
sysmock.assert_called_with(1)
self.assertIn("Error reading file", mock_stderr.getvalue())
def test_load_yaml_empty_file(self):
filename = "/tmp/so-yaml_test-load-empty.yaml"
file = open(filename, "w")
file.close()
result = soyaml.loadYaml(filename)
self.assertEqual(result, {})
def test_load_yaml_whitespace_only(self):
filename = "/tmp/so-yaml_test-load-whitespace.yaml"
file = open(filename, "w")
file.write(" \n\n \n")
file.close()
result = soyaml.loadYaml(filename)
self.assertEqual(result, {})
def test_load_yaml_comments_only(self):
filename = "/tmp/so-yaml_test-load-comments.yaml"
file = open(filename, "w")
file.write("# Just a comment\n# Another comment\n")
file.close()
result = soyaml.loadYaml(filename)
self.assertEqual(result, {})
-16
View File
@@ -1183,18 +1183,6 @@ up_to_3.4.0() {
echo "Removing so-kratos, so-hydra and so-soc so they are recreated on the soauth network."
docker rm -f so-kratos so-hydra so-soc >> $SOUP_LOG 2>&1
# Extract the Sigma rule type for existing detections (Single vs. Correlation)
mkdir -p /opt/so/conf/soc/migrations
echo "0" > /opt/so/conf/soc/migrations/elastalert-migration-3.4.0
chown -R socore:socore /opt/so/conf/soc/migrations
for template in so-metrics-logstash.node so-metrics-logstash.stack_monitoring.node; do
if ! remove_elasticsearch_index_template "$template" "logstash node and node_cel index patterns reversed"; then
FINAL_MESSAGE_QUEUE+=("WARNING: Unable to automatically remove the $template index template. Addon integration templates may fail to load until it is removed:")
FINAL_MESSAGE_QUEUE+=(" - sudo so-elasticsearch-query _index_template/$template -XDELETE && so-checkin")
fi
done
INSTALLEDVERSION=3.4.0
}
@@ -1260,10 +1248,6 @@ valid_soauth_range() {
}
post_to_3.4.0() {
for idx in "metrics-logstash.node-default" "metrics-logstash.stack_monitoring.node-default"; do
rollover_index "$idx"
done
set_postversion 3.4.0
}
### 3.4.0 End ###
+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 {}
+1 -1
View File
@@ -9,7 +9,7 @@
'epel-testing.repo',
'saltstack.repo',
'salt-latest.repo',
'wazuh.repo',
'wazuh.repo'
'Rocky-Base.repo',
'Rocky-CR.repo',
'Rocky-Debuginfo.repo',
+6 -17
View File
@@ -120,30 +120,19 @@ crondetectionsbackup:
socsigmafinalpipeline:
file.managed:
- name: /opt/so/conf/soc/sigma_pipelines/sigma_final_pipeline.yml
- name: /opt/so/conf/soc/sigma_final_pipeline.yaml
- source: salt://soc/files/soc/sigma_final_pipeline.yaml
- user: 939
- group: 939
- mode: 600
- makedirs: True
# sigma-cli loads every *.yml here; clean removes anything else
socsigmapipelines:
file.recurse:
- name: /opt/so/conf/soc/sigma_pipelines
- source: salt://soc/files/soc/sigma_pipelines
socsigmasopipeline:
file.managed:
- name: /opt/so/conf/soc/sigma_so_pipeline.yaml
- source: salt://soc/files/soc/sigma_so_pipeline.yaml
- user: 939
- group: 939
- file_mode: 600
- clean: True
- require:
- file: socsigmafinalpipeline
socsigmapipelinesold:
file.absent:
- names:
- /opt/so/conf/soc/sigma_final_pipeline.yaml
- /opt/so/conf/soc/sigma_so_pipeline.yaml
- mode: 600
socsigmaplaybookpipeline:
file.managed:
+2 -59
View File
@@ -1467,8 +1467,6 @@ soc:
- emerging_threats_addon
useEsql: false
esqlCaseInsensitive: true
esqlQueryDelaySeconds: 30
esqlCorrelationAllowanceSeconds: 600
elastic:
hostUrl:
remoteHostUrls: []
@@ -2681,11 +2679,8 @@ soc:
query: "so_detection.language:suricata | groupby so_detection.ruleset so_detection.isEnabled | groupby so_detection.category"
description: Show all NIDS Detections, which are run with Suricata
- name: "Detection Type - Sigma (Elastalert) - All"
query: "so_detection.language:sigma | groupby so_detection.ruleType | groupby so_detection.ruleset so_detection.isEnabled | groupby so_detection.category | groupby so_detection.product"
query: "so_detection.language:sigma | groupby so_detection.ruleset so_detection.isEnabled | groupby so_detection.category | groupby so_detection.product"
description: Show all Sigma Detections, which are run with Elastalert
- name: "Detection Type - Sigma (Elastalert) - Correlations"
query: "so_detection.ruleType:correlation | groupby so_detection.correlationType so_detection.isEnabled | groupby so_detection.correlationTimespan | groupby so_detection.ruleset"
description: Show Sigma correlation Detections
- name: "Detection Type - YARA (Strelka)"
query: "so_detection.language:yara | groupby so_detection.ruleset so_detection.isEnabled"
description: Show all YARA detections, which are used by Strelka
@@ -2769,7 +2764,7 @@ soc:
elastalert: |
# This is a Sigma rule template, which uses YAML. Replace all template values with your own values.
# The id (UUIDv4) is pregenerated and can safely be used.
# Click "Convert" to convert the Sigma rule to use Security Onion field mappings within a backend query
# Click "Convert" to convert the Sigma rule to use Security Onion field mappings within an EQL query
#
# Rule Creation Guide: https://github.com/SigmaHQ/sigma/wiki/Rule-Creation-High%E2%80%90Level-Guide
# Logsources: https://sigmahq.io/docs/basics/log-sources.html
@@ -2799,58 +2794,6 @@ soc:
- ' -priv'
condition: all of selection_*
level: 'high' # info | low | medium | high | critical
elastalert_correlation: |
# Sigma correlation rule; requires ES|QL (useEsql).
# The first document is the correlation (id, title, severity); the documents after it are the rules it refers to, all in this Detection.
#
# Supported types: event_count, value_count, temporal, value_sum, value_avg, value_percentile, value_median.
# Correlation Guide: https://sigmahq.io/docs/meta/correlations.html
# Logsources: https://sigmahq.io/docs/basics/log-sources.html
title: 'A Short Capitalized Title With Less Than 50 Characters'
id: [publicId]
status: 'experimental'
description: |
Describe what the correlation finds and, importantly, why relating these
events is more meaningful than either of them alone.
references:
- 'https://local.invalid'
author: '@SecurityOnion'
date: '[today]'
tags:
- detection.threat_hunting
- attack.technique_id
correlation:
type: value_count
rules:
- example_base_rule # matches the 'name' of the document below
group-by:
- source.ip
# Xs, Xm, Xh, Xd or Xw.
timespan: 10m
condition:
field: dns.query.name
gte: 40
falsepositives:
- 'Describe the benign activity that also produces this pattern'
# Placeholders: %count%, %start%, %end%, %duration%, and group-by or field values.
summary: '%count% distinct names queried by %source.ip% in %duration%'
level: 'medium' # info | low | medium | high | critical
---
title: 'Base Event'
# The correlation refers to this rule by 'name' (or by 'id').
name: example_base_rule
description: 'The single event that the correlation aggregates.'
logsource:
category: network
service: dns
detection:
selection:
dns.query.name|exists: true
condition: selection
# Carried into the alert.
fields:
- dns.query.name
assistant:
enabled: false
investigationPrompt: Investigate Alert ID {socId}
+2 -2
View File
@@ -47,8 +47,9 @@ so-soc:
{% endif %}
- /opt/so/conf/soc/motd.md:/opt/sensoroni/html/motd.md:ro
- /opt/so/conf/soc/banner.md:/opt/sensoroni/html/login/banner.md:ro
- /opt/so/conf/soc/sigma_pipelines:/opt/sensoroni/sigma_pipelines:ro
- /opt/so/conf/soc/sigma_so_pipeline.yaml:/opt/sensoroni/sigma_so_pipeline.yaml:ro
- /opt/so/conf/soc/sigma_playbook_pipeline.yaml:/opt/sensoroni/sigma_playbook_pipeline.yaml:ro
- /opt/so/conf/soc/sigma_final_pipeline.yaml:/opt/sensoroni/sigma_final_pipeline.yaml:ro
- /opt/so/conf/soc/playbook_placeholder_map.yaml:/opt/sensoroni/playbook_placeholder_map.yaml:ro
- /opt/so/conf/soc/playbook_placeholder_map_custom.yaml:/opt/sensoroni/playbook_placeholder_map_custom.yaml:ro
- /opt/so/conf/soc/custom.js:/opt/sensoroni/html/js/custom.js:ro
@@ -106,7 +107,6 @@ so-soc:
- file: socclientsroles
- file: socplaybookplaceholdermap
- file: socplaybookplaceholdermapcustom
- file: socsigmapipelines
delete_so-soc_so-status.disabled:
file.uncomment:
@@ -1,477 +0,0 @@
name: Security Onion ES|QL Pipeline
# ES|QL query settings
priority: 92
transformations:
- id: esql_default_index
type: set_state
key: index
val: .ds-logs-*
- id: esql_source_metadata
type: set_state
key: metadata
val: "_id, _index, _source"
- id: esql_source_keep
type: set_state
key: keep
val: "_id, _index, _source"
# unmapped fields read as null instead of failing the query
- id: esql_unmapped_fields
type: set_state
key: unmapped_fields
val: nullify
# FROM targets per logsource, any namespace; later entries win, correlations get the union
- id: esql_index_process_creation
type: set_state
key: index
val:
- .ds-logs-endpoint.events.process-*
- .ds-logs-windows.sysmon_operational-*
- .ds-logs-system.security-*
- .ds-logs-windows.powershell-*
- .ds-logs-windows.forwarded-*
- .ds-logs-sysmon_linux.log-*
- .ds-logs-auditd_manager.auditd-*
- .ds-logs-import-so-*
rule_conditions:
- type: logsource
category: process_creation
- id: esql_index_process_creation_windows
type: set_state
key: index
val:
- .ds-logs-endpoint.events.process-*
- .ds-logs-windows.sysmon_operational-*
- .ds-logs-system.security-*
- .ds-logs-windows.powershell-*
- .ds-logs-windows.forwarded-*
- .ds-logs-import-so-*
rule_conditions:
- type: logsource
product: windows
category: process_creation
- id: esql_index_process_creation_linux
type: set_state
key: index
val:
- .ds-logs-endpoint.events.process-*
- .ds-logs-sysmon_linux.log-*
- .ds-logs-auditd_manager.auditd-*
- .ds-logs-import-so-*
rule_conditions:
- type: logsource
product: linux
category: process_creation
- id: esql_index_process_creation_macos
type: set_state
key: index
val:
- .ds-logs-endpoint.events.process-*
- .ds-logs-import-so-*
rule_conditions:
- type: logsource
product: macos
category: process_creation
- id: esql_index_file
type: set_state
key: index
val:
- .ds-logs-endpoint.events.file-*
- .ds-logs-windows.sysmon_operational-*
- .ds-logs-windows.forwarded-*
- .ds-logs-sysmon_linux.log-*
- .ds-logs-import-so-*
rule_cond_op: or
rule_conditions:
- type: logsource
category: file_event
- type: logsource
category: file_delete
- type: logsource
category: file_rename
- type: logsource
category: file_change
- type: logsource
category: file_access
- id: esql_index_file_windows
type: set_state
key: index
val:
- .ds-logs-endpoint.events.file-*
- .ds-logs-windows.sysmon_operational-*
- .ds-logs-windows.forwarded-*
- .ds-logs-import-so-*
rule_cond_op: or
rule_conditions:
- type: logsource
product: windows
category: file_event
- type: logsource
product: windows
category: file_delete
- type: logsource
product: windows
category: file_rename
- type: logsource
product: windows
category: file_change
- type: logsource
product: windows
category: file_access
- id: esql_index_file_linux
type: set_state
key: index
val:
- .ds-logs-endpoint.events.file-*
- .ds-logs-sysmon_linux.log-*
- .ds-logs-import-so-*
rule_cond_op: or
rule_conditions:
- type: logsource
product: linux
category: file_event
- type: logsource
product: linux
category: file_delete
- type: logsource
product: linux
category: file_rename
- type: logsource
product: linux
category: file_change
- type: logsource
product: linux
category: file_access
- id: esql_index_file_macos
type: set_state
key: index
val:
- .ds-logs-endpoint.events.file-*
- .ds-logs-import-so-*
rule_cond_op: or
rule_conditions:
- type: logsource
product: macos
category: file_event
- type: logsource
product: macos
category: file_delete
- type: logsource
product: macos
category: file_rename
- type: logsource
product: macos
category: file_change
- type: logsource
product: macos
category: file_access
- id: esql_index_registry
type: set_state
key: index
val:
- .ds-logs-endpoint.events.registry-*
- .ds-logs-windows.sysmon_operational-*
- .ds-logs-windows.forwarded-*
- .ds-logs-import-so-*
rule_cond_op: or
rule_conditions:
- type: logsource
category: registry_set
- type: logsource
category: registry_add
- type: logsource
category: registry_delete
- type: logsource
category: registry_event
- id: esql_index_library
type: set_state
key: index
val:
- .ds-logs-endpoint.events.library-*
- .ds-logs-windows.sysmon_operational-*
- .ds-logs-windows.forwarded-*
- .ds-logs-import-so-*
rule_cond_op: or
rule_conditions:
- type: logsource
category: image_load
- type: logsource
category: driver_load
- id: esql_index_endpoint_network
type: set_state
key: index
val:
- .ds-logs-endpoint.events.network-*
- .ds-logs-windows.sysmon_operational-*
- .ds-logs-windows.forwarded-*
- .ds-logs-sysmon_linux.log-*
- .ds-logs-import-so-*
rule_cond_op: or
rule_conditions:
- type: logsource
category: network_connection
- type: logsource
category: dns_query
- id: esql_index_endpoint_network_windows
type: set_state
key: index
val:
- .ds-logs-endpoint.events.network-*
- .ds-logs-windows.sysmon_operational-*
- .ds-logs-windows.forwarded-*
- .ds-logs-import-so-*
rule_cond_op: or
rule_conditions:
- type: logsource
product: windows
category: network_connection
- type: logsource
product: windows
category: dns_query
- id: esql_index_endpoint_network_linux
type: set_state
key: index
val:
- .ds-logs-endpoint.events.network-*
- .ds-logs-sysmon_linux.log-*
- .ds-logs-import-so-*
rule_cond_op: or
rule_conditions:
- type: logsource
product: linux
category: network_connection
- type: logsource
product: linux
category: dns_query
- id: esql_index_endpoint_network_macos
type: set_state
key: index
val:
- .ds-logs-endpoint.events.network-*
- .ds-logs-import-so-*
rule_cond_op: or
rule_conditions:
- type: logsource
product: macos
category: network_connection
- type: logsource
product: macos
category: dns_query
- id: esql_index_sysmon_only
type: set_state
key: index
val:
- .ds-logs-windows.sysmon_operational-*
- .ds-logs-windows.forwarded-*
- .ds-logs-import-so-*
rule_cond_op: or
rule_conditions:
- type: logsource
category: process_access
- type: logsource
category: create_remote_thread
- type: logsource
category: pipe_created
- type: logsource
category: create_stream_hash
- type: logsource
category: wmi_event
- type: logsource
category: raw_access_thread
- type: logsource
category: process_tampering
- type: logsource
category: sysmon_status
- type: logsource
category: sysmon_error
- type: logsource
category: file_executable_detected
- type: logsource
category: file_block_executable
- type: logsource
category: file_block_shredding
- type: logsource
category: clipboard_capture
- type: logsource
product: windows
service: sysmon
- id: esql_index_ps_operational
type: set_state
key: index
val:
- .ds-logs-windows.powershell_operational-*
- .ds-logs-windows.forwarded-*
- .ds-logs-import-so-*
rule_cond_op: or
rule_conditions:
- type: logsource
category: ps_script
- type: logsource
category: ps_module
- type: logsource
product: windows
service: powershell
- id: esql_index_ps_classic
type: set_state
key: index
val:
- .ds-logs-windows.powershell-*
- .ds-logs-windows.forwarded-*
- .ds-logs-import-so-*
rule_cond_op: or
rule_conditions:
- type: logsource
category: ps_classic_start
- type: logsource
category: ps_classic_provider_start
- type: logsource
category: ps_classic_script
- type: logsource
product: windows
service: powershell-classic
- id: esql_index_win_security
type: set_state
key: index
val:
- .ds-logs-system.security-*
- .ds-logs-windows.forwarded-*
- .ds-logs-import-so-*
rule_conditions:
- type: logsource
product: windows
service: security
- id: esql_index_win_system
type: set_state
key: index
val:
- .ds-logs-system.system-*
- .ds-logs-windows.forwarded-*
- .ds-logs-import-so-*
rule_conditions:
- type: logsource
product: windows
service: system
- id: esql_index_win_application
type: set_state
key: index
val:
- .ds-logs-system.application-*
- .ds-logs-windows.forwarded-*
- .ds-logs-import-so-*
rule_conditions:
- type: logsource
product: windows
service: application
- id: esql_index_linux_auth
type: set_state
key: index
val:
- .ds-logs-system.auth-*
- .ds-logs-import-so-*
rule_cond_op: or
rule_conditions:
- type: logsource
product: linux
service: auth
- type: logsource
product: linux
service: sshd
- type: logsource
product: linux
service: sudo
- id: esql_index_linux_syslog
type: set_state
key: index
val:
- .ds-logs-system.syslog-*
- .ds-logs-syslog-so-*
- .ds-logs-import-so-*
rule_conditions:
- type: logsource
product: linux
service: syslog
- id: esql_index_linux_auditd
type: set_state
key: index
val:
- .ds-logs-auditd_manager.auditd-*
- .ds-logs-auditd.log-*
- .ds-logs-import-so-*
rule_conditions:
- type: logsource
product: linux
service: auditd
- id: esql_index_network
type: set_state
key: index
val:
- .ds-logs-zeek-so-*
- .ds-logs-suricata-so-*
- .ds-logs-suricata.alerts-so-*
- .ds-logs-endpoint.events.network-*
- .ds-logs-windows.sysmon_operational-*
- .ds-logs-sysmon_linux.log-*
- .ds-logs-windows.forwarded-*
- .ds-logs-import-so-*
rule_conditions:
- type: logsource
category: network
- id: esql_index_so_network
type: set_state
key: index
val:
- .ds-logs-zeek-so-*
- .ds-logs-suricata-so-*
- .ds-logs-import-so-*
rule_cond_op: or
rule_conditions:
- type: logsource
category: network
service: connection
- type: logsource
category: network
service: dns
- type: logsource
category: network
service: http
- type: logsource
category: network
service: file
- type: logsource
category: network
service: x509
- type: logsource
category: network
service: ssl
- type: logsource
category: network
service: ssh
- type: logsource
category: dns
- id: esql_index_zeek
type: set_state
key: index
val:
- .ds-logs-zeek-so-*
- .ds-logs-import-so-*
rule_conditions:
- type: logsource
product: zeek
- id: esql_index_opencanary
type: set_state
key: index
val:
- .ds-logs-idh-so-*
- .ds-logs-import-so-*
rule_conditions:
- type: logsource
product: opencanary
- id: esql_index_kratos
type: set_state
key: index
val:
- .ds-logs-kratos-so-*
- .ds-logs-import-so-*
rule_conditions:
- type: logsource
product: kratos
@@ -2,13 +2,13 @@ name: Security Onion - Playbook Pipeline
priority: 97
transformations:
# Route to lowercase-normalized .caseless subfields for case-insensitive matching.
# file.path.caseless exists on Defend only (Sysmon file events lack it);
# registry.path / dll.path / file.name have no .caseless on any source.
- id: case_insensitive_string_fields
type: field_name_mapping
mapping:
process.executable: process.executable.caseless
process.parent.executable: process.parent.executable.caseless
process.parent.name: process.parent.name.caseless
process.command_line: process.command_line.caseless
process.parent.command_line: process.parent.command_line.caseless
file.path: file.path.caseless
@@ -14,25 +14,18 @@ transformations:
- process.args
- related.ip
- dns.resolved_ip
# Always lowercase, so matched exactly; the backend can then use the indexed ':' operator.
- id: case_sensitive_categorization_fields
- id: esql_default_index
type: set_state
key: case_insensitive_exempt_fields
val:
- tags
- event.category
- event.type
- event.kind
# Not every source maps .caseless; EQL/ES|QL already match case-insensitively.
- id: caseless_to_parent_fields
type: field_name_mapping
mapping:
process.executable.caseless: process.executable
process.name.caseless: process.name
process.parent.executable.caseless: process.parent.executable
process.parent.name.caseless: process.parent.name
target.process.executable.caseless: target.process.executable
target.process.name.caseless: target.process.name
key: index
val: .ds-logs-*
- id: esql_source_metadata
type: set_state
key: metadata
val: "_id, _index, _source"
- id: esql_source_keep
type: set_state
key: keep
val: "_id, _index, _source"
- id: baseline_field_name_mapping
type: field_name_mapping
mapping:
@@ -127,9 +120,6 @@ transformations:
valid_hash_algos: ["MD5", "SHA1", "SHA256", "SHA512", "IMPHASH"]
field_prefix: "file"
drop_algo_prefix: False
# ecs_windows has already renamed Hashes; pySigma 1.5+ only parses the fields listed here
field_to_parse:
- winlog.event_data.Hashes
field_name_conditions:
- type: include_fields
fields:
-17
View File
@@ -8,7 +8,6 @@
{% from 'elasticsearch/config.map.jinja' import ELASTICSEARCH_NODES %}
{% from 'manager/map.jinja' import MANAGERMERGED %}
{% from 'telegraf/map.jinja' import TELEGRAFMERGED %}
{% from 'elastalert/map.jinja' import ELASTALERTMERGED %}
{%- set PG_ENTRY = salt['pillar.get']('telegraf:postgres_creds:' ~ grains.id, {}) %}
{%- set PG_USER = PG_ENTRY.get('user', '') %}
{%- set PG_PASS = PG_ENTRY.get('pass', '') %}
@@ -64,10 +63,6 @@
{% do SOCMERGED.config.server.modules.elastalertengine.update({'enabledSigmaRules': SOCMERGED.config.server.modules.elastalertengine.enabledSigmaRules.default}) %}
{% endif %}
{# correlation schedules follow ElastAlert's run_every #}
{% set run_every = ELASTALERTMERGED.config.run_every %}
{% do SOCMERGED.config.server.modules.elastalertengine.update({'elastAlertRunEverySeconds': run_every.get('minutes', 0) * 60 + run_every.get('seconds', 0)}) %}
{# set elastalertengine.rulesRepos, strelkaengine.rulesRepos, and suricataengine.rulesetSources based on airgap or not #}
{% if GLOBALS.airgap %}
{% do SOCMERGED.config.server.modules.elastalertengine.update({'rulesRepos': SOCMERGED.config.server.modules.elastalertengine.rulesRepos.airgap}) %}
@@ -85,18 +80,6 @@
{% do SOCMERGED.config.server.update({'airgapEnabled': false}) %}
{% endif %}
{# Sigma correlations require ES|QL: load the community correlations and offer correlation authoring only when it is on #}
{% set use_esql = SOCMERGED.config.server.modules.elastalertengine.useEsql %}
{% for repo in SOCMERGED.config.server.modules.elastalertengine.rulesRepos %}
{% if repo.get('rulesetName') == 'securityonion-resources' and repo.get('folder') in ['sigma', 'sigma/stable'] %}
{% do repo.update({'folder': 'sigma' if use_esql else 'sigma/stable'}) %}
{% endif %}
{% endfor %}
{% if not use_esql %}
{% do SOCMERGED.config.server.client.detection.templateDetections.pop('elastalert_correlation', None) %}
{% do SOCMERGED.config.server.client.detections.update({'queries': SOCMERGED.config.server.client.detections.queries | rejectattr('name', 'equalto', 'Detection Type - Sigma (Elastalert) - Correlations') | list}) %}
{% endif %}
{# Define the postgresmetrics module if telegraf is setup to only use Postgres #}
{% if TELEGRAFMERGED.output != 'INFLUXDB' and PG_USER and PG_PASS %}
{% do SOCMERGED.config.server.modules.update({
+10 -21
View File
@@ -401,7 +401,7 @@ soc:
advanced: False
helpLink: sigma
useEsql:
description: "(Pre-release) Use Elasticsearch Piped Query Language (ES|QL) instead of EQL (Elastic Query Language) for Elasticsearch queries. The Sigma converter will output ES|QL instead of EQL, allowing support for correlations. Switching back to EQL is not supported for correlations."
description: "(Pre-release) Use Elasticsearch Piped Query Language (ES|QL) instead of EQL (Elastic Query Language) for Elasticsearch queries. The Sigma converter will output ES|QL instead of EQL, allowing support for correlations."
global: True
advanced: True
forcedType: bool
@@ -410,18 +410,6 @@ soc:
global: True
advanced: True
forcedType: bool
esqlQueryDelaySeconds:
description: "Seconds ES|QL rules search behind now, so unsearchable events aren't missed. Delays alerts by the same amount. Set at least the longest index refresh interval. ES|QL only."
global: True
advanced: True
forcedType: int
helpLink: sigma
esqlCorrelationAllowanceSeconds:
description: "Extra seconds of arrivals each correlation run re-reads beyond its timespan, so a burst whose events arrive spread out is still counted together. ES|QL only."
global: True
advanced: True
forcedType: int
helpLink: sigma
elastic:
index:
description: Comma-separated list of indices or index patterns (wildcard "*" supported) that SOC will search for records.
@@ -874,14 +862,15 @@ soc:
global: True
forcedType: bool
automations:
description: Scheduled automations for the Onion AI assistant, managed from the Agent Studio.
global: True
advanced: True
readonlyUi: True
storage: db
forcedType: string
syntax: json
helpLink: onion-ai
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
+1 -1
View File
@@ -67,7 +67,7 @@ log_has_errors() {
grep -vE "Reading first line of patchfile" | \
grep -vE "Command failed with exit code" | \
grep -vE "Running scope as unit" | \
grep -vE "securityonion-resources/sigma/" | \
grep -vE "securityonion-resources/sigma/stable" | \
grep -vE "remove_failed_vm.sls" | \
grep -vE "failed to copy: httpReadSeeker" | \
grep -vE "Error response from daemon: failed to resolve reference" | \