feat(go): persist Python-compatible incidents

This commit is contained in:
Luna 2026-07-22 01:59:49 -07:00
parent f85c2e831a
commit 538d1d954a
No known key found for this signature in database
13 changed files with 555 additions and 46 deletions

View file

@ -0,0 +1,256 @@
// SPDX-License-Identifier: GPL-3.0-or-later
// Package incident persists Python-compatible incident grouping for retained
// Go alert snapshots. Grouping is process-lineage first and time-window second;
// it adds evidence without hiding or replacing the raw alerts.
package incident
import (
"crypto/rand"
"encoding/hex"
"encoding/json"
"fmt"
"os"
"path/filepath"
"sort"
"sync"
"time"
"codeberg.org/anassaeneroi/enodia-sentinal/go-agent/internal/correlation"
"codeberg.org/anassaeneroi/enodia-sentinal/go-agent/internal/model"
)
const (
Schema = "enodia.incident.v1"
IndexName = "incidents.json"
MaxIncidents = 1000
)
type Record struct {
Schema string `json:"schema"`
ID string `json:"id"`
Host string `json:"host"`
FirstTS float64 `json:"first_ts"`
LastTS float64 `json:"last_ts"`
FirstSeen string `json:"first_seen"`
LastSeen string `json:"last_seen"`
Severity string `json:"severity"`
Signatures []string `json:"signatures"`
SIDs []int `json:"sids"`
PIDs []int `json:"pids"`
Lineage []int `json:"lineage"`
Snapshots []string `json:"snapshots"`
AlertCount int `json:"alert_count"`
Correlations []correlation.Record `json:"correlations"`
}
type Store struct {
Dir string
Window int
Depth int
Enabled bool
mu sync.Mutex
}
// New validates incident settings and preflights the isolated state directory.
func New(dir string, enabled bool, window, depth int) (*Store, error) {
if dir == "" {
return nil, fmt.Errorf("incident directory is required")
}
if window < 0 || depth < 0 {
return nil, fmt.Errorf("incident window and lineage depth must be non-negative")
}
if err := os.MkdirAll(dir, 0o750); err != nil {
return nil, fmt.Errorf("create incident directory: %w", err)
}
return &Store{Dir: dir, Window: window, Depth: depth, Enabled: enabled}, nil
}
// RecordAlert attaches one newly appended alert to a snapshot's incident.
// preferredID keeps all alerts merged into the same snapshot in one incident.
// The snapshot name itself is appended idempotently.
func (s *Store) RecordAlert(snapshotName string, alert model.Alert, lineage map[int]bool,
when time.Time, host, preferredID string) (string, error) {
if !s.Enabled {
return "", nil
}
s.mu.Lock()
defer s.mu.Unlock()
index, err := load(filepath.Join(s.Dir, IndexName))
if err != nil {
return "", err
}
id := ""
if preferredID != "" && index[preferredID] != nil {
id = preferredID
}
if id == "" {
candidates := make([]correlation.IndexIncident, 0, len(index))
for _, item := range index {
candidates = append(candidates, correlation.IndexIncident{
ID: item.ID, Host: item.Host, LastTimestamp: item.LastTS, Lineage: item.Lineage,
})
}
id = correlation.Assign(candidates, lineage, float64(when.UnixNano())/1e9, host, s.Window)
}
whenFloat := float64(when.UnixNano()) / 1e9
whenISO := when.Format(time.RFC3339Nano)
if id == "" {
id, err = newID(when)
if err != nil {
return "", err
}
index[id] = &Record{
Schema: Schema, ID: id, Host: host, FirstTS: whenFloat, LastTS: whenFloat,
FirstSeen: whenISO, LastSeen: whenISO, Severity: alert.Severity,
Signatures: []string{}, SIDs: []int{}, PIDs: []int{}, Lineage: []int{},
Snapshots: []string{}, Correlations: []correlation.Record{},
}
}
item := index[id]
item.Schema = Schema
item.LastTS = whenFloat
item.LastSeen = whenISO
item.Severity = correlation.MaxSeverity(item.Severity, alert.Severity)
item.Signatures = appendUniqueString(item.Signatures, alert.Signature)
if alert.SID != 0 {
item.SIDs = appendUniqueInt(item.SIDs, alert.SID)
}
item.PIDs = sortedUnion(item.PIDs, alert.PIDs)
item.Lineage = sortedUnion(item.Lineage, mapKeys(lineage))
item.Snapshots = appendUniqueString(item.Snapshots, snapshotName)
item.AlertCount++
item.Correlations = correlation.Correlate(correlation.Incident{
FirstTimestamp: item.FirstTS, LastTimestamp: item.LastTS, Signatures: item.Signatures,
}, nil)
for _, match := range item.Correlations {
item.SIDs = appendUniqueInt(item.SIDs, match.SID)
item.Severity = correlation.MaxSeverity(item.Severity, match.Severity)
}
if err := save(filepath.Join(s.Dir, IndexName), index); err != nil {
return "", err
}
return id, nil
}
// LoadIndex reads the durable index for compatibility tests and consumers.
func LoadIndex(dir string) (map[string]*Record, error) {
return load(filepath.Join(dir, IndexName))
}
func load(path string) (map[string]*Record, error) {
raw, err := os.ReadFile(path)
if os.IsNotExist(err) {
return map[string]*Record{}, nil
}
if err != nil {
return nil, fmt.Errorf("read incident index: %w", err)
}
index := map[string]*Record{}
if json.Unmarshal(raw, &index) != nil {
// Python treats malformed legacy state as an empty index.
return map[string]*Record{}, nil
}
for id, item := range index {
if item == nil {
delete(index, id)
continue
}
if item.Schema == "" {
item.Schema = Schema
}
if item.Correlations == nil {
item.Correlations = []correlation.Record{}
}
}
return index, nil
}
func save(path string, index map[string]*Record) error {
if len(index) > MaxIncidents {
items := make([]*Record, 0, len(index))
for _, item := range index {
items = append(items, item)
}
sort.Slice(items, func(i, j int) bool {
if items[i].LastTS == items[j].LastTS {
return items[i].ID < items[j].ID
}
return items[i].LastTS > items[j].LastTS
})
trimmed := make(map[string]*Record, MaxIncidents)
for _, item := range items[:MaxIncidents] {
trimmed[item.ID] = item
}
index = trimmed
}
raw, err := json.MarshalIndent(index, "", " ")
if err != nil {
return fmt.Errorf("encode incident index: %w", err)
}
temporary := path + ".tmp"
if err := os.WriteFile(temporary, append(raw, '\n'), 0o600); err != nil {
return fmt.Errorf("write incident index: %w", err)
}
if err := os.Rename(temporary, path); err != nil {
return fmt.Errorf("replace incident index: %w", err)
}
if err := os.Chmod(path, 0o600); err != nil {
return fmt.Errorf("protect incident index: %w", err)
}
return nil
}
func newID(when time.Time) (string, error) {
suffix := make([]byte, 2)
if _, err := rand.Read(suffix); err != nil {
return "", fmt.Errorf("generate incident id: %w", err)
}
return "inc-" + when.Local().Format("20060102-150405") + "-" + hex.EncodeToString(suffix), nil
}
func appendUniqueString(values []string, value string) []string {
for _, current := range values {
if current == value {
return values
}
}
return append(values, value)
}
func appendUniqueInt(values []int, value int) []int {
for _, current := range values {
if current == value {
return values
}
}
return append(values, value)
}
func sortedUnion(left, right []int) []int {
values := make(map[int]bool, len(left)+len(right))
for _, value := range left {
values[value] = true
}
for _, value := range right {
values[value] = true
}
result := make([]int, 0, len(values))
for value := range values {
result = append(result, value)
}
sort.Ints(result)
return result
}
func mapKeys(values map[int]bool) []int {
result := make([]int, 0, len(values))
for value, enabled := range values {
if enabled {
result = append(result, value)
}
}
return result
}

View file

@ -0,0 +1,89 @@
// SPDX-License-Identifier: GPL-3.0-or-later
package incident
import (
"os"
"path/filepath"
"testing"
"time"
"codeberg.org/anassaeneroi/enodia-sentinal/go-agent/internal/correlation"
"codeberg.org/anassaeneroi/enodia-sentinal/go-agent/internal/model"
)
func TestRecordAlertGroupsLineageAndPersistsCorrelation(t *testing.T) {
directory := t.TempDir()
store, err := New(directory, true, 1800, 8)
if err != nil {
t.Fatal(err)
}
firstTime := time.Unix(1000, 0).UTC()
first, err := store.RecordAlert("alert-1.log", model.Alert{
SID: 100003, Severity: "HIGH", Signature: "exec_rule.web-rce", PIDs: []int{4242},
}, map[int]bool{4242: true, 999: true}, firstTime, "host-a", "")
if err != nil {
t.Fatal(err)
}
second, err := store.RecordAlert("alert-2.log", model.Alert{
SID: 100067, Severity: "HIGH", Signature: "host_rule.suspicious-egress", PIDs: []int{4243},
}, map[int]bool{4243: true, 999: true}, firstTime.Add(50*time.Second), "host-a", "")
if err != nil {
t.Fatal(err)
}
if first == "" || second != first {
t.Fatalf("incident ids first=%q second=%q", first, second)
}
index, err := LoadIndex(directory)
if err != nil {
t.Fatal(err)
}
item := index[first]
if len(index) != 1 || item == nil || item.Schema != Schema || item.AlertCount != 2 ||
item.Severity != "CRITICAL" || len(item.Snapshots) != 2 ||
len(item.Correlations) != 1 || item.Correlations[0].SID != correlation.SIDMultiStageIntrusion {
t.Fatalf("index=%#v", index)
}
if info, err := os.Stat(filepath.Join(directory, IndexName)); err != nil || info.Mode().Perm() != 0o600 {
t.Fatalf("incident index mode info=%v err=%v", info, err)
}
}
func TestPreferredIncidentMakesSnapshotAppendIdempotent(t *testing.T) {
store, err := New(t.TempDir(), true, 1800, 8)
if err != nil {
t.Fatal(err)
}
when := time.Unix(1000, 0).UTC()
id, err := store.RecordAlert("alert-1.log", model.Alert{
SID: 1, Severity: "HIGH", Signature: "one", PIDs: []int{10},
}, map[int]bool{10: true}, when, "host-a", "")
if err != nil {
t.Fatal(err)
}
if _, err := store.RecordAlert("alert-1.log", model.Alert{
SID: 2, Severity: "CRITICAL", Signature: "two", PIDs: []int{20},
}, map[int]bool{20: true}, when.Add(500*time.Millisecond), "host-a", id); err != nil {
t.Fatal(err)
}
index, _ := LoadIndex(store.Dir)
item := index[id]
if len(item.Snapshots) != 1 || item.AlertCount != 2 || len(item.PIDs) != 2 || item.Severity != "CRITICAL" {
t.Fatalf("incident=%#v", item)
}
}
func TestDisabledStoreDoesNotWriteIndex(t *testing.T) {
directory := t.TempDir()
store, err := New(directory, false, 1800, 8)
if err != nil {
t.Fatal(err)
}
id, err := store.RecordAlert("alert.log", model.Alert{SID: 1, Severity: "HIGH", Signature: "test"}, nil, time.Now(), "host-a", "")
if err != nil || id != "" {
t.Fatalf("id=%q err=%v", id, err)
}
if _, err := os.Stat(filepath.Join(directory, IndexName)); !os.IsNotExist(err) {
t.Fatalf("disabled index exists: %v", err)
}
}