enodia-sentinal/go-agent/cmd/enodia-sentinel-go/main.go

646 lines
20 KiB
Go

// SPDX-License-Identifier: GPL-3.0-or-later
// enodia-sentinel-go is the migration validation sidecar. Python remains the
// production default and behavioral oracle until parity is complete.
package main
import (
"context"
"encoding/json"
"errors"
"flag"
"fmt"
"io"
"os"
"os/signal"
"path/filepath"
"sort"
"strings"
"sync"
"syscall"
"time"
"codeberg.org/anassaeneroi/enodia-sentinal/go-agent/internal/agent"
"codeberg.org/anassaeneroi/enodia-sentinal/go-agent/internal/baseline"
"codeberg.org/anassaeneroi/enodia-sentinal/go-agent/internal/config"
"codeberg.org/anassaeneroi/enodia-sentinal/go-agent/internal/correlation"
"codeberg.org/anassaeneroi/enodia-sentinal/go-agent/internal/ebpfsource"
"codeberg.org/anassaeneroi/enodia-sentinal/go-agent/internal/eventlog"
"codeberg.org/anassaeneroi/enodia-sentinal/go-agent/internal/events"
"codeberg.org/anassaeneroi/enodia-sentinal/go-agent/internal/health"
"codeberg.org/anassaeneroi/enodia-sentinal/go-agent/internal/incident"
"codeberg.org/anassaeneroi/enodia-sentinal/go-agent/internal/model"
"codeberg.org/anassaeneroi/enodia-sentinal/go-agent/internal/schema"
"codeberg.org/anassaeneroi/enodia-sentinal/go-agent/internal/sdnotify"
"codeberg.org/anassaeneroi/enodia-sentinal/go-agent/internal/snapshot"
"codeberg.org/anassaeneroi/enodia-sentinal/go-agent/internal/system"
)
func main() {
if err := run(); err != nil {
fmt.Fprintln(os.Stderr, "enodia-sentinel-go:", err)
os.Exit(1)
}
}
func run() error {
once := flag.Bool("once", false, "emit one sweep and exit")
configPath := flag.String("config", "", "Sentinel TOML path (default: ENODIA_CONFIG or /etc/enodia-sentinel.toml)")
procRoot := flag.String("proc-root", "/proc", "procfs root")
stateDir := flag.String("state-dir", "", "enable baseline lifecycle in this directory")
suidRoot := flag.String("suid-root", "/", "SUID scan root (testing only)")
fixture := flag.String("fixture", "", "JSON SystemState fixture (parity/testing only)")
execEvents := flag.String("exec-events", "", "replay a JSON array of exec events and exit")
syscallEvents := flag.String("syscall-events", "", "replay a JSON array of syscall events and exit")
hostEvents := flag.String("host-events", "", "replay a JSON array of typed host events and exit")
eventStream := flag.String("event-stream", "", "match mixed JSONL events from a file or - for stdin")
eventLogPath := flag.String("event-log", "", "append live event envelopes to a bounded JSONL file")
eventLogMaxBytes := flag.Int64("event-log-max-bytes", 64<<20, "rotate the retained event log before this size")
eventsTail := flag.Int("events-tail", 0, "print the newest N retained event records and exit")
snapshotDir := flag.String("snapshot-dir", "", "retain bounded alert snapshot pairs in this directory")
rulesList := flag.Bool("rules-list", false, "emit the active event-rule catalog as JSON and exit")
rulesShow := flag.Int("rules-show", 0, "emit one event rule by SID as JSON and exit")
correlateFile := flag.String("correlate", "", "correlate a JSON array of incident summaries and exit")
incidentsList := flag.Bool("incidents-list", false, "list retained Go incidents as JSON and exit")
incidentShow := flag.String("incident-show", "", "show one retained Go incident timeline as JSON and exit")
ebpfExec := flag.Bool("ebpf-exec", false, "enable the native execve eBPF source (fails open to polling)")
ebpfSyscall := flag.Bool("ebpf-syscall", false, "enable the native security-syscall eBPF source (fails open to polling)")
checkHealth := flag.Bool("health", false, "check the state heartbeat as JSON and exit")
showVersion := flag.Bool("version", false, "print version and exit")
host := flag.String("host", "", "override hostname (parity/testing only)")
timestamp := flag.String("timestamp", "", "override RFC3339 timestamp (parity/testing only)")
flag.Parse()
if *showVersion {
fmt.Println(agent.Version)
return nil
}
if *eventsTail < 0 {
return fmt.Errorf("events-tail must be non-negative")
}
if *eventsTail > 0 {
path := *eventLogPath
if path == "" {
eventsDir := *stateDir
if eventsDir == "" {
eventsDir = "/var/lib/enodia-sentinel-go"
}
path = eventsDir + "/events.jsonl"
}
rows, err := eventlog.ReadRecent(path, *eventsTail)
if err != nil {
return err
}
for _, row := range rows {
if _, err := os.Stdout.Write(append(row, '\n')); err != nil {
return err
}
}
return nil
}
if *incidentsList || *incidentShow != "" {
directory := *snapshotDir
if directory == "" {
directory = *stateDir
}
if directory == "" {
directory = "/var/lib/enodia-sentinel-go"
}
if *incidentsList {
return writeIncidentList(directory, os.Stdout)
}
return writeIncidentView(directory, *incidentShow, os.Stdout)
}
cfg, err := config.Load(*configPath)
if err != nil {
return err
}
if *checkHealth {
healthDir := *stateDir
if healthDir == "" {
healthDir = "/var/lib/enodia-sentinel-go"
}
status := health.Check(healthDir, time.Now(), time.Duration(cfg.HeartbeatMaxAge)*time.Second)
encoder := json.NewEncoder(os.Stdout)
encoder.SetEscapeHTML(false)
if err := encoder.Encode(status); err != nil {
return err
}
if !status.Healthy {
return errors.New(status.Detail)
}
return nil
}
if *rulesList || *rulesShow != 0 {
engine, err := events.LoadExecRuleEngine(cfg.ExecRulesFile)
if err != nil {
return err
}
catalog := events.RuleCatalog(engine)
payload := any(catalog)
if *rulesShow != 0 {
rule := events.FindRule(catalog, *rulesShow)
if rule == nil {
return fmt.Errorf("no such rule sid: %d", *rulesShow)
}
payload = rule
}
encoder := json.NewEncoder(os.Stdout)
encoder.SetEscapeHTML(false)
return encoder.Encode(payload)
}
if *correlateFile != "" {
raw, err := os.ReadFile(*correlateFile)
if err != nil {
return err
}
var incidents []correlation.Incident
if err := json.Unmarshal(raw, &incidents); err != nil {
return fmt.Errorf("correlation incidents: %w", err)
}
results := make([][]correlation.Record, 0, len(incidents))
for _, incident := range incidents {
results = append(results, correlation.Correlate(incident, nil))
}
encoder := json.NewEncoder(os.Stdout)
encoder.SetEscapeHTML(false)
return encoder.Encode(results)
}
if *execEvents != "" {
return replayExecEvents(cfg, *execEvents, *host, *timestamp, os.Stdout)
}
if *syscallEvents != "" {
return replaySyscallEvents(*syscallEvents, *host, *timestamp, os.Stdout)
}
if *hostEvents != "" {
return replayHostEvents(*hostEvents, *host, *timestamp, os.Stdout)
}
if *eventStream != "" {
return replayEventStream(cfg, *eventStream, *host, *timestamp, os.Stdout)
}
capture := func() (model.State, error) { return system.Capture(*procRoot) }
if *fixture != "" {
capture = func() (model.State, error) {
var state model.State
raw, err := os.ReadFile(*fixture)
if err != nil {
return state, err
}
err = json.Unmarshal(raw, &state)
return state, err
}
}
runner := agent.New(cfg, capture)
runner.SetExecProbeStatus("off (not requested)")
runner.SetSyscallProbeStatus("off (not requested)")
if *fixture == "" && *stateDir != "" {
// Stateful mode is explicit so an opt-in validation sidecar
// cannot overwrite the production Python agent's baseline files merely
// by being run for a smoke test.
cfg.LogDir = *stateDir
runner.Config = cfg
runner.Lifecycle = baseline.New(cfg, func() []string {
return system.ScanSUIDBinaries(*suidRoot, cfg.SUIDScanExtraDirs)
})
}
if *host != "" {
runner.Host = func() (string, error) { return *host, nil }
}
if *timestamp != "" {
fixed, err := time.Parse(time.RFC3339Nano, *timestamp)
if err != nil {
return fmt.Errorf("timestamp: %w", err)
}
runner.Now = func() time.Time { return fixed }
}
if err := runner.Initialize(); err != nil {
return err
}
var execSource *ebpfsource.ExecSource
var execEngine events.ExecRuleEngine
if *ebpfExec {
execEngine, err = events.LoadExecRuleEngine(cfg.ExecRulesFile)
if err != nil {
return err
}
execSource, err = ebpfsource.OpenExecSource()
if err != nil {
runner.SetExecProbeStatus("disabled (" + probeFailureStatus(err) + ")")
fmt.Fprintln(os.Stderr, "enodia-sentinel-go: eBPF exec monitor disabled:", err)
} else {
runner.SetExecProbeStatus("enabled")
defer execSource.Close()
}
}
var syscallSource *ebpfsource.SyscallSource
if *ebpfSyscall {
syscallSource, err = ebpfsource.OpenSyscallSource(*procRoot)
if err != nil {
runner.SetSyscallProbeStatus("disabled (" + probeFailureStatus(err) + ")")
fmt.Fprintln(os.Stderr, "enodia-sentinel-go: eBPF syscall monitor disabled:", err)
} else {
runner.SetSyscallProbeStatus("enabled")
defer syscallSource.Close()
}
}
var retainedEvents *eventlog.Writer
if *eventLogPath != "" {
retainedEvents, err = eventlog.New(*eventLogPath, *eventLogMaxBytes)
if err != nil {
return err
}
}
var snapshotStore *snapshot.Store
if *snapshotDir != "" {
snapshotStore, err = snapshot.New(
*snapshotDir, *procRoot, cfg.MaxSnapshots, cfg.MaxSnapshotAgeDays,
)
if err != nil {
return err
}
if err := snapshotStore.ConfigureIncidents(
cfg.IncidentTracking, cfg.IncidentWindow, cfg.IncidentLineageDepth,
); err != nil {
return err
}
snapshotStore.EnableSlowEnrichment()
defer snapshotStore.Close()
}
notifier := sdnotify.FromEnvironment()
if err := notifier.Notify("READY=1\nSTATUS=initialization complete; poll loop ready"); err != nil {
return fmt.Errorf("systemd readiness notification: %w", err)
}
defer notifier.Notify("STOPPING=1\nSTATUS=stopping") //nolint:errcheck -- process exit is already committed
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()
encoder := json.NewEncoder(os.Stdout)
encoder.SetEscapeHTML(false)
var encoderMu sync.Mutex
emit := func(event map[string]any) error {
encoderMu.Lock()
defer encoderMu.Unlock()
if snapshotStore != nil && event["event_type"] == "alert" {
if _, err := snapshotStore.CaptureEvent(event); err != nil {
return err
}
}
if snapshotStore != nil && event["event_type"] == "status" {
if status, ok := event["status"].(schema.Status); ok {
stats, err := snapshotStore.SnapshotStats()
if err != nil {
return err
}
status.TotalAlerts = stats.Total
status.Counts = stats.Counts
status.LastAlert = stats.LastAlert
event["status"] = status
}
}
if retainedEvents != nil {
if err := retainedEvents.Append(event); err != nil {
return err
}
}
if err := encoder.Encode(event); err != nil {
return err
}
if *stateDir != "" && event["event_type"] == "status" {
if err := health.WriteHeartbeat(*stateDir, time.Now()); err != nil {
return err
}
}
return nil
}
if execSource != nil && !*once {
go monitorExecSource(ctx, execSource, execEngine, runner, emit)
}
if syscallSource != nil && !*once {
go monitorSyscallSource(
ctx, syscallSource, events.DefaultSyscallRuleEngine(),
events.DefaultHostRuleEngine(), runner, emit,
)
}
return runner.Run(ctx, *once, emit)
}
func writeIncidentList(directory string, output io.Writer) error {
index, err := incident.LoadIndex(directory)
if err != nil {
return err
}
items := make([]*incident.Record, 0, len(index))
for _, item := range index {
items = append(items, item)
}
sort.Slice(items, func(i, j int) bool { return items[i].LastTS > items[j].LastTS })
return json.NewEncoder(output).Encode(items)
}
func writeIncidentView(directory, id string, output io.Writer) error {
index, err := incident.LoadIndex(directory)
if err != nil {
return err
}
item := index[id]
if item == nil {
return fmt.Errorf("no such incident: %s", id)
}
timeline := make([]map[string]any, 0, len(item.Snapshots))
snapshots := make([]snapshot.Report, 0, len(item.Snapshots))
for _, name := range item.Snapshots {
path := filepath.Join(directory, strings.TrimSuffix(name, ".log")+".json")
report, err := snapshot.LoadReport(path)
if err != nil {
timeline = append(timeline, map[string]any{"snapshot": name, "time": "?", "missing": true})
continue
}
snapshots = append(snapshots, report)
signatures := make([]string, 0, len(report.Alerts))
pids := make([]int, 0, len(report.Processes))
seen := map[string]bool{}
for _, alert := range report.Alerts {
if !seen[alert.Signature] {
seen[alert.Signature] = true
signatures = append(signatures, alert.Signature)
}
}
for _, process := range report.Processes {
pids = append(pids, process.PID)
}
timeline = append(timeline, map[string]any{
"snapshot": name, "time": report.Time, "severity": report.Severity,
"signatures": signatures, "pids": pids,
})
}
sort.Slice(timeline, func(i, j int) bool { return timeline[i]["time"].(string) < timeline[j]["time"].(string) })
return json.NewEncoder(output).Encode(map[string]any{
"schema": "enodia.incident.view.v1", "incident": item,
"timeline": timeline, "snapshots": snapshots,
})
}
func monitorExecSource(
ctx context.Context,
source execEventReader,
engine events.ExecRuleEngine,
runner *agent.Agent,
emit agent.EmitFunc,
) {
for {
execEvent, err := source.Read()
if err != nil {
var lost ebpfsource.LostSamplesError
if errors.As(err, &lost) {
fmt.Fprintln(os.Stderr, "enodia-sentinel-go:", lost.Error())
continue
}
if ctx.Err() == nil {
runner.SetExecProbeStatus("disabled (reader failure)")
fmt.Fprintln(os.Stderr, "enodia-sentinel-go: eBPF exec monitor stopped:", err)
}
return
}
records, err := runner.AlertEvents(engine.Match(execEvent))
if err != nil {
runner.SetExecProbeStatus("disabled (event encoding failure)")
fmt.Fprintln(os.Stderr, "enodia-sentinel-go: eBPF exec monitor stopped:", err)
return
}
for _, record := range records {
if err := emit(record); err != nil {
runner.SetExecProbeStatus("disabled (output failure)")
fmt.Fprintln(os.Stderr, "enodia-sentinel-go: eBPF exec monitor stopped:", err)
return
}
}
}
}
func monitorSyscallSource(
ctx context.Context,
source syscallEventReader,
syscallEngine events.SyscallRuleEngine,
hostEngine events.HostRuleEngine,
runner *agent.Agent,
emit agent.EmitFunc,
) {
for {
securityEvent, err := source.Read()
if err != nil {
var lost ebpfsource.SyscallLostSamplesError
if errors.As(err, &lost) {
fmt.Fprintln(os.Stderr, "enodia-sentinel-go:", lost.Error())
continue
}
if ctx.Err() == nil {
runner.SetSyscallProbeStatus("disabled (reader failure)")
fmt.Fprintln(os.Stderr, "enodia-sentinel-go: eBPF syscall monitor stopped:", err)
}
return
}
var alerts []model.Alert
if securityEvent.Syscall != nil {
alerts = syscallEngine.Match(*securityEvent.Syscall)
} else if securityEvent.Host != nil {
alerts = hostEngine.Match(*securityEvent.Host)
}
records, err := runner.AlertEvents(alerts)
if err != nil {
runner.SetSyscallProbeStatus("disabled (event encoding failure)")
fmt.Fprintln(os.Stderr, "enodia-sentinel-go: eBPF syscall monitor stopped:", err)
return
}
for _, record := range records {
if err := emit(record); err != nil {
runner.SetSyscallProbeStatus("disabled (output failure)")
fmt.Fprintln(os.Stderr, "enodia-sentinel-go: eBPF syscall monitor stopped:", err)
return
}
}
}
}
type execEventReader interface {
Read() (events.ExecEvent, error)
}
type syscallEventReader interface {
Read() (ebpfsource.SecurityEvent, error)
}
func probeFailureStatus(err error) string {
if errors.Is(err, syscall.EPERM) || errors.Is(err, syscall.EACCES) {
return "permission denied"
}
line, _, _ := strings.Cut(err.Error(), "\n")
const limit = 160
characters := []rune(line)
if len(characters) > limit {
line = string(characters[:limit-1]) + "…"
}
return line
}
func replayExecEvents(cfg config.Config, path, host, timestamp string, output *os.File) error {
raw, err := os.ReadFile(path)
if err != nil {
return err
}
var replay []events.ExecEvent
if err := json.Unmarshal(raw, &replay); err != nil {
return fmt.Errorf("exec events: %w", err)
}
engine, err := events.LoadExecRuleEngine(cfg.ExecRulesFile)
if err != nil {
return err
}
if host == "" {
host, err = os.Hostname()
if err != nil {
return err
}
}
if timestamp == "" {
timestamp = time.Now().Local().Format(time.RFC3339Nano)
} else if _, err := time.Parse(time.RFC3339Nano, timestamp); err != nil {
return fmt.Errorf("timestamp: %w", err)
}
encoder := json.NewEncoder(output)
encoder.SetEscapeHTML(false)
for _, execEvent := range replay {
for _, alert := range engine.Match(execEvent) {
record, err := schema.Build("alert", alert, host, timestamp)
if err != nil {
return err
}
if err := encoder.Encode(record); err != nil {
return err
}
}
}
return nil
}
func replaySyscallEvents(path, host, timestamp string, output *os.File) error {
raw, err := os.ReadFile(path)
if err != nil {
return err
}
var replay []events.SyscallEvent
if err := json.Unmarshal(raw, &replay); err != nil {
return fmt.Errorf("syscall events: %w", err)
}
if host == "" {
host, err = os.Hostname()
if err != nil {
return err
}
}
if timestamp == "" {
timestamp = time.Now().Local().Format(time.RFC3339Nano)
} else if _, err := time.Parse(time.RFC3339Nano, timestamp); err != nil {
return fmt.Errorf("timestamp: %w", err)
}
engine := events.DefaultSyscallRuleEngine()
encoder := json.NewEncoder(output)
encoder.SetEscapeHTML(false)
for _, syscallEvent := range replay {
for _, alert := range engine.Match(syscallEvent) {
record, err := schema.Build("alert", alert, host, timestamp)
if err != nil {
return err
}
if err := encoder.Encode(record); err != nil {
return err
}
}
}
return nil
}
func replayHostEvents(path, host, timestamp string, output *os.File) error {
raw, err := os.ReadFile(path)
if err != nil {
return err
}
var replay []events.HostEvent
if err := json.Unmarshal(raw, &replay); err != nil {
return fmt.Errorf("host events: %w", err)
}
if host == "" {
host, err = os.Hostname()
if err != nil {
return err
}
}
if timestamp == "" {
timestamp = time.Now().Local().Format(time.RFC3339Nano)
} else if _, err := time.Parse(time.RFC3339Nano, timestamp); err != nil {
return fmt.Errorf("timestamp: %w", err)
}
engine := events.DefaultHostRuleEngine()
encoder := json.NewEncoder(output)
encoder.SetEscapeHTML(false)
for _, hostEvent := range replay {
for _, alert := range engine.Match(hostEvent) {
record, err := schema.Build("alert", alert, host, timestamp)
if err != nil {
return err
}
if err := encoder.Encode(record); err != nil {
return err
}
}
}
return nil
}
func replayEventStream(cfg config.Config, path, host, timestamp string, output *os.File) error {
engine, err := events.LoadExecRuleEngine(cfg.ExecRulesFile)
if err != nil {
return err
}
var input *os.File
if path == "-" {
input = os.Stdin
} else {
input, err = os.Open(path)
if err != nil {
return err
}
defer input.Close()
}
if host == "" {
host, err = os.Hostname()
if err != nil {
return err
}
}
if timestamp == "" {
timestamp = time.Now().Local().Format(time.RFC3339Nano)
} else if _, err := time.Parse(time.RFC3339Nano, timestamp); err != nil {
return fmt.Errorf("timestamp: %w", err)
}
rules := events.NewRuleSet(engine)
encoder := json.NewEncoder(output)
encoder.SetEscapeHTML(false)
return events.ScanJSONL(input, func(raw []byte) error {
alerts, err := rules.MatchJSON(raw)
if err != nil {
return err
}
for _, alert := range alerts {
record, err := schema.Build("alert", alert, host, timestamp)
if err != nil {
return err
}
if err := encoder.Encode(record); err != nil {
return err
}
}
return nil
})
}