mirror of
https://github.com/Security-Onion-Solutions/securityonion.git
synced 2026-10-06 22:44:48 +02:00
Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
cad18f5bfc |
No files matched your search
@@ -131,6 +131,8 @@ def beacon(config): # noqa: C901
|
||||
'setting_id': setting_id,
|
||||
'node_id': node_id,
|
||||
})
|
||||
log.info('postgres_pillar_beacon: audit_settings id=%d setting_id=%s node_id=%s',
|
||||
row_id, setting_id, node_id)
|
||||
if row_id > max_id:
|
||||
max_id = row_id
|
||||
|
||||
|
||||
@@ -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()
|
||||
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
]
|
||||
}
|
||||
@@ -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"
|
||||
}
|
||||
}
|
||||
@@ -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,22 @@ 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
|
||||
|
||||
# Lead-in salt puts on the comment of an orchestration step that raised.
|
||||
STEP_RAISED = 'An exception occurred in this state:'
|
||||
|
||||
# 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 +130,189 @@ 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 _failed_steps(ret):
|
||||
steps = []
|
||||
for job in ret.values() if isinstance(ret, dict) else []:
|
||||
job_ret = job.get('return') if isinstance(job, dict) else None
|
||||
data = job_ret.get('data') if isinstance(job_ret, dict) else None
|
||||
for group in data.values() if isinstance(data, dict) else []:
|
||||
if isinstance(group, dict):
|
||||
steps.extend(step for step in group.values() if isinstance(step, dict) and step.get('result') is False)
|
||||
return steps
|
||||
|
||||
|
||||
def _result_unknown(ret):
|
||||
# A step that raised (e.g. salt-master restarted while it waited on a queued state run)
|
||||
# never collected the minion's return, so the state run may still have completed.
|
||||
steps = _failed_steps(ret)
|
||||
return bool(steps) and all(str(step.get('comment', '')).startswith(STEP_RAISED) for step in steps)
|
||||
|
||||
|
||||
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 and _result_unknown(ret):
|
||||
log.warning('push result unknown jid=%s paths=%s; the orchestration lost track of the state run, which may '
|
||||
'still have completed (if not, the change will be applied at the next scheduled highstate): %s',
|
||||
jid, paths, ' | '.join(failures))
|
||||
elif 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 +335,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 +403,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,501 @@
|
||||
# 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)
|
||||
|
||||
RAISED = ('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')
|
||||
|
||||
RAISED_RET = _orch_ret(dict(REFRESH_STEP, **{
|
||||
'salt_|-apply_hydra_1_|-apply_hydra_1_|-state': {
|
||||
'__id__': 'apply_hydra_1', 'result': False, 'changes': {}, 'comment': RAISED,
|
||||
},
|
||||
}), 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 test_result_unknown(self):
|
||||
self.assertTrue(drainer._result_unknown(RAISED_RET))
|
||||
for ret in (CONFLICT_RET, STATE_FAIL_RET, SUCCESS_RET, ['No minions matched'], {MASTER: 'odd'},
|
||||
{MASTER: {'return': {'data': {MASTER: ['Rendering SLS failed']}}, 'success': False}}):
|
||||
self.assertFalse(drainer._result_unknown(ret), ret)
|
||||
mixed = _orch_ret(dict(RAISED_RET[MASTER]['return']['data'][MASTER],
|
||||
**CONFLICT_RET[MASTER]['return']['data'][MASTER]), success=False)
|
||||
self.assertFalse(drainer._result_unknown(mixed))
|
||||
|
||||
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,
|
||||
'6_unknown': RAISED_RET,
|
||||
}
|
||||
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', '6_unknown'):
|
||||
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'))
|
||||
self.assertIn('push result unknown jid=6_unknown', self.logged('warning'))
|
||||
self.assertIn('AuthenticationError: Authentication error occurred.', self.logged('warning'))
|
||||
self.assertNotIn('6_unknown', self.logged('error'))
|
||||
|
||||
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()
|
||||
@@ -1183,11 +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:")
|
||||
|
||||
@@ -3,6 +3,8 @@
|
||||
{% set BATCH = AUTOAPPLY.batch %}
|
||||
{% set BATCH_WAIT = AUTOAPPLY.batch_wait %}
|
||||
|
||||
{# queue must be a top-level salt.state arg (kwarg is ignored); an int is max_queue and still fails on conflict #}
|
||||
|
||||
{% for action in actions %}
|
||||
{% if action.get('highstate') %}
|
||||
apply_highstate_{{ loop.index }}:
|
||||
@@ -12,8 +14,7 @@ apply_highstate_{{ loop.index }}:
|
||||
- highstate: True
|
||||
- batch: {{ action.get('batch', BATCH) }}
|
||||
- batch_wait: {{ action.get('batch_wait', BATCH_WAIT) }}
|
||||
- kwarg:
|
||||
queue: 2
|
||||
- queue: True
|
||||
{% else %}
|
||||
refresh_pillar_{{ loop.index }}:
|
||||
salt.function:
|
||||
@@ -29,8 +30,7 @@ apply_{{ action.state | replace('.', '_') }}_{{ loop.index }}:
|
||||
- {{ action.state }}
|
||||
- batch: {{ action.get('batch', BATCH) }}
|
||||
- batch_wait: {{ action.get('batch_wait', BATCH_WAIT) }}
|
||||
- kwarg:
|
||||
queue: 2
|
||||
- queue: True
|
||||
- require:
|
||||
- salt: refresh_pillar_{{ loop.index }}
|
||||
{% endif %}
|
||||
|
||||
@@ -138,6 +138,7 @@ def run():
|
||||
# top level so the reactor is robust to either shape.
|
||||
event = data.get('data', data) # noqa: F821 -- data provided by reactor
|
||||
setting_id = event.get('setting_id', '')
|
||||
audit_id = event.get('id')
|
||||
node_id = (event.get('node_id') or '').strip()
|
||||
|
||||
app = _app_from_setting(setting_id)
|
||||
@@ -150,8 +151,8 @@ def run():
|
||||
if not entry:
|
||||
LOG.warning(
|
||||
'push_pillar: app "%s" is not in pillar_push_map.yaml; change will be '
|
||||
'picked up at the next scheduled highstate (setting_id=%s)',
|
||||
app, setting_id,
|
||||
'picked up at the next scheduled highstate (setting_id=%s audit_id=%s)',
|
||||
app, setting_id, audit_id,
|
||||
)
|
||||
return {}
|
||||
|
||||
@@ -165,12 +166,12 @@ def run():
|
||||
'node_{}_{}'.format(node_id, app), actions,
|
||||
'audit:{}@{}'.format(setting_id, node_id),
|
||||
)
|
||||
LOG.info('push_pillar: per-node intent updated for %s on %s (setting_id=%s)',
|
||||
app, node_id, setting_id)
|
||||
LOG.info('push_pillar: per-node intent updated for %s on %s (setting_id=%s audit_id=%s)',
|
||||
app, node_id, setting_id, audit_id)
|
||||
return {}
|
||||
|
||||
# Branch B: grid-wide app change -> use the map entry's actions as-is.
|
||||
actions = list(entry) # copy to avoid mutating the cache
|
||||
_write_intent('pillar_{}'.format(app), actions, 'audit:{}'.format(setting_id))
|
||||
LOG.info('push_pillar: app intent updated for %s (setting_id=%s)', app, setting_id)
|
||||
LOG.info('push_pillar: app intent updated for %s (setting_id=%s audit_id=%s)', app, setting_id, audit_id)
|
||||
return {}
|
||||
+4
-61
@@ -1445,7 +1445,7 @@ soc:
|
||||
default:
|
||||
- repo: https://github.com/Security-Onion-Solutions/securityonion-resources
|
||||
license: Elastic-2.0
|
||||
folder: sigma
|
||||
folder: sigma/stable
|
||||
community: true
|
||||
rulesetName: securityonion-resources
|
||||
- repo: file:///nsm/rules/custom-local-repos/local-sigma
|
||||
@@ -1455,7 +1455,7 @@ soc:
|
||||
airgap:
|
||||
- repo: file:///nsm/rules/detect-sigma/repos/securityonion-resources
|
||||
license: Elastic-2.0
|
||||
folder: sigma
|
||||
folder: sigma/stable
|
||||
community: true
|
||||
rulesetName: securityonion-resources
|
||||
- repo: file:///nsm/rules/custom-local-repos/local-sigma
|
||||
@@ -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}
|
||||
|
||||
@@ -14,15 +14,6 @@ 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
|
||||
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
|
||||
@@ -127,9 +118,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:
|
||||
|
||||
@@ -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,12 +80,6 @@
|
||||
{% do SOCMERGED.config.server.update({'airgapEnabled': false}) %}
|
||||
{% endif %}
|
||||
|
||||
{# correlation authoring requires ES|QL #}
|
||||
{% if not SOCMERGED.config.server.modules.elastalertengine.useEsql %}
|
||||
{% 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({
|
||||
|
||||
+1
-13
@@ -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.
|
||||
|
||||
+1
-1
@@ -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" | \
|
||||
|
||||
Reference in new issue
Block a user