diff --git a/README.md b/README.md
index 86bad64..623aaec 100644
--- a/README.md
+++ b/README.md
@@ -1,3 +1,71 @@
+# SIEM Backend – Metadata-first Variante
+
+## Warum die Datenbank bisher wächst
+
+Die ursprüngliche Implementierung speichert jedes Windows-Event zweimal: als breite normalisierte Zeile in `event_logs` und zusätzlich als vollständiges XML in `event_log_raw`. Außerdem enthielt die Partitionswartung einen Tabellennamenfehler (`event_logs_raw` statt `event_log_raw`). Der Wartungslauf brach deshalb ab, bevor alte Partitionen gelöscht wurden. Auch `event_count_buckets` und `ueba_context_buckets` waren nicht partitioniert und konnten unbegrenzt wachsen.
+
+## Neues Speichermodell
+
+- `event_logs`: kurzlebige Hot-Daten für Regeln, die eine genaue Reihenfolge einzelner Events benötigen. Standard: 72 Stunden.
+- `event_occurrences`: langfristige Metadaten-Aggregate pro Zeit-Bucket, Host, Channel, Event-ID, Benutzer, IP und Workstation. Enthält `cnt`, `first_event_ts` und `last_event_ts`. Standard: 180 Tage.
+- `event_catalog`: kleine, dauerhafte Landkarte des ersten und letzten Auftretens jeder Event-ID pro Host und Channel.
+- `event_log_raw`: optionales Raw-XML. Standardmäßig deaktiviert; bei Aktivierung nur 24 Stunden Aufbewahrung.
+- `event_count_buckets` und `ueba_context_buckets`: jetzt partitioniert und automatisch bereinigt.
+
+Für einen User-Lockout (typischerweise Security Event 4740) kann die Anwendung damit direkt beantworten: wann er zuerst/zuletzt im Bucket auftrat, auf welchem Host, für welchen Benutzer und wie oft. Das vollständige XML ist dafür nicht erforderlich. `CallerComputerName` wird dabei als Gerät/Workstation übernommen.
+
+## Bevorzugtes Metadaten-Ingest
+
+Bestehende Collector dürfen weiterhin `msg` mit XML senden. Neue oder angepasste Collector können stattdessen direkt ein `meta`-Objekt senden; dann muss das Backend das XML weder speichern noch parsen:
+
+```json
+[
+ {
+ "host": "DC01.example.local",
+ "channel": "Security",
+ "id": 4740,
+ "source": "Microsoft-Windows-Security-Auditing",
+ "ts": "2026-07-18T12:34:56Z",
+ "meta": {
+ "target_user": "alice",
+ "subject_user": "DC01$",
+ "device": "CLIENT-42",
+ "provider": "Microsoft-Windows-Security-Auditing"
+ }
+ }
+]
+```
+
+Mindestens `msg` oder `meta` ist erforderlich. Für Regeln, die auf `msg`/Volltext prüfen, muss Raw-XML weiterhin vom Collector geliefert werden; im reinen Metadatenmodus sollten solche Regeln auf strukturierte Felder umgestellt werden.
+
+## Wichtige Einstellungen
+
+```env
+STORE_RAW_XML=false
+METADATA_BUCKET=1m
+EVENT_RETENTION=72h
+RAW_RETENTION=24h
+METADATA_RETENTION=4320h
+PARTITION_MAINTENANCE_ENABLED=true
+```
+
+`METADATA_BUCKET=1m` ist ein guter Ausgangspunkt. Bei sehr hohen Raten kann auf `5m` erhöht werden; dadurch sinkt die Zeilenzahl weiter, die zeitliche Auflösung wird aber gröber.
+
+## Bestehende Installation migrieren
+
+1. Datenbank sichern und Ingest vorübergehend stoppen.
+2. `deploy/mariadb/migrations/002-metadata-first.sql` in einem Wartungsfenster ausführen. Das Script wärmt `event_catalog` aus den kleineren Baseline-Buckets vor, damit bekannte Event-IDs nicht als neu alarmiert werden. Die beiden `ALTER TABLE ... PARTITION BY`-Operationen können große Tabellen neu aufbauen und sperren.
+3. Neue Umgebungsvariablen aus `dot_env` übernehmen.
+4. Backend aktualisieren und starten.
+5. Im Log kontrollieren, dass `partition maintenance completed` erscheint. Die frühere Meldung zur nicht vorhandenen Tabelle `event_logs_raw` darf nicht mehr auftreten.
+6. Nach Ablauf von `EVENT_RETENTION` prüfen, ob alte `event_logs`-Partitionen verschwinden. Raw-XML kann nach erfolgreicher Validierung separat gelöscht oder durch die kurze `RAW_RETENTION` automatisch entfernt werden. `event_occurrences` wird nicht rückwirkend aus der alten Volltabelle aufgebaut; die langfristige Metadatenhistorie beginnt mit dem neuen Backend.
+
+## Datenbankwahl
+
+MariaDB bleibt für dieses Modell sinnvoll: Die Abfragen sind überwiegend zeitbasierte Filter, gruppierte Zähler und kleine relationale Dimensionen. Ein Wechsel zu ClickHouse oder OpenSearch lohnt sich erst, wenn trotz Aggregation sehr hohe Eventraten, ad-hoc Volltextsuchen oder jahrelange Rohdatenhaltung erforderlich sind. Für die hier beschriebene Metadatenanforderung wäre ein sofortiger Engine-Wechsel zusätzlicher Betriebsaufwand ohne zwingenden Nutzen.
+
+---
+
# SIEM-lite Admin-Handbuch: Anlernphase, manuelle Bewertung und Lerneffekt
Stand: 2026-04-27
diff --git a/compose.yml b/compose.yml
index 1469953..2022a76 100644
--- a/compose.yml
+++ b/compose.yml
@@ -62,6 +62,17 @@ services:
NEW_SOURCE_IP_LOOKBACK: ${NEW_SOURCE_IP_LOOKBACK}
NEW_SOURCE_IP_WINDOW: ${NEW_SOURCE_IP_WINDOW}
DETECTIONS_LIMIT: ${DETECTIONS_LIMIT}
+ STORE_RAW_XML: ${STORE_RAW_XML}
+ METADATA_BUCKET: ${METADATA_BUCKET}
+ EVENT_RETENTION: ${EVENT_RETENTION}
+ RAW_RETENTION: ${RAW_RETENTION}
+ METADATA_RETENTION: ${METADATA_RETENTION}
+ PARTITION_MAINTENANCE_ENABLED: ${PARTITION_MAINTENANCE_ENABLED}
+ PARTITION_MAINTENANCE_INTERVAL: ${PARTITION_MAINTENANCE_INTERVAL}
+ PARTITION_INTERVAL: ${PARTITION_INTERVAL}
+ PARTITION_AHEAD: ${PARTITION_AHEAD}
+ PARTITION_BEHIND: ${PARTITION_BEHIND}
+ PARTITION_RETENTION: ${PARTITION_RETENTION}
TZ: ${TZ}
depends_on:
mariadb:
diff --git a/deploy/mariadb/init/001-schema.sql b/deploy/mariadb/init/001-schema.sql
index 3d6c7b2..5620e57 100644
--- a/deploy/mariadb/init/001-schema.sql
+++ b/deploy/mariadb/init/001-schema.sql
@@ -22,6 +22,8 @@ SET FOREIGN_KEY_CHECKS = 0;
DROP PROCEDURE IF EXISTS ensure_siem_partitions;
+DROP TABLE IF EXISTS event_catalog;
+DROP TABLE IF EXISTS event_occurrences;
DROP TABLE IF EXISTS event_count_buckets;
DROP TABLE IF EXISTS ueba_context_buckets;
DROP TABLE IF EXISTS host_risk_scores;
@@ -181,6 +183,63 @@ PARTITION BY RANGE COLUMNS(ts) (
PARTITION pmax VALUES LESS THAN (MAXVALUE)
);
+-- ---------------------------------------------------------------------
+-- Kleine Event-Landkarte: erstes/letztes Auftreten pro Host/Channel/EventID.
+-- Verhindert teure NOT EXISTS-Scans über die Hot-Event-Tabelle.
+-- ---------------------------------------------------------------------
+
+CREATE TABLE event_catalog (
+ hostname VARCHAR(191) NOT NULL,
+ channel_name VARCHAR(128) NOT NULL,
+ event_id INT UNSIGNED NOT NULL,
+ first_seen DATETIME(6) NOT NULL,
+ last_seen DATETIME(6) NOT NULL,
+ total_count BIGINT UNSIGNED NOT NULL DEFAULT 0,
+ updated_at DATETIME(6) NOT NULL DEFAULT CURRENT_TIMESTAMP(6) ON UPDATE CURRENT_TIMESTAMP(6),
+ PRIMARY KEY (hostname, channel_name, event_id),
+ KEY idx_event_catalog_first_seen (first_seen),
+ KEY idx_event_catalog_last_seen (last_seen)
+) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci;
+
+-- ---------------------------------------------------------------------
+-- Langfristige, kompakte Event-Metadaten
+-- Eine Zeile pro Zeit-Bucket und Dimensionskombination statt einer Zeile pro XML.
+-- ---------------------------------------------------------------------
+
+CREATE TABLE event_occurrences (
+ bucket_start DATETIME(6) NOT NULL,
+ bucket_end DATETIME(6) NOT NULL,
+ dimension_key BINARY(16) NOT NULL,
+
+ hostname VARCHAR(191) NOT NULL,
+ channel_name VARCHAR(128) NOT NULL,
+ event_id INT UNSIGNED NOT NULL,
+ provider_name VARCHAR(191) NOT NULL DEFAULT '',
+
+ target_user VARCHAR(191) NOT NULL DEFAULT '',
+ subject_user VARCHAR(191) NOT NULL DEFAULT '',
+ src_ip VARCHAR(64) NOT NULL DEFAULT '',
+ workstation VARCHAR(191) NOT NULL DEFAULT '',
+ logon_type VARCHAR(32) NOT NULL DEFAULT '',
+ status_text VARCHAR(128) NOT NULL DEFAULT '',
+ failure_reason VARCHAR(255) NOT NULL DEFAULT '',
+
+ cnt BIGINT UNSIGNED NOT NULL DEFAULT 0,
+ first_event_ts DATETIME(6) NOT NULL,
+ last_event_ts DATETIME(6) NOT NULL,
+ updated_at DATETIME(6) NOT NULL DEFAULT CURRENT_TIMESTAMP(6) ON UPDATE CURRENT_TIMESTAMP(6),
+
+ PRIMARY KEY (bucket_start, dimension_key),
+ KEY idx_occurrences_time_host_event (bucket_start, hostname, channel_name, event_id),
+ KEY idx_occurrences_host_event_time (hostname, channel_name, event_id, bucket_start),
+ KEY idx_occurrences_target_user_time (target_user, bucket_start),
+ KEY idx_occurrences_subject_user_time (subject_user, bucket_start),
+ KEY idx_occurrences_src_ip_time (src_ip, bucket_start)
+) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci
+PARTITION BY RANGE COLUMNS(bucket_start) (
+ PARTITION pmax VALUES LESS THAN (MAXVALUE)
+);
+
-- ---------------------------------------------------------------------
-- Detection-Regeln
-- ---------------------------------------------------------------------
@@ -361,7 +420,10 @@ CREATE TABLE event_count_buckets (
bucket_start,
bucket_end
)
-) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci;
+) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci
+PARTITION BY RANGE COLUMNS(bucket_start) (
+ PARTITION pmax VALUES LESS THAN (MAXVALUE)
+);
-- ---------------------------------------------------------------------
-- Baseline Stats
@@ -518,7 +580,10 @@ CREATE TABLE ueba_context_buckets (
workstation,
bucket_start
)
-) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci;
+) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci
+PARTITION BY RANGE COLUMNS(bucket_start) (
+ PARTITION pmax VALUES LESS THAN (MAXVALUE)
+);
-- ---------------------------------------------------------------------
-- Privileged Users
diff --git a/dot_env b/dot_env
index 7b73975..4e96437 100644
--- a/dot_env
+++ b/dot_env
@@ -50,4 +50,20 @@ BASELINE_SUPPRESS_FOR=1h
#BASELINE_MIN_SAMPLES=84
#BASELINE_MEDIUM_Z=3.0
#BASELINE_HIGH_Z=5.0
-#BASELINE_MIN_COUNT=20
\ No newline at end of file
+#BASELINE_MIN_COUNT=20
+# Metadata-first storage
+# Raw XML is disabled by default. Enable only for short forensic retention.
+STORE_RAW_XML=false
+METADATA_BUCKET=1m
+EVENT_RETENTION=72h
+RAW_RETENTION=24h
+METADATA_RETENTION=4320h
+
+# Partition maintenance
+PARTITION_MAINTENANCE_ENABLED=true
+PARTITION_MAINTENANCE_INTERVAL=15m
+PARTITION_INTERVAL=3h
+PARTITION_AHEAD=24h
+PARTITION_BEHIND=6h
+# Compatibility fallback; table-specific retention values above take precedence.
+PARTITION_RETENTION=720h
diff --git a/main.go b/main.go
index 5492707..d5ba649 100644
--- a/main.go
+++ b/main.go
@@ -632,19 +632,19 @@ a:hover {
-
- | Zeit | Host | Channel | EventID | User | IP | Nachricht |
+ | Zuletzt | Host | Channel | EventID | User | IP | Anzahl |
{{range .RecentEvents}}
- | {{fmtTime .Time}} |
+ {{fmtTime .LastSeen}} |
{{.Hostname}} |
{{.Channel}} |
{{.EventID}} |
{{if .TargetUser}}{{.TargetUser}}{{else}}{{.SubjectUser}}{{end}} |
{{.SrcIP}} |
- {{short .Message 120}} |
+ {{.Count}} |
{{end}}
@@ -1021,11 +1021,13 @@ a:hover {
- | Zeit | Host | Channel | EventID | Target User | Subject User | IP | Workstation | Detail |
+ Erstes Event | Letztes Event | Anzahl | Host | Channel | EventID | Target User | Subject User | IP | Workstation | Status |
{{range .Events}}
- | {{fmtTime .Time}} |
+ {{fmtTime .FirstSeen}} |
+ {{fmtTime .LastSeen}} |
+ {{.Count}} |
{{.Hostname}} |
{{.Channel}} |
{{.EventID}} |
@@ -1033,7 +1035,7 @@ a:hover {
{{.SubjectUser}} |
{{.SrcIP}} |
{{.Workstation}} |
- öffnen |
+ {{.StatusText}} |
{{end}}
@@ -1115,7 +1117,7 @@ a:hover {
SubStatus{{.Event.SubStatusText}}
-
Rohes Event XML
+
Rohes Event XML (nur bei aktivierter Kurzzeit-Speicherung)
{{.Event.Message}}
{{template "footer" .}}
{{end}}
@@ -1169,16 +1171,47 @@ type Config struct {
PartitionInterval time.Duration
PartitionAhead time.Duration
PartitionBehind time.Duration
- PartitionRetention time.Duration
+ PartitionRetention time.Duration // compatibility fallback
+
+ // Metadata-first storage: raw XML is optional, compact event rows are short-lived,
+ // and aggregated occurrences remain queryable for a much longer period.
+ StoreRawXML bool
+ MetadataBucket time.Duration
+ EventRetention time.Duration
+ RawRetention time.Duration
+ MetadataRetention time.Duration
}
type LogPayload struct {
- Hostname string `json:"host"`
- Channel string `json:"channel"`
- EventID uint32 `json:"id"`
- Source string `json:"source"`
- Time time.Time `json:"ts"`
- Message string `json:"msg"`
+ Hostname string `json:"host"`
+ Channel string `json:"channel"`
+ EventID uint32 `json:"id"`
+ Source string `json:"source"`
+ Time time.Time `json:"ts"`
+ Message string `json:"msg,omitempty"`
+ Metadata *EventMetadataPayload `json:"meta,omitempty"`
+}
+
+// EventMetadataPayload lets collectors send already-extracted fields without the
+// original XML. This is the preferred ingestion format for metadata-first mode.
+type EventMetadataPayload struct {
+ Computer string `json:"computer,omitempty"`
+ ProviderName string `json:"provider,omitempty"`
+ TargetUser string `json:"target_user,omitempty"`
+ TargetDomain string `json:"target_domain,omitempty"`
+ SubjectUser string `json:"subject_user,omitempty"`
+ SubjectDomain string `json:"subject_domain,omitempty"`
+ Workstation string `json:"workstation,omitempty"`
+ Device string `json:"device,omitempty"`
+ SrcIP string `json:"src_ip,omitempty"`
+ SrcPort string `json:"src_port,omitempty"`
+ LogonType string `json:"logon_type,omitempty"`
+ ProcessName string `json:"process_name,omitempty"`
+ AuthenticationPackage string `json:"authentication_package,omitempty"`
+ LogonProcess string `json:"logon_process,omitempty"`
+ StatusText string `json:"status,omitempty"`
+ SubStatusText string `json:"sub_status,omitempty"`
+ FailureReason string `json:"failure_reason,omitempty"`
}
type NormalizedEvent struct {
@@ -1294,6 +1327,10 @@ type EventRow struct {
Time time.Time
ReceivedAt time.Time
Message string
+ Count uint64
+ FirstSeen time.Time
+ LastSeen time.Time
+ Aggregated bool
}
type DashboardStats struct {
@@ -1482,9 +1519,10 @@ type PrivilegedUsersPageData struct {
}
type RawEventInsert struct {
- Message string
- SHA256 string
- Time time.Time
+ EventOffset int
+ Message string
+ SHA256 string
+ Time time.Time
}
type EventCountBucketAgg struct {
@@ -1498,9 +1536,30 @@ type EventCountBucketAgg struct {
LastTS time.Time
}
+type MetadataOccurrenceAgg struct {
+ BucketStart time.Time
+ BucketEnd time.Time
+ DimensionKey []byte
+ Hostname string
+ Channel string
+ EventID uint32
+ ProviderName string
+ TargetUser string
+ SubjectUser string
+ SrcIP string
+ Workstation string
+ LogonType string
+ StatusText string
+ FailureReason string
+ Count uint64
+ FirstTS time.Time
+ LastTS time.Time
+}
+
type partitionedTable struct {
Name string
TimeColumn string
+ Retention time.Duration
}
var (
@@ -1838,6 +1897,99 @@ ON DUPLICATE KEY UPDATE
return nil
}
+func upsertEventCatalogTx(ctx context.Context, tx *sql.Tx, buckets map[string]*EventCountBucketAgg) error {
+ if len(buckets) == 0 {
+ return nil
+ }
+ var sb strings.Builder
+ args := make([]any, 0, len(buckets)*7)
+ sb.WriteString(`
+INSERT INTO event_catalog
+(hostname, channel_name, event_id, first_seen, last_seen, total_count, updated_at)
+VALUES
+`)
+ i := 0
+ for _, b := range buckets {
+ if i > 0 {
+ sb.WriteString(",")
+ }
+ i++
+ sb.WriteString("(?,?,?,?,?,?,UTC_TIMESTAMP(6))")
+ args = append(args, b.Hostname, b.Channel, b.EventID, b.FirstTS.UTC(), b.LastTS.UTC(), b.Count)
+ }
+ sb.WriteString(`
+ON DUPLICATE KEY UPDATE
+ first_seen = LEAST(first_seen, VALUES(first_seen)),
+ last_seen = GREATEST(last_seen, VALUES(last_seen)),
+ total_count = total_count + VALUES(total_count),
+ updated_at = UTC_TIMESTAMP(6)
+`)
+ if _, err := tx.ExecContext(ctx, sb.String(), args...); err != nil {
+ return fmt.Errorf("upsert event_catalog: %w", err)
+ }
+ return nil
+}
+
+func metadataDimensionKey(parts ...string) []byte {
+ h := sha256.New()
+ for _, part := range parts {
+ _, _ = io.WriteString(h, part)
+ _, _ = h.Write([]byte{0})
+ }
+ sum := h.Sum(nil)
+ key := make([]byte, 16)
+ copy(key, sum[:16])
+ return key
+}
+
+func upsertMetadataOccurrencesTx(ctx context.Context, tx *sql.Tx, occurrences map[string]*MetadataOccurrenceAgg) error {
+ if len(occurrences) == 0 {
+ return nil
+ }
+
+ var sb strings.Builder
+ args := make([]any, 0, len(occurrences)*16)
+ sb.WriteString(`
+INSERT INTO event_occurrences (
+ bucket_start, bucket_end, dimension_key,
+ hostname, channel_name, event_id, provider_name,
+ target_user, subject_user, src_ip, workstation,
+ logon_type, status_text, failure_reason,
+ cnt, first_event_ts, last_event_ts
+) VALUES
+`)
+
+ i := 0
+ for _, o := range occurrences {
+ if i > 0 {
+ sb.WriteString(",")
+ }
+ i++
+ sb.WriteString("(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)")
+ args = append(args,
+ o.BucketStart.UTC(), o.BucketEnd.UTC(), o.DimensionKey,
+ o.Hostname, o.Channel, o.EventID, o.ProviderName,
+ o.TargetUser, o.SubjectUser, o.SrcIP, o.Workstation,
+ o.LogonType, o.StatusText, o.FailureReason,
+ o.Count, o.FirstTS.UTC(), o.LastTS.UTC(),
+ )
+ }
+
+ sb.WriteString(`
+ON DUPLICATE KEY UPDATE
+ cnt = cnt + VALUES(cnt),
+ first_event_ts = LEAST(first_event_ts, VALUES(first_event_ts)),
+ last_event_ts = GREATEST(last_event_ts, VALUES(last_event_ts)),
+ bucket_end = GREATEST(bucket_end, VALUES(bucket_end)),
+ updated_at = UTC_TIMESTAMP(6)
+`)
+
+ if _, err := tx.ExecContext(ctx, sb.String(), args...); err != nil {
+ return fmt.Errorf("upsert event_occurrences: %w", err)
+ }
+ return nil
+}
+
func (s *server) runSOCLoop() {
ticker := time.NewTicker(1 * time.Minute)
defer ticker.Stop()
@@ -2247,8 +2399,13 @@ func parsePartitionName(name string) (time.Time, bool) {
func (s *server) ensureConfiguredPartitions(ctx context.Context) error {
tables := []partitionedTable{
- {Name: "event_logs", TimeColumn: "ts"},
- {Name: "event_logs_raw", TimeColumn: "ts"},
+ {Name: "event_logs", TimeColumn: "ts", Retention: s.cfg.EventRetention},
+ // Fixed: the actual table is event_log_raw (singular "log"). The old typo
+ // aborted every maintenance run and therefore prevented all partition drops.
+ {Name: "event_log_raw", TimeColumn: "ts", Retention: s.cfg.RawRetention},
+ {Name: "event_occurrences", TimeColumn: "bucket_start", Retention: s.cfg.MetadataRetention},
+ {Name: "event_count_buckets", TimeColumn: "bucket_start", Retention: s.cfg.MetadataRetention},
+ {Name: "ueba_context_buckets", TimeColumn: "bucket_start", Retention: s.cfg.MetadataRetention},
}
for _, tbl := range tables {
@@ -2260,8 +2417,12 @@ func (s *server) ensureConfiguredPartitions(ctx context.Context) error {
return err
}
- if s.cfg.PartitionRetention > 0 {
- if err := s.dropOldPartitions(ctx, tbl.Name, s.cfg.PartitionRetention); err != nil {
+ retention := tbl.Retention
+ if retention <= 0 {
+ retention = s.cfg.PartitionRetention
+ }
+ if retention > 0 {
+ if err := s.dropOldPartitions(ctx, tbl.Name, retention); err != nil {
return err
}
}
@@ -2540,6 +2701,9 @@ LIMIT 1
if err != nil {
return "", err
}
+ if strings.TrimSpace(msg) == "" {
+ return "Roh-XML wurde nicht gespeichert oder ist gemäß RAW_RETENTION bereits abgelaufen.", nil
+ }
return msg, nil
}
@@ -3592,7 +3756,7 @@ func (s *server) getDashboardStats(ctx context.Context) (DashboardStats, error)
if err := s.db.QueryRowContext(ctx, `SELECT COUNT(*) FROM agents WHERE is_enabled = 1 AND last_seen >= ?`, time.Now().UTC().Add(-s.cfg.OfflineAfter)).Scan(&stats.AgentsActive); err != nil {
return stats, err
}
- if err := s.db.QueryRowContext(ctx, `SELECT COUNT(*) FROM event_logs WHERE ts >= ?`, time.Now().UTC().Add(-24*time.Hour)).Scan(&stats.Events24h); err != nil {
+ if err := s.db.QueryRowContext(ctx, `SELECT COALESCE(SUM(cnt), 0) FROM event_occurrences WHERE bucket_start >= ?`, time.Now().UTC().Add(-24*time.Hour)).Scan(&stats.Events24h); err != nil {
return stats, err
}
if err := s.db.QueryRowContext(ctx, `SELECT COUNT(*) FROM detections WHERE created_at >= ?`, time.Now().UTC().Add(-24*time.Hour)).Scan(&stats.Detections24h); err != nil {
@@ -3643,18 +3807,17 @@ func (s *server) listEvents(ctx context.Context, f EventFilter) ([]EventRow, err
f.Limit = 100
}
+ // Deliberately return bucket rows directly. This keeps ORDER BY + LIMIT on the
+ // partition key and avoids a large GROUP BY across the selected time range.
query := `
-SELECT id, hostname, channel_name, event_id, source, computer, provider_name,
- target_user, target_domain, subject_user, subject_domain,
- workstation, src_ip, src_port, logon_type, process_name,
- authentication_package, logon_process, status_text, sub_status_text,
- failure_reason, ts, received_at
-FROM event_logs
+SELECT hostname, channel_name, event_id, provider_name,
+ target_user, subject_user, src_ip, workstation,
+ logon_type, status_text, failure_reason,
+ cnt, first_event_ts, last_event_ts
+FROM event_occurrences
WHERE 1=1
`
-
args := make([]any, 0, 16)
-
if f.Host != "" {
query += ` AND hostname = ?`
args = append(args, f.Host)
@@ -3669,7 +3832,8 @@ WHERE 1=1
}
if f.User != "" {
query += ` AND (target_user = ? OR subject_user = ?)`
- args = append(args, f.User, f.User)
+ u := normalizeUsername(f.User)
+ args = append(args, u, u)
}
if f.SrcIP != "" {
query += ` AND src_ip = ?`
@@ -3677,26 +3841,22 @@ WHERE 1=1
}
if f.TimeFrom != "" {
if t, err := parseUIRFC3339(f.TimeFrom); err == nil {
- query += ` AND ts >= ?`
- args = append(args, t)
+ query += ` AND bucket_start >= ?`
+ args = append(args, bucketStart(t, s.cfg.MetadataBucket))
}
}
if f.TimeTo != "" {
if t, err := parseUIRFC3339(f.TimeTo); err == nil {
- query += ` AND ts <= ?`
+ query += ` AND bucket_start <= ?`
args = append(args, t)
}
}
-
if f.Rule != "" || f.Severity != "" {
query += ` AND EXISTS (
- SELECT 1
- FROM detections d
- WHERE d.hostname = event_logs.hostname
- AND event_logs.ts >= d.window_start
- AND event_logs.ts <= d.window_end
- AND d.created_at >= DATE_SUB(event_logs.ts, INTERVAL 1 DAY)
- `
+ SELECT 1 FROM detections d
+ WHERE d.hostname = event_occurrences.hostname
+ AND event_occurrences.last_event_ts >= d.window_start
+ AND event_occurrences.first_event_ts <= d.window_end`
if f.Rule != "" {
query += ` AND d.rule_name = ?`
args = append(args, f.Rule)
@@ -3707,8 +3867,7 @@ WHERE 1=1
}
query += ` )`
}
-
- query += ` ORDER BY ts DESC LIMIT ?`
+ query += ` ORDER BY bucket_start DESC LIMIT ?`
args = append(args, f.Limit)
rows, err := s.db.QueryContext(ctx, query, args...)
@@ -3721,23 +3880,20 @@ WHERE 1=1
for rows.Next() {
var ev EventRow
if err := rows.Scan(
- &ev.ID, &ev.Hostname, &ev.Channel, &ev.EventID, &ev.Source,
- &ev.Computer, &ev.ProviderName,
- &ev.TargetUser, &ev.TargetDomain, &ev.SubjectUser, &ev.SubjectDomain,
- &ev.Workstation, &ev.SrcIP, &ev.SrcPort, &ev.LogonType, &ev.ProcessName,
- &ev.AuthenticationPackage, &ev.LogonProcess, &ev.StatusText, &ev.SubStatusText,
- &ev.FailureReason, &ev.Time, &ev.ReceivedAt,
+ &ev.Hostname, &ev.Channel, &ev.EventID, &ev.ProviderName,
+ &ev.TargetUser, &ev.SubjectUser, &ev.SrcIP, &ev.Workstation,
+ &ev.LogonType, &ev.StatusText, &ev.FailureReason,
+ &ev.Count, &ev.FirstSeen, &ev.LastSeen,
); err != nil {
return nil, err
}
-
- ev.Time = normalizeTime(ev.Time)
- ev.ReceivedAt = normalizeTime(ev.ReceivedAt)
- ev.Message = eventListSummary(ev)
-
+ ev.FirstSeen = normalizeTime(ev.FirstSeen)
+ ev.LastSeen = normalizeTime(ev.LastSeen)
+ ev.Time = ev.LastSeen
+ ev.Aggregated = true
+ ev.Message = fmt.Sprintf("%s EventID %d: %d Vorkommen", ev.Channel, ev.EventID, ev.Count)
out = append(out, ev)
}
-
return out, rows.Err()
}
@@ -3832,6 +3988,12 @@ func loadConfig() Config {
PartitionAhead: getenvDuration("PARTITION_AHEAD", 24*time.Hour),
PartitionBehind: getenvDuration("PARTITION_BEHIND", 6*time.Hour),
PartitionRetention: getenvDuration("PARTITION_RETENTION", 30*24*time.Hour),
+
+ StoreRawXML: getenvBool("STORE_RAW_XML", false),
+ MetadataBucket: getenvDuration("METADATA_BUCKET", time.Minute),
+ EventRetention: getenvDuration("EVENT_RETENTION", 72*time.Hour),
+ RawRetention: getenvDuration("RAW_RETENTION", 24*time.Hour),
+ MetadataRetention: getenvDuration("METADATA_RETENTION", 180*24*time.Hour),
}
}
@@ -4109,13 +4271,14 @@ func (s *server) insertBatch(ctx context.Context, agentID uint64, batch []LogPay
defer func() { _ = tx.Rollback() }()
var sb strings.Builder
- args := make([]any, 0, len(batch)*28)
+ args := make([]any, 0, len(batch)*29)
sb.WriteString(`
INSERT INTO event_logs (
agent_id, hostname, channel_name, event_id, source, computer, provider_name,
level_value, task_value, opcode_value, keywords,
- target_user, target_domain, subject_user, subject_domain,
+ target_user, target_user_norm, target_domain,
+ subject_user, subject_user_norm, subject_domain,
workstation, src_ip, src_port, logon_type, process_name,
authentication_package, logon_process, status_text, sub_status_text,
failure_reason, ts, received_at, msg, msg_sha256
@@ -4123,60 +4286,51 @@ INSERT INTO event_logs (
`)
bucketAgg := make(map[string]*EventCountBucketAgg)
+ occurrenceAgg := make(map[string]*MetadataOccurrenceAgg)
rawEvents := make([]RawEventInsert, 0, len(batch))
for i, item := range batch {
if i > 0 {
sb.WriteString(",")
}
- sb.WriteString("(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,UTC_TIMESTAMP(6),'',?)")
+ sb.WriteString("(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,UTC_TIMESTAMP(6),'',?)")
- msgHash := sha256Hex(item.Message)
- rawEvents = append(rawEvents, RawEventInsert{
- Message: item.Message,
- SHA256: msgHash,
- Time: item.Time.UTC(),
- })
-
- norm := NormalizeEventXML(item.Message)
-
- realHost := firstNonEmpty(norm.Computer, item.Hostname)
- targetUser := normalizeUsername(norm.TargetUser)
- subjectUser := normalizeUsername(norm.SubjectUser)
-
- bucketSize := s.cfg.BaselineWindow
- if bucketSize <= 0 {
- bucketSize = 5 * time.Minute
+ msgHash := payloadFingerprint(item)
+ if s.cfg.StoreRawXML && strings.TrimSpace(item.Message) != "" {
+ rawEvents = append(rawEvents, RawEventInsert{
+ EventOffset: i,
+ Message: item.Message,
+ SHA256: msgHash,
+ Time: item.Time.UTC(),
+ })
}
- bs := bucketStart(item.Time, bucketSize)
- be := bs.Add(bucketSize)
-
- bucketKey := fmt.Sprintf(
- "%s|%s|%s|%d",
- bs.Format(time.RFC3339Nano),
- realHost,
- item.Channel,
- item.EventID,
- )
+ norm := NormalizeEventXML(item.Message)
+ if item.Metadata != nil {
+ norm = overlayMetadata(norm, *item.Metadata)
+ }
+ realHost := firstNonEmpty(norm.Computer, item.Hostname)
+ targetUserNorm := normalizeUsername(norm.TargetUser)
+ subjectUserNorm := normalizeUsername(norm.SubjectUser)
+ providerName := firstNonEmpty(norm.ProviderName, item.Source)
+ eventTS := item.Time.UTC()
+ baselineBucketSize := s.cfg.BaselineWindow
+ if baselineBucketSize <= 0 {
+ baselineBucketSize = 5 * time.Minute
+ }
+ bs := bucketStart(eventTS, baselineBucketSize)
+ be := bs.Add(baselineBucketSize)
+ bucketKey := fmt.Sprintf("%s|%s|%s|%d", bs.Format(time.RFC3339Nano), realHost, item.Channel, item.EventID)
agg := bucketAgg[bucketKey]
if agg == nil {
agg = &EventCountBucketAgg{
- BucketStart: bs,
- BucketEnd: be,
- Hostname: realHost,
- Channel: item.Channel,
- EventID: item.EventID,
- FirstTS: item.Time.UTC(),
- LastTS: item.Time.UTC(),
+ BucketStart: bs, BucketEnd: be, Hostname: realHost, Channel: item.Channel,
+ EventID: item.EventID, FirstTS: eventTS, LastTS: eventTS,
}
bucketAgg[bucketKey] = agg
}
-
agg.Count++
-
- eventTS := item.Time.UTC()
if eventTS.Before(agg.FirstTS) {
agg.FirstTS = eventTS
}
@@ -4184,64 +4338,76 @@ INSERT INTO event_logs (
agg.LastTS = eventTS
}
- // Erfolgreicher Logon
- if item.Channel == "Security" && item.EventID == 4624 && targetUser != "" {
- priv, err := s.detector.isPrivilegedUser(ctx, targetUser)
- if err != nil {
- s.logger.Printf("privileged user check failed for %s: %v", targetUser, err)
- } else if priv {
- s.detector.privilegedLogonsTotal.WithLabelValues(targetUser, realHost).Inc()
+ metadataBucketSize := s.cfg.MetadataBucket
+ if metadataBucketSize <= 0 {
+ metadataBucketSize = time.Minute
+ }
+ mbs := bucketStart(eventTS, metadataBucketSize)
+ mbe := mbs.Add(metadataBucketSize)
+ failureReason := norm.FailureReason
+ if len(failureReason) > 255 {
+ failureReason = failureReason[:255]
+ }
+ dimKey := metadataDimensionKey(
+ realHost, item.Channel, strconv.FormatUint(uint64(item.EventID), 10), providerName,
+ targetUserNorm, subjectUserNorm, norm.SrcIP, norm.Workstation,
+ norm.LogonType, norm.StatusText, failureReason,
+ )
+ occKey := mbs.Format(time.RFC3339Nano) + "|" + hex.EncodeToString(dimKey)
+ occ := occurrenceAgg[occKey]
+ if occ == nil {
+ occ = &MetadataOccurrenceAgg{
+ BucketStart: mbs, BucketEnd: mbe, DimensionKey: dimKey,
+ Hostname: realHost, Channel: item.Channel, EventID: item.EventID,
+ ProviderName: providerName, TargetUser: targetUserNorm, SubjectUser: subjectUserNorm,
+ SrcIP: norm.SrcIP, Workstation: norm.Workstation, LogonType: norm.LogonType,
+ StatusText: norm.StatusText, FailureReason: failureReason,
+ FirstTS: eventTS, LastTS: eventTS,
}
+ occurrenceAgg[occKey] = occ
+ }
+ occ.Count++
+ if eventTS.Before(occ.FirstTS) {
+ occ.FirstTS = eventTS
+ }
+ if eventTS.After(occ.LastTS) {
+ occ.LastTS = eventTS
}
- // Fehlgeschlagener Logon
- if item.Channel == "Security" && item.EventID == 4625 && targetUser != "" {
- priv, err := s.detector.isPrivilegedUser(ctx, targetUser)
+ // Privileged-user metrics are based only on compact normalized metadata.
+ if item.Channel == "Security" && item.EventID == 4624 && targetUserNorm != "" {
+ priv, err := s.detector.isPrivilegedUser(ctx, targetUserNorm)
if err != nil {
- s.logger.Printf("privileged user check failed for %s: %v", targetUser, err)
+ s.logger.Printf("privileged user check failed for %s: %v", targetUserNorm, err)
} else if priv {
- s.detector.privilegedLogonFailuresTotal.WithLabelValues(targetUser, realHost).Inc()
+ s.detector.privilegedLogonsTotal.WithLabelValues(targetUserNorm, realHost).Inc()
}
}
-
- // Special Privileged Logon
- if item.Channel == "Security" && item.EventID == 4672 && subjectUser != "" {
- priv, err := s.detector.isPrivilegedUser(ctx, subjectUser)
+ if item.Channel == "Security" && item.EventID == 4625 && targetUserNorm != "" {
+ priv, err := s.detector.isPrivilegedUser(ctx, targetUserNorm)
if err != nil {
- s.logger.Printf("privileged user check failed for %s: %v", subjectUser, err)
+ s.logger.Printf("privileged user check failed for %s: %v", targetUserNorm, err)
} else if priv {
- s.detector.privilegedLogonsTotal.WithLabelValues(subjectUser, realHost).Inc()
+ s.detector.privilegedLogonFailuresTotal.WithLabelValues(targetUserNorm, realHost).Inc()
+ }
+ }
+ if item.Channel == "Security" && item.EventID == 4672 && subjectUserNorm != "" {
+ priv, err := s.detector.isPrivilegedUser(ctx, subjectUserNorm)
+ if err != nil {
+ s.logger.Printf("privileged user check failed for %s: %v", subjectUserNorm, err)
+ } else if priv {
+ s.detector.privilegedLogonsTotal.WithLabelValues(subjectUserNorm, realHost).Inc()
}
}
args = append(args,
- agentID,
- realHost,
- item.Channel,
- item.EventID,
- item.Source,
- norm.Computer,
- firstNonEmpty(norm.ProviderName, item.Source),
- norm.LevelValue,
- norm.TaskValue,
- norm.OpcodeValue,
- norm.Keywords,
- norm.TargetUser,
- norm.TargetDomain,
- norm.SubjectUser,
- norm.SubjectDomain,
- norm.Workstation,
- norm.SrcIP,
- norm.SrcPort,
- norm.LogonType,
- norm.ProcessName,
- norm.AuthenticationPackage,
- norm.LogonProcess,
- norm.StatusText,
- norm.SubStatusText,
- norm.FailureReason,
- item.Time.UTC(),
- msgHash,
+ agentID, realHost, item.Channel, item.EventID, item.Source, norm.Computer, providerName,
+ norm.LevelValue, norm.TaskValue, norm.OpcodeValue, norm.Keywords,
+ norm.TargetUser, targetUserNorm, norm.TargetDomain,
+ norm.SubjectUser, subjectUserNorm, norm.SubjectDomain,
+ norm.Workstation, norm.SrcIP, norm.SrcPort, norm.LogonType, norm.ProcessName,
+ norm.AuthenticationPackage, norm.LogonProcess, norm.StatusText, norm.SubStatusText,
+ norm.FailureReason, eventTS, msgHash,
)
}
@@ -4249,29 +4415,32 @@ INSERT INTO event_logs (
if err != nil {
return err
}
-
firstID, err := res.LastInsertId()
if err != nil {
return fmt.Errorf("event_logs LastInsertId: %w", err)
}
-
affected, err := res.RowsAffected()
if err != nil {
return fmt.Errorf("event_logs RowsAffected: %w", err)
}
-
if affected != int64(len(batch)) {
return fmt.Errorf("event_logs insert affected %d rows, expected %d", affected, len(batch))
}
- if err := insertRawEventsTx(ctx, tx, uint64(firstID), rawEvents); err != nil {
- return err
+ if s.cfg.StoreRawXML {
+ if err := insertRawEventsTx(ctx, tx, uint64(firstID), rawEvents); err != nil {
+ return err
+ }
}
-
if err := upsertEventCountBucketsTx(ctx, tx, bucketAgg); err != nil {
return err
}
-
+ if err := upsertEventCatalogTx(ctx, tx, bucketAgg); err != nil {
+ return err
+ }
+ if err := upsertMetadataOccurrencesTx(ctx, tx, occurrenceAgg); err != nil {
+ return err
+ }
if err := tx.Commit(); err != nil {
return err
}
@@ -4351,7 +4520,7 @@ VALUES
sb.WriteString("(?,?,?,?,UTC_TIMESTAMP(6))")
args = append(args,
- firstEventID+uint64(i),
+ firstEventID+uint64(raw.EventOffset),
raw.Time.UTC(),
raw.Message,
raw.SHA256,
@@ -5333,15 +5502,14 @@ func (d *detector) runUEBABaselineUpdate(ctx context.Context) error {
windowStart := windowEnd.Add(-d.cfg.UEBANewContextWindow)
rows, err := d.db.QueryContext(ctx, `
-SELECT target_user, hostname, src_ip, workstation, COUNT(*)
+SELECT target_user_norm, hostname, src_ip, workstation, COUNT(*)
FROM event_logs
WHERE channel_name = 'Security'
AND event_id = 4624
AND ts >= ? AND ts < ?
- AND target_user <> ''
- AND target_user <> '-'
- AND target_user NOT LIKE '%$'
-GROUP BY target_user, hostname, src_ip, workstation
+ AND target_user_norm <> ''
+ AND target_user_norm NOT LIKE '%$'
+GROUP BY target_user_norm, hostname, src_ip, workstation
`, windowStart, windowEnd)
if err != nil {
return err
@@ -5391,7 +5559,7 @@ func (d *detector) runUEBANewUserContextRule(ctx context.Context) error {
const q = `
SELECT
e.hostname,
- e.target_user,
+ e.target_user_norm,
e.src_ip,
e.workstation,
MIN(e.ts) AS first_seen,
@@ -5400,10 +5568,9 @@ FROM event_logs e
WHERE e.channel_name = 'Security'
AND e.event_id = 4624
AND e.ts >= ? AND e.ts < ?
- AND e.target_user <> ''
- AND e.target_user <> '-'
- AND e.target_user NOT LIKE '%$'
-GROUP BY e.hostname, e.target_user, e.src_ip, e.workstation
+ AND e.target_user_norm <> ''
+ AND e.target_user_norm NOT LIKE '%$'
+GROUP BY e.hostname, e.target_user_norm, e.src_ip, e.workstation
`
rows, err := d.db.QueryContext(ctx, q, windowStart, windowEnd)
@@ -5810,15 +5977,15 @@ func (d *detector) runFailedLogonSpikeRule(ctx context.Context) error {
windowStart := windowEnd.Add(-d.cfg.FailedLogonWindow)
const q = `
-SELECT hostname, COUNT(*) AS cnt
-FROM event_logs
+SELECT hostname, SUM(cnt) AS cnt
+FROM event_occurrences
WHERE channel_name = 'Security'
AND event_id = 4625
- AND ts >= ? AND ts < ?
+ AND bucket_start >= ? AND bucket_start < ?
GROUP BY hostname
-HAVING COUNT(*) >= ?
+HAVING SUM(cnt) >= ?
`
- rows, err := d.db.QueryContext(ctx, q, windowStart, windowEnd, d.cfg.FailedLogonThreshold)
+ rows, err := d.db.QueryContext(ctx, q, bucketStart(windowStart, d.cfg.MetadataBucket), windowEnd, d.cfg.FailedLogonThreshold)
if err != nil {
return err
}
@@ -5867,15 +6034,15 @@ func (d *detector) runRebootSpikeRule(ctx context.Context) error {
windowStart := windowEnd.Add(-d.cfg.RebootWindow)
const q = `
-SELECT hostname, COUNT(*) AS cnt
-FROM event_logs
+SELECT hostname, SUM(cnt) AS cnt
+FROM event_occurrences
WHERE channel_name = 'System'
AND event_id IN (1074, 6005, 6006)
- AND ts >= ? AND ts < ?
+ AND bucket_start >= ? AND bucket_start < ?
GROUP BY hostname
-HAVING COUNT(*) >= ?
+HAVING SUM(cnt) >= ?
`
- rows, err := d.db.QueryContext(ctx, q, windowStart, windowEnd, d.cfg.RebootThreshold)
+ rows, err := d.db.QueryContext(ctx, q, bucketStart(windowStart, d.cfg.MetadataBucket), windowEnd, d.cfg.RebootThreshold)
if err != nil {
return err
}
@@ -5923,20 +6090,12 @@ func (d *detector) runNewEventIDRule(ctx context.Context) error {
windowStart := windowEnd.Add(-d.cfg.DetectionInterval)
const q = `
-SELECT e.hostname, e.channel_name, e.event_id, COUNT(*) AS cnt
-FROM event_logs e
-WHERE e.ts >= ? AND e.ts < ?
- AND NOT EXISTS (
- SELECT 1
- FROM event_logs old
- WHERE old.hostname = e.hostname
- AND old.channel_name = e.channel_name
- AND old.event_id = e.event_id
- AND old.ts < ?
- )
-GROUP BY e.hostname, e.channel_name, e.event_id
+SELECT hostname, channel_name, event_id, total_count
+FROM event_catalog
+WHERE first_seen >= ? AND first_seen < ?
+ORDER BY first_seen ASC
`
- rows, err := d.db.QueryContext(ctx, q, windowStart, windowEnd, windowStart)
+ rows, err := d.db.QueryContext(ctx, q, windowStart, windowEnd)
if err != nil {
return err
}
@@ -5988,17 +6147,17 @@ func (d *detector) runPasswordSprayRule(ctx context.Context) error {
windowStart := windowEnd.Add(-d.cfg.PasswordSprayWindow)
const q = `
-SELECT hostname, src_ip, COUNT(*) AS attempts, COUNT(DISTINCT target_user) AS users
-FROM event_logs
+SELECT hostname, src_ip, SUM(cnt) AS attempts, COUNT(DISTINCT target_user) AS users
+FROM event_occurrences
WHERE channel_name = 'Security'
AND event_id = 4625
- AND ts >= ? AND ts < ?
+ AND bucket_start >= ? AND bucket_start < ?
AND src_ip <> '' AND src_ip <> '-' AND src_ip <> '::1' AND src_ip <> '127.0.0.1'
AND target_user <> '' AND target_user <> '-'
GROUP BY hostname, src_ip
-HAVING COUNT(*) >= ? AND COUNT(DISTINCT target_user) >= ?
+HAVING SUM(cnt) >= ? AND COUNT(DISTINCT target_user) >= ?
`
- rows, err := d.db.QueryContext(ctx, q, windowStart, windowEnd, d.cfg.PasswordSprayMinAttempts, d.cfg.PasswordSprayMinUsers)
+ rows, err := d.db.QueryContext(ctx, q, bucketStart(windowStart, d.cfg.MetadataBucket), windowEnd, d.cfg.PasswordSprayMinAttempts, d.cfg.PasswordSprayMinUsers)
if err != nil {
return err
}
@@ -6048,19 +6207,19 @@ func (d *detector) runSuccessAfterFailuresRule(ctx context.Context) error {
windowStart := windowEnd.Add(-d.cfg.SuccessAfterFailureWindow)
const q = `
-SELECT s.hostname, s.target_user, s.src_ip, COUNT(*) AS success_count
+SELECT s.hostname, s.target_user_norm, s.src_ip, COUNT(*) AS success_count
FROM event_logs s
WHERE s.channel_name = 'Security'
AND s.event_id = 4624
AND s.ts >= ? AND s.ts < ?
- AND s.target_user <> '' AND s.target_user <> '-'
+ AND s.target_user_norm <> ''
AND EXISTS (
SELECT 1
FROM event_logs f
WHERE f.hostname = s.hostname
AND f.channel_name = 'Security'
AND f.event_id = 4625
- AND f.target_user = s.target_user
+ AND f.target_user_norm = s.target_user_norm
AND (
(f.src_ip = s.src_ip AND s.src_ip <> '' AND s.src_ip <> '-')
OR (s.src_ip = '' OR s.src_ip = '-' OR f.src_ip = '' OR f.src_ip = '-')
@@ -6068,7 +6227,7 @@ WHERE s.channel_name = 'Security'
AND f.ts >= DATE_SUB(s.ts, INTERVAL ? SECOND)
AND f.ts < s.ts
)
-GROUP BY s.hostname, s.target_user, s.src_ip
+GROUP BY s.hostname, s.target_user_norm, s.src_ip
`
rows, err := d.db.QueryContext(ctx, q, windowStart, windowEnd, int(d.cfg.SuccessAfterFailureWindow.Seconds()))
if err != nil {
@@ -6121,26 +6280,25 @@ func (d *detector) runNewSourceIPForUserRule(ctx context.Context) error {
windowStart := windowEnd.Add(-d.cfg.NewSourceIPWindow)
const q = `
-SELECT e.hostname, e.target_user, e.src_ip, MIN(e.ts) AS first_seen, COUNT(*) AS cnt
+SELECT e.hostname, e.target_user_norm, e.src_ip, MIN(e.ts) AS first_seen, COUNT(*) AS cnt
FROM event_logs e
WHERE e.channel_name = 'Security'
AND e.event_id = 4624
AND e.ts >= ? AND e.ts < ?
- AND e.target_user <> ''
- AND e.target_user <> '-'
- AND e.target_user NOT LIKE '%$'
+ AND e.target_user_norm <> ''
+ AND e.target_user_norm NOT LIKE '%$'
AND e.src_ip <> ''
AND e.src_ip <> '-'
AND e.src_ip <> '::1'
AND e.src_ip <> '127.0.0.1'
- AND LOWER(e.target_user) NOT IN (
+ AND e.target_user_norm NOT IN (
'system',
'localsystem',
'local service',
'network service',
'anonymous logon'
)
-GROUP BY e.hostname, e.target_user, e.src_ip
+GROUP BY e.hostname, e.target_user_norm, e.src_ip
`
rows, err := d.db.QueryContext(ctx, q, windowStart, windowEnd)
@@ -6361,6 +6519,44 @@ VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, UTC_TIMESTAMP(6), ?)
return affected > 0, nil
}
+func payloadFingerprint(p LogPayload) string {
+ if strings.TrimSpace(p.Message) != "" {
+ return sha256Hex(p.Message)
+ }
+ b, err := json.Marshal(struct {
+ Host string `json:"host"`
+ Channel string `json:"channel"`
+ ID uint32 `json:"id"`
+ Source string `json:"source"`
+ Time time.Time `json:"ts"`
+ Meta *EventMetadataPayload `json:"meta"`
+ }{p.Hostname, p.Channel, p.EventID, p.Source, p.Time.UTC(), p.Metadata})
+ if err != nil {
+ return sha256Hex(fmt.Sprintf("%s|%s|%d|%s", p.Hostname, p.Channel, p.EventID, p.Time.UTC().Format(time.RFC3339Nano)))
+ }
+ return sha256Hex(string(b))
+}
+
+func overlayMetadata(base NormalizedEvent, m EventMetadataPayload) NormalizedEvent {
+ base.Computer = firstNonEmpty(m.Computer, base.Computer)
+ base.ProviderName = firstNonEmpty(m.ProviderName, base.ProviderName)
+ base.TargetUser = firstNonEmpty(m.TargetUser, base.TargetUser)
+ base.TargetDomain = firstNonEmpty(m.TargetDomain, base.TargetDomain)
+ base.SubjectUser = firstNonEmpty(m.SubjectUser, base.SubjectUser)
+ base.SubjectDomain = firstNonEmpty(m.SubjectDomain, base.SubjectDomain)
+ base.Workstation = firstNonEmpty(m.Workstation, m.Device, base.Workstation)
+ base.SrcIP = firstNonEmpty(m.SrcIP, base.SrcIP)
+ base.SrcPort = firstNonEmpty(m.SrcPort, base.SrcPort)
+ base.LogonType = firstNonEmpty(m.LogonType, base.LogonType)
+ base.ProcessName = firstNonEmpty(m.ProcessName, base.ProcessName)
+ base.AuthenticationPackage = firstNonEmpty(m.AuthenticationPackage, base.AuthenticationPackage)
+ base.LogonProcess = firstNonEmpty(m.LogonProcess, base.LogonProcess)
+ base.StatusText = firstNonEmpty(m.StatusText, base.StatusText)
+ base.SubStatusText = firstNonEmpty(m.SubStatusText, base.SubStatusText)
+ base.FailureReason = firstNonEmpty(m.FailureReason, base.FailureReason)
+ return base
+}
+
func NormalizeEventXML(xmlStr string) NormalizedEvent {
var out NormalizedEvent
if strings.TrimSpace(xmlStr) == "" {
@@ -6447,7 +6643,7 @@ func NormalizeEventXML(xmlStr string) NormalizedEvent {
out.SubjectUser = v
case "SubjectDomainName":
out.SubjectDomain = v
- case "WorkstationName":
+ case "WorkstationName", "CallerComputerName":
out.Workstation = v
case "IpAddress":
out.SrcIP = v
@@ -6561,8 +6757,8 @@ func validatePayload(p *LogPayload) error {
if p.EventID == 0 {
return errors.New("event id is required")
}
- if p.Message == "" {
- return errors.New("msg is required")
+ if strings.TrimSpace(p.Message) == "" && p.Metadata == nil {
+ return errors.New("either msg or meta is required")
}
if len(p.Message) > 2*1024*1024 {
return errors.New("msg too large")
@@ -6732,23 +6928,22 @@ func (d *detector) runOffHoursLoginRule(ctx context.Context) error {
windowStart := windowEnd.Add(-5 * time.Minute)
const q = `
-SELECT hostname, target_user, COUNT(*) AS cnt
+SELECT hostname, target_user_norm, COUNT(*) AS cnt
FROM event_logs
WHERE channel_name = 'Security'
AND event_id = 4624
AND ts >= ? AND ts < ?
- AND target_user <> ''
- AND target_user <> '-'
- AND target_user NOT LIKE '%$'
+ AND target_user_norm <> ''
+ AND target_user_norm NOT LIKE '%$'
AND logon_type IN ('2', '7', '10', '11')
- AND LOWER(target_user) NOT IN (
+ AND target_user_norm NOT IN (
'system',
'localsystem',
'local service',
'network service',
'anonymous logon'
)
-GROUP BY hostname, target_user
+GROUP BY hostname, target_user_norm
`
rows, err := d.db.QueryContext(ctx, q, windowStart, windowEnd)
@@ -6857,15 +7052,14 @@ func (d *detector) runFirstTimePrivilegedRule(ctx context.Context) error {
windowStart := windowEnd.Add(-10 * time.Minute)
const q = `
-SELECT hostname, subject_user, COUNT(*) AS cnt
+SELECT hostname, subject_user_norm, COUNT(*) AS cnt
FROM event_logs
WHERE channel_name = 'Security'
AND event_id = 4672
AND ts >= ? AND ts < ?
- AND subject_user <> ''
- AND subject_user <> '-'
- AND subject_user NOT LIKE '%$'
-GROUP BY hostname, subject_user
+ AND subject_user_norm <> ''
+ AND subject_user_norm NOT LIKE '%$'
+GROUP BY hostname, subject_user_norm
`
rows, err := d.db.QueryContext(ctx, q, windowStart, windowEnd)
@@ -7010,23 +7204,22 @@ func (d *detector) runAdminNewHostRule(ctx context.Context) error {
lookbackStart := windowEnd.Add(-d.cfg.UEBALookback)
rows, err := d.db.QueryContext(ctx, `
-SELECT e.hostname, e.target_user, COUNT(*) AS cnt
+SELECT e.hostname, e.target_user_norm, COUNT(*) AS cnt
FROM event_logs e
WHERE e.channel_name = 'Security'
AND e.event_id = 4624
AND e.ts >= ? AND e.ts < ?
- AND e.target_user <> ''
- AND e.target_user <> '-'
- AND e.target_user NOT LIKE '%$'
+ AND e.target_user_norm <> ''
+ AND e.target_user_norm NOT LIKE '%$'
AND NOT EXISTS (
SELECT 1
FROM ueba_user_baseline b
- WHERE b.username = e.target_user
+ WHERE b.username = e.target_user_norm
AND b.hostname = e.hostname
AND b.last_seen >= ?
AND b.first_seen < ?
)
-GROUP BY e.hostname, e.target_user
+GROUP BY e.hostname, e.target_user_norm
`, windowStart, windowEnd, lookbackStart, windowStart)
if err != nil {
return err
diff --git a/schema.sql b/schema.sql
index 64cba70..f007780 100644
--- a/schema.sql
+++ b/schema.sql
@@ -1,101 +1,11 @@
-CREATE DATABASE IF NOT EXISTS eventcollector
- CHARACTER SET utf8mb4
- COLLATE utf8mb4_unicode_ci;
-
-USE eventcollector;
-
-CREATE TABLE IF NOT EXISTS agents (
- id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT,
- hostname VARCHAR(255) NOT NULL,
- api_key_hash CHAR(64) NOT NULL,
- first_seen DATETIME(6) NOT NULL DEFAULT CURRENT_TIMESTAMP(6),
- last_seen DATETIME(6) NOT NULL DEFAULT CURRENT_TIMESTAMP(6),
- last_ip VARCHAR(64) NOT NULL DEFAULT '',
- is_enabled TINYINT(1) NOT NULL DEFAULT 1,
- PRIMARY KEY (id),
- UNIQUE KEY ux_agents_hostname (hostname),
- KEY ix_agents_last_seen (last_seen)
-) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci;
-
-CREATE TABLE IF NOT EXISTS event_logs (
- id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT,
- agent_id BIGINT UNSIGNED NOT NULL,
- hostname VARCHAR(255) NOT NULL,
- channel_name VARCHAR(128) NOT NULL,
- event_id INT UNSIGNED NOT NULL,
- source VARCHAR(255) NOT NULL,
- ts DATETIME(6) NOT NULL,
- received_at DATETIME(6) NOT NULL DEFAULT CURRENT_TIMESTAMP(6),
- msg LONGTEXT NOT NULL,
- msg_sha256 CHAR(64) NOT NULL,
- PRIMARY KEY (id),
- KEY ix_event_logs_ts (ts),
- KEY ix_event_logs_received_at (received_at),
- KEY ix_event_logs_agent_ts (agent_id, ts),
- KEY ix_event_logs_eventid_ts (event_id, ts),
- KEY ix_event_logs_hostname_ts (hostname, ts),
- KEY ix_event_logs_channel_event_ts (channel_name, event_id, ts),
- CONSTRAINT fk_event_logs_agent
- FOREIGN KEY (agent_id) REFERENCES agents(id)
- ON DELETE RESTRICT
- ON UPDATE RESTRICT
-) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci;
-
-CREATE TABLE IF NOT EXISTS detections (
- id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT,
- rule_name VARCHAR(128) NOT NULL,
- severity VARCHAR(32) NOT NULL,
- hostname VARCHAR(255) NOT NULL,
- channel_name VARCHAR(128) NOT NULL DEFAULT '',
- event_id INT UNSIGNED NOT NULL DEFAULT 0,
- score DOUBLE NOT NULL DEFAULT 0,
- window_start DATETIME(6) NOT NULL,
- window_end DATETIME(6) NOT NULL,
- summary VARCHAR(512) NOT NULL,
- details_json JSON NOT NULL,
- created_at DATETIME(6) NOT NULL DEFAULT CURRENT_TIMESTAMP(6),
- PRIMARY KEY (id),
- UNIQUE KEY ux_detection_dedupe (rule_name, hostname, channel_name, event_id, window_start, window_end),
- KEY ix_detections_created (created_at),
- KEY ix_detections_rule_host_time (rule_name, hostname, created_at),
- KEY ix_detections_severity_time (severity, created_at)
-) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci;
-
-USE eventcollector;
-
-INSERT INTO agents (hostname, api_key_hash)
-VALUES
- ('client01.domain.local', SHA2('SUPER-LANGER-AGENT-KEY-01', 256)),
- ('client02.domain.local', SHA2('SUPER-LANGER-AGENT-KEY-02', 256));
-
- #V2
-
- ALTER TABLE event_logs
- ADD COLUMN computer VARCHAR(255) NOT NULL DEFAULT '' AFTER source,
- ADD COLUMN provider_name VARCHAR(255) NOT NULL DEFAULT '' AFTER computer,
- ADD COLUMN level_value INT UNSIGNED NOT NULL DEFAULT 0 AFTER provider_name,
- ADD COLUMN task_value INT UNSIGNED NOT NULL DEFAULT 0 AFTER level_value,
- ADD COLUMN opcode_value INT UNSIGNED NOT NULL DEFAULT 0 AFTER task_value,
- ADD COLUMN keywords VARCHAR(255) NOT NULL DEFAULT '' AFTER opcode_value,
- ADD COLUMN target_user VARCHAR(255) NOT NULL DEFAULT '' AFTER keywords,
- ADD COLUMN target_domain VARCHAR(255) NOT NULL DEFAULT '' AFTER target_user,
- ADD COLUMN subject_user VARCHAR(255) NOT NULL DEFAULT '' AFTER target_domain,
- ADD COLUMN subject_domain VARCHAR(255) NOT NULL DEFAULT '' AFTER subject_user,
- ADD COLUMN workstation VARCHAR(255) NOT NULL DEFAULT '' AFTER subject_domain,
- ADD COLUMN src_ip VARCHAR(64) NOT NULL DEFAULT '' AFTER workstation,
- ADD COLUMN src_port VARCHAR(32) NOT NULL DEFAULT '' AFTER src_ip,
- ADD COLUMN logon_type VARCHAR(32) NOT NULL DEFAULT '' AFTER src_port,
- ADD COLUMN process_name VARCHAR(512) NOT NULL DEFAULT '' AFTER logon_type,
- ADD COLUMN authentication_package VARCHAR(128) NOT NULL DEFAULT '' AFTER process_name,
- ADD COLUMN logon_process VARCHAR(128) NOT NULL DEFAULT '' AFTER authentication_package,
- ADD COLUMN status_text VARCHAR(64) NOT NULL DEFAULT '' AFTER logon_process,
- ADD COLUMN sub_status_text VARCHAR(64) NOT NULL DEFAULT '' AFTER status_text,
- ADD COLUMN failure_reason VARCHAR(512) NOT NULL DEFAULT '' AFTER sub_status_text;
-
-ALTER TABLE event_logs
- ADD KEY ix_event_logs_target_user_ts (target_user, ts),
- ADD KEY ix_event_logs_src_ip_ts (src_ip, ts),
- ADD KEY ix_event_logs_target_user_src_ip_ts (target_user, src_ip, ts),
- ADD KEY ix_event_logs_eventid_srcip_ts (event_id, src_ip, ts),
- ADD KEY ix_event_logs_eventid_targetuser_ts (event_id, target_user, ts),
- ADD KEY ix_event_logs_eventid_logontype_ts (event_id, logon_type, ts);
\ No newline at end of file
+-- Canonical schema notice
+--
+-- New empty database:
+-- deploy/mariadb/init/001-schema.sql
+--
+-- Existing database:
+-- deploy/mariadb/migrations/002-metadata-first.sql
+--
+-- The former schema.sql was an obsolete, non-idempotent development schema and
+-- has been retained as schema.legacy.sql only for historical reference. Do not
+-- run schema.legacy.sql against a production database.