From 10541fb22d5a2f1120464621a725134a88bc1c82 Mon Sep 17 00:00:00 2001 From: jbergner Date: Sat, 18 Jul 2026 17:16:10 +0200 Subject: [PATCH] Update to META-Data Only --- README.md | 68 ++++ compose.yml | 11 + deploy/mariadb/init/001-schema.sql | 69 +++- dot_env | 18 +- main.go | 627 +++++++++++++++++++---------- schema.sql | 112 +----- 6 files changed, 584 insertions(+), 321 deletions(-) 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 { -

Neueste Events

+

Neueste Event-Metadaten

- + {{range .RecentEvents}} - + - + {{end}}
ZeitHostChannelEventIDUserIPNachricht
ZuletztHostChannelEventIDUserIPAnzahl
{{fmtTime .Time}}{{fmtTime .LastSeen}} {{.Hostname}} {{.Channel}} {{.EventID}} {{if .TargetUser}}{{.TargetUser}}{{else}}{{.SubjectUser}}{{end}} {{.SrcIP}}{{short .Message 120}}{{.Count}}
@@ -1021,11 +1021,13 @@ a:hover {
- + {{range .Events}} - + + + @@ -1033,7 +1035,7 @@ a:hover { - + {{end}}
ZeitHostChannelEventIDTarget UserSubject UserIPWorkstationDetailErstes EventLetztes EventAnzahlHostChannelEventIDTarget UserSubject UserIPWorkstationStatus
{{fmtTime .Time}}{{fmtTime .FirstSeen}}{{fmtTime .LastSeen}}{{.Count}} {{.Hostname}} {{.Channel}} {{.EventID}}{{.SubjectUser}} {{.SrcIP}} {{.Workstation}}öffnen{{.StatusText}}
@@ -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.