Compare commits

..
Author SHA1 Message Date
Josh Patterson 9732e1c639 Trim tracebacks in push failure log lines
When an orchestration step raises, salt returns the full traceback as the
step comment, and the drainer wrote it verbatim, putting ~70 lines into
so-push-drainer.log per failure. Collapse comments to one line and, for
tracebacks, keep only the lead-in and the raised exception, e.g.
"apply_soc_1: An exception occurred in this state:
salt.exceptions.AuthenticationError: Authentication error occurred."

Seen on a standalone when a pushed highstate restarted salt-master while
two queued pushes were waiting: their orchestrations lost the master
connection and failed with AuthenticationError, although the minion
completed both state runs.
2026-09-30 16:18:37 -04:00
Josh Patterson 53f9ebcd46 FIX: queue auto-applied state runs instead of failing on conflict
orch.push_batch passed `kwarg: {queue: 2}` to salt.state, but in Salt
3006 queue is a top-level salt.state argument and salt.state always sets
the minion's queue kwarg from it (default False), so the kwarg block was
silently dropped and every pushed state ran with queue=False. The drainer
dispatches a separate async orchestration each 15s pass, so settings saved
more than ~15s apart overlap on the same minion and every run after the
first fails immediately with 'The function "state.sls" is running as PID
...'. The change then waits for the next scheduled highstate.

Seen on a 3.4.0 standalone: hydra.enabled, telegraf.output, and two soc
settings (including soc.config.licenseKey) were saved within 30s. The soc
state was dispatched while the telegraf state was still running and was
rejected, so the license key was not applied.

Use `queue: True`, as orch.deploy_newnode already does. An int is treated
as max_queue and still falls through to the conflict error once that many
state runs are active.

The failure was only visible in the master log, since the drainer
dispatches with --async and logged only "dispatch accepted". The drainer
now:
  - logs each dispatched action
  - parses the orchestration jid from salt-run's stderr (the only place
    --async reports it) and records it under /opt/so/state/push_dispatched
  - on later passes looks each jid up with jobs.lookup_jid and logs either
    "push succeeded" or an ERROR with the failed step, the per-minion
    failed states or rejection text, and the triggering paths
Lookups run outside the pending-intent lock since the reactors share it.

The beacon now logs each audit_settings row it emits and the reactor logs
the audit row id, so a single change can be traced from audit_settings to
its push result.

Adds so-push-drainer_test.py; the drainer is now held to the 100% coverage
requirement in python-test.

Verified on the standalone: a soc push dispatched while a 90s state run
was in progress queued behind it (queue=True in the job args), completed,
and the drainer logged "push succeeded" for its jid. The new result
parsing reports the original soc conflict and the hydra license failure
from the job cache.
2026-09-30 16:18:37 -04:00
coreyogburn 47d74f1ae1 Merge pull request #16270 from Security-Onion-Solutions/cogburn/automation
New Automation Fields
2026-09-30 11:10:51 -06:00
Corey Ogburn 855716846a New Automation Fields 2026-09-29 16:50:04 -06:00
Mike Reeves 8e35d70595 Merge pull request #16266 from Security-Onion-Solutions/mreeves/soai-context-1m
Raise SOAI Sonnet default small context limit to 1M
2026-09-29 12:19:20 -04:00
Jason Ertel e4625cfcae Merge pull request #16267 from Security-Onion-Solutions/jertel/wip
resolve startup errors
2026-09-29 12:07:26 -04:00
Jason Ertel eb803dce0e resolve startup errors 2026-09-29 12:02:48 -04:00
Mike Reeves a06f08217a Raise SOAI Sonnet default small context limit to 1M
Context is now flat-priced, so match contextLimitSmall to contextLimitLarge.
With equal limits the SOC assistant hides the increase-context toggle.
2026-09-29 11:42:29 -04:00
Josh Patterson 29d27cf255 Merge pull request #16261 from Security-Onion-Solutions/fix/telegraf-drop-docker-socket
FIX: remove the docker socket from so-telegraf
2026-09-29 10:46:54 -04:00
Jason Ertel 235a60e587 Merge pull request #16264 from Security-Onion-Solutions/jertel/wip
Metric alarms and more NTF annotations
2026-09-29 08:31:28 -04:00
Jason Ertel a8f7c46b0d Merge branch '3/dev' into jertel/wip 2026-09-28 13:47:19 -04:00
Jason Ertel 26d895ccb7 alarms and ntf 2026-09-28 13:47:16 -04:00
Josh Patterson b43efc458f Merge pull request #16263 from Security-Onion-Solutions/fix/service-account-nologin
FIX: use /sbin/nologin for service accounts
2026-09-28 13:05:26 -04:00
Josh Patterson 21222ff119 FIX: use /sbin/nologin for service accounts
These accounts existed only for container UID mapping and filesystem
ownership, but user.present omitted shell:, so Salt fell through to the
platform useradd default and every one of them got /bin/bash. Pin them to
/sbin/nologin so none can be used as an interactive login or `su -` target.

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

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

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

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

Closes #16256
2026-09-25 09:23:09 -04:00
Mike Reeves 88fa7e7fb4 Merge pull request #16262 from Security-Onion-Solutions/TOoSmOotH-patch-4
Add openai_embeddings to the YAML configuration
2026-09-24 15:29:17 -04:00
Mike Reeves efe0581892 Add openai_embeddings to the YAML configuration 2026-09-24 15:27:51 -04:00
Josh Patterson 72f60fcaa9 FIX: remove the docker socket from so-telegraf
so-telegraf mounted /var/run/docker.sock and joined the host docker group on
every node type. The :ro flag blocks write() to the inode, not connect() plus
HTTP over the socket, so any code execution inside the container could reach
POST /containers/create with Privileged:true and become root on the host.

No telegraf script used the socket; the only consumer was the native
[[inputs.docker]] plugin, and group_add 920 existed solely to feed it.
Container metrics now come from so-container-stats, a collector that runs on
the host from cron and writes influx line protocol to a file telegraf already
had mounted. This is the pattern so-status, so-raid-status and
so-elasticagent-status already use, so the privileged docker access stays on
the host side where root cron already ran it.

The collector runs as somon, a service account in the docker group with no
login shell and a locked password, rather than root. Docker group membership
is still root-equivalent on the host, so this is defense in depth rather than
a privilege boundary.

The telegraf scripts were root:939 mode 770, letting socore rewrite them for
code execution inside the container; they are now 750, which still allows the
read and execute telegraf needs. The container also ran with no group, giving
it gid 0, and now runs as 939:939. That alone would have broken
lasthighstate.sh, which reached /opt/so/log/salt only via the root group and
could not tell an unreadable file from a missing one, so it silently reported
a 56 year highstate age. Only the lasthighstate file is bind mounted now, and
the script tests readability instead of existence.

Everything inputs.docker collected beyond the five fields the shipped
dashboards query is available per stat under telegraf:container_stats,
annotated for SOC so an operator can enable it without editing files. Defaults
reproduce the previous output exactly. With every stat enabled the emitted
field set matches what inputs.docker wrote, verified by running the plugin
against the live socket and diffing: 53 fields, no type mismatches, no field
present on one side only.

Two deliberate differences: max_usage carries the real cgroup peak where the
daemon reports 0 on cgroup v2, and host-network containers emit no
docker_container_net row, matching inputs.docker. Docker label tags are not
restored, since nothing queries them and they cost significant cardinality.

so-status, the influxdb size cron, so-elasticagent-status, so-raid-status and
so-common-status-check truncated their output in place while telegraf read it,
so telegraf periodically saw an empty file and logged a parse error or emitted
empty values. They now write aside and rename. Measured on a live manager, the
old so-status cron left status.log empty for 215 of 10997 reads.

Tested on a fresh install, a converted grid and a 3.0 upgrade.
2026-09-24 14:30:25 -04:00
Jorge Reyes 2684a5ca95 Merge pull request #16260 from Security-Onion-Solutions/reyesj2/16254
FIX: Fleet scripts failing when endpoints-initial is missing
2026-09-23 08:58:15 -05:00
reyesj2 bcee63bde5 remove policy precheck 2026-09-23 08:51:11 -05:00
reyesj2 6ce89eb323 jq -e exits 0 with empty input 2026-09-22 21:25:50 -05:00
reyesj2 ebab4b0d90 split between missing token and multiple enrollment tokens 2026-09-22 16:11:56 -05:00
reyesj2 db60c27da2 fail cleanly when endpoints-initial is missing or has multiple enrollment tokens 2026-09-22 12:24:23 -05:00
Jason Ertel f4defdfde0 Merge pull request #16258 from Security-Onion-Solutions/jertel/wip
upgrade whoislookup analyzer deps; notification annotations
2026-09-22 08:41:04 -04:00
Jason Ertel 36652e8f23 fixed malformed desc annotation 2026-09-22 08:30:44 -04:00
Jason Ertel 7bef194540 force transitive dep version 2026-09-22 08:28:44 -04:00
Jason Ertel e2bf2837fe notification annotations 2026-09-22 08:18:59 -04:00
Jason Ertel 06704dad22 upgrade whoislookup analyzer deps 2026-09-22 08:17:18 -04:00
Mike Reeves 47fe0758d0 Merge pull request #16255 from Security-Onion-Solutions/mreeves/fix-so-user-stdin-drain
Fix user creation failing with a valid password
2026-09-18 17:12:09 -04:00
58 changed files with 2432 additions and 98 deletions

No files matched your search

+25
View File
@@ -6,6 +6,9 @@ on:
- "salt/sensoroni/files/analyzers/**"
- "salt/manager/tools/sbin/**"
- "salt/_beacons/**"
- "salt/telegraf/tools/sbin_jinja/**"
- "salt/telegraf/defaults.yaml"
- "salt/telegraf/soc_telegraf.yaml"
jobs:
build:
@@ -34,3 +37,25 @@ jobs:
- name: Test with pytest
run: |
PYTHONPATH=${{ matrix.python-code-path }} pytest ${{ matrix.python-code-path }} --cov=${{ matrix.python-code-path }} --doctest-modules --cov-report=term --cov-fail-under=100 --cov-config=pytest.ini
telegraf-collector:
# so-container-stats is a jinja template rather than an importable module, so it gets its
# own job: the test renders it the way salt does, then drives it with a faked docker engine
# and cgroup tree. No container runtime is needed.
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v3
- name: Set up Python
uses: actions/setup-python@v3
with:
python-version: "3.14"
- name: Install dependencies
run: |
python -m pip install --upgrade pip
python -m pip install flake8 pytest jinja2 pyyaml
- name: Lint with flake8
run: |
flake8 salt/telegraf/tools/sbin_jinja/so-container-stats_test.py --config=pytest.ini
- name: Test with pytest
run: |
pytest salt/telegraf/tools/sbin_jinja/so-container-stats_test.py -v
+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
+3 -1
View File
@@ -217,9 +217,11 @@ sostatus_log:
- replace: False
# Install sostatus check cron. This is used to populate Grid.
# telegraf reads status.log on the same minute boundary this runs, so write aside and rename
# rather than truncating the file it is reading
so-status_check_cron:
cron.present:
- name: '/usr/sbin/so-status -j > /opt/so/log/sostatus/status.log 2>&1'
- name: '/usr/sbin/so-status -j > /opt/so/log/sostatus/status.log.tmp 2>&1; mv -f /opt/so/log/sostatus/status.log.tmp /opt/so/log/sostatus/status.log'
- identifier: so-status_check_cron
- user: root
- minute: '*/1'
+19 -6
View File
@@ -9,6 +9,7 @@ import sys
import subprocess
import os
import json
import tempfile
sys.path.append('/opt/saltstack/salt/lib/python3.10/site-packages/')
import salt.config
@@ -17,6 +18,21 @@ import salt.loader
__opts__ = salt.config.minion_config('/etc/salt/minion')
__grains__ = salt.loader.grains(__opts__)
def write_atomic(path, value):
# telegraf reads these files on its own schedule; replacing them by rename means it never
# reads a truncated file and reports an empty value as if it were real
directory = os.path.dirname(path)
handle, temp = tempfile.mkstemp(dir=directory)
try:
with os.fdopen(handle, 'w') as f:
f.write(str(value))
os.chmod(temp, 0o644)
os.replace(temp, path)
except Exception:
os.path.exists(temp) and os.unlink(temp)
raise
def check_needs_restarted():
osfam = __grains__['os_family']
val = '0'
@@ -34,8 +50,7 @@ def check_needs_restarted():
else:
fail("Unsupported OS")
with open(outfile, 'w') as f:
f.write(val)
write_atomic(outfile, val)
def check_for_fps():
feat = 'fps'
@@ -56,8 +71,7 @@ def check_for_fps():
# Unknown, so assume 0
fps = 0
with open('/opt/so/log/sostatus/fps_enabled', 'w') as f:
f.write(str(fps))
write_atomic('/opt/so/log/sostatus/fps_enabled', fps)
def check_for_lks():
feat = 'Lks'
@@ -80,8 +94,7 @@ def check_for_lks():
lks = 1
if lks:
break
with open('/opt/so/log/sostatus/lks_enabled', 'w') as f:
f.write(str(lks))
write_atomic('/opt/so/log/sostatus/lks_enabled', lks)
def fail(msg):
print(msg, file=sys.stderr)
+1
View File
@@ -177,6 +177,7 @@ if [[ $EXCLUDE_FALSE_POSITIVE_ERRORS == 'Y' ]]; then
EXCLUDED_ERRORS="$EXCLUDED_ERRORS|Unexpected authorization header" # expected WARN log lines indicating invalid auth header
EXCLUDED_ERRORS="$EXCLUDED_ERRORS|Missing ory_kratos_session cookie" # expected WARN log lines indicating invalid auth header
EXCLUDED_ERRORS="$EXCLUDED_ERRORS|Static assets preprocessor only supports GET and HEAD requests" # expected WARN log lines indicating invalid auth header
EXCLUDED_ERRORS="$EXCLUDED_ERRORS|respondError" # respondError is a function name, output via http middleware as standard request logging
fi
if [[ $EXCLUDE_KNOWN_ERRORS == 'Y' ]]; then
+3 -1
View File
@@ -125,4 +125,6 @@ else
RAIDSTATUS=1
fi
echo "nsmraid=$RAIDSTATUS" > /opt/so/log/raid/status.log
# telegraf reads this file; write aside and rename so it never sees a half-written file
echo "nsmraid=$RAIDSTATUS" > /opt/so/log/raid/status.log.tmp
mv -f /opt/so/log/raid/status.log.tmp /opt/so/log/raid/status.log
+1
View File
@@ -21,6 +21,7 @@ elastalert:
- gid: 933
- home: /opt/so/conf/elastalert
- createhome: False
- shell: /sbin/nologin
elastalogdir:
file.directory:
@@ -19,6 +19,7 @@ elastic-agent-pr:
- gid: 948
- home: /opt/so/conf/elastic-fleet-pr
- createhome: False
- shell: /sbin/nologin
{% else %}
+1
View File
@@ -20,6 +20,7 @@ elastic-agent:
- gid: 949
- home: /opt/so/conf/elastic-agent
- createhome: False
- shell: /sbin/nologin
elasticagentconfdir:
file.directory:
+1
View File
@@ -26,6 +26,7 @@ elastic-fleet:
- gid: 947
- home: /opt/so/conf/elastic-fleet
- createhome: False
- shell: /sbin/nologin
elasticfleet_sbin:
file.recurse:
@@ -30,6 +30,56 @@ fleet_api() {
curl -sK /opt/so/conf/elasticsearch/curl.config -L "localhost:5601/api/fleet/${QUERYPATH}" "$@" --retry 3 --retry-delay 10 --fail 2>/dev/null
}
elastic_fleet_require_agent_policy() {
local AGENT_POLICY=$1
local POLICY_JSON
if ! POLICY_JSON=$(fleet_api "agent_policies/$AGENT_POLICY") || [ -z "$POLICY_JSON" ]; then
echo "Error: Agent policy '$AGENT_POLICY' was not found or is not visible in the current Kibana space." >&2
return 1
fi
if ! jq -e '.item.package_policies | type == "array"' <<<"$POLICY_JSON" >/dev/null 2>&1; then
echo "Error: Agent policy '$AGENT_POLICY' was not found or is not visible in the current Kibana space." >&2
return 1
fi
echo "$POLICY_JSON"
}
# Print the single active enrollment token for POLICY_ID.
# Exit 1: retryable (API failure, invalid response, no active token)
# Exit 2: multiple active tokens - Shouldn't get into this state without manual intervention
elastic_fleet_active_enrollment_token() {
local POLICY_ID=$1
local RESP TOKEN_COUNT API_KEY
if ! RESP=$(fleet_api "enrollment_api_keys?perPage=100" -H 'kbn-xsrf: true' -H 'Content-Type: application/json'); then
echo "Error: Failed to retrieve enrollment tokens for agent policy '$POLICY_ID'." >&2
return 1
fi
if ! jq -e '.list' <<<"$RESP" >/dev/null 2>&1; then
echo "Error: Invalid enrollment token response for agent policy '$POLICY_ID'." >&2
return 1
fi
TOKEN_COUNT=$(jq --arg pid "$POLICY_ID" '[.list[] | select(.policy_id == $pid and .active == true)] | length' <<<"$RESP")
if [ "${TOKEN_COUNT:-0}" -eq 0 ]; then
echo "Error: No active enrollment token found for agent policy '$POLICY_ID'." >&2
return 1
fi
if [ "$TOKEN_COUNT" -gt 1 ]; then
echo "Error: Found $TOKEN_COUNT active enrollment tokens for agent policy '$POLICY_ID'; expected exactly one." >&2
return 2
fi
API_KEY=$(jq -r --arg pid "$POLICY_ID" '.list[] | select(.policy_id == $pid and .active == true) | .api_key' <<<"$RESP")
echo "$API_KEY"
}
# Max number of concurrent Fleet write jobs (create/update). Override via env if needed.
MAX_FLEET_JOBS=${MAX_FLEET_JOBS:-10}
@@ -62,15 +112,7 @@ elastic_fleet_load_integrations_dir() {
i=0
# Fetch the agent policy a single time; we look up integration ids locally below.
if ! POLICY_JSON=$(fleet_api "agent_policies/$AGENT_POLICY"); then
echo "Error: Failed to retrieve agent policy '$AGENT_POLICY'."
rm -f "$FAIL_FILE"
rm -rf "$OUT_DIR"
return 1
fi
if ! jq -e '.item.package_policies' <<<"$POLICY_JSON" >/dev/null 2>&1; then
echo "Error: Invalid agent policy response for '$AGENT_POLICY'."
if ! POLICY_JSON=$(elastic_fleet_require_agent_policy "$AGENT_POLICY"); then
rm -f "$FAIL_FILE"
rm -rf "$OUT_DIR"
return 1
@@ -124,9 +166,15 @@ elastic_fleet_integration_check() {
JSON_STRING=$2
NAME=$(jq -r .name $JSON_STRING)
NAME=$(jq -r .name "$JSON_STRING")
INTEGRATION_ID=""
INTEGRATION_ID=$(/usr/sbin/so-elastic-fleet-agent-policy-view "$AGENT_POLICY" | jq -r '.item.package_policies[] | select(.name=="'"$NAME"'") | .id')
local POLICY_JSON
if ! POLICY_JSON=$(elastic_fleet_require_agent_policy "$AGENT_POLICY"); then
return 1
fi
INTEGRATION_ID=$(jq -r --arg n "$NAME" '.item.package_policies[]? | select(.name==$n) | .id' <<<"$POLICY_JSON")
}
@@ -148,7 +196,16 @@ elastic_fleet_integration_remove() {
NAME=$2
INTEGRATION_ID=$(/usr/sbin/so-elastic-fleet-agent-policy-view "$AGENT_POLICY" | jq -r '.item.package_policies[] | select(.name=="'"$NAME"'") | .id')
local POLICY_JSON
if ! POLICY_JSON=$(elastic_fleet_require_agent_policy "$AGENT_POLICY"); then
return 1
fi
INTEGRATION_ID=$(jq -r --arg n "$NAME" '.item.package_policies[]? | select(.name==$n) | .id' <<<"$POLICY_JSON")
if [ -z "$INTEGRATION_ID" ]; then
echo "Error: Integration '$NAME' was not found in agent policy '$AGENT_POLICY'." >&2
return 1
fi
JSON_STRING=$( jq -n \
--arg INTEGRATIONID "$INTEGRATION_ID" \
@@ -13,7 +13,10 @@ ERROR=false
for INTEGRATION in /opt/so/conf/elastic-fleet/integrations/elastic-defend/*.json
do
printf "\n\nInitial Endpoints Policy - Loading $INTEGRATION\n"
elastic_fleet_integration_check "endpoints-initial" "$INTEGRATION"
if ! elastic_fleet_integration_check "endpoints-initial" "$INTEGRATION"; then
ERROR=true
continue
fi
if [ -n "$INTEGRATION_ID" ]; then
printf "\n\nIntegration $NAME exists - Upgrading integration policy\n"
if ! elastic_fleet_integration_policy_upgrade "$INTEGRATION_ID"; then
@@ -7,20 +7,35 @@
. /usr/sbin/so-elastic-fleet-common
# Get all the fleet policies
json_output=$(curl -s -K /opt/so/conf/elasticsearch/curl.config -L -X GET "localhost:5601/api/fleet/agent_policies" -H 'kbn-xsrf: true')
if ! json_output=$(fleet_api "agent_policies" -H 'kbn-xsrf: true'); then
echo "Error: Failed to retrieve Fleet agent policies." >&2
exit 1
fi
if ! jq -e '.items' <<<"$json_output" >/dev/null 2>&1; then
echo "Error: Invalid Fleet agent policies response." >&2
exit 1
fi
# Extract the IDs that start with "FleetServer_"
POLICY=$(echo "$json_output" | jq -r '.items[] | select(.id | startswith("FleetServer_")) | .id')
POLICY=$(jq -r '.items[] | select(.id | startswith("FleetServer_")) | .id' <<<"$json_output")
# Iterate over each ID in the POLICY variable
for POLICYNAME in $POLICY; do
printf "\nUpdating Policy: $POLICYNAME\n"
# First get the Integration ID
INTEGRATION_ID=$(/usr/sbin/so-elastic-fleet-agent-policy-view "$POLICYNAME" | jq -r '.item.package_policies[] | select(.package.name == "fleet_server") | .id')
if ! POLICY_JSON=$(elastic_fleet_require_agent_policy "$POLICYNAME"); then
exit 1
fi
INTEGRATION_ID=$(jq -r '.item.package_policies[]? | select(.package.name == "fleet_server") | .id' <<<"$POLICY_JSON")
if [ -z "$INTEGRATION_ID" ]; then
echo "Error: fleet_server integration was not found in agent policy '$POLICYNAME'." >&2
exit 1
fi
# Modify the default integration policy to update the policy_id and an with the correct naming
UPDATED_INTEGRATION_POLICY=$(jq --arg policy_id "$POLICYNAME" --arg name "fleet_server-$POLICYNAME" '
UPDATED_INTEGRATION_POLICY=$(jq --arg policy_id "$POLICYNAME" --arg name "fleet_server-$POLICYNAME" '
.policy_id = $policy_id |
.name = $name' /opt/so/conf/elastic-fleet/integrations/fleet-server/fleet-server.json)
@@ -22,12 +22,19 @@ NUM_RUNNING=$(pgrep -cf "/bin/bash /sbin/so-elastic-agent-gen-installers")
for i in {1..30}
do
ENROLLMENTOKEN=$(curl -K /opt/so/conf/elasticsearch/curl.config -L "localhost:5601/api/fleet/enrollment_api_keys?perPage=100" -H 'kbn-xsrf: true' -H 'Content-Type: application/json' | jq .list | jq -r -c '.[] | select(.policy_id | contains("endpoints-initial")) | .api_key')
ENROLLMENTOKEN=$(elastic_fleet_active_enrollment_token "endpoints-initial")
TOKEN_RC=$?
if [ "$TOKEN_RC" -eq 2 ]; then
exit 1
fi
FLEETHOST=$(curl -K /opt/so/conf/elasticsearch/curl.config 'http://localhost:5601/api/fleet/fleet_server_hosts/grid-default' | jq -r '.item.host_urls[]' | paste -sd ',')
if [[ $FLEETHOST ]] && [[ $ENROLLMENTOKEN ]]; then break; else sleep 10; fi
if [[ -n "$FLEETHOST" ]] && [[ -n "$ENROLLMENTOKEN" ]]; then
break
fi
sleep 10
done
if [[ -z $FLEETHOST ]] || [[ -z $ENROLLMENTOKEN ]]; then
if [[ -z "$FLEETHOST" ]] || [[ -z "$ENROLLMENTOKEN" ]]; then
printf "\nFleet Host URL, Enrollment Token or Elastic Version empty - exiting..."
printf "\nFleet Host: $FLEETHOST, Enrollment Token: $ENROLLMENTOKEN\n"
exit 1
@@ -67,19 +74,25 @@ for GOOS in "${GOTARGETOS[@]}"; do
GOARCH="amd64"
if [[ $GOOS == 'darwin/arm64' ]]; then GOOS="darwin" && GOARCH="arm64"; fi
printf "\n\n### Generating $GOOS/$GOARCH Installer...\n"
docker run -e CGO_ENABLED=0 -e GOOS=$GOOS -e GOARCH=$GOARCH \
if ! docker run -e CGO_ENABLED=0 -e GOOS=$GOOS -e GOARCH=$GOARCH \
--mount type=bind,source=/etc/pki/tls/certs/,target=/workspace/files/cert/ \
--mount type=bind,source=/nsm/elastic-agent-workspace/,target=/workspace/files/elastic-agent/ \
--mount type=bind,source=/opt/so/saltstack/local/salt/elasticfleet/files/,target=/output/ \
{{ GLOBALS.registry_host }}:5000/{{ GLOBALS.image_repo }}/so-elastic-agent-builder:{{ GLOBALS.so_version }} go build -ldflags "-X main.fleetHostURLsList=$FLEETHOST -X main.enrollmentToken=$ENROLLMENTOKEN" -o /output/so-elastic-agent_${GOOS}_${GOARCH}
{{ GLOBALS.registry_host }}:5000/{{ GLOBALS.image_repo }}/so-elastic-agent-builder:{{ GLOBALS.so_version }} go build -ldflags "-X main.fleetHostURLsList=$FLEETHOST -X main.enrollmentToken=$ENROLLMENTOKEN" -o /output/so-elastic-agent_${GOOS}_${GOARCH}; then
printf "\n### ERROR: Failed to generate $GOOS/$GOARCH installer. Exiting...\n"
exit 1
fi
printf "\n### $GOOS/$GOARCH Installer Generated...\n"
done
printf "\n\n### Generating MSI...\n"
cp /opt/so/saltstack/local/salt/elasticfleet/files/so-elastic-agent_windows_amd64 /opt/so/saltstack/local/salt/elasticfleet/files/so-elastic-agent_windows_amd64.exe
docker run \
if ! docker run \
--mount type=bind,source=/opt/so/saltstack/local/salt/elasticfleet/files/,target=/output/ -w /output \
{{ GLOBALS.registry_host }}:5000/{{ GLOBALS.image_repo }}/so-elastic-agent-builder:{{ GLOBALS.so_version }} wixl -o so-elastic-agent_windows_amd64_msi --arch x64 /workspace/so-elastic-agent.wxs
{{ GLOBALS.registry_host }}:5000/{{ GLOBALS.image_repo }}/so-elastic-agent-builder:{{ GLOBALS.so_version }} wixl -o so-elastic-agent_windows_amd64_msi --arch x64 /workspace/so-elastic-agent.wxs; then
printf "\n### ERROR: Failed to generate MSI. Exiting...\n"
exit 1
fi
printf "\n### MSI Generated...\n"
# Verify installers were created
@@ -202,26 +202,9 @@ fi
### Finalization ###
# Query for Enrollment Tokens for default policies
if ENDPOINTSENROLLMENTOKEN_RAW=$(fleet_api "enrollment_api_keys" -H 'kbn-xsrf: true' -H 'Content-Type: application/json'); then
ENDPOINTSENROLLMENTOKEN=$(echo "$ENDPOINTSENROLLMENTOKEN_RAW" | jq .list | jq -r -c '.[] | select(.policy_id | contains("endpoints-initial")) | .api_key')
else
echo -e "\nFailed to query for Endpoints enrollment token"
exit 1
fi
if GRIDNODESENROLLMENTOKENGENERAL_RAW=$(fleet_api "enrollment_api_keys" -H 'kbn-xsrf: true' -H 'Content-Type: application/json'); then
GRIDNODESENROLLMENTOKENGENERAL=$(echo "$GRIDNODESENROLLMENTOKENGENERAL_RAW" | jq .list | jq -r -c '.[] | select(.policy_id | contains("so-grid-nodes_general")) | .api_key')
else
echo -e "\nFailed to query for Grid nodes - General enrollment token"
exit 1
fi
if GRIDNODESENROLLMENTOKENHEAVY_RAW=$(fleet_api "enrollment_api_keys" -H 'kbn-xsrf: true' -H 'Content-Type: application/json'); then
GRIDNODESENROLLMENTOKENHEAVY=$(echo "$GRIDNODESENROLLMENTOKENHEAVY_RAW" | jq .list | jq -r -c '.[] | select(.policy_id | contains("so-grid-nodes_heavy")) | .api_key')
else
echo -e "\nFailed to query for Grid nodes - Heavy enrollment token"
exit 1
fi
ENDPOINTSENROLLMENTOKEN=$(elastic_fleet_active_enrollment_token "endpoints-initial") || exit 1
GRIDNODESENROLLMENTOKENGENERAL=$(elastic_fleet_active_enrollment_token "so-grid-nodes_general") || exit 1
GRIDNODESENROLLMENTOKENHEAVY=$(elastic_fleet_active_enrollment_token "so-grid-nodes_heavy") || exit 1
# Store needed data in minion pillar
pillar_file=/opt/so/saltstack/local/pillar/minions/{{ GLOBALS.minion_id }}.sls
+1
View File
@@ -32,6 +32,7 @@ elasticsearch:
- gid: 930
- home: /opt/so/conf/elasticsearch
- createhome: False
- shell: /sbin/nologin
elasticsearch_sbin:
file.recurse:
+3 -1
View File
@@ -94,9 +94,11 @@ metrics_link_file:
- docker_container: so-influxdb
# Install cron job to determine size of influxdb for telegraf
# telegraf reads this while the cron rewrites it, so write aside and rename rather than
# truncating in place. tgraflogdir recurses ownership, so the temp file is chowned to match
get_influxdb_size:
cron.present:
- name: 'du -s -k /nsm/influxdb | cut -f1 > /opt/so/log/telegraf/influxdb_size.log 2>&1'
- name: 'du -s -k /nsm/influxdb | cut -f1 > /opt/so/log/telegraf/influxdb_size.log.tmp 2>&1; chown 939:939 /opt/so/log/telegraf/influxdb_size.log.tmp; mv -f /opt/so/log/telegraf/influxdb_size.log.tmp /opt/so/log/telegraf/influxdb_size.log'
- identifier: get_influxdb_size
- user: root
- minute: '*/1'
+1
View File
@@ -21,6 +21,7 @@ kafka_user:
- gid: 960
- home: /opt/so/conf/kafka
- createhome: False
- shell: /sbin/nologin
kafka_home_dir:
file.absent:
+1
View File
@@ -22,6 +22,7 @@ kibana:
- gid: 932
- home: /opt/so/conf/kibana
- createhome: False
- shell: /sbin/nologin
# Drop the correct nginx config based on role
+1
View File
@@ -27,6 +27,7 @@ kratos:
- uid: 928
- gid: 928
- home: /opt/so/conf/kratos
- shell: /sbin/nologin
kratosdir:
file.directory:
+1
View File
@@ -35,6 +35,7 @@ logstash:
- uid: 931
- gid: 931
- home: /opt/so/conf/logstash
- shell: /sbin/nologin
logstash_sbin:
file.recurse:
+1 -1
View File
@@ -166,7 +166,7 @@ so-repo-sync:
so_fleetagent_status:
cron.present:
- name: /usr/sbin/so-elasticagent-status > /opt/so/log/agents/agentstatus.log 2>&1
- name: '/usr/sbin/so-elasticagent-status > /opt/so/log/agents/agentstatus.log.tmp 2>&1; mv -f /opt/so/log/agents/agentstatus.log.tmp /opt/so/log/agents/agentstatus.log'
- identifier: so_fleetagent_status
- user: root
- minute: '*/5'
+143 -6
View File
@@ -19,6 +19,8 @@ is older than debounce_seconds, this script:
* dispatches a single `salt-run state.orchestrate orch.push_batch --async`
with the deduped actions list passed as pillar kwargs
* deletes the contributed intent files on successful dispatch
* records the orchestration jid under /opt/so/state/push_dispatched and, on
later passes, looks up its result and logs success or per-minion failures
Reactor sls files (push_files, push_pillar) write intents
but never dispatch directly
@@ -30,6 +32,7 @@ import json
import logging
import logging.handlers
import os
import re
import subprocess
import sys
import time
@@ -40,8 +43,18 @@ PENDING_DIR = '/opt/so/state/push_pending'
LOCK_FILE = os.path.join(PENDING_DIR, '.lock')
LOG_FILE = '/opt/so/log/salt/so-push-drainer.log'
DISPATCHED_DIR = '/opt/so/state/push_dispatched'
HIGHSTATE_SENTINEL = '__highstate__'
RESULT_CHECK_DELAY = 30
RESULT_MAX_AGE = 7200
RESULT_CHECKS_PER_PASS = 5
TEXT_LIMIT = 500
# salt-run --async reports the jid only in a log line on stderr.
JID_RE = re.compile(r'salt/run/(\d{20})')
def _make_logger():
logger = logging.getLogger('so-push-drainer')
@@ -113,15 +126,130 @@ def _dispatch(actions, log):
except subprocess.CalledProcessError as exc:
log.error('dispatch failed (rc=%s): stdout=%s stderr=%s',
exc.returncode, exc.stdout, exc.stderr)
return False
return None
except subprocess.TimeoutExpired:
log.error('dispatch timed out after 60s')
return False
return None
except Exception:
log.exception('dispatch raised')
return False
log.info('dispatch accepted: %s', (result.stdout or '').strip())
return True
return None
match = JID_RE.search(result.stderr or '')
if not match:
log.warning('dispatch accepted but no jid found, result will not be tracked: stderr=%s',
_trim(result.stderr))
return ''
log.info('dispatch accepted: jid=%s', match.group(1))
return match.group(1)
def _trim(value):
text = value if isinstance(value, str) else json.dumps(value, default=str)
lines = [line.strip() for line in text.splitlines() if line.strip()]
if 'Traceback (most recent call last):' in text:
# Keep the lead-in and the raised exception; the frames are noise in a log line.
lines = [text.split('Traceback (most recent call last):', 1)[0].strip(), lines[-1]]
text = ' '.join(line for line in lines if line)
return text if len(text) <= TEXT_LIMIT else text[:TEXT_LIMIT] + '...'
def _unlink(path, log):
try:
os.unlink(path)
except OSError:
log.exception('failed to remove %s', path)
def _record_dispatch(jid, actions, paths, log):
record = {'jid': jid, 'dispatched_at': time.time(), 'actions': actions, 'paths': paths}
path = os.path.join(DISPATCHED_DIR, '{}.json'.format(jid))
try:
os.makedirs(DISPATCHED_DIR, exist_ok=True)
tmp_path = path + '.tmp'
with open(tmp_path, 'w') as f:
json.dump(record, f)
os.rename(tmp_path, path)
except OSError:
log.exception('failed to record dispatch %s', jid)
def _lookup_jid(jid, log):
"""Returns the job cache entry for jid, {} while it is still running, or None on error."""
cmd = ['salt-run', 'jobs.lookup_jid', jid, '--out=json']
try:
result = subprocess.run(cmd, check=True, capture_output=True, text=True, timeout=60)
return json.loads(result.stdout or '{}')
except (subprocess.CalledProcessError, subprocess.TimeoutExpired, ValueError) as exc:
log.warning('lookup of jid %s failed: %s', jid, exc)
return None
def _minion_failure(minion_ret):
if isinstance(minion_ret, dict):
return '; '.join(
'{}: {}'.format(state.get('__id__', state_key), _trim(state.get('comment', '')))
for state_key, state in minion_ret.items()
if isinstance(state, dict) and state.get('result') is False
)
# A state run rejected before it starts (e.g. another state run is in
# progress) returns a list of error strings instead of state results.
if isinstance(minion_ret, (list, str)):
return _trim(minion_ret)
return ''
def _orch_failures(ret):
failures = []
for job in ret.values():
if not isinstance(job, dict):
continue
job_ret = job.get('return')
data = (job_ret.get('data') or {}) if isinstance(job_ret, dict) else {}
for steps in data.values():
if not isinstance(steps, dict):
failures.append(_trim(steps))
continue
for step in steps.values():
if not isinstance(step, dict) or step.get('result') is not False:
continue
failures.append('{}: {}'.format(step.get('__id__', step.get('name')), _trim(step.get('comment', ''))))
minion_rets = (step.get('changes') or {}).get('ret') or {}
for minion, minion_ret in minion_rets.items():
text = _minion_failure(minion_ret)
if text:
failures.append('{}: {}'.format(minion, text))
if job.get('success') is False and not failures:
failures.append('orchestration reported failure: {}'.format(_trim(job.get('return'))))
return failures
def _check_dispatched(log, now):
checked = 0
for path in sorted(glob.glob(os.path.join(DISPATCHED_DIR, '*.json'))):
if checked >= RESULT_CHECKS_PER_PASS:
break
record = _read_intent(path, log)
if not isinstance(record, dict) or not record.get('jid'):
_unlink(path, log)
continue
age = now - record.get('dispatched_at', 0)
if age < RESULT_CHECK_DELAY:
continue
checked += 1
jid = record['jid']
paths = record.get('paths', [])
ret = _lookup_jid(jid, log)
if not ret:
if age > RESULT_MAX_AGE:
log.warning('no result for jid=%s after %ds, no longer tracking; paths=%s', jid, age, paths)
_unlink(path, log)
continue
failures = _orch_failures(ret)
if failures:
log.error('push failed jid=%s paths=%s; change will be applied at the next scheduled highstate: %s',
jid, paths, ' | '.join(failures))
else:
log.info('push succeeded jid=%s paths=%s', jid, paths)
_unlink(path, log)
def main():
@@ -143,6 +271,9 @@ def main():
debounce_seconds = int(push.get('debounce_seconds', 30))
# Outside the lock: lookups are slow and the reactors take the same lock.
_check_dispatched(log, time.time())
os.makedirs(PENDING_DIR, exist_ok=True)
lock_fd = os.open(LOCK_FILE, os.O_CREAT | os.O_RDWR, 0o644)
try:
@@ -208,10 +339,16 @@ def main():
len(ready), len(deduped), len(combined_actions),
debounce_duration, all_paths[:20],
)
for action in deduped:
log.info('action: %s tgt=%s', 'highstate' if action.get('highstate') else action.get('state'),
action.get('tgt'))
if not _dispatch(deduped, log):
jid = _dispatch(deduped, log)
if jid is None:
log.warning('dispatch failed; leaving intent files in place for retry')
return 1
if jid:
_record_dispatch(jid, deduped, all_paths[:20], log)
for path, _ in ready:
try:
@@ -0,0 +1,400 @@
# Copyright Security Onion Solutions LLC and/or licensed to Security Onion Solutions LLC under one
# or more contributor license agreements. Licensed under the Elastic License 2.0 as shown at
# https://securityonion.net/license; you may not use this file except in compliance with the
# Elastic License 2.0.
import importlib.util
import json
import logging
import os
import shutil
import subprocess
import sys
import tempfile
import time
import unittest
from importlib.machinery import SourceFileLoader
from unittest.mock import MagicMock, patch
# salt is not installed where these tests run; the drainer only needs salt.client.Caller.
_salt = MagicMock()
sys.modules.setdefault('salt', _salt)
sys.modules.setdefault('salt.client', _salt.client)
HERE = os.path.dirname(os.path.abspath(__file__))
SCRIPT = os.path.join(HERE, 'so-push-drainer')
_loader = SourceFileLoader('so_push_drainer', SCRIPT)
_spec = importlib.util.spec_from_loader('so_push_drainer', _loader)
drainer = importlib.util.module_from_spec(_spec)
_loader.exec_module(drainer)
MASTER = 'manager.localdomain_master'
JID = '20260930171554259426'
ASYNC_STDERR = ('[WARNING ] Running in asynchronous mode. Results of this execution may be collected '
'by attaching to the master event bus or by examining the master job cache, if '
'configured. This execution is running under tag salt/run/{}\n'.format(JID))
CONFLICT = ('The function "state.sls" is running as PID 372218 and was started at '
'2026, Sep 30 17:15:40.466233 with jid 20260930171540466233')
def _orch_ret(steps, success=True):
return {MASTER: {
'fun': 'runner.state.orchestrate',
'jid': JID,
'return': {'data': {MASTER: steps}, 'outputter': 'highstate', 'retcode': 0 if success else 1},
'success': success,
}}
REFRESH_STEP = {
'salt_|-refresh_pillar_1_|-saltutil.refresh_pillar_|-function': {
'__id__': 'refresh_pillar_1', 'result': True,
'changes': {'ret': {'manager_standalone': True}},
'comment': 'Function ran successfully.',
},
}
CONFLICT_RET = _orch_ret(dict(REFRESH_STEP, **{
'salt_|-apply_soc_1_|-apply_soc_1_|-state': {
'__id__': 'apply_soc_1', 'result': False,
'changes': {'out': 'highstate', 'ret': {'manager_standalone': [CONFLICT]}},
'comment': 'Run failed on minions: manager_standalone',
},
}), success=False)
STATE_FAIL_RET = _orch_ret({
'salt_|-apply_hydra_1_|-apply_hydra_1_|-state': {
'__id__': 'apply_hydra_1', 'result': False,
'changes': {'out': 'highstate', 'ret': {'manager_standalone': {
'test_|-no_license_|-no_license_|-fail_without_changes': {
'__id__': 'hydra.enabled_no_license_detected', 'result': False,
'comment': 'This is a feature supported only for customers with a valid license.',
},
'file_|-hydra_conf_|-/opt/so/conf/hydra_|-managed': {'result': True, 'comment': 'ok'},
}}},
'comment': 'Run failed on minions: manager_standalone',
},
}, success=False)
SUCCESS_RET = _orch_ret(dict(REFRESH_STEP, **{
'salt_|-apply_telegraf_1_|-apply_telegraf_1_|-state': {
'__id__': 'apply_telegraf_1', 'result': True,
'changes': {'out': 'highstate', 'ret': {'manager_standalone': {
'file_|-tgrafconf_|-/opt/so/conf/telegraf/etc/telegraf.conf_|-managed': {'result': True},
}}},
'comment': 'States ran successfully.',
},
}))
class DrainerTestCase(unittest.TestCase):
def setUp(self):
self.tmpdir = tempfile.mkdtemp()
self.pending = os.path.join(self.tmpdir, 'push_pending')
self.dispatched = os.path.join(self.tmpdir, 'push_dispatched')
os.makedirs(self.pending)
for name, value in (
('PENDING_DIR', self.pending),
('LOCK_FILE', os.path.join(self.pending, '.lock')),
('DISPATCHED_DIR', self.dispatched),
('LOG_FILE', os.path.join(self.tmpdir, 'log', 'so-push-drainer.log')),
):
patcher = patch.object(drainer, name, value)
patcher.start()
self.addCleanup(patcher.stop)
self.log = MagicMock()
def tearDown(self):
shutil.rmtree(self.tmpdir, ignore_errors=True)
def write_json(self, directory, name, data):
os.makedirs(directory, exist_ok=True)
path = os.path.join(directory, name)
with open(path, 'w') as f:
if isinstance(data, str):
f.write(data)
else:
json.dump(data, f)
return path
def logged(self, level):
return ' '.join(c.args[0] % c.args[1:] for c in getattr(self.log, level).call_args_list)
class TestHelpers(DrainerTestCase):
def test_make_logger_adds_handler_once(self):
logger = logging.getLogger('so-push-drainer')
self.addCleanup(logger.handlers.clear)
logger.handlers.clear()
self.assertIs(drainer._make_logger(), logger)
drainer._make_logger()
self.assertEqual(len(logger.handlers), 1)
self.assertTrue(os.path.isdir(os.path.dirname(drainer.LOG_FILE)))
def test_load_push_cfg(self):
with patch.object(drainer.salt.client, 'Caller') as caller:
caller.return_value.cmd.return_value = {'enabled': False}
self.assertEqual(drainer._load_push_cfg(), {'enabled': False})
caller.return_value.cmd.return_value = 'garbage'
self.assertEqual(drainer._load_push_cfg(), {})
def test_read_intent(self):
good = self.write_json(self.pending, 'good.json', {'a': 1})
bad = self.write_json(self.pending, 'bad.json', '{nope')
self.assertEqual(drainer._read_intent(good, self.log), {'a': 1})
self.assertIsNone(drainer._read_intent(bad, self.log))
with patch('builtins.open', side_effect=RuntimeError('boom')):
self.assertIsNone(drainer._read_intent(good, self.log))
self.log.exception.assert_called_once()
def test_dedupe_actions(self):
actions = [
'not a dict',
{'state': 'soc'},
{'state': 'soc', 'tgt': '*'},
{'state': 'soc', 'tgt': '*', 'tgt_type': 'compound'},
{'highstate': True, 'tgt': '*'},
{'state': 'soc', 'tgt': 'node1', 'tgt_type': 'glob'},
]
self.assertEqual(drainer._dedupe_actions(actions), [actions[2], actions[4], actions[5]])
def test_trim(self):
self.assertEqual(drainer._trim(' text \n'), 'text')
self.assertEqual(drainer._trim(['a']), '["a"]')
self.assertEqual(drainer._trim(None), 'null')
self.assertEqual(drainer._trim('x' * 600), 'x' * drainer.TEXT_LIMIT + '...')
def test_trim_traceback(self):
comment = ('An exception occurred in this state: Traceback (most recent call last):\n'
' File "salt/client/__init__.py", line 1934, in pub\n'
' raise AuthenticationError(err_msg)\n'
'salt.exceptions.AuthenticationError: Authentication error occurred.\n')
self.assertEqual(drainer._trim(comment), 'An exception occurred in this state: '
'salt.exceptions.AuthenticationError: Authentication error occurred.')
self.assertEqual(drainer._trim('line one\n line two\n'), 'line one line two')
def test_unlink_missing_logs(self):
drainer._unlink(os.path.join(self.tmpdir, 'missing'), self.log)
self.log.exception.assert_called_once()
class TestDispatch(DrainerTestCase):
def run_dispatch(self, **kwargs):
with patch.object(drainer.subprocess, 'run', **kwargs) as run:
jid = drainer._dispatch([{'state': 'soc', 'tgt': '*'}], self.log)
return jid, run
def test_jid_parsed_from_stderr(self):
jid, run = self.run_dispatch(return_value=MagicMock(stdout='', stderr=ASYNC_STDERR))
self.assertEqual(jid, JID)
cmd = run.call_args[0][0]
self.assertEqual(cmd[:3], ['salt-run', 'state.orchestrate', 'orch.push_batch'])
self.assertIn('--async', cmd)
def test_no_jid(self):
jid, _ = self.run_dispatch(return_value=MagicMock(stdout='', stderr=None))
self.assertEqual(jid, '')
self.log.warning.assert_called_once()
def test_failures_return_none(self):
for exc in (subprocess.CalledProcessError(1, 'salt-run', 'out', 'err'),
subprocess.TimeoutExpired('salt-run', 60),
RuntimeError('boom')):
jid, _ = self.run_dispatch(side_effect=exc)
self.assertIsNone(jid)
def test_record_dispatch(self):
drainer._record_dispatch(JID, [{'state': 'soc'}], ['audit:soc.config.licenseKey'], self.log)
with open(os.path.join(self.dispatched, JID + '.json')) as f:
record = json.load(f)
self.assertEqual(record['jid'], JID)
self.assertEqual(record['paths'], ['audit:soc.config.licenseKey'])
self.assertIn('dispatched_at', record)
def test_record_dispatch_oserror(self):
with patch.object(drainer.os, 'makedirs', side_effect=OSError('ro')):
drainer._record_dispatch(JID, [], [], self.log)
self.log.exception.assert_called_once()
class TestResults(DrainerTestCase):
def test_lookup_jid(self):
with patch.object(drainer.subprocess, 'run') as run:
run.return_value = MagicMock(stdout=json.dumps(SUCCESS_RET))
self.assertEqual(drainer._lookup_jid(JID, self.log), SUCCESS_RET)
self.assertEqual(run.call_args[0][0], ['salt-run', 'jobs.lookup_jid', JID, '--out=json'])
run.return_value = MagicMock(stdout='')
self.assertEqual(drainer._lookup_jid(JID, self.log), {})
run.return_value = MagicMock(stdout='not json')
self.assertIsNone(drainer._lookup_jid(JID, self.log))
run.side_effect = subprocess.TimeoutExpired('salt-run', 60)
self.assertIsNone(drainer._lookup_jid(JID, self.log))
def test_minion_failure_shapes(self):
self.assertEqual(drainer._minion_failure([CONFLICT]), json.dumps([CONFLICT]))
self.assertEqual(drainer._minion_failure('Rendering SLS failed'), 'Rendering SLS failed')
self.assertEqual(drainer._minion_failure(True), '')
self.assertEqual(drainer._minion_failure({'a': {'result': True}}), '')
def test_orch_failures_conflict(self):
failures = drainer._orch_failures(CONFLICT_RET)
self.assertEqual(failures[0], 'apply_soc_1: Run failed on minions: manager_standalone')
self.assertIn('manager_standalone', failures[1])
self.assertIn('is running as PID 372218', failures[1])
self.assertEqual(len(failures), 2)
def test_orch_failures_failed_state(self):
failures = drainer._orch_failures(STATE_FAIL_RET)
self.assertEqual(len(failures), 2)
self.assertIn('hydra.enabled_no_license_detected: This is a feature', failures[1])
self.assertNotIn('hydra_conf', failures[1])
def test_orch_failures_success(self):
self.assertEqual(drainer._orch_failures(SUCCESS_RET), [])
def test_orch_failures_render_error(self):
ret = {MASTER: {'return': {'data': {MASTER: ['Rendering SLS failed']}}, 'success': False}}
self.assertEqual(drainer._orch_failures(ret), ['["Rendering SLS failed"]'])
def test_orch_failures_unparsed(self):
self.assertEqual(drainer._orch_failures({MASTER: 'odd'}), [])
ret = {MASTER: {'return': 'Exception occurred', 'success': False}}
self.assertEqual(drainer._orch_failures(ret), ['orchestration reported failure: Exception occurred'])
def record(self, jid, age, now):
return self.write_json(self.dispatched, jid + '.json', {
'jid': jid, 'dispatched_at': now - age, 'actions': [], 'paths': ['audit:' + jid],
})
def test_check_dispatched(self):
now = time.time()
results = {
'1_failed': CONFLICT_RET,
'2_ok': SUCCESS_RET,
'3_pending': {},
'4_expired': None,
}
young = self.record('0_young', 5, now)
paths = {jid: self.record(jid, 60, now) for jid in results}
paths['4_expired'] = self.record('4_expired', drainer.RESULT_MAX_AGE + 1, now)
bad = self.write_json(self.dispatched, '5_bad.json', '{nope')
with patch.object(drainer, '_lookup_jid', side_effect=lambda jid, log: results[jid]):
drainer._check_dispatched(self.log, now)
self.assertTrue(os.path.exists(young))
self.assertTrue(os.path.exists(paths['3_pending']))
for jid in ('1_failed', '2_ok', '4_expired'):
self.assertFalse(os.path.exists(paths[jid]), jid)
self.assertFalse(os.path.exists(bad))
self.assertIn('push failed jid=1_failed', self.logged('error'))
self.assertIn('is running as PID 372218', self.logged('error'))
self.assertIn('push succeeded jid=2_ok', self.logged('info'))
self.assertIn('no result for jid=4_expired', self.logged('warning'))
def test_check_dispatched_limit(self):
now = time.time()
for i in range(drainer.RESULT_CHECKS_PER_PASS + 2):
self.record('{:02d}'.format(i), 60, now)
with patch.object(drainer, '_lookup_jid', return_value={}) as lookup:
drainer._check_dispatched(self.log, now)
self.assertEqual(lookup.call_count, drainer.RESULT_CHECKS_PER_PASS)
class TestMain(DrainerTestCase):
def setUp(self):
super().setUp()
self.cfg = {'enabled': True, 'debounce_seconds': 30}
for name, kwargs in (
('_make_logger', {'return_value': self.log}),
('_load_push_cfg', {'side_effect': lambda: self.cfg}),
('_check_dispatched', {}),
):
patcher = patch.object(drainer, name, **kwargs)
setattr(self, name, patcher.start())
self.addCleanup(patcher.stop)
def intent(self, name, age=60, actions=None, paths=None):
now = time.time()
return self.write_json(self.pending, name, {
'first_touch': now - age - 5, 'last_touch': now - age,
'actions': [{'state': 'soc', 'tgt': '*'}] if actions is None else actions,
'paths': paths or ['audit:soc.config.licenseKey'],
})
def test_no_pending_dir(self):
shutil.rmtree(self.pending)
self.assertEqual(drainer.main(), 0)
self._load_push_cfg.assert_not_called()
def test_cfg_error(self):
self._load_push_cfg.side_effect = RuntimeError('no salt')
self.assertEqual(drainer.main(), 1)
def test_disabled(self):
self.cfg['enabled'] = False
self.assertEqual(drainer.main(), 0)
self._check_dispatched.assert_not_called()
def test_no_intents_still_checks_results(self):
self.assertEqual(drainer.main(), 0)
self._check_dispatched.assert_called_once()
def test_debounce_and_broken(self):
young = self.intent('young.json', age=1)
broken = self.write_json(self.pending, 'broken.json', '{nope')
with patch.object(drainer, '_dispatch') as dispatch:
self.assertEqual(drainer.main(), 0)
dispatch.assert_not_called()
self.assertTrue(os.path.exists(young))
self.assertFalse(os.path.exists(broken))
def test_broken_unlink_error_ignored(self):
self.write_json(self.pending, 'broken.json', '{nope')
with patch.object(drainer.os, 'unlink', side_effect=OSError('busy')):
self.assertEqual(drainer.main(), 0)
def test_no_usable_actions(self):
path = self.intent('empty.json', actions=[{'state': 'soc'}])
self.assertEqual(drainer.main(), 0)
self.assertFalse(os.path.exists(path))
self.intent('empty.json', actions=[{'state': 'soc'}])
with patch.object(drainer.os, 'unlink', side_effect=OSError('busy')):
self.assertEqual(drainer.main(), 0)
def test_dispatch_failure_keeps_intents(self):
path = self.intent('pillar_soc.json')
with patch.object(drainer, '_dispatch', return_value=None):
self.assertEqual(drainer.main(), 1)
self.assertTrue(os.path.exists(path))
def test_dispatch_records_jid(self):
soc = self.intent('pillar_soc.json')
hs = self.intent('pillar_global.json', actions=[{'highstate': True, 'tgt': '*'}], paths=['audit:global.x'])
with patch.object(drainer, '_dispatch', return_value=JID) as dispatch, \
patch.object(drainer, '_record_dispatch') as record:
self.assertEqual(drainer.main(), 0)
self.assertEqual(len(dispatch.call_args[0][0]), 2)
record.assert_called_once()
self.assertEqual(record.call_args[0][0], JID)
self.assertEqual(sorted(record.call_args[0][2]), ['audit:global.x', 'audit:soc.config.licenseKey'])
self.assertFalse(os.path.exists(soc))
self.assertFalse(os.path.exists(hs))
self.assertIn('action: highstate tgt=*', self.logged('info'))
def test_dispatch_without_jid_not_recorded(self):
self.intent('pillar_soc.json')
with patch.object(drainer, '_dispatch', return_value=''), \
patch.object(drainer, '_record_dispatch') as record, \
patch.object(drainer.os, 'unlink', side_effect=OSError('busy')):
self.assertEqual(drainer.main(), 0)
record.assert_not_called()
self.log.exception.assert_called_once()
if __name__ == '__main__':
unittest.main()
+4 -4
View File
@@ -3,6 +3,8 @@
{% set BATCH = AUTOAPPLY.batch %}
{% set BATCH_WAIT = AUTOAPPLY.batch_wait %}
{# queue must be a top-level salt.state arg (kwarg is ignored); an int is max_queue and still fails on conflict #}
{% for action in actions %}
{% if action.get('highstate') %}
apply_highstate_{{ loop.index }}:
@@ -12,8 +14,7 @@ apply_highstate_{{ loop.index }}:
- highstate: True
- batch: {{ action.get('batch', BATCH) }}
- batch_wait: {{ action.get('batch_wait', BATCH_WAIT) }}
- kwarg:
queue: 2
- queue: True
{% else %}
refresh_pillar_{{ loop.index }}:
salt.function:
@@ -29,8 +30,7 @@ apply_{{ action.state | replace('.', '_') }}_{{ loop.index }}:
- {{ action.state }}
- batch: {{ action.get('batch', BATCH) }}
- batch_wait: {{ action.get('batch_wait', BATCH_WAIT) }}
- kwarg:
queue: 2
- queue: True
- require:
- salt: refresh_pillar_{{ loop.index }}
{% endif %}
+6 -5
View File
@@ -138,6 +138,7 @@ def run():
# top level so the reactor is robust to either shape.
event = data.get('data', data) # noqa: F821 -- data provided by reactor
setting_id = event.get('setting_id', '')
audit_id = event.get('id')
node_id = (event.get('node_id') or '').strip()
app = _app_from_setting(setting_id)
@@ -150,8 +151,8 @@ def run():
if not entry:
LOG.warning(
'push_pillar: app "%s" is not in pillar_push_map.yaml; change will be '
'picked up at the next scheduled highstate (setting_id=%s)',
app, setting_id,
'picked up at the next scheduled highstate (setting_id=%s audit_id=%s)',
app, setting_id, audit_id,
)
return {}
@@ -165,12 +166,12 @@ def run():
'node_{}_{}'.format(node_id, app), actions,
'audit:{}@{}'.format(setting_id, node_id),
)
LOG.info('push_pillar: per-node intent updated for %s on %s (setting_id=%s)',
app, node_id, setting_id)
LOG.info('push_pillar: per-node intent updated for %s on %s (setting_id=%s audit_id=%s)',
app, node_id, setting_id, audit_id)
return {}
# Branch B: grid-wide app change -> use the map entry's actions as-is.
actions = list(entry) # copy to avoid mutating the cache
_write_intent('pillar_{}'.format(app), actions, 'audit:{}'.format(setting_id))
LOG.info('push_pillar: app intent updated for %s (setting_id=%s)', app, setting_id)
LOG.info('push_pillar: app intent updated for %s (setting_id=%s audit_id=%s)', app, setting_id, audit_id)
return {}
@@ -1,2 +1,3 @@
requests>=2.31.0
whoisit>=2.7.0
requests>=2.34.0
whoisit>=4.0.5
anyio>=4.15.1
+12 -1
View File
@@ -1494,6 +1494,9 @@ soc:
org: Security Onion
bucket: telegraf/so_short_term
verifyCert: false
notification:
dismissedPruneDays: 30
enabled: true
playbook:
autoUpdateEnabled: true
playbookImportFrequencySeconds: 86400
@@ -1558,6 +1561,14 @@ soc:
reconcilePersona: ""
toolUseTurnAttempts: 12
toolUseTurnDelayMs: 175
agentSessionMaxTurns: 20
agentStreamFlushIntervalMs: 1000
agentStreamIdleTimeoutSeconds: 300
automationSettings:
tickIntervalSeconds: 60
maxConcurrentItems: 4
maxQueuedItems: 0
alertTriageEpoch: "2026-09-24T00:00:00Z"
tools:
filterEventFields:
- "@timestamp"
@@ -2796,7 +2807,7 @@ soc:
- id: sonnet
displayName: Claude Sonnet
origin: USA
contextLimitSmall: 200000
contextLimitSmall: 1000000
contextLimitLarge: 1000000
lowBalanceColorAlert: 500000
enabled: true
+108
View File
@@ -160,6 +160,7 @@ soc:
description: Schedules that are shared across the Security Onion product. Modify via one of the SOC Schedules view.
readonlyUi: True
global: True
advanced: True
forcedType: string
syntax: json
storage: db
@@ -495,9 +496,23 @@ soc:
description: JSON list of notifications. Modify via the SOC Notifications view.
readonlyUi: True
global: True
advanced: True
forcedType: string
syntax: json
storage: db
dismissedPruneDays:
title: Dismissed Retention Days
description: The number of days to retain dismissed notifications. When a notification is dismissed, it will be pruned after this many days. Only one user need dismiss a notification for it to be pruned.
forcedType: int
global: True
maxListLimit:
description: Maximum number of notifications to display.
forcedType: int
global: True
enabled:
description: Enables or disables the SOC notification module.
forcedType: bool
global: True
postgres:
host:
description: Hostname or IP address of the PostgreSQL server used by SOC. Defaults to the manager hostname.
@@ -524,6 +539,48 @@ soc:
global: True
sensitive: True
advanced: True
postgresmetrics:
host:
description: Hostname or IP address of the PostgreSQL server used by Telegraf. Defaults to the manager hostname.
global: True
advanced: True
port:
description: Port of the PostgreSQL server used by Telegraf.
global: True
advanced: True
sslMode:
description: "Use encrypted connections to the PostgreSQL server used by Telegraf. Must be one of the following values: disable, allow, prefer, require, verify-ca, verify-full."
global: True
advanced: True
database:
description: Database to authenticate to on the PostgreSQL server.
global: True
advanced: True
user:
description: Username to authenticate to the PostgreSQL server used by Telegraf.
global: True
advanced: True
password:
description: Password used to authenticate to the PostgreSQL server used by Telegraf.
global: True
sensitive: True
advanced: True
cacheExpirationMs:
description: The interval (in milliseconds) to wait before querying the DB for updated metrics.
global: True
advanced: True
maxMetricAgeSeconds:
description: The maximum age (in seconds) of metrics to display in the SOC Grid Metrics view. Metrics older than this value will not be displayed.
global: True
advanced: True
alarms:
description: JSON list of metric alarms. Modify via the SOC Grid Alarms view.
readonlyUi: True
advanced: True
global: True
forcedType: string
syntax: json
storage: db
salt:
longRelayTimeoutMs:
description: Duration (in milliseconds) to wait for a response from the Salt API when executing tasks known for being long running before giving up and showing an error on the SOC UI.
@@ -782,6 +839,7 @@ soc:
- gemini
- openai_responses
- openai_chat
- openai_embeddings
- field: apiUrl
label: API URL
required: False
@@ -803,6 +861,16 @@ soc:
description: Indicates if the Assistant Module should operate in agentic mode or not. If true, agents can work together to solve tasks.
global: True
forcedType: bool
automations:
template:
description: Scheduled automations for the Onion AI assistant, managed from the Agent Studio. Each automation is stored under its own generated ID so its history and rollback are independent of every other automation.
global: True
advanced: False
readonlyUi: True
duplicates: True
forcedType: string
syntax: json
helpLink: onion-ai
agents:
description: Agent definitions for the Onion AI assistant, managed from the Agent Studio. An entry naming a system agent overrides only the fields an admin may change; everything else comes from the built-in definition.
global: True
@@ -835,6 +903,9 @@ soc:
- field: persona
label: Persona
multiline: True
- field: maxConcurrentInstances
label: Max Concurrent Instances
forcedType: int
skills:
description: Skill definitions for the Onion AI assistant, managed from the Agent Studio. An entry naming a system skill overrides only its enabled state and persona addendum; its tool set comes from the built-in definition.
global: True
@@ -938,6 +1009,43 @@ soc:
description: The number of times to retry extracting memories from a session if errors occur.
global: True
advanced: True
agentSessionMaxTurns:
description: Maximum number of model turns a headless agent session, such as one started by an automation, may take before it is stopped. Turns taken by delegated sub-agents count toward this limit. A session that reaches the limit is recorded as failed.
global: True
advanced: True
forcedType: int
agentStreamFlushIntervalMs:
description: Milliseconds between writes of a streaming headless agent turn to the database. Lower values show progress sooner in the Agent Studio at the cost of more frequent Elasticsearch updates.
global: True
advanced: True
forcedType: int
agentStreamIdleTimeoutSeconds:
description: Seconds a streaming headless agent turn may go without receiving any output before it is abandoned and the session is recorded as failed. Set to 0 to disable the timeout.
global: True
advanced: True
forcedType: int
automationSettings:
tickIntervalSeconds:
description: How often, in seconds, the automation scheduler checks for automations that are due to run. Must be greater than 0.
global: True
advanced: True
forcedType: int
maxConcurrentItems:
description: Maximum number of automation work items that may run at the same time. Additional work items wait in the queue until a running item finishes. User chat sessions count toward this limit but are never held back by it. Set to 0 to disable the limit.
global: True
advanced: True
forcedType: int
maxQueuedItems:
description: Maximum number of automation work items that may wait to start. Once the queue is full, no new work items are created until the backlog drains. Set to 0 to disable the limit.
global: True
advanced: True
forcedType: int
alertTriageEpoch:
description: The earliest alert time the Alert Triage automation will consider. Alerts before this time are never triaged, which keeps a first run on an existing deployment from working through old history. Must be in UTC format (2026-09-24T00:00:00Z).
regex: '^(\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}(\.\d+)?Z)?$'
regexFailureMessage: Expecting date in RFC3339 format (2026-09-24T00:00:00Z)
global: True
advanced: True
tools:
filterEventFields:
description: A whitelist of fields to return when OnionAI uses the query_events tool. All other fields are removed. One field per line.
+1
View File
@@ -64,6 +64,7 @@ suricata:
- gid: 940
- home: /nsm/suricata
- createhome: False
- shell: /sbin/nologin
socoregroupwithsuricata:
group.present:
+90 -10
View File
@@ -36,7 +36,7 @@ tgraf_sync_script_{{script}}:
- name: /opt/so/conf/telegraf/scripts/{{script}}
- user: root
- group: 939
- mode: 770
- mode: 750
- template: jinja
- source: salt://telegraf/scripts/{{script}}
- defaults:
@@ -49,7 +49,7 @@ tgraf_sync_script_esindexsize.sh:
- name: /opt/so/conf/telegraf/scripts/esindexsize.sh
- user: root
- group: 939
- mode: 770
- mode: 750
- source: salt://telegraf/scripts/esindexsize.sh
{# Copy conf/elasticsearch/curl.config for telegraf to use with esindexsize.sh #}
tgraf_sync_escurl_conf:
@@ -61,6 +61,80 @@ tgraf_sync_escurl_conf:
- source: salt://elasticsearch/curl.config
{% endif %}
# so-container-stats runs on the host as somon, a docker group member, so the container does
# not need the docker socket
somongroup:
group.present:
- name: somon
- gid: 961
# cron chdirs to $HOME before running a job, so home must exist
somon:
user.present:
- uid: 961
- gid: 961
- home: /opt/so/log/somon
- createhome: False
- shell: /sbin/nologin
- groups:
- docker
# renumbering an existing somon is a no-op on a fresh host and lets a host created
# before the id changed converge instead of failing the whole telegraf state
- allow_uid_change: True
- allow_gid_change: True
- require:
- group: somongroup
somonlogdir:
file.directory:
- name: /opt/so/log/somon
- user: 961
- group: 961
- mode: 755
# the lock file is not otherwise managed; recurse so a renumber rechowns it too
- recurse:
- user
- group
- require:
- user: somon
containers_log:
file.managed:
- name: /opt/so/log/somon/containers.log
- user: 961
- group: 961
- mode: 644
- replace: False
- require:
- file: somonlogdir
# telegraf reads on the same minute boundary the collector runs, and docker stats takes
# seconds, so write aside and rename rather than truncating the file telegraf is reading.
# ; not && so a failed run replaces the file instead of leaving stale metrics behind.
# flock -n keeps a run that outlives its minute from racing the next one over the same tmp
# file; the skipped run leaves a stale containers.log, which containers.sh discards by age
so-container-stats_cron:
cron.present:
- name: "flock -n /opt/so/log/somon/containers.lock -c '/usr/sbin/so-container-stats > /opt/so/log/somon/containers.log.tmp 2>&1; mv -f /opt/so/log/somon/containers.log.tmp /opt/so/log/somon/containers.log'"
- identifier: so-container-stats_cron
- user: somon
- minute: '*/1'
- hour: '*'
- daymonth: '*'
- month: '*'
- dayweek: '*'
- require:
- user: somon
# salt.lasthighstate touches this at order 9001, after the container starts; pre-create it so
# docker does not create a directory at the bind mount source
lasthighstate_placeholder:
file.managed:
- name: /opt/so/log/salt/lasthighstate
- mode: 644
- replace: False
- makedirs: True
telegraf_sbin:
file.recurse:
- name: /usr/sbin
@@ -69,14 +143,20 @@ telegraf_sbin:
- group: root
- file_mode: 755
#telegraf_sbin_jinja:
# file.recurse:
# - name: /usr/sbin
# - source: salt://telegraf/tools/sbin_jinja
# - user: 939
# - group: 939
# - file_mode: 755
# - template: jinja
# so-container-stats needs the per-stat toggles, so it renders instead of copying
tgraf_sbin_jinja:
file.recurse:
- name: /usr/sbin
- source: salt://telegraf/tools/sbin_jinja
- user: root
- group: root
- file_mode: 755
# the unit test lives beside the script; it must not ship or be rendered as jinja
- exclude_pat:
- "*_test.py"
- template: jinja
- defaults:
CONTAINER_STATS: {{ TELEGRAFMERGED.container_stats }}
tgrafconf:
file.managed:
+77
View File
@@ -10,6 +10,69 @@ telegraf:
flush_jitter: '0s'
debug: false
quiet: false
container_stats:
tags:
identity: False
engine:
n_containers: False
n_containers_running: False
n_containers_stopped: False
n_containers_paused: False
n_images: False
n_cpus: False
n_goroutines: False
n_used_file_descriptors: False
n_listener_events: False
memory_total: False
cpu:
usage_percent: True
usage_total: False
usage_in_usermode: False
usage_in_kernelmode: False
usage_system: False
throttling_periods: False
throttling_throttled_periods: False
throttling_throttled_time: False
container_id: False
mem:
usage_percent: True
usage: False
limit: False
max_usage: False
active_anon: False
active_file: False
inactive_anon: False
inactive_file: False
unevictable: False
pgfault: False
pgmajfault: False
container_id: False
net:
rx_bytes: True
rx_packets: False
rx_errors: False
rx_dropped: False
tx_bytes: False
tx_packets: False
tx_errors: False
tx_dropped: False
container_id: False
blkio:
io_service_bytes_recursive_read: False
io_service_bytes_recursive_write: False
container_id: False
status:
uptime_ns: True
oomkilled: True
pid: False
exitcode: False
restart_count: False
started_at: False
finished_at: False
container_id: False
health:
health_status: False
failing_streak: False
scripts:
eval:
- agentstatus.sh
@@ -19,6 +82,7 @@ telegraf:
- oldpcap.sh
- os.sh
- raid.sh
- containers.sh
- sostatus.sh
- suriloss.sh
- surirules.sh
@@ -34,6 +98,7 @@ telegraf:
- os.sh
- raid.sh
- redis.sh
- containers.sh
- sostatus.sh
- suriloss.sh
- surirules.sh
@@ -47,6 +112,7 @@ telegraf:
- os.sh
- raid.sh
- redis.sh
- containers.sh
- sostatus.sh
- features.sh
managerhype:
@@ -56,6 +122,7 @@ telegraf:
- os.sh
- raid.sh
- redis.sh
- containers.sh
- sostatus.sh
- features.sh
managersearch:
@@ -66,12 +133,14 @@ telegraf:
- os.sh
- raid.sh
- redis.sh
- containers.sh
- sostatus.sh
- features.sh
import:
- influxdbsize.sh
- lasthighstate.sh
- os.sh
- containers.sh
- sostatus.sh
sensor:
- checkfiles.sh
@@ -79,6 +148,7 @@ telegraf:
- oldpcap.sh
- os.sh
- raid.sh
- containers.sh
- sostatus.sh
- suriloss.sh
- surirules.sh
@@ -93,6 +163,7 @@ telegraf:
- os.sh
- raid.sh
- redis.sh
- containers.sh
- sostatus.sh
- suriloss.sh
- surirules.sh
@@ -101,12 +172,14 @@ telegraf:
idh:
- lasthighstate.sh
- os.sh
- containers.sh
- sostatus.sh
searchnode:
- eps.sh
- lasthighstate.sh
- os.sh
- raid.sh
- containers.sh
- sostatus.sh
- features.sh
receiver:
@@ -115,16 +188,20 @@ telegraf:
- os.sh
- raid.sh
- redis.sh
- containers.sh
- sostatus.sh
fleet:
- lasthighstate.sh
- os.sh
- containers.sh
- sostatus.sh
hypervisor:
- lasthighstate.sh
- os.sh
- containers.sh
- sostatus.sh
desktop:
- lasthighstate.sh
- os.sh
- containers.sh
- sostatus.sh
+5
View File
@@ -13,6 +13,11 @@ so-telegraf:
docker_container.absent:
- force: True
so-container-stats_cron:
cron.absent:
- identifier: so-container-stats_cron
- user: somon
so-telegraf_so-status.disabled:
file.comment:
- name: /opt/so/conf/so-status/so-status.conf
+7 -4
View File
@@ -19,8 +19,7 @@ so-telegraf:
docker_container.running:
- image: {{ GLOBALS.registry_host }}:5000/{{ GLOBALS.image_repo }}/so-telegraf:{{ GLOBALS.so_version }}
- restart_policy: unless-stopped
- user: 939
- group_add: 939,920
- user: 939:939
- environment:
- HOST_ETC=/host/etc
- HOST_SYS=/host/sys
@@ -38,7 +37,6 @@ so-telegraf:
- /opt/so/conf/telegraf/etc/telegraf.conf:/etc/telegraf/telegraf.conf:ro
- /opt/so/conf/telegraf/node_config.json:/etc/telegraf/node_config.json:ro
- /var/run/utmp:/var/run/utmp:ro
- /var/run/docker.sock:/var/run/docker.sock:ro
- /:/host:ro
- /sys:/host/sys:ro
- /proc:/host/proc:ro
@@ -51,7 +49,8 @@ so-telegraf:
- /opt/so/log/suricata:/var/log/suricata:ro
- /opt/so/log/raid:/var/log/raid:ro
- /opt/so/log/sostatus:/var/log/sostatus:ro
- /opt/so/log/salt:/var/log/salt:ro
- /opt/so/log/somon:/var/log/somon:ro
- /opt/so/log/salt/lasthighstate:/var/log/salt/lasthighstate:ro
- /opt/so/log/agents:/var/log/agents:ro
{% if GLOBALS.is_manager or GLOBALS.role == 'so-heavynode' %}
- /opt/so/conf/telegraf/etc/escurl.config:/etc/telegraf/elasticsearch.config:ro
@@ -74,6 +73,8 @@ so-telegraf:
{% endfor %}
{% endif %}
- watch:
- file: tgraf_sbin_jinja
- file: lasthighstate_placeholder
- file: trusttheca
- x509: telegraf_crt
- x509: telegraf_key
@@ -83,6 +84,8 @@ so-telegraf:
- file: tgraf_sync_script_{{script}}
{% endfor %}
- require:
- file: lasthighstate_placeholder
- file: somonlogdir
- file: trusttheca
- x509: telegraf_crt
- x509: telegraf_key
+11 -7
View File
@@ -228,13 +228,6 @@
# ## bond interfaces.
# # bond_interfaces = ["bond0"]
# # Read metrics about docker containers
[[inputs.docker]]
# ## Docker Endpoint
# ## To use TCP, set endpoint = "tcp://[ip]:[port]"
# ## To use environment variables (ie, docker-machine), set endpoint = "ENV"
endpoint = "unix:///var/run/docker.sock"
#
# # Read stats from one or more Elasticsearch servers or clusters
{%- if GLOBALS.is_manager or GLOBALS.role == 'so-heavynode' %}
@@ -342,6 +335,17 @@
interval = "60s"
{%- endif %}
{%- if 'containers.sh' in TELEGRAFMERGED.scripts[GLOBALS.role.split('-')[1]] %}
{%- do TELEGRAFMERGED.scripts[GLOBALS.role.split('-')[1]].remove('containers.sh') %}
[[inputs.exec]]
commands = [
["/scripts/containers.sh"]
]
data_format = "influx"
timeout = "15s"
interval = "60s"
{%- endif %}
{%- if TELEGRAFMERGED.scripts[GLOBALS.role.split('-')[1]] | length > 0 %}
[[inputs.exec]]
commands = [
+29
View File
@@ -0,0 +1,29 @@
#!/bin/bash
#
# 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.
# if this script isn't already running
if [[ ! "`pidof -x $(basename $0) -o %PPID`" ]]; then
CONTAINERSLOG=/var/log/somon/containers.log
# the collector rewrites this every minute; report nothing rather than repeating a stale
# file as if it were current, in case a run was skipped or the collector is wedged
MAXAGE=150
if [ -r "$CONTAINERSLOG" ]; then
AGE=$(( $(date +%s) - $(stat -c %Y "$CONTAINERSLOG") ))
if [ "$AGE" -le "$MAXAGE" ]; then
cat $CONTAINERSLOG
fi
fi
exit 0
fi
exit 0
+6 -4
View File
@@ -8,10 +8,12 @@
# if this script isn't already running
if [[ ! "`pidof -x $(basename $0) -o %PPID`" ]]; then
LAST_HIGHSTATE_END=$([ -e "/var/log/salt/lasthighstate" ] && date -r /var/log/salt/lasthighstate +%s || echo 0)
NOW=$(date +%s)
HIGHSTATE_AGE_SECONDS=$((NOW-LAST_HIGHSTATE_END))
echo "salt highstate_age_seconds=$HIGHSTATE_AGE_SECONDS"
if [ -r "/var/log/salt/lasthighstate" ]; then
LAST_HIGHSTATE_END=$(date -r /var/log/salt/lasthighstate +%s)
NOW=$(date +%s)
HIGHSTATE_AGE_SECONDS=$((NOW-LAST_HIGHSTATE_END))
echo "salt highstate_age_seconds=$HIGHSTATE_AGE_SECONDS"
fi
fi
+313
View File
@@ -53,6 +53,319 @@ telegraf:
forcedType: bool
advanced: True
helpLink: influxdb
container_stats:
tags:
identity:
description: Adds the container_image, container_version, engine_host and server_version tags to every container measurement, as the retired inputs.docker plugin did. Increases series cardinality. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
engine:
n_containers:
description: Total number of containers known to the Docker engine, running or not. Part of the engine-level docker measurement. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
n_containers_running:
description: Number of containers currently running. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
n_containers_stopped:
description: Number of containers currently stopped. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
n_containers_paused:
description: Number of containers currently paused. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
n_images:
description: Number of container images held by the Docker engine. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
n_cpus:
description: Number of CPUs the Docker engine reports for this host. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
n_goroutines:
description: Number of goroutines inside the Docker daemon. Diagnostic detail for the daemon itself. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
n_used_file_descriptors:
description: Number of file descriptors held open by the Docker daemon. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
n_listener_events:
description: Number of event listeners subscribed to the Docker daemon. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
memory_total:
description: Total physical memory the Docker engine reports for this host, in bytes. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
cpu:
usage_percent:
description: Percentage of host CPU consumed by the container. Required by the Container CPU Usage chart on the Security Onion Performance dashboard.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
usage_total:
description: Cumulative CPU time consumed by the container, in nanoseconds. Read from the cgroup. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
usage_in_usermode:
description: Cumulative CPU time consumed in user mode, in nanoseconds. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
usage_in_kernelmode:
description: Cumulative CPU time consumed in kernel mode, in nanoseconds. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
usage_system:
description: Cumulative host-wide CPU time, in nanoseconds. Used as the denominator when calculating container CPU percentage. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
throttling_periods:
description: Number of CPU enforcement periods the container has seen. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
throttling_throttled_periods:
description: Number of periods in which the container was throttled against its CPU limit. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
throttling_throttled_time:
description: Total time the container spent throttled, in nanoseconds. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
container_id: &containerid
description: Full 64 character container ID, emitted as a field on this measurement. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
mem:
usage_percent:
description: Percentage of its memory limit the container is using. Required by the Container Memory Usage chart on the Security Onion Performance dashboard.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
usage:
description: Container memory usage in bytes, excluding reclaimable page cache. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
limit:
description: Memory limit for the container in bytes. Reports total host memory when the container is unlimited. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
max_usage:
description: Peak memory usage for the container in bytes, read from the cgroup peak counter. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
active_anon:
description: Anonymous memory on the active LRU list, in bytes. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
active_file:
description: Page cache on the active LRU list, in bytes. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
inactive_anon:
description: Anonymous memory on the inactive LRU list, in bytes. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
inactive_file:
description: Page cache on the inactive LRU list, in bytes. This is the reclaimable cache subtracted from usage. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
unevictable:
description: Memory that cannot be reclaimed, in bytes. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
pgfault:
description: Cumulative number of page faults taken by the container. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
pgmajfault:
description: Cumulative number of major page faults, those requiring disk access. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
container_id: *containerid
net:
rx_bytes:
description: Bytes received by the container across all interfaces except loopback. Required by the Container Traffic - Inbound chart on the Security Onion Performance dashboard.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
rx_packets:
description: Packets received by the container. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
rx_errors:
description: Receive errors counted on the container interfaces. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
rx_dropped:
description: Received packets dropped by the container interfaces. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
tx_bytes:
description: Bytes transmitted by the container. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
tx_packets:
description: Packets transmitted by the container. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
tx_errors:
description: Transmit errors counted on the container interfaces. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
tx_dropped:
description: Transmitted packets dropped by the container interfaces. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
container_id: *containerid
blkio:
io_service_bytes_recursive_read:
description: Cumulative bytes read from block devices by the container. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
io_service_bytes_recursive_write:
description: Cumulative bytes written to block devices by the container. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
container_id: *containerid
status:
uptime_ns:
description: How long the container has been running, in nanoseconds. Required by the Container Uptime chart on the Security Onion Performance dashboard.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
oomkilled:
description: Whether the container was killed by the kernel out-of-memory handler. Required by the Most Recent Container Events table on the Security Onion Performance dashboard.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
pid:
description: Host process ID of the container main process. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
exitcode:
description: Exit code of the container main process. Meaningful once the container has stopped. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
restart_count:
description: Number of times the Docker engine has restarted this container. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
started_at:
description: Time the container last started, as a Unix timestamp in nanoseconds. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
finished_at:
description: Time the container last exited, as a Unix timestamp in nanoseconds. Absent for a container that has never exited. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
container_id: *containerid
health:
health_status:
description: Result of the container healthcheck, such as healthy, unhealthy or starting. Only emitted for containers that define a healthcheck. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
failing_streak:
description: Number of consecutive failed healthchecks. Only emitted for containers that define a healthcheck. Defaults to off.
forcedType: bool
global: True
advanced: True
helpLink: influxdb
scripts:
eval: &telegrafscripts
description: List of input.exec scripts to run for this node type. The script must be present in salt/telegraf/scripts.
@@ -0,0 +1,398 @@
#!/usr/bin/env python3
# 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.
# Runs from cron as somon, a docker group member, and emits influx line protocol for telegraf to
# read. This exists so so-telegraf does not need the docker socket; the container reads the output
# file instead.
# Measurement/tag/field names match telegraf's inputs.docker plugin because the InfluxDB
# dashboards query them directly. Which fields are emitted is set per stat in SOC; see
# telegraf.container_stats in defaults.yaml.
import json
SETTINGS = json.loads('''{{ CONTAINER_STATS | tojson }}''')
{% raw %}
import os
import re
import subprocess
import sys
from datetime import datetime, timezone
CGROUP_ROOT = '/sys/fs/cgroup'
NANOSEC = 10**9
# cgroup and docker report cpu time in microseconds; inputs.docker published nanoseconds
USEC_TO_NSEC = 1000
def want(group, field):
return bool(SETTINGS.get(group, {}).get(field))
def wants_any(group, fields):
return any(want(group, field) for field in fields)
def docker(args):
proc = subprocess.run(['docker'] + args, stdout=subprocess.PIPE, stderr=subprocess.DEVNULL, encoding='utf-8')
if proc.returncode != 0:
print('Container system error; unable to query docker', file=sys.stderr)
sys.exit(1)
return proc.stdout
def escape_tag(value):
return value.replace(',', '\\,').replace(' ', '\\ ').replace('=', '\\=')
def quote(value):
return '"{}"'.format(str(value).replace('\\', '\\\\').replace('"', '\\"'))
def integer(value):
return '{}i'.format(int(value))
def unsigned(value):
# inputs.docker published the cgroup and network counters as uint64; influx stores u and i
# as different field types, so matching it keeps historical series readable
return '{}u'.format(max(int(value), 0))
def read_text(path):
try:
with open(path) as handle:
return handle.read()
except OSError:
return ''
def read_pairs(path):
# cgroup files such as memory.stat and cpu.stat are "key value" per line
values = {}
for line in read_text(path).splitlines():
parts = line.split()
if len(parts) == 2:
try:
values[parts[0]] = int(parts[1])
except ValueError:
pass
return values
def read_value(path):
raw = read_text(path).strip()
try:
return int(raw)
except ValueError:
return None
def cgroup_path(pid):
# "0::/system.slice/docker-<id>.scope" on cgroup v2
for line in read_text('/proc/{}/cgroup'.format(pid)).splitlines():
parts = line.split(':', 2)
if len(parts) == 3 and parts[1] == '':
return CGROUP_ROOT + parts[2]
return None
def host_memory_total():
match = re.search(r'MemTotal:\s+(\d+) kB', read_text('/proc/meminfo'))
return int(match.group(1)) * 1024 if match else 0
def host_cpu_nanoseconds():
# inputs.docker's usage_system is the host-wide cpu time the daemon reads from /proc/stat
for line in read_text('/proc/stat').splitlines():
if line.startswith('cpu '):
ticks = sum(int(value) for value in line.split()[1:])
return int(ticks * NANOSEC / os.sysconf('SC_CLK_TCK'))
return 0
def parse_image(image):
# mirrors telegraf's internal/docker ParseImage so the tags match what inputs.docker emitted
domain = ''
remainder = image
if '/' in image:
head, _, tail = image.partition('/')
if '.' in head or ':' in head or head == 'localhost':
domain, remainder = head + '/', tail
if ':' in remainder:
name, _, version = remainder.rpartition(':')
return domain + name, version
return domain + remainder, 'unknown'
def to_nanoseconds(stamp):
# docker emits 9 fractional digits; fromisoformat takes at most 6 before python 3.11
stamp = stamp.rstrip('Z')
if '.' in stamp:
whole, _, frac = stamp.partition('.')
stamp = whole + '.' + frac[:6]
try:
parsed = datetime.fromisoformat(stamp).replace(tzinfo=timezone.utc)
except ValueError:
return None
if parsed.year <= 1:
return None
return int(parsed.timestamp() * NANOSEC)
def percent(value):
try:
return float(value.strip().rstrip('%'))
except ValueError:
return 0.0
def net_counters(pid):
# summed across interfaces except loopback, matching the daemon's per-container totals
rows = read_text('/proc/{}/net/dev'.format(pid)).splitlines()[2:]
if not rows:
return None
names = ['rx_bytes', 'rx_packets', 'rx_errors', 'rx_dropped']
totals = dict.fromkeys(names + ['tx_bytes', 'tx_packets', 'tx_errors', 'tx_dropped'], 0)
found = False
for row in rows:
iface, _, rest = row.partition(':')
if iface.strip() == 'lo':
continue
columns = rest.split()
if len(columns) < 12:
continue
found = True
for index, name in enumerate(names):
totals[name] += int(columns[index])
for index, name in enumerate(['tx_bytes', 'tx_packets', 'tx_errors', 'tx_dropped']):
totals[name] += int(columns[8 + index])
return totals if found else None
def blkio_counters(path):
# inputs.docker emitted the device=total row unconditionally, so a container that has done
# no block io reports zeros rather than dropping out of the measurement entirely
totals = {'io_service_bytes_recursive_read': 0, 'io_service_bytes_recursive_write': 0}
rows = read_text(path + '/io.stat').splitlines()
for row in rows:
for token in row.split()[1:]:
key, _, value = token.partition('=')
try:
if key == 'rbytes':
totals['io_service_bytes_recursive_read'] += int(value)
elif key == 'wbytes':
totals['io_service_bytes_recursive_write'] += int(value)
except ValueError:
pass
return totals
def collect_stats():
stats = {}
for line in docker(['stats', '--no-stream', '--format', '{{json .}}']).splitlines():
line = line.strip()
if not line:
continue
try:
entry = json.loads(line)
except json.JSONDecodeError:
continue
stats[entry.get('Name', '')] = entry
return stats
def collect_inspect():
ids = docker(['ps', '-aq']).split()
if not ids:
return []
try:
return json.loads(docker(['inspect'] + ids))
except json.JSONDecodeError:
return []
def collect_info():
try:
return json.loads(docker(['info', '--format', '{{json .}}']))
except json.JSONDecodeError:
return {}
class Emitter:
def __init__(self):
self.lines = []
def add(self, measurement, tags, fields):
if not fields:
return
tagset = ','.join('{}={}'.format(key, escape_tag(value)) for key, value in tags)
body = ','.join('{}={}'.format(key, fields[key]) for key in sorted(fields))
self.lines.append('{},{} {}'.format(measurement, tagset, body))
def main():
cpu_extra = ['usage_total', 'usage_in_usermode', 'usage_in_kernelmode',
'throttling_periods', 'throttling_throttled_periods', 'throttling_throttled_time']
mem_extra = ['usage', 'limit', 'max_usage', 'active_anon', 'active_file', 'inactive_anon',
'inactive_file', 'unevictable', 'pgfault', 'pgmajfault']
net_fields = ['rx_bytes', 'rx_packets', 'rx_errors', 'rx_dropped',
'tx_bytes', 'tx_packets', 'tx_errors', 'tx_dropped']
engine_fields = ['n_containers', 'n_containers_running', 'n_containers_stopped',
'n_containers_paused', 'n_images', 'n_cpus', 'n_goroutines',
'n_used_file_descriptors', 'n_listener_events']
identity = want('tags', 'identity')
need_stats = want('cpu', 'usage_percent')
need_cgroup_cpu = wants_any('cpu', cpu_extra)
need_cgroup_mem = wants_any('mem', mem_extra + ['usage_percent'])
need_blkio = wants_any('blkio', ['io_service_bytes_recursive_read',
'io_service_bytes_recursive_write', 'container_id'])
need_net = wants_any('net', net_fields + ['container_id'])
need_info = identity or wants_any('engine', engine_fields + ['memory_total'])
emitter = Emitter()
stats = collect_stats() if need_stats else {}
info = collect_info() if need_info else {}
engine_tags = [('engine_host', info.get('Name', '')),
('server_version', info.get('ServerVersion', ''))]
system_ns = host_cpu_nanoseconds() if want('cpu', 'usage_system') else 0
mem_total = host_memory_total() if wants_any('mem', ['limit', 'usage_percent']) else 0
if info:
counts = {'n_containers': 'Containers', 'n_containers_running': 'ContainersRunning',
'n_containers_stopped': 'ContainersStopped', 'n_containers_paused': 'ContainersPaused',
'n_images': 'Images', 'n_cpus': 'NCPU', 'n_goroutines': 'NGoroutines',
'n_used_file_descriptors': 'NFd', 'n_listener_events': 'NEventsListener'}
fields = {name: integer(info[key]) for name, key in counts.items()
if want('engine', name) and key in info}
emitter.add('docker', engine_tags, fields)
# inputs.docker published memory_total as its own point
if want('engine', 'memory_total') and 'MemTotal' in info:
emitter.add('docker', engine_tags, {'memory_total': integer(info['MemTotal'])})
for container in collect_inspect():
name = container.get('Name', '').lstrip('/')
if not name:
continue
state = container.get('State', {})
status = state.get('Status', 'unknown')
pid = state.get('Pid')
container_id = container.get('Id', '')
tags = [('container_name', name), ('container_status', status)]
if identity:
image, version = parse_image(container.get('Config', {}).get('Image', ''))
tags += [('container_image', image), ('container_version', version)] + engine_tags
status_fields = {}
if want('status', 'uptime_ns') or want('status', 'started_at') or want('status', 'finished_at'):
started = to_nanoseconds(state.get('StartedAt', ''))
finished = to_nanoseconds(state.get('FinishedAt', ''))
if started is not None:
if want('status', 'started_at'):
status_fields['started_at'] = integer(started)
if want('status', 'uptime_ns'):
end = finished if finished is not None and finished >= started else int(
datetime.now(timezone.utc).timestamp() * NANOSEC)
status_fields['uptime_ns'] = integer(end - started)
if finished is not None and want('status', 'finished_at'):
status_fields['finished_at'] = integer(finished)
if want('status', 'oomkilled'):
status_fields['oomkilled'] = 'true' if state.get('OOMKilled') else 'false'
if want('status', 'pid'):
status_fields['pid'] = integer(pid or 0)
if want('status', 'exitcode'):
status_fields['exitcode'] = integer(state.get('ExitCode', 0))
if want('status', 'restart_count'):
status_fields['restart_count'] = integer(container.get('RestartCount', 0))
if want('status', 'container_id'):
status_fields['container_id'] = quote(container_id)
emitter.add('docker_container_status', tags, status_fields)
health = state.get('Health')
if health:
health_fields = {}
if want('health', 'health_status'):
health_fields['health_status'] = quote(health.get('Status', ''))
if want('health', 'failing_streak'):
health_fields['failing_streak'] = integer(health.get('FailingStreak', 0))
emitter.add('docker_container_health', tags, health_fields)
if status != 'running' or not pid:
continue
path = cgroup_path(pid)
entry = stats.get(name, {})
cpu_fields = {}
if want('cpu', 'usage_percent'):
cpu_fields['usage_percent'] = percent(entry.get('CPUPerc', '0%'))
if want('cpu', 'usage_system'):
cpu_fields['usage_system'] = unsigned(system_ns)
if want('cpu', 'container_id'):
cpu_fields['container_id'] = quote(container_id)
if need_cgroup_cpu and path:
cpu = read_pairs(path + '/cpu.stat')
mapping = {'usage_total': 'usage_usec', 'usage_in_usermode': 'user_usec',
'usage_in_kernelmode': 'system_usec',
'throttling_periods': 'nr_periods',
'throttling_throttled_periods': 'nr_throttled',
'throttling_throttled_time': 'throttled_usec'}
for field, key in mapping.items():
if want('cpu', field) and key in cpu:
scale = 1 if field.startswith('throttling_') and field != 'throttling_throttled_time' else USEC_TO_NSEC
cpu_fields[field] = unsigned(cpu[key] * scale)
emitter.add('docker_container_cpu', tags + [('cpu', 'cpu-total')], cpu_fields)
mem_fields = {}
if want('mem', 'container_id'):
mem_fields['container_id'] = quote(container_id)
if need_cgroup_mem and path:
memory = read_pairs(path + '/memory.stat')
for field in ['active_anon', 'active_file', 'inactive_anon', 'inactive_file',
'unevictable', 'pgfault', 'pgmajfault']:
if want('mem', field) and field in memory:
mem_fields[field] = unsigned(memory[field])
current = read_value(path + '/memory.current')
# inputs.docker reports usage net of reclaimable page cache
usage = max(current - memory.get('inactive_file', 0), 0) if current is not None else None
raw_limit = read_value(path + '/memory.max')
limit = raw_limit if raw_limit is not None else mem_total
if want('mem', 'usage') and usage is not None:
mem_fields['usage'] = unsigned(usage)
if want('mem', 'limit'):
mem_fields['limit'] = unsigned(limit)
if want('mem', 'usage_percent') and usage is not None:
# same ratio inputs.docker computes, from the same two values
mem_fields['usage_percent'] = usage / limit * 100.0 if limit else 0.0
if want('mem', 'max_usage'):
peak = read_value(path + '/memory.peak')
if peak is not None:
mem_fields['max_usage'] = unsigned(peak)
emitter.add('docker_container_mem', tags, mem_fields)
# inputs.docker emitted nothing for host-network containers; its Networks map was empty
if need_net and container.get('HostConfig', {}).get('NetworkMode', '') != 'host':
counters = net_counters(pid)
if counters is not None:
net_out = {field: unsigned(counters[field]) for field in net_fields if want('net', field)}
if want('net', 'container_id'):
net_out['container_id'] = quote(container_id)
emitter.add('docker_container_net', tags + [('network', 'total')], net_out)
if need_blkio and path:
counters = blkio_counters(path)
if counters is not None:
blkio_out = {field: unsigned(value) for field, value in counters.items() if want('blkio', field)}
if want('blkio', 'container_id'):
blkio_out['container_id'] = quote(container_id)
emitter.add('docker_container_blkio', tags + [('device', 'total')], blkio_out)
print('\n'.join(emitter.lines))
if __name__ == '__main__':
main()
{% endraw %}
@@ -0,0 +1,635 @@
# 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.
# so-container-stats replaces telegraf's inputs.docker plugin, which was dropped along with the
# docker socket. The InfluxDB dashboards query its output by measurement, tag and field name, so
# these tests pin that contract: the shipped defaults, the per-stat toggles, the field types, and
# the value semantics copied from the plugin.
#
# The collector is a jinja template, so every test renders it the way salt does and imports the
# result. Docker and the cgroup filesystem are faked, so nothing here needs a container runtime.
import contextlib
import importlib.util
import io
import json
import os
import tempfile
import unittest
import jinja2
import yaml
HERE = os.path.dirname(os.path.abspath(__file__))
TEMPLATE = os.path.join(HERE, 'so-container-stats')
DEFAULTS = os.path.join(HERE, '..', '..', 'defaults.yaml')
# a container id is a 64 character hex string; the plugin published it in full
SOC_ID = 'a' * 64
TELEGRAF_ID = 'b' * 64
NGINX_ID = 'c' * 64
IDSTOOLS_ID = 'd' * 64
def shipped_defaults():
with open(DEFAULTS) as handle:
return yaml.safe_load(handle)['telegraf']['container_stats']
def all_enabled():
return {group: {field: True for field in fields} for group, fields in shipped_defaults().items()}
def render(settings):
"""Render the template as salt does, import it, and hand back the module."""
with open(TEMPLATE) as handle:
source = handle.read()
rendered = jinja2.Template(source, keep_trailing_newline=True).render(CONTAINER_STATS=settings)
path = os.path.join(tempfile.mkdtemp(), 'so_container_stats.py')
with open(path, 'w') as handle:
handle.write(rendered)
spec = importlib.util.spec_from_file_location('so_container_stats', path)
module = importlib.util.module_from_spec(spec)
spec.loader.exec_module(module)
return module
def split_escaped(text, sep):
"""Split on an unescaped, unquoted separator, the way influx line protocol is written."""
parts, current, escaped, quoted = [], '', False, False
for char in text:
if escaped:
current += char
escaped = False
elif char == '\\':
current += char
escaped = True
elif char == '"':
quoted = not quoted
current += char
elif char == sep and not quoted:
parts.append(current)
current = ''
else:
current += char
parts.append(current)
return parts
def split_on_space(line):
escaped = quoted = False
for index, char in enumerate(line):
if escaped:
escaped = False
elif char == '\\':
escaped = True
elif char == '"':
quoted = not quoted
elif char == ' ' and not quoted:
return line[:index], line[index + 1:]
return line, ''
def parse(output):
"""Parse line protocol into {(measurement, container): (tags, fields)}."""
points = {}
for line in output.splitlines():
line = line.strip()
if not line:
continue
head, fieldpart = split_on_space(line)
pieces = split_escaped(head, ',')
measurement, tags = pieces[0], {}
for piece in pieces[1:]:
key, _, value = piece.partition('=')
tags[key] = value
fields = {}
for piece in split_escaped(fieldpart, ','):
key, _, value = piece.partition('=')
fields[key] = value
key = (measurement, tags.get('container_name', ''))
if key in points:
# the engine measurement is published as two points, the way inputs.docker did
points[key][1].update(fields)
else:
points[key] = (tags, fields)
return points
def field_type(value):
if value.endswith('u'):
return 'unsigned'
if value.endswith('i'):
return 'integer'
if value.startswith('"'):
return 'string'
if value in ('true', 'false'):
return 'boolean'
return 'float'
def container(name, cid, pid, status='running', network='bridge', image='registry:5000/repo/img:3.4.0',
started='2026-09-24T10:00:00.123456789Z', finished='0001-01-01T00:00:00Z', health=None,
oomkilled=False, exitcode=0, restarts=0):
state = {'Status': status, 'Pid': pid, 'StartedAt': started, 'FinishedAt': finished,
'OOMKilled': oomkilled, 'ExitCode': exitcode}
if health is not None:
state['Health'] = health
return {'Id': cid, 'Name': '/' + name, 'State': state, 'RestartCount': restarts,
'Config': {'Image': image}, 'HostConfig': {'NetworkMode': network}}
class CollectorTestCase(unittest.TestCase):
"""Builds a fake docker engine and cgroup tree, then runs the collector against it."""
def setUp(self):
self.inspect = [
container('so-soc', SOC_ID, 1001),
container('so-telegraf', TELEGRAF_ID, 1002, network='host'),
container('so-nginx', NGINX_ID, 1003, health={'Status': 'healthy', 'FailingStreak': 0}),
]
self.stats = {
'so-soc': {'Name': 'so-soc', 'CPUPerc': '1.25%', 'MemPerc': '6.39%', 'NetIO': '1kB / 2kB'},
'so-telegraf': {'Name': 'so-telegraf', 'CPUPerc': '0.04%', 'MemPerc': '0.58%', 'NetIO': '0B / 0B'},
'so-nginx': {'Name': 'so-nginx', 'CPUPerc': '0.00%', 'MemPerc': '0.09%', 'NetIO': '3kB / 4kB'},
}
self.info = {'Name': 'sohost', 'ServerVersion': '29.2.1', 'Containers': 3, 'ContainersRunning': 3,
'ContainersStopped': 0, 'ContainersPaused': 0, 'Images': 9, 'NCPU': 8,
'NGoroutines': 42, 'NFd': 77, 'NEventsListener': 1, 'MemTotal': 16000000000}
# one cgroup per container, addressed through /proc/<pid>/cgroup exactly as the collector does
self.files = {
'/proc/meminfo': 'MemTotal: 15625000 kB\n',
'/proc/stat': 'cpu 100 200 300 400\ncpu0 1 2 3 4\n',
}
for pid in (1001, 1002, 1003):
self.files['/proc/%d/cgroup' % pid] = '0::/scope%d\n' % pid
self.cgroup(pid, 'cpu.stat',
'usage_usec 1000\nuser_usec 600\nsystem_usec 400\nnr_periods 5\nnr_throttled 2\nthrottled_usec 700\n')
self.cgroup(pid, 'memory.stat',
'active_anon 300\nactive_file 40\ninactive_anon 200\ninactive_file 50\nunevictable 0\npgfault 1234\npgmajfault 56\n')
self.cgroup(pid, 'memory.current', '1000\n')
self.cgroup(pid, 'memory.max', '4000\n')
self.cgroup(pid, 'memory.peak', '2500\n')
self.cgroup(pid, 'io.stat', '8:0 rbytes=100 wbytes=200 rios=1 wios=2\n252:0 rbytes=10 wbytes=20 rios=1 wios=1\n')
self.files['/proc/%d/net/dev' % pid] = (
'Inter-| Receive | Transmit\n'
' face |bytes packets errs drop fifo frame compressed multicast|bytes packets errs drop fifo colls carrier compressed\n'
' lo: 9999 99 9 9 0 0 0 0 9999 99 9 9 0 0 0 0\n'
' eth0: 1000 10 1 2 0 0 0 0 2000 20 3 4 0 0 0 0\n'
' eth1: 500 5 0 0 0 0 0 0 1000 10 0 0 0 0 0 0\n')
def cgroup(self, pid, name, contents):
self.files['/sys/fs/cgroup/scope%d/%s' % (pid, name)] = contents
def run_collector(self, settings=None, module=None):
module = module or render(settings if settings is not None else all_enabled())
def fake_docker(args):
if args[0] == 'stats':
return ''.join(json.dumps(entry) + '\n' for entry in self.stats.values())
if args[0] == 'ps':
return ' '.join(entry['Id'] for entry in self.inspect)
if args[0] == 'inspect':
return json.dumps(self.inspect)
if args[0] == 'info':
return json.dumps(self.info)
raise AssertionError('unexpected docker call: %s' % args)
module.docker = fake_docker
module.read_text = lambda path: self.files.get(path, '')
module.CGROUP_ROOT = '/sys/fs/cgroup'
buffer = io.StringIO()
with contextlib.redirect_stdout(buffer):
module.main()
self.output = buffer.getvalue()
return parse(self.output)
class TestShippedDefaults(CollectorTestCase):
def test_defaults_emit_only_the_dashboard_fields(self):
# the Security Onion Performance dashboard queries exactly these five
points = self.run_collector(shipped_defaults())
emitted = {(measurement, field) for (measurement, _), (_, fields) in points.items() for field in fields}
self.assertEqual(emitted, {
('docker_container_cpu', 'usage_percent'),
('docker_container_mem', 'usage_percent'),
('docker_container_net', 'rx_bytes'),
('docker_container_status', 'uptime_ns'),
('docker_container_status', 'oomkilled'),
})
def test_defaults_do_not_emit_the_opt_in_measurements(self):
points = self.run_collector(shipped_defaults())
measurements = {measurement for measurement, _ in points}
self.assertNotIn('docker', measurements)
self.assertNotIn('docker_container_blkio', measurements)
self.assertNotIn('docker_container_health', measurements)
def test_defaults_carry_the_tags_the_dashboard_filters_on(self):
points = self.run_collector(shipped_defaults())
tags, _ = points[('docker_container_cpu', 'so-soc')]
self.assertEqual(tags['container_status'], 'running')
self.assertEqual(tags['cpu'], 'cpu-total')
# identity tags are opt in, so they must be absent by default
self.assertNotIn('container_image', tags)
self.assertNotIn('engine_host', tags)
class TestToggles(CollectorTestCase):
def test_enabling_one_stat_adds_only_that_field(self):
settings = shipped_defaults()
settings['cpu']['usage_total'] = True
_, fields = self.run_collector(settings)[('docker_container_cpu', 'so-soc')]
self.assertIn('usage_total', fields)
self.assertNotIn('usage_in_usermode', fields)
def test_disabling_one_stat_leaves_its_neighbours(self):
settings = all_enabled()
settings['blkio']['io_service_bytes_recursive_read'] = False
_, fields = self.run_collector(settings)[('docker_container_blkio', 'so-soc')]
self.assertNotIn('io_service_bytes_recursive_read', fields)
self.assertIn('io_service_bytes_recursive_write', fields)
def test_identity_tags_are_added_when_enabled(self):
tags, _ = self.run_collector(all_enabled())[('docker_container_cpu', 'so-soc')]
self.assertEqual(tags['container_image'], 'registry:5000/repo/img')
self.assertEqual(tags['container_version'], '3.4.0')
self.assertEqual(tags['engine_host'], 'sohost')
self.assertEqual(tags['server_version'], '29.2.1')
def test_everything_off_emits_nothing(self):
settings = {group: {field: False for field in fields} for group, fields in shipped_defaults().items()}
self.assertEqual(self.run_collector(settings), {})
class TestFieldTypes(CollectorTestCase):
"""inputs.docker wrote the cgroup and network counters as unsigned; influx treats u and i as
different field types, so a mismatch breaks queries spanning the change."""
def test_counter_fields_are_unsigned(self):
points = self.run_collector(all_enabled())
for measurement, field in (('docker_container_cpu', 'usage_total'),
('docker_container_cpu', 'throttling_periods'),
('docker_container_mem', 'usage'),
('docker_container_mem', 'limit'),
('docker_container_mem', 'pgfault'),
('docker_container_net', 'rx_bytes'),
('docker_container_blkio', 'io_service_bytes_recursive_read')):
_, fields = points[(measurement, 'so-soc')]
self.assertEqual(field_type(fields[field]), 'unsigned', '%s.%s' % (measurement, field))
def test_status_and_engine_fields_are_signed(self):
points = self.run_collector(all_enabled())
_, status = points[('docker_container_status', 'so-soc')]
for field in ('uptime_ns', 'pid', 'exitcode', 'restart_count', 'started_at'):
self.assertEqual(field_type(status[field]), 'integer', field)
_, engine = points[('docker', '')]
self.assertEqual(field_type(engine['n_containers']), 'integer')
def test_percentages_are_floats_and_ids_are_quoted_strings(self):
points = self.run_collector(all_enabled())
_, cpu = points[('docker_container_cpu', 'so-soc')]
self.assertEqual(field_type(cpu['usage_percent']), 'float')
self.assertEqual(field_type(cpu['container_id']), 'string')
self.assertEqual(cpu['container_id'], '"%s"' % SOC_ID)
_, status = points[('docker_container_status', 'so-soc')]
self.assertEqual(field_type(status['oomkilled']), 'boolean')
_, health = points[('docker_container_health', 'so-nginx')]
self.assertEqual(health['health_status'], '"healthy"')
class TestMemorySemantics(CollectorTestCase):
def test_usage_subtracts_reclaimable_page_cache(self):
# inputs.docker reports usage net of inactive_file: 1000 - 50
_, fields = self.run_collector(all_enabled())[('docker_container_mem', 'so-soc')]
self.assertEqual(fields['usage'], '950u')
def test_usage_percent_is_usage_over_limit(self):
_, fields = self.run_collector(all_enabled())[('docker_container_mem', 'so-soc')]
self.assertAlmostEqual(float(fields['usage_percent']), 950 / 4000 * 100.0)
def test_limit_falls_back_to_host_memory_when_unlimited(self):
for pid in (1001, 1002, 1003):
self.cgroup(pid, 'memory.max', 'max\n')
_, fields = self.run_collector(all_enabled())[('docker_container_mem', 'so-soc')]
self.assertEqual(fields['limit'], '%du' % (15625000 * 1024))
def test_max_usage_reports_the_cgroup_peak(self):
# the daemon reports 0 on cgroup v2, so this deliberately carries the real peak
_, fields = self.run_collector(all_enabled())[('docker_container_mem', 'so-soc')]
self.assertEqual(fields['max_usage'], '2500u')
class TestCpuSemantics(CollectorTestCase):
def test_microsecond_counters_are_published_as_nanoseconds(self):
_, fields = self.run_collector(all_enabled())[('docker_container_cpu', 'so-soc')]
self.assertEqual(fields['usage_total'], '1000000u')
self.assertEqual(fields['usage_in_usermode'], '600000u')
self.assertEqual(fields['usage_in_kernelmode'], '400000u')
self.assertEqual(fields['throttling_throttled_time'], '700000u')
def test_throttling_counts_are_not_scaled(self):
_, fields = self.run_collector(all_enabled())[('docker_container_cpu', 'so-soc')]
self.assertEqual(fields['throttling_periods'], '5u')
self.assertEqual(fields['throttling_throttled_periods'], '2u')
def test_usage_system_is_host_wide_cpu_time(self):
_, fields = self.run_collector(all_enabled())[('docker_container_cpu', 'so-soc')]
self.assertEqual(field_type(fields['usage_system']), 'unsigned')
self.assertNotEqual(fields['usage_system'], '0u')
class TestNetwork(CollectorTestCase):
def test_counters_sum_interfaces_and_ignore_loopback(self):
_, fields = self.run_collector(all_enabled())[('docker_container_net', 'so-soc')]
self.assertEqual(fields['rx_bytes'], '1500u')
self.assertEqual(fields['rx_packets'], '15u')
self.assertEqual(fields['tx_bytes'], '3000u')
self.assertEqual(fields['rx_dropped'], '2u')
def test_host_network_containers_emit_no_row(self):
# inputs.docker skipped these: its Networks map is empty for --net=host
points = self.run_collector(all_enabled())
self.assertNotIn(('docker_container_net', 'so-telegraf'), points)
self.assertIn(('docker_container_net', 'so-soc'), points)
def test_total_tag_is_present(self):
tags, _ = self.run_collector(all_enabled())[('docker_container_net', 'so-soc')]
self.assertEqual(tags['network'], 'total')
class TestBlkio(CollectorTestCase):
def test_counters_sum_devices(self):
_, fields = self.run_collector(all_enabled())[('docker_container_blkio', 'so-soc')]
self.assertEqual(fields['io_service_bytes_recursive_read'], '110u')
self.assertEqual(fields['io_service_bytes_recursive_write'], '220u')
def test_container_with_no_io_reports_zero_rather_than_disappearing(self):
for pid in (1001, 1002, 1003):
self.cgroup(pid, 'io.stat', '')
tags, fields = self.run_collector(all_enabled())[('docker_container_blkio', 'so-soc')]
self.assertEqual(fields['io_service_bytes_recursive_read'], '0u')
self.assertEqual(tags['device'], 'total')
class TestStatus(CollectorTestCase):
def test_running_container_uptime_counts_from_start(self):
_, fields = self.run_collector(all_enabled())[('docker_container_status', 'so-soc')]
self.assertGreater(int(fields['uptime_ns'].rstrip('i')), 0)
self.assertNotIn('finished_at', fields)
def test_exited_container_reports_its_lifetime_and_finished_at(self):
self.inspect.append(container('so-idstools', IDSTOOLS_ID, 0, status='exited',
started='2026-09-24T10:00:00.000000000Z',
finished='2026-09-24T10:00:02.000000000Z', exitcode=3, restarts=1))
points = self.run_collector(all_enabled())
tags, fields = points[('docker_container_status', 'so-idstools')]
self.assertEqual(tags['container_status'], 'exited')
self.assertEqual(fields['uptime_ns'], '2000000000i')
self.assertEqual(fields['finished_at'], '1790244002000000000i')
self.assertEqual(fields['exitcode'], '3i')
self.assertEqual(fields['restart_count'], '1i')
# a stopped container has no live stats, so only the status row is emitted
self.assertNotIn(('docker_container_cpu', 'so-idstools'), points)
def test_health_is_emitted_only_for_containers_with_a_healthcheck(self):
points = self.run_collector(all_enabled())
self.assertIn(('docker_container_health', 'so-nginx'), points)
self.assertNotIn(('docker_container_health', 'so-soc'), points)
def test_oomkilled_is_reported(self):
self.inspect[0]['State']['OOMKilled'] = True
_, fields = self.run_collector(all_enabled())[('docker_container_status', 'so-soc')]
self.assertEqual(fields['oomkilled'], 'true')
class TestEngineMeasurement(CollectorTestCase):
def test_engine_counts_come_from_docker_info(self):
_, fields = self.run_collector(all_enabled())[('docker', '')]
self.assertEqual(fields['n_containers'], '3i')
self.assertEqual(fields['n_cpus'], '8i')
self.assertEqual(fields['n_used_file_descriptors'], '77i')
def test_engine_is_published_as_two_points(self):
# inputs.docker emitted memory_total on its own point, so keep that shape
self.run_collector(all_enabled())
engine = [line for line in self.output.splitlines() if line.startswith('docker,')]
self.assertEqual(len(engine), 2)
self.assertTrue(any('memory_total=' in line for line in engine))
self.assertTrue(any('n_containers=' in line for line in engine))
def test_engine_rows_are_tagged_with_host_and_version(self):
tags, _ = self.run_collector(all_enabled())[('docker', '')]
self.assertEqual(tags['engine_host'], 'sohost')
self.assertEqual(tags['server_version'], '29.2.1')
class TestLineProtocol(CollectorTestCase):
def test_tag_values_are_escaped(self):
self.inspect[0]['Name'] = '/odd name,with=chars'
self.stats['odd name,with=chars'] = self.stats.pop('so-soc')
self.stats['odd name,with=chars']['Name'] = 'odd name,with=chars'
output = self.run_collector(all_enabled())
self.assertIn(('docker_container_cpu', 'odd\\ name\\,with\\=chars'), output)
def test_image_without_a_tag_reports_version_unknown(self):
self.inspect[0]['Config']['Image'] = 'busybox'
tags, _ = self.run_collector(all_enabled())[('docker_container_cpu', 'so-soc')]
self.assertEqual(tags['container_image'], 'busybox')
self.assertEqual(tags['container_version'], 'unknown')
def test_every_line_has_a_measurement_tagset_and_fieldset(self):
module = render(all_enabled())
def fake_docker(args):
if args[0] == 'stats':
return ''.join(json.dumps(entry) + '\n' for entry in self.stats.values())
if args[0] == 'ps':
return ' '.join(entry['Id'] for entry in self.inspect)
if args[0] == 'inspect':
return json.dumps(self.inspect)
return json.dumps(self.info)
module.docker = fake_docker
module.read_text = lambda path: self.files.get(path, '')
buffer = io.StringIO()
with contextlib.redirect_stdout(buffer):
module.main()
lines = [line for line in buffer.getvalue().splitlines() if line.strip()]
self.assertTrue(lines)
for line in lines:
head, fieldpart = split_on_space(line)
self.assertIn(',', head, line)
self.assertIn('=', fieldpart, line)
self.assertFalse(fieldpart.endswith(','), line)
# group -> (measurement, the container whose row carries it)
GROUP_TARGET = {
'engine': ('docker', ''),
'cpu': ('docker_container_cpu', 'so-soc'),
'mem': ('docker_container_mem', 'so-soc'),
'net': ('docker_container_net', 'so-soc'),
'blkio': ('docker_container_blkio', 'so-soc'),
'status': ('docker_container_status', 'so-soc'),
'health': ('docker_container_health', 'so-nginx'),
}
class TestEverySetting(CollectorTestCase):
"""Whatever is offered in defaults.yaml has to actually be collectable. These tests are driven
off that file, so a new setting that is never wired up fails here rather than shipping."""
def setUp(self):
super().setUp()
# an exited container so status.finished_at has a value to report
self.inspect.append(container('so-idstools', IDSTOOLS_ID, 0, status='exited',
started='2026-09-24T10:00:00.000000000Z',
finished='2026-09-24T10:00:02.000000000Z'))
def test_every_setting_emits_its_field_when_enabled(self):
points = self.run_collector(all_enabled())
for group, fields in shipped_defaults().items():
if group == 'tags':
continue
measurement, name = GROUP_TARGET[group]
if group == 'status':
# finished_at only exists for a container that has actually exited
_, exited = points[(measurement, 'so-idstools')]
self.assertIn('finished_at', exited)
_, emitted = points[(measurement, name)]
for field in fields:
if group == 'status' and field == 'finished_at':
continue
self.assertIn(field, emitted, '%s.%s is offered but never emitted' % (group, field))
def test_every_setting_is_individually_wired(self):
# enabling one stat on its own must produce exactly that field, proving each toggle is
# read rather than riding along with a neighbour
for group, fields in shipped_defaults().items():
if group == 'tags':
continue
measurement, name = GROUP_TARGET[group]
for field in fields:
settings = {other: {key: False for key in values} for other, values in shipped_defaults().items()}
settings[group][field] = True
target = 'so-idstools' if (group == 'status' and field == 'finished_at') else name
points = self.run_collector(settings)
self.assertIn((measurement, target), points, '%s.%s emitted no row' % (group, field))
_, emitted = points[(measurement, target)]
self.assertEqual(sorted(emitted), [field], '%s.%s did not emit itself alone' % (group, field))
def test_all_enabled_values_are_exact(self):
points = self.run_collector(all_enabled())
clock = os.sysconf('SC_CLK_TCK')
expected = {
('docker_container_cpu', 'so-soc'): {
'usage_percent': '1.25', 'usage_total': '1000000u', 'usage_in_usermode': '600000u',
'usage_in_kernelmode': '400000u', 'usage_system': '%du' % int(1000 * 10**9 / clock),
'throttling_periods': '5u', 'throttling_throttled_periods': '2u',
'throttling_throttled_time': '700000u', 'container_id': '"%s"' % SOC_ID,
},
('docker_container_mem', 'so-soc'): {
'usage': '950u', 'limit': '4000u', 'max_usage': '2500u', 'active_anon': '300u',
'active_file': '40u', 'inactive_anon': '200u', 'inactive_file': '50u',
'unevictable': '0u', 'pgfault': '1234u', 'pgmajfault': '56u',
'usage_percent': '23.75', 'container_id': '"%s"' % SOC_ID,
},
('docker_container_net', 'so-soc'): {
'rx_bytes': '1500u', 'rx_packets': '15u', 'rx_errors': '1u', 'rx_dropped': '2u',
'tx_bytes': '3000u', 'tx_packets': '30u', 'tx_errors': '3u', 'tx_dropped': '4u',
'container_id': '"%s"' % SOC_ID,
},
('docker_container_blkio', 'so-soc'): {
'io_service_bytes_recursive_read': '110u',
'io_service_bytes_recursive_write': '220u', 'container_id': '"%s"' % SOC_ID,
},
('docker', ''): {
'n_containers': '3i', 'n_containers_running': '3i', 'n_containers_stopped': '0i',
'n_containers_paused': '0i', 'n_images': '9i', 'n_cpus': '8i', 'n_goroutines': '42i',
'n_used_file_descriptors': '77i', 'n_listener_events': '1i',
'memory_total': '16000000000i',
},
('docker_container_health', 'so-nginx'): {
'health_status': '"healthy"', 'failing_streak': '0i',
},
}
for key, fields in expected.items():
_, emitted = points[key]
for field, value in fields.items():
self.assertEqual(emitted[field], value, '%s %s' % (key[0], field))
def test_status_values_are_exact(self):
# uptime is relative to now, so it is checked separately from the fixed fields
points = self.run_collector(all_enabled())
_, running = points[('docker_container_status', 'so-soc')]
self.assertEqual(running['pid'], '1001i')
self.assertEqual(running['exitcode'], '0i')
self.assertEqual(running['restart_count'], '0i')
self.assertEqual(running['oomkilled'], 'false')
self.assertEqual(running['container_id'], '"%s"' % SOC_ID)
self.assertEqual(int(running['started_at'].rstrip('i')) // 10**9, 1790244000)
self.assertGreater(int(running['uptime_ns'].rstrip('i')), 0)
_, exited = points[('docker_container_status', 'so-idstools')]
self.assertEqual(exited['uptime_ns'], '2000000000i')
self.assertEqual(exited['finished_at'], '1790244002000000000i')
def test_tags_identity_toggle_controls_the_identity_tags(self):
settings = shipped_defaults()
settings['tags']['identity'] = False
tags, _ = self.run_collector(settings)[('docker_container_cpu', 'so-soc')]
self.assertEqual(sorted(tags), ['container_name', 'container_status', 'cpu'])
settings['tags']['identity'] = True
tags, _ = self.run_collector(settings)[('docker_container_cpu', 'so-soc')]
self.assertEqual(sorted(tags), ['container_image', 'container_name', 'container_status',
'container_version', 'cpu', 'engine_host', 'server_version'])
def test_identity_tags_cover_every_documented_tag(self):
tags, _ = self.run_collector(all_enabled())[('docker_container_cpu', 'so-soc')]
for tag in ('container_image', 'container_version', 'engine_host', 'server_version'):
self.assertIn(tag, tags)
class TestTemplate(unittest.TestCase):
def test_template_renders_to_valid_python_for_the_shipped_defaults(self):
with open(TEMPLATE) as handle:
source = handle.read()
rendered = jinja2.Template(source, keep_trailing_newline=True).render(CONTAINER_STATS=shipped_defaults())
compile(rendered, 'so-container-stats', 'exec')
self.assertNotIn('{%', rendered)
# the docker format strings must survive rendering untouched
self.assertIn('{{json .}}', rendered)
def test_every_annotated_setting_exists_in_defaults(self):
# SOC reads both trees; an annotation without a default cannot be reverted in the UI
with open(os.path.join(HERE, '..', '..', 'soc_telegraf.yaml')) as handle:
annotated = yaml.safe_load(handle)['telegraf']['container_stats']
defaults = shipped_defaults()
for group, fields in annotated.items():
self.assertIn(group, defaults)
for field in fields:
self.assertIn(field, defaults[group], '%s.%s annotated but missing from defaults' % (group, field))
def test_every_default_setting_is_annotated_for_soc(self):
with open(os.path.join(HERE, '..', '..', 'soc_telegraf.yaml')) as handle:
annotated = yaml.safe_load(handle)['telegraf']['container_stats']
for group, fields in shipped_defaults().items():
self.assertIn(group, annotated)
for field in fields:
self.assertIn(field, annotated[group], '%s.%s missing a SOC annotation' % (group, field))
if __name__ == '__main__':
unittest.main()
+1
View File
@@ -23,6 +23,7 @@ zeek:
- gid: 937
- home: /opt/so/conf/zeek
- createhome: False
- shell: /sbin/nologin
# Create some directories
zeekpolicydir:
+1
View File
@@ -1611,6 +1611,7 @@ reserve_group_ids() {
logCmd "groupadd -g 949 elastic-agent"
logCmd "groupadd -g 947 elastic-fleet"
logCmd "groupadd -g 960 kafka"
logCmd "groupadd -g 961 somon"
}
reserve_ports() {