From 827cf50e50bb85c5236a700986cbd6c602248ab0 Mon Sep 17 00:00:00 2001 From: dborup Date: Mon, 28 Sep 2026 14:31:10 +0200 Subject: [PATCH 1/5] test(anomaly): reproduce narrow-range period miss and mixed-burst label Two deterministic reproductions for the periodic detector in the merged (not yet integrated) internal/anomaly library. Both fail on master: - A 5m train with gaps alternating 294s/306s, MinPeriod=299s, MaxPeriod=301s, tolerance 6s: 300s explains every gap and meets every criterion (asserted as a fixture guard), but every seed g/k lies outside the range and is dropped before chainRefined can refine it into 300s. - A periodic chain whose bursts merge an expected and an unclassified event is labelled unclassified instead of mixed when the newest pulse starts with an unclassified event, or when the only mixed burst sits among unclassified pulses. Controls cover the label cases that are already correct. Co-Authored-By: Claude Opus 5.5 --- internal/anomaly/periodic_followup_test.go | 115 +++++++++++++++++++++ 1 file changed, 115 insertions(+) create mode 100644 internal/anomaly/periodic_followup_test.go diff --git a/internal/anomaly/periodic_followup_test.go b/internal/anomaly/periodic_followup_test.go new file mode 100644 index 000000000..0e05b3e22 --- /dev/null +++ b/internal/anomaly/periodic_followup_test.go @@ -0,0 +1,115 @@ +package anomaly + +import ( + "fmt" + "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. +func TestPeriodicNarrowRangeFindsPeriodBehindOutOfRangeSeeds(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}} + 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, tolerance=6s") + } + 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) + } + }) + } +} From 74f86e81c9f4387895dbf300b4f4004f5cf8f2ec Mon Sep 17 00:00:00 2001 From: dborup Date: Mon, 28 Sep 2026 14:44:40 +0200 Subject: [PATCH 2/5] fix(anomaly): refine out-of-range seeds; label mixed bursts as mixed Periodic detector of the merged, not yet integrated internal/anomaly library. Narrow period range: estimate() dropped every seed g/k outside [MinPeriod, MaxPeriod] before chainRefined could refine it, so a range narrower than the jitter (gaps 294s/306s, range 299s..301s, tol 6s) found no period although 300s meets every criterion. A seed is now measured while some period in range may explain its gap as k periods (seedInReach: k*MinPeriod - tol(MinPeriod) <= g <= k*MaxPeriod + tol(MaxPeriod), sound for absolute and relative jitter since k*P - tol(P) never decreases in P). An out-of-range seed only feeds its refinement and is never the estimate itself. The seed count per pulse is still at most recentGaps*(MaxMissing+1), so searchCells (Chance) and periodicWorkBound are unchanged; measuredSteps is refreshed (+0..6.8% gaps visited). Mixed bursts: the chain label counted only all-expected pulses, so a chain whose only expected events sat in bursts together with unclassified events came out unclassified. Such a chain is now mixed; monitored still wins, and all-expected chains stay expected. Adds a property test for seedInReach (soundness under absolute, relative and combined jitter, and tight bounds) and an invariant test that an out-of-range seed is measured but never reported. Co-Authored-By: Claude Opus 5.5 --- internal/anomaly/periodic.go | 38 +++++++++-- internal/anomaly/periodic_bench_test.go | 8 +-- internal/anomaly/periodic_followup_test.go | 78 ++++++++++++++++++++++ 3 files changed, 115 insertions(+), 9 deletions(-) diff --git a/internal/anomaly/periodic.go b/internal/anomaly/periodic.go index 77db75578..cb8b6b4df 100644 --- a/internal/anomaly/periodic.go +++ b/internal/anomaly/periodic.go @@ -321,7 +321,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 +329,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 +353,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,21 +372,24 @@ 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. + // + // A seed outside [MinPeriod, MaxPeriod] is measured only for its + // refinement (see seedInReach); it 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 { + if !seedInReach(g, k, r) { continue } + cand := g / int64(k) if slices.Contains(seen, cand) { continue } seen = append(seen, cand) c, ref := p.chainRefined(cand, r, sc) - if c.gaps > 0 && (best.gaps == 0 || c.preferred(best, r)) { + if cand >= minP && cand <= maxP && c.gaps > 0 && (best.gaps == 0 || c.preferred(best, r)) { best = c } if ref != cand && !slices.Contains(seenRef, ref) && !slices.Contains(seen, ref) { @@ -416,6 +427,23 @@ func (p *periodicState) estimate(r *PeriodicRule, sc *periodicScratch) chainInfo return best } +// seedInReach reports whether the seed g/k is worth measuring: whether some +// period P in [MinPeriod, MaxPeriod] may explain gap g as k periods +// (|g - k*P| <= tol(P)). A seed outside the range can still be refined into +// it by chainRefined: gaps of 294s, 306s... with tol 6s and range +// [299s, 301s] give only out-of-range seeds, which refine to 300s. +// +// The P nearest g/k are the range ends, and k*P - tol(P) never decreases in +// P (tol grows by at most JitterRel*dP, JitterRel <= 0.25 < k), so such a P +// exists only if k*MinPeriod - tol(MinPeriod) <= g <= k*MaxPeriod + +// tol(MaxPeriod). It admits at most one seed per gap and k, 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. // 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 index 0e05b3e22..ebe66ca38 100644 --- a/internal/anomaly/periodic_followup_test.go +++ b/internal/anomaly/periodic_followup_test.go @@ -2,6 +2,8 @@ package anomaly import ( "fmt" + "math/rand" + "slices" "testing" "time" ) @@ -113,3 +115,79 @@ func TestPeriodicBurstTrafficLabel(t *testing.T) { }) } } + +// seedInReach keeps every seed g/k for which some period in range explains g +// as k periods with its own tolerance, under absolute and relative jitter, +// and its bounds are tight. +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) + } + lo, hi := int64(k)*minP-tolNanos(&r, minP), int64(k)*maxP+tolNanos(&r, maxP) + if !seedInReach(lo, k, &r) || !seedInReach(hi, k, &r) || seedInReach(lo-1, k, &r) || seedInReach(hi+1, k, &r) { + t.Fatalf("trial %d: reach of k=%d is not exactly [%v, %v]", trial, k, time.Duration(lo), time.Duration(hi)) + } + } + 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) + } + } + } +} From 9899b9f695c0022faa6acdbe4995b1a52b33b0e8 Mon Sep 17 00:00:00 2001 From: dborup Date: Mon, 28 Sep 2026 14:55:08 +0200 Subject: [PATCH 3/5] reviewfix(anomaly): always measure in-range seeds; prove seedInReach exhaustively From the independent review of 74f86e81 (verdict: approve, one should-fix). - With a tolerance below k-1 ns, seedInReach alone rejected a gap in (k*MaxPeriod+tol, k*MaxPeriod+k-1] whose seed floor(g/k) = MaxPeriod is in range, so the fix could drop a seed master measured and lower the estimate. estimate() now measures every in-range seed and, beyond them, the out-of-range seeds seedInReach keeps: the seed pool is a superset of master's and the in-range ranking is unchanged. Regression test: TestPeriodicInRangeSeedIsAlwaysMeasured (fails with seedInReach-only admission). - seedInReach's comment states that the test is necessary, not sufficient, and gives the monotonicity argument including the 1 ns truncation step of tolNanos. - TestSeedInReachIsSoundExhaustively checks every explained gap over small integer periods, with absolute, relative and combined jitter and every k; the sampled test no longer restates the formula as a tightness check. - The narrow-range reproduction also runs with a relative tolerance. Work counters of TestPeriodicSearchWorkIsBounded are unchanged by this commit. Co-Authored-By: Claude Opus 5.5 --- internal/anomaly/periodic.go | 36 +++-- internal/anomaly/periodic_followup_test.go | 166 +++++++++++++++------ 2 files changed, 140 insertions(+), 62 deletions(-) diff --git a/internal/anomaly/periodic.go b/internal/anomaly/periodic.go index cb8b6b4df..f55ec1427 100644 --- a/internal/anomaly/periodic.go +++ b/internal/anomaly/periodic.go @@ -373,23 +373,25 @@ func (p *periodicState) estimate(r *PeriodicRule, sc *periodicScratch) chainInfo // alone would rank higher, and until a refined period would signal it // changes nothing. // - // A seed outside [MinPeriod, MaxPeriod] is measured only for its - // refinement (see seedInReach); it is never the estimate itself. + // 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++ { - if !seedInReach(g, k, r) { + cand := g / int64(k) + inRange := cand >= minP && cand <= maxP + if !inRange && !seedInReach(g, k, r) { continue } - cand := g / int64(k) if slices.Contains(seen, cand) { continue } seen = append(seen, cand) c, ref := p.chainRefined(cand, r, sc) - if cand >= minP && cand <= maxP && 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) { @@ -427,18 +429,20 @@ func (p *periodicState) estimate(r *PeriodicRule, sc *periodicScratch) chainInfo return best } -// seedInReach reports whether the seed g/k is worth measuring: whether some -// period P in [MinPeriod, MaxPeriod] may explain gap g as k periods -// (|g - k*P| <= tol(P)). A seed outside the range can still be refined into -// it by chainRefined: gaps of 294s, 306s... with tol 6s and range -// [299s, 301s] give only out-of-range seeds, which refine to 300s. +// 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 P nearest g/k are the range ends, and k*P - tol(P) never decreases in -// P (tol grows by at most JitterRel*dP, JitterRel <= 0.25 < k), so such a P -// exists only if k*MinPeriod - tol(MinPeriod) <= g <= k*MaxPeriod + -// tol(MaxPeriod). It admits at most one seed per gap and k, so there are -// still at most recentGaps*(MaxMissing+1) seeds, as searchCells charges and -// periodicWorkBound counts. +// 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) diff --git a/internal/anomaly/periodic_followup_test.go b/internal/anomaly/periodic_followup_test.go index ebe66ca38..9d5a25aa1 100644 --- a/internal/anomaly/periodic_followup_test.go +++ b/internal/anomaly/periodic_followup_test.go @@ -15,50 +15,59 @@ import ( // 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. +// the range instead of being dropped before refinement, with an absolute and +// with a relative tolerance. func TestPeriodicNarrowRangeFindsPeriodBehindOutOfRangeSeeds(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}} - 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)) - } + 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) + // 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, tolerance=6s") - } - ev := a[0].Periodic - if ev == nil || ev.Period != 5*time.Minute { - t.Fatalf("detected period %+v, want 5m0s", ev) + 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) + } + }) } } @@ -117,8 +126,8 @@ func TestPeriodicBurstTrafficLabel(t *testing.T) { } // seedInReach keeps every seed g/k for which some period in range explains g -// as k periods with its own tolerance, under absolute and relative jitter, -// and its bounds are tight. +// 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 @@ -150,10 +159,6 @@ func TestSeedInReachKeepsEverySeedAPeriodInRangeExplains(t *testing.T) { 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) } - lo, hi := int64(k)*minP-tolNanos(&r, minP), int64(k)*maxP+tolNanos(&r, maxP) - if !seedInReach(lo, k, &r) || !seedInReach(hi, k, &r) || seedInReach(lo-1, k, &r) || seedInReach(hi+1, k, &r) { - t.Fatalf("trial %d: reach of k=%d is not exactly [%v, %v]", trial, k, time.Duration(lo), time.Duration(hi)) - } } if checked < 19000 { t.Fatalf("only %d of 20000 trials checked", checked) @@ -191,3 +196,72 @@ func TestPeriodicOutOfRangeSeedIsNeverReported(t *testing.T) { } } } + +// 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) + } + } + } +} From 0a30e8b4100b8b3884d7fdb1189f9c73dd56940a Mon Sep 17 00:00:00 2001 From: dborup Date: Mon, 28 Sep 2026 16:33:53 +0200 Subject: [PATCH 4/5] test(anomaly): reproduce relative-jitter refinement miss and lost ConfidenceLow The two remaining findings from the review of the merged (not yet integrated) anomaly library. Every new test fails on 9899b9f6 at its intended assertion: - TestPeriodicRelativeJitterRefinesWithTheCandidateTolerance: JitterRel 2%, gaps repeating 305.95s, 305.95s, 294.05s. 300s explains every gap within its own 6s tolerance and meets every criterion (asserted as a fixture guard), but chainRefined builds the common range with the seed's tolerance, so no seed refines to a period that explains the train and it is never detected. - TestCensoredNewStreamKeepsLowConfidenceOn{QuietEnd,EvictionEnd,Reopen}: a new-stream episode opened from censored evidence is reported low, but its quiet end, its eviction end and a reopen within the cooldown (by an uncensored reactivation) come out as route_observed. Controls (pass before and after): an uncensored stream is never lowered, and the low confidence does not carry into a new episode after the cooldown. Co-Authored-By: Claude Opus 5.5 --- internal/anomaly/confidence_followup_test.go | 138 +++++++++++++++++++ internal/anomaly/periodic_followup_test.go | 52 +++++++ 2 files changed, 190 insertions(+) create mode 100644 internal/anomaly/confidence_followup_test.go 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/periodic_followup_test.go b/internal/anomaly/periodic_followup_test.go index 9d5a25aa1..2b0cba9c1 100644 --- a/internal/anomaly/periodic_followup_test.go +++ b/internal/anomaly/periodic_followup_test.go @@ -265,3 +265,55 @@ func TestPeriodicInRangeSeedIsAlwaysMeasured(t *testing.T) { } } } + +// 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) + } +} From cfbc7c8472f3d6882e5e9d637503e580e50f7313 Mon Sep 17 00:00:00 2001 From: dborup Date: Mon, 28 Sep 2026 16:34:14 +0200 Subject: [PATCH 5/5] fix(anomaly): refine with each period's own tolerance; keep ConfidenceLow per episode Relative jitter (P2): chainRefined intersected the ranges [(g-tol)/k, (g+tol)/k] with the seed's tolerance, although a period P explains gap g with its own, tol(P) = max(JitterAbs, JitterRel*P). The new periodRange returns exactly {P : |g - k*P| <= tol(P)}: the integer ends as before when JitterRel is 0 (proven against the old formula as an oracle), and otherwise float estimates of g/(k+JitterRel) and g/(k-JitterRel) stepped to the exact ends (k*P -+ tol(P) never decrease in P). Still at most one refined period per seed, so searchCells and the Chance union bound are unchanged. The refined period is now clamped into the common range intersected with [MinPeriod, MaxPeriod]. With exact ranges a common range can straddle MaxPeriod; unclamped, the phase estimate above it was rejected although periods in range explain every gap (15 of 10,000 narrow-range edge trains found by 9899b9f6 were lost without it). Under absolute jitter this only adds detections at the range ends. periodRange results are memoized in the Detector-owned scratch, one slot per (pulse mod recentGaps, k), at most recentGaps*(maxMissingLimit+1) entries; an entry is used only when g, JitterAbs and JitterRel match. Lost ConfidenceLow (P3): the episode now records that one of its hot signals came from censored new-stream evidence (like its traffic label and coverage), every transition carries it, and emit lowers the confidence for any transition of such an episode: quiet end, eviction end and a reopen within the cooldown. A new episode starts without it. Golden: exactly one field changes, the quiet end of censored episode 7771b168 from route_observed to low (sha256 3d6fbf08...). The golden guard keeps its pinned hash and maps that one documented field back. The pinned signals of Monte Carlo train 1240 become 57/2/69: the train, broken at gap 58, is found again as 172.48s (residual 2.561s <= 2.587s, Chance 1.1e-11), which the seed-tolerance range could not reach. Co-Authored-By: Claude Opus 5.5 --- internal/anomaly/chance_test.go | 32 +++- internal/anomaly/detector.go | 10 +- internal/anomaly/periodic.go | 93 ++++++++-- internal/anomaly/periodic_followup_test.go | 179 +++++++++++++++++++ internal/anomaly/review_fixes_test.go | 11 +- internal/anomaly/state.go | 17 +- internal/anomaly/testdata/replay_golden.json | 2 +- 7 files changed, 317 insertions(+), 27 deletions(-) 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/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 f55ec1427..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 @@ -454,8 +455,9 @@ func seedInReach(g int64, k int, r *PeriodicRule) bool { // 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 @@ -496,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 @@ -520,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_followup_test.go b/internal/anomaly/periodic_followup_test.go index 2b0cba9c1..cb35931d6 100644 --- a/internal/anomaly/periodic_followup_test.go +++ b/internal/anomaly/periodic_followup_test.go @@ -317,3 +317,182 @@ func TestPeriodicRelativeJitterRefinesWithTheCandidateTolerance(t *testing.T) { 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": [],