diff --git a/internal/anomaly/chance_test.go b/internal/anomaly/chance_test.go index 9ab1bd6ac..cf9b4cc79 100644 --- a/internal/anomaly/chance_test.go +++ b/internal/anomaly/chance_test.go @@ -243,8 +243,36 @@ func TestPeriodicEvidenceJSONIsCanonical(t *testing.T) { // unchanged; only the documented Chance values changed: rounded to // ChanceDigits, then multiplied by 97/49 (searchCells of the fixture rule) // when the union bound began to charge the refined candidates. +// +// Since then one documented field changed as well: the censored new-stream +// episode below keeps ConfidenceLow on its quiet end (it was route_observed, +// losing the confidence the episode opened with). The test maps that one +// field back before hashing, so the hash is unchanged and any other change +// to the golden still fails. const goldenWithoutChance = "2538014b83e1976a614da59c3db2eeb8d403ba76de63f3f58d55aad221eaf959" +const censoredQuietEndEpisode = "7771b1680b2f2d77043ed1601a6e2dbb" + +// unlowerCensoredQuietEnd checks that the one documented confidence change is +// in the golden and restores the value the hash was taken with. +func unlowerCensoredQuietEnd(t *testing.T, v interface{}) { + t.Helper() + found := 0 + for _, x := range v.(map[string]interface{})["Candidates"].([]interface{}) { + c := x.(map[string]interface{}) + if c["EventID"] == censoredQuietEndEpisode && c["Reason"] == "quiet" { + if c["Confidence"] != "low" { + t.Fatalf("episode %s ends with confidence %v, want low", censoredQuietEndEpisode, c["Confidence"]) + } + c["Confidence"] = "route_observed" + found++ + } + } + if found != 1 { + t.Fatalf("%d quiet ends of episode %s in the golden, want 1", found, censoredQuietEndEpisode) + } +} + func stripChance(v interface{}) interface{} { switch x := v.(type) { case map[string]interface{}: @@ -299,7 +327,9 @@ func TestGoldenChangedOnlyInTheDocumentedChanceValues(t *testing.T) { if len(chances) != 4 || n["9.09399e-8"] != 2 || n["0.000478992"] != 2 { t.Fatalf("golden Chance values %v, want 9.09399e-8 x2 and 0.000478992 x2", chances) } - b, err := json.Marshal(stripChance(decode())) + v := stripChance(decode()) + unlowerCensoredQuietEnd(t, v) + b, err := json.Marshal(v) if err != nil { t.Fatal(err) } diff --git a/internal/anomaly/confidence_followup_test.go b/internal/anomaly/confidence_followup_test.go new file mode 100644 index 000000000..9052231ca --- /dev/null +++ b/internal/anomaly/confidence_followup_test.go @@ -0,0 +1,138 @@ +package anomaly + +import ( + "testing" + "time" +) + +// A new-stream episode opened from left-censored evidence is reported with +// ConfidenceLow ("never as a certain new stream", newstream.go). Every later +// transition of that episode, including those without a trigger event (quiet +// end, eviction end) and a reopen by an uncensored signal within the +// cooldown, keeps that confidence. + +// episodes groups the candidates of rule by episode, in emission order. +func episodes(cs []Candidate, rule string) (ids []string, by map[string][]Candidate) { + by = map[string][]Candidate{} + for _, c := range cs { + if c.Rule != rule { + continue + } + if _, ok := by[c.EventID]; !ok { + ids = append(ids, c.EventID) + } + by[c.EventID] = append(by[c.EventID], c) + } + return ids, by +} + +// censoredEpisode returns the only episode of rule "ns", after checking that +// it was opened from censored evidence with ConfidenceLow and that it +// contains a transition to `to` for reason `reason`. +func censoredEpisode(t *testing.T, cs []Candidate, to State, reason Reason) []Candidate { + t.Helper() + ids, by := episodes(cs, "ns") + if len(ids) != 1 { + t.Fatalf("fixture: %d new-stream episodes, want 1:\n%s", len(ids), describe(cs)) + } + ep := by[ids[0]] + if open := ep[0]; open.NewStream == nil || !open.NewStream.Censored || open.Confidence != ConfidenceLow { + t.Fatalf("fixture: episode opened %s->%s with confidence %v, evidence %+v; want a censored, low-confidence open", + open.From, open.To, open.Confidence, open.NewStream) + } + for _, c := range ep { + if c.To == to && c.Reason == reason { + return ep + } + } + t.Fatalf("fixture: episode has no ->%v (%v) transition:\n%s", to, reason, describe(ep)) + return nil +} + +func assertAllLow(t *testing.T, ep []Candidate) { + t.Helper() + for _, c := range ep { + if c.Confidence != ConfidenceLow { + t.Errorf("%s->%s (%s) at %v carries confidence %v, want low (the episode opened from censored evidence)", + c.From, c.To, c.Reason, c.At.Sub(t0), c.Confidence) + } + } +} + +// censoredStream is a stream that starts one minute after the history start +// (inside CensorWindow) and reaches MinCount 30 after 14m30s. +func censoredStream(t *testing.T, prefix string, start time.Time) []Event { + return series(t, prefix, start, 30, 30*time.Second, nil, 0x11, 0x22) +} + +func TestCensoredNewStreamKeepsLowConfidenceOnQuietEnd(t *testing.T) { + d := newDet(t, Config{NewStream: []NewStreamRule{nsRule(30)}, HistoryStart: t0}) + cs := run(t, d, censoredStream(t, "n", t0.Add(time.Minute))) + cs = append(cs, d.Advance(t0.Add(3*time.Hour)).Candidates...) + assertAllLow(t, censoredEpisode(t, cs, StateEnded, ReasonQuiet)) +} + +func TestCensoredNewStreamKeepsLowConfidenceOnEvictionEnd(t *testing.T) { + cfg := Config{NewStream: []NewStreamRule{nsRule(30)}, HistoryStart: t0, Limits: experimentalLimits()} + cfg.Limits.MaxKeysPerScope = 1 + evs := append(censoredStream(t, "n", t0.Add(time.Minute)), relayEvent(t, "other", t0.Add(16*time.Minute), 0x33, 0x44)) + cs := run(t, newDet(t, cfg), evs) + assertAllLow(t, censoredEpisode(t, cs, StateEnded, ReasonEvicted)) +} + +// The stream goes quiet for more than QuietPeriod and comes back: the new +// stream state is reactivated (never censored) and fires again within the +// cooldown, which reopens the censored episode. +func TestCensoredNewStreamKeepsLowConfidenceOnReopen(t *testing.T) { + r := nsRule(30) + r.State.Cooldown = 24 * time.Hour + d := newDet(t, Config{NewStream: []NewStreamRule{r}, HistoryStart: t0}) + evs := append(censoredStream(t, "a", t0.Add(time.Minute)), censoredStream(t, "b", t0.Add(13*time.Hour))...) + cs := run(t, d, evs) + cs = append(cs, d.Advance(t0.Add(15*time.Hour)).Candidates...) + ep := censoredEpisode(t, cs, StateActive, ReasonReopened) + for _, c := range ep { + if c.Reason == ReasonReopened && (c.NewStream == nil || c.NewStream.Censored) { + t.Fatalf("fixture: the reopening signal must be uncensored, got %+v", c.NewStream) + } + } + assertAllLow(t, ep) +} + +// Controls: an uncensored stream is never lowered, and the low confidence +// does not outlive its episode (after the cooldown, a reactivated stream +// opens a new episode with the key's confidence). +func TestNewStreamConfidenceIsNotLoweredWithoutCensoring(t *testing.T) { + t.Run("uncensored stream", func(t *testing.T) { + d := newDet(t, historyThen(nsRule(30))) + cs := run(t, d, series(t, "n", t0, 30, 30*time.Second, nil, 0x11, 0x22)) + cs = append(cs, d.Advance(t0.Add(3*time.Hour)).Candidates...) + ids, by := episodes(cs, "ns") + if len(ids) != 1 || len(by[ids[0]]) < 3 { + t.Fatalf("fixture: want one episode with an end:\n%s", describe(cs)) + } + for _, c := range by[ids[0]] { + if c.Confidence != ConfidenceRouteObserved { + t.Errorf("%s->%s (%s) carries %v, want route_observed", c.From, c.To, c.Reason, c.Confidence) + } + } + }) + t.Run("new episode after the cooldown", func(t *testing.T) { + d := newDet(t, Config{NewStream: []NewStreamRule{nsRule(30)}, HistoryStart: t0}) + evs := append(censoredStream(t, "a", t0.Add(time.Minute)), censoredStream(t, "b", t0.Add(13*time.Hour))...) + cs := run(t, d, evs) + cs = append(cs, d.Advance(t0.Add(15*time.Hour)).Candidates...) + ids, by := episodes(cs, "ns") + if len(ids) != 2 { + t.Fatalf("fixture: want two episodes:\n%s", describe(cs)) + } + if first := by[ids[0]][0]; first.Confidence != ConfidenceLow || first.NewStream == nil || !first.NewStream.Censored { + t.Fatalf("fixture: first episode opened with %v", first.Confidence) + } + for _, c := range by[ids[1]] { + if c.Confidence != ConfidenceRouteObserved { + t.Errorf("second episode %s->%s (%s) carries %v, want route_observed", c.From, c.To, c.Reason, c.Confidence) + } + } + }) +} diff --git a/internal/anomaly/detector.go b/internal/anomaly/detector.go index a451eb083..9261a2acd 100644 --- a/internal/anomaly/detector.go +++ b/internal/anomaly/detector.go @@ -703,7 +703,7 @@ func (d *Detector) params(tb *keyTable, slot int) (*StateParams, string, Kind) { func (d *Detector) hot(tb *keyTable, k *keyState, slot int, p *StateParams, ctx emitCtx) { ep := &k.eps[slot] d.trBuf = ep.onHot(ctx.at, p, d.trBuf[:0]) - ep.note(ctx.label, ctx.coverage) + ep.note(ctx.label, ctx.coverage, ctx.ns != nil && ctx.ns.Censored) d.emit(tb, k, slot, d.trBuf, ctx) d.schedule(tb, k, slot, p) } @@ -766,9 +766,6 @@ func (d *Detector) emit(tb *keyTable, k *keyState, slot int, trs []transition, c if ctx.ns != nil { ev := *ctx.ns c.NewStream = &ev - if ev.Censored { - c.Confidence = ConfidenceLow - } } if ctx.per != nil { ev := *ctx.per @@ -776,6 +773,11 @@ func (d *Detector) emit(tb *keyTable, k *keyState, slot int, trs []transition, c c.Periodic = &ev } } + // censored evidence is never a certain new stream: its signal, and + // every other transition of its episode, are low confidence + if tr.low || ctx.ns != nil && ctx.ns.Censored { + c.Confidence = ConfidenceLow + } c.Suppressed = d.cfg.Expected.Mode == ExpectedSuppress && c.Traffic == LabelExpected d.push(c) } diff --git a/internal/anomaly/periodic.go b/internal/anomaly/periodic.go index 77db75578..0f1a03a4e 100644 --- a/internal/anomaly/periodic.go +++ b/internal/anomaly/periodic.go @@ -99,6 +99,7 @@ func classFlags(c TrafficClass) pulseFlags { // buffers per key. type periodicScratch struct { resid, med, seen, seenRef []int64 + ranges []rangeMemo // see memoRange // work counters (chains measured, gaps visited), for bound tests and // benchmarks chains, steps uint64 @@ -321,7 +322,7 @@ func (p *periodicState) evaluate(now int64, r *PeriodicRule, sc *periodicScratch p.alts = p.alternatives(best, r, sc) } // traffic label and burst sizes over the chain's pulses - var total, exp, mon uint64 + var total, exp, anyExp, mon uint64 events := 0 for i := p.n - 1 - best.gaps; i < p.n; i++ { fl := p.flagAt(i) @@ -329,11 +330,19 @@ func (p *periodicState) evaluate(now int64, r *PeriodicRule, sc *periodicScratch if fl&pulseAllExpected != 0 { exp++ } + if fl&pulseAnyExpected != 0 { + anyExp++ + } if fl&pulseAnyMonitored != 0 { mon++ } events += int(p.countAt(i)) } + label := labelFromCounts(total, exp, mon) + if label == LabelUnclassified && anyExp > 0 { + // a burst merged expected and unclassified events + label = LabelMixed + } ev := PeriodicEvidence{ Period: time.Duration(best.period), Tolerance: time.Duration(best.tol), @@ -345,7 +354,7 @@ func (p *periodicState) evaluate(now int64, r *PeriodicRule, sc *periodicScratch SufficientAt: time.Unix(0, p.sufficientAt).UTC(), Alternatives: append([]PeriodAlternative(nil), p.alts...), } - return periodicSignal{ev: ev, label: labelFromCounts(total, exp, mon)}, true + return periodicSignal{ev: ev, label: label}, true } // estimate searches candidate periods derived from recent gaps. @@ -364,13 +373,18 @@ func (p *periodicState) estimate(r *PeriodicRule, sc *periodicScratch) chainInfo // strictly outranks it. So the search never returns a chain the seeds // alone would rank higher, and until a refined period would signal it // changes nothing. + // + // Every seed in [MinPeriod, MaxPeriod] is measured. A seed outside it is + // measured only for its refinement (see seedInReach) and is never the + // estimate itself. var bestRef chainInfo seen, seenRef := sc.seen[:0], sc.seenRef[:0] for i := p.n - 1; i >= 1 && i >= p.n-recentGaps; i-- { g := p.at(i) - p.at(i-1) for k := 1; k <= maxK; k++ { cand := g / int64(k) - if cand < minP || cand > maxP { + inRange := cand >= minP && cand <= maxP + if !inRange && !seedInReach(g, k, r) { continue } if slices.Contains(seen, cand) { @@ -378,7 +392,7 @@ func (p *periodicState) estimate(r *PeriodicRule, sc *periodicScratch) chainInfo } seen = append(seen, cand) c, ref := p.chainRefined(cand, r, sc) - if c.gaps > 0 && (best.gaps == 0 || c.preferred(best, r)) { + if inRange && c.gaps > 0 && (best.gaps == 0 || c.preferred(best, r)) { best = c } if ref != cand && !slices.Contains(seenRef, ref) && !slices.Contains(seen, ref) { @@ -416,14 +430,34 @@ func (p *periodicState) estimate(r *PeriodicRule, sc *periodicScratch) chainInfo return best } +// seedInReach reports whether some period P in [MinPeriod, MaxPeriod] may +// explain gap g as k periods (|g - k*P| <= tol(P)), so that the seed g/k is +// worth measuring even when it lies outside the range: chainRefined can +// refine it into the range (gaps of 294s, 306s... with tol 6s and range +// [299s, 301s] give only out-of-range seeds, which refine to 300s). +// +// The test is necessary, not sufficient. k*P + tol(P) increases with P, and +// so does k*P - tol(P): from P to P+1, tolNanos grows by at most 1ns (its +// JitterRel*P by at most 0.25 before truncation) while k*P grows by k >= 1. +// So an in-range P can explain g only if +// k*MinPeriod - tol(MinPeriod) <= g <= k*MaxPeriod + tol(MaxPeriod). +// estimate() measures at most one seed per gap and k either way, so there +// are still at most recentGaps*(MaxMissing+1) seeds, as searchCells charges +// and periodicWorkBound counts. +func seedInReach(g int64, k int, r *PeriodicRule) bool { + minP, maxP := int64(r.MinPeriod), int64(r.MaxPeriod) + return g >= int64(k)*minP-tolNanos(r, minP) && g <= int64(k)*maxP+tolNanos(r, maxP) +} + // chainRefined measures seed's chain exactly as chainFor does and returns it // with the period this search cell tests: seed, or a refinement of it. // // A single gap fixes P only to within tol/k. When the jitter of neighbouring // gaps cancels (gaps of 294s, 306s, 294s... around 300s with tol 6s) no seed // taken from one gap explains the next one, so no chain forms. Every period -// in the intersection of the ranges [(g-tol)/k, (g+tol)/k] of the newest gaps -// explains all of them. In the same pass as the chain, over the recentGaps +// in the intersection of the ranges of the newest gaps (periodRange: the P +// with |g - k*P| <= tol(P), each with its own tolerance) explains all of +// them. In the same pass as the chain, over the recentGaps // newest gaps (the window the seeds come from) and with the multiples k the // seed assigns them, chainRefined intersects those ranges. If the common // range covers at least two gaps and more than the seed's chain, it also @@ -464,14 +498,12 @@ func (p *periodicState) chainRefined(seed int64, r *PeriodicRule, sc *periodicSc } if inRange { k := max((g+seed/2)/seed, 1) - l, h := lo, min(hi, (g+c.tol)/k) - if g-c.tol > 0 { - l = max(l, (g-c.tol+k-1)/k) - } - if k > int64(maxK) || l > h { + if k > int64(maxK) { + inRange = false + } else if gl, gh := sc.memoRange(p.evaluated-uint64(p.n-1-i), g, k, r); max(lo, gl) > min(hi, gh) { inRange = false } else { - lo, hi = l, h + lo, hi = max(lo, gl), min(hi, gh) n++ span += g slots += k @@ -488,16 +520,85 @@ func (p *periodicState) chainRefined(seed int64, r *PeriodicRule, sc *periodicSc if n < 2 || n <= c.gaps { return c, seed } + // Every period in the common range explains its n gaps with its own + // tolerance. ref is the phase estimate clamped into the part of that + // range inside [MinPeriod, MaxPeriod] (a train whose estimate lies just + // beyond MaxPeriod may still have periods in range that explain it). + // The recount keeps ref only when it explains more of the newest gaps + // than the seed's chain. + lo, hi = max(lo, int64(r.MinPeriod)), min(hi, int64(r.MaxPeriod)) + if lo > hi { + return c, seed + } ref := min(max(span/slots, lo), hi) - // The range uses the seed's tolerance, but the tolerance grows with the - // period (JitterRel): keep ref only if, with its own tolerance, it - // explains more of the newest gaps than the seed does. - if ref < int64(r.MinPeriod) || ref > int64(r.MaxPeriod) || p.recentExplained(ref, r, sc) <= c.gaps { + if p.recentExplained(ref, r, sc) <= c.gaps { return c, seed } return c, ref } +// periodRange returns the periods P that explain gap g as k periods with +// their own tolerance, |g - k*P| <= tol(P), as [lo, hi] (empty if lo > hi). +// Both k*P - tol(P) and k*P + tol(P) never decrease in P (see seedInReach), +// so these P form one range: hi is the last P with k*P - tol(P) <= g, lo the +// first with k*P + tol(P) >= g. +// +// Without JitterRel, tol is JitterAbs and the ends are exact integer +// divisions. With JitterRel, the ends of the relative part, g/(k+JitterRel) +// and g/(k-JitterRel), are estimated in floating point and then stepped to +// the exact ends under tolNanos's truncation (a few steps at most). +func periodRange(g, k int64, r *PeriodicRule) (lo, hi int64) { + abs := int64(r.JitterAbs) + lo, hi = max((g-abs+k-1)/k, 1), (g+abs)/k + if r.JitterRel == 0 { + return lo, hi + } + lo = max(min(lo, int64(float64(g)/(float64(k)+r.JitterRel))), 1) + hi = max(hi, int64(float64(g)/(float64(k)-r.JitterRel))) + for (hi+1)*k-tolNanos(r, hi+1) <= g { + hi++ + } + for hi >= 1 && hi*k-tolNanos(r, hi) > g { + hi-- + } + for lo > 1 && (lo-1)*k+tolNanos(r, lo-1) >= g { + lo-- + } + for lo <= hi && lo*k+tolNanos(r, lo) < g { + lo++ + } + return lo, hi +} + +// rangeMemo is one cached periodRange result, with every input it depends on. +type rangeMemo struct { + g, abs, lo, hi int64 + rel float64 +} + +// memoRange is periodRange for the gap ending at the key's pulse number +// pulse. The seeds of one search mostly assign the same multiple k to the +// same gap, and the next pulses' searches see that gap again, so with +// JitterRel the result is cached in one slot per (pulse mod recentGaps, k): +// at most recentGaps*(maxMissingLimit+1) entries, allocated once per +// Detector. The slot fixes k, and an entry is used only when g, JitterAbs +// and JitterRel match too, so it is never stale, whichever rule, key or ring +// asks. +func (sc *periodicScratch) memoRange(pulse uint64, g, k int64, r *PeriodicRule) (lo, hi int64) { + if r.JitterRel == 0 { + return periodRange(g, k, r) // two integer divisions: not worth a lookup + } + if sc.ranges == nil { + sc.ranges = make([]rangeMemo, recentGaps*(maxMissingLimit+1)) + } + m := &sc.ranges[int(pulse%recentGaps)*(maxMissingLimit+1)+int(k-1)] + if m.g != g || m.abs != int64(r.JitterAbs) || m.rel != r.JitterRel { + m.lo, m.hi = periodRange(g, k, r) + m.g, m.abs, m.rel = g, int64(r.JitterAbs), r.JitterRel + } + return m.lo, m.hi +} + // recentExplained counts the consecutive newest gaps, at most recentGaps, // that period explains with its own tolerance. func (p *periodicState) recentExplained(period int64, r *PeriodicRule, sc *periodicScratch) int { diff --git a/internal/anomaly/periodic_bench_test.go b/internal/anomaly/periodic_bench_test.go index ebeff832e..0547d8bd0 100644 --- a/internal/anomaly/periodic_bench_test.go +++ b/internal/anomaly/periodic_bench_test.go @@ -124,10 +124,10 @@ func periodicWorkBound(r *PeriodicRule) (chains, steps uint64) { // not count (more CPU per visited gap, say) is not caught; only // BenchmarkPeriodicSearch measures it. var measuredSteps = map[string]uint64{ - "periodic/0.6": 236631, "alternating/0.6": 1434346, "jitter/0.6": 246825, - "missing/0.6": 243414, "bursty/0.6": 196561, "noise/0.6": 144778, - "periodic/1": 236631, "alternating/1": 1434346, "jitter/1": 5846059, - "missing/1": 1174373, "bursty/1": 196561, "noise/1": 144778, + "periodic/0.6": 236631, "alternating/0.6": 1466800, "jitter/0.6": 247574, + "missing/0.6": 243414, "bursty/0.6": 199730, "noise/0.6": 148177, + "periodic/1": 236631, "alternating/1": 1466800, "jitter/1": 6242089, + "missing/1": 1174373, "bursty/1": 199730, "noise/1": 148177, } // Every single pulse stays within periodicWorkBound, at the largest allowed diff --git a/internal/anomaly/periodic_followup_test.go b/internal/anomaly/periodic_followup_test.go new file mode 100644 index 000000000..cb35931d6 --- /dev/null +++ b/internal/anomaly/periodic_followup_test.go @@ -0,0 +1,498 @@ +package anomaly + +import ( + "fmt" + "math/rand" + "slices" + "testing" + "time" +) + +// Follow-up to the periodic detector: a period search range narrower than +// the jitter, and the traffic label of bursts that mix expected and +// unclassified events. + +// A period range narrower than the tolerance: gaps alternate 294s and 306s, +// so every seed g/k lies outside [299s, 301s], yet 300s explains every gap +// within 6s and meets every criterion. The seeds must still be refined into +// the range instead of being dropped before refinement, with an absolute and +// with a relative tolerance. +func TestPeriodicNarrowRangeFindsPeriodBehindOutOfRangeSeeds(t *testing.T) { + for _, j := range []struct { + name string + abs time.Duration + rel float64 + }{{"absolute 6s", 6 * time.Second, 0}, {"relative 2%", 0, 0.02}} { + t.Run(j.name, func(t *testing.T) { + rule := experimentalPeriodic() + rule.MinPeriod, rule.MaxPeriod = 299*time.Second, 301*time.Second + rule.JitterAbs, rule.JitterRel, rule.MaxJitterFraction = j.abs, j.rel, 1 + cfg := Config{Limits: experimentalLimits(), Expected: ExpectedPolicy{Mode: ExpectedInclude}, Periodic: []PeriodicRule{rule}} + if err := cfg.Validate(); err != nil { + t.Fatal(err) + } + const pulses = 32 + offs := make([]time.Duration, pulses) + var events []Event + for i := range offs { + offs[i] = time.Duration(i) * 5 * time.Minute + if i%2 == 1 { + offs[i] -= 6 * time.Second + } + events = append(events, relayEvent(t, fmt.Sprintf("narrow-%02d", i), at(offs[i]), 1, 2)) + } + + // The fixture is valid: 300s explains the whole chain, meets the + // criteria and passes MaxChance, and no seed g/k is in range. + p := ringOf(&rule, offs...) + p.evaluated = pulses + var sc periodicScratch + ideal := p.chainFor(int64(300*time.Second), &rule, &sc) + if ideal.gaps != pulses-1 || !ideal.meets(&rule) || p.chance(ideal, &rule) > rule.MaxChance { + t.Fatalf("invalid repro: ideal=%+v chance=%g", ideal, p.chance(ideal, &rule)) + } + for _, g := range []time.Duration{294 * time.Second, 306 * time.Second} { + for k := 1; k <= rule.MaxMissing+1; k++ { + if s := g / time.Duration(k); s >= rule.MinPeriod && s <= rule.MaxPeriod { + t.Fatalf("invalid repro: seed %v/%d = %v is in range", g, k, s) + } + } + } + + a := filter(run(t, newDet(t, cfg), events), rule.Name, StateActive) + if len(a) == 0 { + t.Fatal("5m train with gaps 294s/306s was not detected with MinPeriod=299s, MaxPeriod=301s") + } + ev := a[0].Periodic + if ev == nil || ev.Period != 5*time.Minute { + t.Fatalf("detected period %+v, want 5m0s", ev) + } + }) + } +} + +// The traffic label of a periodic chain follows the events of its pulses: +// a burst that merges an expected and an unclassified event is mixed, in +// either order, exactly as two single-event pulses with those classes are. +// +// A chain is evaluated when its newest pulse starts, so that pulse holds only +// its first event then; the older pulses are complete. +func TestPeriodicBurstTrafficLabel(t *testing.T) { + x, u, m := TrafficExpected, TrafficUnclassified, TrafficMonitored + every := func(cls ...TrafficClass) func(int) []TrafficClass { + return func(int) []TrafficClass { return cls } + } + oneAmong := func(burst []TrafficClass, other TrafficClass) func(int) []TrafficClass { + return func(i int) []TrafficClass { + if i == 3 { + return burst + } + return []TrafficClass{other} + } + } + for _, c := range []struct { + name string + pulse func(i int) []TrafficClass // classes of the events in pulse i + want TrafficLabel + }{ + {"bursts unclassified then expected", every(u, x), LabelMixed}, + {"one expected-then-unclassified burst among unclassified pulses", oneAmong([]TrafficClass{x, u}, u), LabelMixed}, + {"one unclassified-then-expected burst among unclassified pulses", oneAmong([]TrafficClass{u, x}, u), LabelMixed}, + // controls + // (mixed already before the fix, only because the newest pulse + // holds just its expected first event when the chain is evaluated) + {"bursts expected then unclassified", every(x, u), LabelMixed}, + {"one mixed burst among expected pulses", oneAmong([]TrafficClass{x, u}, x), LabelMixed}, + {"single events alternating", func(i int) []TrafficClass { return []TrafficClass{[]TrafficClass{x, u}[i%2]} }, LabelMixed}, + {"expected bursts", every(x, x), LabelExpected}, + {"unclassified bursts", every(u, u), LabelUnclassified}, + {"mixed bursts with a monitored event", every(x, u, m), LabelMonitored}, + } { + t.Run(c.name, func(t *testing.T) { + var evs []Event + for i := 0; i < 12; i++ { + for j, cls := range c.pulse(i) { + e := relayEvent(t, fmt.Sprintf("p%02d-%d", i, j), at(time.Duration(i)*5*time.Minute+time.Duration(j)*time.Second), 1, 2) + e.Traffic = cls + evs = append(evs, e) + } + } + a := firstActive(t, run(t, perDet(t, ExpectedInclude), evs)) + if a.Traffic != c.want { + t.Fatalf("traffic label %v, want %v", a.Traffic, c.want) + } + }) + } +} + +// seedInReach keeps every seed g/k for which some period in range explains g +// as k periods with its own tolerance, under absolute, relative and combined +// jitter, at realistic magnitudes (sampled). +func TestSeedInReachKeepsEverySeedAPeriodInRangeExplains(t *testing.T) { + rng := rand.New(rand.NewSource(94)) + checked := 0 + for trial := 0; trial < 20000; trial++ { + r := experimentalPeriodic() + r.MinPeriod = time.Duration(30+rng.Intn(7200)) * time.Second + r.MaxPeriod = r.MinPeriod + time.Duration(1+rng.Int63n(int64(2*r.MinPeriod))) + r.JitterAbs, r.JitterRel = 0, 0 + switch trial % 3 { + case 0: + r.JitterAbs = time.Duration(1 + rng.Int63n(int64(r.MinPeriod/4))) + case 1: + r.JitterRel = 0.25 * (1 - rng.Float64()) + default: + r.JitterAbs, r.JitterRel = time.Duration(1+rng.Int63n(int64(r.MinPeriod/8))), 0.25*rng.Float64() + } + minP, maxP := int64(r.MinPeriod), int64(r.MaxPeriod) + k := 1 + rng.Intn(maxMissingLimit+1) + // periods and residuals at the edges are where a bound that is off + // (or uses the wrong period's tolerance) drops a seed + per := []int64{minP, maxP, minP + rng.Int63n(maxP-minP+1)}[rng.Intn(3)] + tol := tolNanos(&r, per) + g := []int64{int64(k)*per - tol, int64(k)*per + tol, int64(k)*per - tol + rng.Int63n(2*tol+1)}[rng.Intn(3)] + if kk, _ := explain(g, per, tol, k); kk != k { + continue // tol >= per/2: another k explains g + } + checked++ + if !seedInReach(g, k, &r) { + t.Fatalf("trial %d: %v explains gap %v as %d periods (tol %v), but the seed is dropped; rule %v..%v abs %v rel %g", + trial, time.Duration(per), time.Duration(g), k, time.Duration(tol), r.MinPeriod, r.MaxPeriod, r.JitterAbs, r.JitterRel) + } + } + if checked < 19000 { + t.Fatalf("only %d of 20000 trials checked", checked) + } +} + +// A seed outside the range is measured for its refinement but never +// reported: a steady 294s or 306s train is measured (its seed is in the +// seed pool) and never yields a period outside [299s, 301s]. +func TestPeriodicOutOfRangeSeedIsNeverReported(t *testing.T) { + rule := experimentalPeriodic() + rule.MinPeriod, rule.MaxPeriod = 299*time.Second, 301*time.Second + rule.JitterAbs, rule.JitterRel, rule.MaxJitterFraction = 6*time.Second, 0, 1 + cfg := Config{Limits: experimentalLimits(), Expected: ExpectedPolicy{Mode: ExpectedInclude}, Periodic: []PeriodicRule{rule}} + for _, step := range []time.Duration{294 * time.Second, 306 * time.Second} { + offs := make([]time.Duration, 32) + for i := range offs { + offs[i] = time.Duration(i) * step + } + p := ringOf(&rule, offs...) + p.evaluated = uint64(len(offs)) + var sc periodicScratch + est := p.estimate(&rule, &sc) + if !slices.Contains(sc.seen, int64(step)) { + t.Fatalf("%v: fixture: the out-of-range seed was not measured (seeds %v)", step, sc.seen) + } + if est.gaps > 0 && (est.period < int64(rule.MinPeriod) || est.period > int64(rule.MaxPeriod)) { + t.Fatalf("%v: estimate %v is outside the range", step, time.Duration(est.period)) + } + evs := series(t, "steady", t0, len(offs), step, nil, 1, 2) + for _, c := range run(t, newDet(t, cfg), evs) { + if c.Periodic != nil && (c.Periodic.Period < rule.MinPeriod || c.Periodic.Period > rule.MaxPeriod) { + t.Fatalf("%v: reported period %v outside [%v, %v]", step, c.Periodic.Period, rule.MinPeriod, rule.MaxPeriod) + } + } + } +} + +// The same, exhaustively over small integer periods, where the truncation of +// tolNanos and the rounding of explain matter most: every gap that some +// period in range explains as k periods is in reach. +func TestSeedInReachIsSoundExhaustively(t *testing.T) { + checked := 0 + for minP := int64(2); minP <= 40; minP++ { + for width := int64(1); width <= 12; width++ { + for _, abs := range []int64{0, 1, 2, 3} { + for _, rel := range []float64{0, 0.1, 0.13, 0.25} { + r := PeriodicRule{MinPeriod: time.Duration(minP), MaxPeriod: time.Duration(minP + width), + JitterAbs: time.Duration(abs), JitterRel: rel} + for k := 1; k <= maxMissingLimit+1; k++ { + for per := minP; per <= minP+width; per++ { + tol := tolNanos(&r, per) + for g := int64(k)*per - tol; g <= int64(k)*per+tol; g++ { + if kk, _ := explain(g, per, tol, k); kk != k { + continue + } + checked++ + if !seedInReach(g, k, &r) { + t.Fatalf("period %d explains gap %d as %d periods (tol %d) but the seed is dropped; range %d..%d abs %d rel %g", + per, g, k, tol, minP, minP+width, abs, rel) + } + } + } + } + } + } + } + } + if checked < 100000 { + t.Fatalf("only %d cases checked", checked) + } +} + +// Every seed in range is measured, also where seedInReach alone would drop +// it: with a 1ns tolerance the gap 3*60s+2ns gives the in-range seed 60s +// (floor of g/3), beyond k*MaxPeriod+tol. That seed explains the three +// newest gaps with full coverage and must stay the estimate. +func TestPeriodicInRangeSeedIsAlwaysMeasured(t *testing.T) { + r := experimentalPeriodic() + r.MinPeriod, r.MaxPeriod = 30*time.Second, time.Minute + r.JitterAbs, r.JitterRel = 1, 0 + cfg := Config{Limits: experimentalLimits(), Expected: ExpectedPolicy{Mode: ExpectedInclude}, Periodic: []PeriodicRule{r}} + if err := cfg.Validate(); err != nil { + t.Fatal(err) + } + const m = time.Minute + g0 := 3*m + 2 + if seedInReach(int64(g0), 3, &r) { + t.Fatal("fixture: seedInReach alone keeps the seed") + } + p := ringOf(&r, 0, g0, g0+m-1, g0+2*m, g0+3*m-1) + p.evaluated = uint64(p.n) + var sc periodicScratch + est := p.estimate(&r, &sc) + if est.period != int64(m) || est.gaps != 3 || est.coverage != 1 { + t.Fatalf("estimate %v over %d gaps (coverage %g), want 1m0s over 3 (coverage 1)", time.Duration(est.period), est.gaps, est.coverage) + } + for i := p.n - 1; i >= 1; i-- { + g := p.at(i) - p.at(i-1) + for k := 1; k <= r.MaxMissing+1; k++ { + if s := g / int64(k); s >= int64(r.MinPeriod) && s <= int64(r.MaxPeriod) && !slices.Contains(sc.seen, s) { + t.Fatalf("in-range seed %v (gap %v / %d) was not measured", time.Duration(s), time.Duration(g), k) + } + } + } +} + +// Relative jitter: the refinement must use the tolerance of the period it +// refines to, not the seed's. JitterRel 2% gives tol(300s) = 6s; gaps repeat +// 305.95s, 305.95s, 294.05s, which 300s explains within 5.95s. With the +// seed's tolerance, the low seed's range (tol 5.881s) is empty, and the high +// seed's range (tol 6.119s) reaches 300.169s, which misses the 294.05s gaps +// with its own tolerance, so no seed refined to a period that explains the +// train and it was never detected. +func TestPeriodicRelativeJitterRefinesWithTheCandidateTolerance(t *testing.T) { + rule := experimentalPeriodic() + rule.JitterAbs, rule.JitterRel, rule.MaxJitterFraction = 0, 0.02, 1 + cfg := Config{Limits: experimentalLimits(), Expected: ExpectedPolicy{Mode: ExpectedInclude}, Periodic: []PeriodicRule{rule}} + if err := cfg.Validate(); err != nil { + t.Fatal(err) + } + hi, lo := 305950*time.Millisecond, 294050*time.Millisecond + offs := []time.Duration{0} + for i := 1; i < 32; i++ { + g := hi + if i%3 == 0 { + g = lo + } + offs = append(offs, offs[i-1]+g) + } + p := ringOf(&rule, offs...) + p.evaluated = uint64(len(offs)) + var sc periodicScratch + ideal := p.chainFor(int64(300*time.Second), &rule, &sc) + if ideal.gaps != len(offs)-1 || !ideal.meets(&rule) || p.chance(ideal, &rule) > rule.MaxChance { + t.Fatalf("invalid repro: ideal=%+v chance=%g", ideal, p.chance(ideal, &rule)) + } + + var events []Event + for i, o := range offs { + events = append(events, relayEvent(t, fmt.Sprintf("rel-%02d", i), at(o), 1, 2)) + } + a := filter(run(t, newDet(t, cfg), events), rule.Name, StateActive) + if len(a) == 0 { + t.Fatal("train with gaps 305.95s, 305.95s, 294.05s was not detected; 300s explains every gap with JitterRel 2%") + } + ev := a[0].Periodic + if ev == nil { + t.Fatal("activation without periodic evidence") + } + // the detected period explains every gap of the chain with its own + // tolerance, and is within that tolerance of 300s + c := p.chainFor(int64(ev.Period), &rule, &sc) + if c.gaps < int(ev.Support)-1 || ev.Tolerance != time.Duration(tolNanos(&rule, int64(ev.Period))) || + (ev.Period-300*time.Second).Abs() > ev.Tolerance { + t.Fatalf("detected %v (tolerance %v, support %d), whose chain explains %d gaps", ev.Period, ev.Tolerance, ev.Support, c.gaps) + } +} + +// oldAbsRange is the refinement range chainRefined used before periodRange, +// for one gap g as k periods with tolerance tol (then always the seed's): +// [ceil((g-tol)/k), floor((g+tol)/k)], with no lower end while g-tol <= 0. +func oldAbsRange(g, k, tol int64) (lo, hi int64) { + lo, hi = 1, (g+tol)/k + if g-tol > 0 { + lo = max(lo, (g-tol+k-1)/k) + } + return lo, hi +} + +// With absolute jitter only, periodRange gives exactly the integers of the +// old computation: exhaustively over small gaps and tolerances, and sampled +// at nanosecond magnitudes up to the largest gaps a rule can bridge. +func TestPeriodRangeWithAbsoluteJitterIsTheOldRange(t *testing.T) { + check := func(g, k, abs int64) { + r := PeriodicRule{JitterAbs: time.Duration(abs)} + lo, hi := periodRange(g, k, &r) + if olo, ohi := oldAbsRange(g, k, abs); lo != olo || hi != ohi { + t.Fatalf("g=%d k=%d JitterAbs=%d: periodRange [%d, %d], old range [%d, %d]", g, k, abs, lo, hi, olo, ohi) + } + } + for g := int64(1); g <= 2000; g++ { + for k := int64(1); k <= maxMissingLimit+1; k++ { + for abs := int64(0); abs <= 40; abs++ { + check(g, k, abs) + } + } + } + rng := rand.New(rand.NewSource(108)) + for i := 0; i < 200000; i++ { + k := 1 + rng.Int63n(maxMissingLimit+1) + g := 1 + rng.Int63n(k*int64(maxWindow)) + check(g, k, rng.Int63n(int64(maxWindow/2))) + } +} + +// With relative (and combined) jitter, periodRange is exactly the set of P +// with |g - k*P| <= tol(P): by brute force over small integers, and by its +// defining ends at nanosecond magnitudes. +func TestPeriodRangeIsExactlyThePeriodsThatExplainTheGap(t *testing.T) { + for _, abs := range []int64{0, 1, 3, 7} { + for _, rel := range []float64{0.01, 0.1, 0.13, 0.25} { + r := PeriodicRule{JitterAbs: time.Duration(abs), JitterRel: rel} + for k := int64(1); k <= maxMissingLimit+1; k++ { + for g := int64(1); g <= 600; g++ { + lo, hi := periodRange(g, k, &r) + for per := int64(1); per <= g+abs+1; per++ { + d := g - k*per + in := max(d, -d) <= tolNanos(&r, per) + if in != (per >= lo && per <= hi) { + t.Fatalf("abs=%d rel=%g k=%d g=%d: P=%d explains=%v, but periodRange is [%d, %d]", abs, rel, k, g, per, in, lo, hi) + } + } + } + } + } + } + rng := rand.New(rand.NewSource(1081)) + for i := 0; i < 200000; i++ { + r := PeriodicRule{JitterAbs: time.Duration(rng.Int63n(int64(10 * time.Second))), JitterRel: 0.25 * rng.Float64()} + k := 1 + rng.Int63n(maxMissingLimit+1) + g := int64(time.Second) + rng.Int63n(k*int64(maxWindow)) + lo, hi := periodRange(g, k, &r) + f := func(per int64) int64 { return k*per - tolNanos(&r, per) } + h := func(per int64) int64 { return k*per + tolNanos(&r, per) } + if lo > hi || f(hi) > g || f(hi+1) <= g || h(lo) < g || (lo > 1 && h(lo-1) >= g) { + t.Fatalf("abs=%v rel=%g k=%d g=%d: [%d, %d] are not the exact ends", r.JitterAbs, r.JitterRel, k, g, lo, hi) + } + } +} + +// The refined period is clamped into the part of the common range that lies +// in [MinPeriod, MaxPeriod]. Both trains have a true period of 300.98s in +// the range [299s, 301s], residuals of 0.90..0.99 tol and two of three gaps +// long, so their phase estimate (mean gap) lies beyond MaxPeriod. The common +// range straddles MaxPeriod: clamped only into the common range, the refined +// period would be above 301s and rejected, although the periods from its +// lower end up to 301s explain every gap. +func TestPeriodicRefinementIsClampedIntoTheSearchRange(t *testing.T) { + for _, c := range []struct { + name string + abs time.Duration + rel float64 + period int64 + gaps []int64 + }{ + {"relative 2%", 0, 0.02, 300983680322, []int64{ + 295433727560, 295081034588, 306545473781, 306861519212, 306631836381, 306839668100, + 295481269048, 306684537919, 306797954449, 306605636817, 306836785762, 306717842243, + 295064704066, 306841202094, 306918734081, 306475916213, 306738216837, 306599618627, + 295439541190, 306701648665, 306655004215, 295172194564, 295422882548, 306816575831, + }}, + {"absolute 6s", 6 * time.Second, 0, 300983815273, []int64{ + 306831828975, 306822681220, 295457880758, 306702249706, 306592244168, 306845246623, + 295130477268, 306753857882, 306756812033, 306408553938, 295364746775, 306677382170, + 306511141502, 306590311144, 295567381355, 306481296874, 306868730214, 306614050442, + 295348035963, 306797789043, 306408756668, 295272100140, 295171304764, 306667020981, + }}, + } { + t.Run(c.name, func(t *testing.T) { + r := experimentalPeriodic() + r.MinPeriod, r.MaxPeriod = 299*time.Second, 301*time.Second + r.JitterAbs, r.JitterRel, r.MaxJitterFraction = c.abs, c.rel, 1 + if err := (&Config{Limits: experimentalLimits(), Expected: ExpectedPolicy{Mode: ExpectedInclude}, Periodic: []PeriodicRule{r}}).Validate(); err != nil { + t.Fatal(err) + } + // fixture: the true period explains the train and would signal, + // while the mean gap is above MaxPeriod + offs := []time.Duration{0} + var sum int64 + for _, g := range c.gaps { + offs = append(offs, offs[len(offs)-1]+time.Duration(g)) + sum += g + } + ring := ringOf(&r, offs...) + ring.evaluated = uint64(len(offs)) + var sc periodicScratch + ideal := ring.chainFor(c.period, &r, &sc) + if ideal.gaps != len(c.gaps) || !ideal.meets(&r) || ring.chance(ideal, &r) > r.MaxChance || sum/int64(len(c.gaps)) <= int64(r.MaxPeriod) { + t.Fatalf("invalid fixture: ideal=%+v chance=%g mean gap %v", ideal, ring.chance(ideal, &r), time.Duration(sum/int64(len(c.gaps)))) + } + + p := newPeriodicState(r.HistoryLen) + cur := t0.UnixNano() + for i, g := range c.gaps { + cur += g + s, ok := p.observe(cur, TrafficUnclassified, &r, &sc) + if !ok { + continue + } + per := int64(s.ev.Period) + if per < int64(r.MinPeriod) || per > int64(r.MaxPeriod) || max(per-c.period, c.period-per) > tolNanos(&r, c.period) { + t.Fatalf("gap %d: signalled %v, want a period in range within tol of %v", i, s.ev.Period, time.Duration(c.period)) + } + return + } + t.Fatalf("no signal over %d gaps; %v and periods up to 301s explain every gap", len(c.gaps), time.Duration(c.period)) + }) + } +} + +// The periodRange memo is never stale: one scratch answers interleaved +// queries for different gaps, multiples, rules (the Detector shares its +// scratch between rules) and pulse numbers that share a slot, always with +// periodRange's result. +func TestRangeMemoIsNeverStale(t *testing.T) { + rules := []PeriodicRule{ + {JitterAbs: 5 * time.Second, JitterRel: 0.02}, + {JitterAbs: 1 * time.Second, JitterRel: 0.02}, + {JitterAbs: 5 * time.Second, JitterRel: 0.03}, + {JitterAbs: 6 * time.Second}, + } + gaps := []int64{int64(294 * time.Second), int64(300 * time.Second), int64(306 * time.Second), int64(600 * time.Second), int64(887 * time.Second)} + rng := rand.New(rand.NewSource(1082)) + var sc periodicScratch + hits := 0 + for i := 0; i < 200000; i++ { + r := &rules[rng.Intn(len(rules))] + g, k := gaps[rng.Intn(len(gaps))], int64(1+rng.Intn(maxMissingLimit+1)) + pulse := uint64(rng.Intn(3 * recentGaps)) + if sc.ranges != nil && r.JitterRel != 0 { + if m := sc.ranges[int(pulse%recentGaps)*(maxMissingLimit+1)+int(k-1)]; m.g == g && m.abs == int64(r.JitterAbs) && m.rel == r.JitterRel { + hits++ + } + } + lo, hi := sc.memoRange(pulse, g, k, r) + if wlo, whi := periodRange(g, k, r); lo != wlo || hi != whi { + t.Fatalf("query %d: memo [%d, %d], periodRange [%d, %d] (g=%d k=%d rule %+v pulse %d)", i, lo, hi, wlo, whi, g, k, *r, pulse) + } + } + if hits < 1000 || hits > 199000 { + t.Fatalf("fixture: %d memo hits of 200000 queries; want both hits and misses", hits) + } + if len(sc.ranges) != recentGaps*(maxMissingLimit+1) { + t.Fatalf("memo holds %d entries, want the fixed %d", len(sc.ranges), recentGaps*(maxMissingLimit+1)) + } +} diff --git a/internal/anomaly/review_fixes_test.go b/internal/anomaly/review_fixes_test.go index 93d17b225..f4e9e4570 100644 --- a/internal/anomaly/review_fixes_test.go +++ b/internal/anomaly/review_fixes_test.go @@ -88,7 +88,12 @@ func TestRefinementDoesNotHideAMeetingSeedChain(t *testing.T) { // MaxChance exists. The gaps are two trains from an independent Monte Carlo // (rule: JitterAbs 1s, JitterRel 0.03, MaxChance 1e-9); without the Chance // condition the first signalled 27 pulses late and the second never. -// ad3ab30f's signals are the reference. +// ad3ab30f's signals are the reference, except that train 1240 now also +// signals at gap 69: since refinement ranges use each period's own +// tolerance, the train, broken by the 187.5s gap at 58, is found again as +// 172.48s over 12 pulses (median residual 2.561s <= 0.5*5.174s, Chance +// 1.1e-11), where the seed-tolerance range only reached 171.75s, whose +// residual 2.84s fails MaxJitterFraction. func TestRefinedCandidateReplacesOnlyWhenItWouldSignal(t *testing.T) { for _, c := range []struct { name string @@ -122,7 +127,7 @@ func TestRefinedCandidateReplacesOnlyWhenItWouldSignal(t *testing.T) { 174342474654, 169534334907, 175681725697, 173915952191, 187517102381, 174359296774, 169014698377, 174123854812, 175505387461, 175828494557, 173758510591, 168915340027, 171752353392, 167303471723, 175039029128, 174776307324, - }, 57, 1, 57}, + }, 57, 2, 69}, } { t.Run(c.name, func(t *testing.T) { r := experimentalPeriodic() @@ -141,7 +146,7 @@ func TestRefinedCandidateReplacesOnlyWhenItWouldSignal(t *testing.T) { } } if first != c.first || count != c.count || last != c.last { - t.Fatalf("signals first=%d count=%d last=%d, want %d %d %d (as on ad3ab30f)", first, count, last, c.first, c.count, c.last) + t.Fatalf("signals first=%d count=%d last=%d, want %d %d %d", first, count, last, c.first, c.count, c.last) } }) } diff --git a/internal/anomaly/state.go b/internal/anomaly/state.go index 920a1d122..d76f0d4bf 100644 --- a/internal/anomaly/state.go +++ b/internal/anomaly/state.go @@ -15,17 +15,21 @@ type episode struct { endedAt int64 hot uint32 version uint32 // bumped on every change; invalidates old deadlines - // label and cov summarize the hot signals of the episode, so that - // time-driven and eviction transitions carry the same traffic label and - // coverage flags as the transitions that opened it. + // label, cov and low summarize the hot signals of the episode, so that + // time-driven and eviction transitions, and a reopen, carry the same + // traffic label, coverage flags and low confidence as the transitions + // that opened it. label TrafficLabel labelSet bool cov CoverageFlags + low bool // some hot signal came from incomplete (censored) evidence } -// note merges one hot signal's label and coverage into the episode. -func (ep *episode) note(l TrafficLabel, c CoverageFlags) { +// note merges one hot signal's label, coverage and incomplete-evidence flag +// into the episode. +func (ep *episode) note(l TrafficLabel, c CoverageFlags, low bool) { ep.cov |= c + ep.low = ep.low || low if !ep.labelSet { ep.label, ep.labelSet = l, true return @@ -41,6 +45,7 @@ type transition struct { summary *EpisodeSummary label TrafficLabel cov CoverageFlags + low bool } func (ep *episode) summary() *EpisodeSummary { @@ -55,7 +60,7 @@ func (ep *episode) move(to State, r Reason, at int64, out []transition, withSumm if !legalTransition(ep.state, to, r) { panic("anomaly: illegal state transition " + ep.state.String() + "->" + to.String() + " (" + r.String() + ")") } - tr := transition{from: ep.state, to: to, reason: r, at: at, start: ep.start, label: ep.label, cov: ep.cov} + tr := transition{from: ep.state, to: to, reason: r, at: at, start: ep.start, label: ep.label, cov: ep.cov, low: ep.low} if withSummary { tr.summary = ep.summary() } diff --git a/internal/anomaly/testdata/replay_golden.json b/internal/anomaly/testdata/replay_golden.json index 8e47bdda3..f0e9aa4dd 100644 --- a/internal/anomaly/testdata/replay_golden.json +++ b/internal/anomaly/testdata/replay_golden.json @@ -339,7 +339,7 @@ "At": "2030-01-07T13:15:00Z", "DecidedAt": "2030-01-07T13:15:00Z", "EpisodeStart": "2030-01-07T12:45:00Z", - "Confidence": "route_observed", + "Confidence": "low", "Traffic": "unclassified", "Suppressed": false, "Coverage": [],