correlation

package
v0.2.3 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 27, 2026 License: Apache-2.0 Imports: 10 Imported by: 0

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

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

func SameSignal(left, right Signal) bool

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.

func (Config) Normalize

func (c Config) Normalize() Config

Normalize applies conservative operational defaults shared by batch and incremental callers.

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 Phase

type Phase string
const (
	PhaseSource  Phase = "source"
	PhaseConsume Phase = "consume"
)

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).

func (Signal) Validate

func (s Signal) Validate() error

Validate enforces a well-formed signal.

type SignalPosition

type SignalPosition struct {
	OccurredAt time.Time
	ID         shared.ID
}

SignalPosition is a deterministic staged-signal event-time cursor.

type SourcePosition

type SourcePosition struct {
	RecordedAt time.Time
	ID         shared.ID
}

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.

func (State) Validate

func (s State) Validate() error

Validate rejects corrupt persisted state before it can influence incident assignment.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL