Documentation
¶
Overview ¶
Package correlation is the pure-domain, deterministic engine that folds runtime detection signals into incidents for Phase C of the EDR data plane (#594, C2 #676). It groups related signals by asset + entity within an event-time session window into a single incident, emitting the incident.IncidentEvent log that C1's Project folds and C7 persists. Correlation is EVENT-TIME, never ingest-order: signals are ordered by their OccurredAt (with a stable id tiebreak) before folding, so an out-of-order input yields the same incidents. Deduplication (by signal id), a bounded session window, and anti-storm suppression keep a flood of repeats from exploding an incident while staying coverage-honest (the suppressed count is recorded, never silently dropped).
Index ¶
- func Correlate(cfg Config, signals []Signal) ([]incident.IncidentEvent, error)
- func SameSignal(left, right Signal) bool
- type ActiveSession
- type Assignment
- type AssignmentOutcome
- type Checkpoint
- type Config
- type IncrementalPlan
- type Phase
- type Signal
- type SignalPosition
- type SourcePosition
- type State
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func Correlate ¶
func Correlate(cfg Config, signals []Signal) ([]incident.IncidentEvent, error)
Correlate folds detection signals into incidents and returns the incident.IncidentEvent log to persist and project. It is deterministic and event-time-based: signals are ordered by (OccurredAt, ID) so an out-of-order input yields the same incidents; duplicates (by signal ID) are collapsed; each (asset, entity) session within Window becomes one incident; a storm beyond MaxPerIncident is suppressed to a single recorded note.
Output contract: the returned slice contains the events for MULTIPLE incidents, contiguous per incident and ordered by incident id. Each incident's sub-slice is internally time-ordered and Created-first, so a consumer (C7 persistence) appends the flat stream as-is, but a caller that PROJECTS must first partition by IncidentID and fold each incident separately — folding the whole flat slice through one incident.Project would fail on the second Created.
This is BATCH correlation over a complete signal set. Incident identity is seeded on a session's earliest signal, so callers requiring durable event-time cursors, watermarks, and bounded late-event handling use CorrelateIncrementalPage rather than re-running a backfilled batch window.
func SameSignal ¶
SameSignal compares staged signal values including TimelineRef contents rather than pointer identity.
Types ¶
type ActiveSession ¶
type ActiveSession struct {
AssetID shared.ID
EntityID shared.ID
IncidentID shared.ID
MinOccurredAt time.Time
MaxOccurredAt time.Time
ReflectedCount int
MaxSeverity shared.Severity
}
ActiveSession is the bounded, mutable summary needed to match future event-time signals.
func (ActiveSession) ValidateForStore ¶
func (s ActiveSession) ValidateForStore() error
ValidateForStore validates an active-session summary before persistence.
type Assignment ¶
type Assignment struct {
SignalID shared.ID
IncidentID shared.ID
AssetID shared.ID
EntityID shared.ID
OccurredAt time.Time
Severity shared.Severity
Outcome AssignmentOutcome
}
Assignment is the durable, immutable dedupe and graph edge from one signal to one incident.
func (Assignment) ValidateForStore ¶
func (a Assignment) ValidateForStore() error
ValidateForStore validates an immutable ledger row before persistence.
type AssignmentOutcome ¶
type AssignmentOutcome string
AssignmentOutcome records how a signal was represented by the correlation graph.
const ( AssignmentAttached AssignmentOutcome = "attached" AssignmentSuppressed AssignmentOutcome = "suppressed" AssignmentTooLate AssignmentOutcome = "too_late" )
type Checkpoint ¶
type Checkpoint struct {
Revision uint64
MaxObservedAt time.Time
Watermark time.Time
Phase Phase
Completed SourcePosition
Snapshot SourcePosition
RetentionAsOf time.Time
SourceCursor SourcePosition
StagedCursor SignalPosition
PolicyDigest string
}
Checkpoint is the durable two-phase position for one engagement's correlation stream.
type Config ¶
type Config struct {
// Window is the session gap: within one (asset, entity) key, a signal more than Window after the
// previous signal starts a NEW incident; otherwise it joins the current one. Must be > 0.
Window time.Duration
// MaxPerIncident caps how many signals an incident reflects individually (the Created signal plus
// attaches). Beyond it, further signals are suppressed as a storm and recorded as a single note
// (coverage-honest). Must be > 0.
MaxPerIncident int
// AllowedLateness is subtracted from the maximum observed event time to produce the
// monotonic watermark used by incremental correlation. It must not be negative.
AllowedLateness time.Duration
// Actor is the attribution for emitted events; defaults to "correlator".
Actor string
// EngagementID is the authoritative scope stamped on newly created incidents.
EngagementID shared.ID
// PageSize bounds each source-materialization and staged-consumption invocation.
// Zero uses the conservative default.
PageSize int
// MaxActiveSessions bounds the mutable state loaded into one correlation transaction.
// Zero uses the conservative default.
MaxActiveSessions int
// MaxTimelineRefsPerDetection and MaxTimelineRefsPerPage bound causal fanout during source materialization.
// Zero values use conservative defaults.
MaxTimelineRefsPerDetection int
MaxTimelineRefsPerPage int
}
Config tunes correlation. Zero fields take documented defaults.
type IncrementalPlan ¶
type IncrementalPlan struct {
Events []incident.IncidentEvent
Added []Assignment
ActiveSessions []ActiveSession
Next Checkpoint
TooLate int
Suppressed int
Evicted int
}
IncrementalPlan is one deterministic state transition. Events, Added, ActiveSessions, and Next form one durable commit.
func CorrelateIncremental ¶
func CorrelateIncremental(cfg Config, state State, signals []Signal, processedAt time.Time) (IncrementalPlan, error)
CorrelateIncremental assigns unseen signals using bounded session summaries and exact ledger lookups.
func CorrelateIncrementalPage ¶
func CorrelateIncrementalPage(cfg Config, state State, signals []Signal, processedAt time.Time, finalPage bool) (IncrementalPlan, error)
CorrelateIncrementalPage applies a globally event-ordered staged page. Only the final page advances the finalized watermark and prunes active sessions.
func (IncrementalPlan) Changed ¶
func (p IncrementalPlan) Changed(previous Checkpoint) bool
type Signal ¶
type Signal struct {
ID shared.ID
AssetID shared.ID
EntityID shared.ID
OccurredAt time.Time
Severity shared.Severity
RuleID string
Title string
// Timeline is set only for a causally referenced endpoint transition.
Timeline *incident.TimelineRef
}
Signal is one thing to correlate — a runtime detection with the identity, entity, and event time needed to group it. ID is the detection id and the dedupe key; EntityID is the process/network/file entity the detection concerns (zero for an asset-level detection).
type SignalPosition ¶
SignalPosition is a deterministic staged-signal event-time cursor.
type SourcePosition ¶
SourcePosition is the immutable recorded-order position of a source detection.
type State ¶
type State struct {
Checkpoint Checkpoint
KnownSignalIDs map[shared.ID]struct{}
ActiveSessions []ActiveSession
}
State contains a checkpoint, exact dedupe answers for the current input, and bounded active sessions.