diff --git a/grafana-alertcheck/.changeset/v0.1.9.md b/grafana-alertcheck/.changeset/v0.1.9.md new file mode 100644 index 000000000..f4b11ec7a --- /dev/null +++ b/grafana-alertcheck/.changeset/v0.1.9.md @@ -0,0 +1,2 @@ +- Add support for Grafana's `Recovering` instance state ("keep firing for"): a first-seen recovering instance is classified preexisting, and `check` observes it through its own recovery deadline so a genuine recovery reads `recovered` while an unresolved, paused, or absent one stays `still_failing`. +- Remove `watch --poll-interval`. Every rule now polls at exactly half its evaluation interval, so a state that lasts a full interval can never fall between polls; the flag could only widen `maxGap` and hide transitions. diff --git a/grafana-alertcheck/cmd/check.go b/grafana-alertcheck/cmd/check.go index 57f64996b..7c41c21a4 100644 --- a/grafana-alertcheck/cmd/check.go +++ b/grafana-alertcheck/cmd/check.go @@ -34,7 +34,7 @@ func runCheck(args []string, stdin io.Reader, stdout, stderr io.Writer) int { pidfile := fs.String("pidfile", "", "pidfile of the recorder to stop before reading --in (default .pid)") from := fs.String("from", "", "the moment the deploy finished, RFC3339 (required with --in)") to := fs.String("to", "", "the end of the window to classify, RFC3339 (required)") - states := fs.String("states", "", "comma-separated bad states to classify against (default: firing)") + states := fs.String("states", "", "comma-separated bad states to classify against (default: firing,recovering)") preexisting := fs.String("preexisting", "", "how to judge an instance already bad at `from` (default: fail-unless-recovered)") minObserved := fs.Int("min-observed", 0, "minimum rules that must be observed (default: every resolved rule)") allowPaused := fs.Bool("allow-paused", false, "do not count a rule paused before the window against --min-observed") diff --git a/grafana-alertcheck/cmd/common.go b/grafana-alertcheck/cmd/common.go index 443b986be..4b0ea0594 100644 --- a/grafana-alertcheck/cmd/common.go +++ b/grafana-alertcheck/cmd/common.go @@ -12,11 +12,9 @@ import ( ) // commonFlags is registerCommon's result: the flags watch and check share. -// Connection details are never flags, and states / poll-interval are -// deliberately NOT here — states is check-only because recording is -// unfiltered, and poll-interval is watch-only because check reads the cadence -// from the log header. Putting either here would give both commands an opinion -// about a value only one of them may set. +// Connection details are never flags, and states is deliberately NOT here: +// states is check-only because recording is unfiltered, and putting it here +// would give both commands an opinion about a value only one of them may set. type commonFlags struct { folder *string concurrency *int @@ -101,13 +99,13 @@ func readAlerts(stdin io.Reader, flagName, path string) ([]string, error) { // parseStates parses check's --states flag: a comma-separated list of the // "bad" state vocabulary Config.States matches against (classify.go's // badStateSet). An empty string is not resolved here — it means "use the -// library default of {firing}" — so this returns nil, nil for "" rather than -// an error. +// library default of {firing, recovering}" — so this returns nil, nil for "" +// rather than an error. // // normal is deliberately NOT accepted. The vocabulary is fixed to -// firing | pending | nodata | error precisely because "normal" is the good -// state, never a bad one to classify against: --states normal would turn every -// healthy instance into a violation and fail every healthy fleet. +// firing | pending | recovering | nodata | error precisely because "normal" is +// the good state, never a bad one to classify against: --states normal would +// turn every healthy instance into a violation and fail every healthy fleet. func parseStates(s string) ([]gate.State, error) { if strings.TrimSpace(s) == "" { return nil, nil @@ -119,10 +117,10 @@ func parseStates(s string) ([]gate.State, error) { continue } switch gate.State(part) { - case gate.StateFiring, gate.StatePending, gate.StateNodata, gate.StateError: + case gate.StateFiring, gate.StatePending, gate.StateRecovering, gate.StateNodata, gate.StateError: out = append(out, gate.State(part)) default: - return nil, fmt.Errorf("--states: unknown state %q (want any of: firing, pending, nodata, error)", part) + return nil, fmt.Errorf("--states: unknown state %q (want any of: firing, pending, recovering, nodata, error)", part) } } if len(out) == 0 { diff --git a/grafana-alertcheck/cmd/watch.go b/grafana-alertcheck/cmd/watch.go index 3a12c5ce4..c8441dc98 100644 --- a/grafana-alertcheck/cmd/watch.go +++ b/grafana-alertcheck/cmd/watch.go @@ -13,7 +13,7 @@ import ( const watchUsage = "usage: grafana-alertcheck watch --out [--pidfile F] [--daemon-log F] " + "(--alerts [--folder F] | --include-labels k=v,... [--exclude-labels k=v,...]) [--exclude-alerts ] " + - "[--poll-interval D] [--concurrency N] [--until RFC3339]" + "[--concurrency N] [--until RFC3339]" // runWatch is the record step's entire CLI surface, split in two by one flag // set — gate.DaemonChildFlag ("--daemon-child") and gate.ReadyFDFlag @@ -44,7 +44,6 @@ func runWatch(args []string, stdin io.Reader, stdout, stderr io.Writer) int { pidfile := fs.String("pidfile", "", "pidfile path (default .pid)") daemonLog := fs.String("daemon-log", "", "stdout/stderr sink for the detached recorder (default .daemon.log)") until := fs.String("until", "", "optional hard stop, RFC3339 (default: run until check stops it)") - pollInterval := fs.String("poll-interval", "", "override every rule's poll cadence (default: half its own evaluation interval)") // Hidden: never in watchUsage, never typed by an operator (see doc comment). daemonChild := fs.Bool(gate.DaemonChildFlag[2:], false, "") @@ -124,14 +123,6 @@ func runWatch(args []string, stdin io.Reader, stdout, stderr io.Writer) int { } cfg.Until = t } - if *pollInterval != "" { - d, err := time.ParseDuration(*pollInterval) - if err != nil { - fmt.Fprintf(stderr, "--poll-interval: %v\n", err) - return 2 - } - cfg.PollEvery = d - } if err := gate.Watch(context.Background(), cfg); err != nil { fmt.Fprintln(stderr, err) diff --git a/grafana-alertcheck/cmd/watch_test.go b/grafana-alertcheck/cmd/watch_test.go index e434e88c7..6ca4488c8 100644 --- a/grafana-alertcheck/cmd/watch_test.go +++ b/grafana-alertcheck/cmd/watch_test.go @@ -32,9 +32,6 @@ func TestRunWatch_FlagValidation(t *testing.T) { {"until in the past", true, func(t *testing.T) []string { return []string{"--out", t.TempDir() + "/log.jsonl", "--alerts", writeTempAlerts(t), "--until", "2000-01-01T00:00:00Z"} }, "not in the future"}, - {"bad poll-interval", true, func(t *testing.T) []string { - return []string{"--out", t.TempDir() + "/log.jsonl", "--alerts", writeTempAlerts(t), "--poll-interval", "not-a-duration"} - }, "--poll-interval"}, {"alerts and labels", true, func(t *testing.T) []string { return []string{"--out", t.TempDir() + "/log.jsonl", "--alerts", writeTempAlerts(t), "--include-labels", "team=bcm"} }, "cannot be combined with label selection"}, diff --git a/grafana-alertcheck/docs/advanced.md b/grafana-alertcheck/docs/advanced.md index fa6f0be23..a68c92532 100644 --- a/grafana-alertcheck/docs/advanced.md +++ b/grafana-alertcheck/docs/advanced.md @@ -10,7 +10,7 @@ description: Why grafana-alertcheck schedules per rule, how the request budget w ## Per-rule schedules, never a global cycle -Each rule polls at its **own** cadence (default: half the rule's own evaluation interval). There is deliberately no single global minimum-interval cycle. Overwrite with `--poll-interval`. +Each rule polls at its **own** cadence: half the rule's own evaluation interval, always. The cadence is not configurable — polling faster cannot reveal more (Grafana state only changes on evaluations) and polling slower could let a state fall between polls. One rule at `intervalSeconds=10` beside twenty at `300` keeps a 5 s cadence for itself and 150 s for the other twenty — not a 5 s cycle for all of them, which would be a 60× request bloat at ~1.8 s per request and would fail to start on a reasonable fleet. @@ -25,7 +25,7 @@ The gate records one observation of every rule up front and checks the schedule - **Burst bound** — the slowest request exceeds the fleet's tightest cadence, which can open a mid-run gap. - **Startup handoff** — draining the first-observation pass's backlog at `--concurrency` would leave some rule unpolled past its own `maxGap`. A rule the pass observed early is seeded overdue, and a tight rule observed late can queue behind every rule due before it. The gate simulates the poller's first cycles from the recorded observation times and measured latencies — each wake takes every rule due at that instant, polls the batch at `--concurrency`, and wakes again when it ends — and refuses if any rule's first poll would land past its `maxGap`. Steady-state utilization cannot see this — a long pass at low concurrency is exactly the case it passes. -The error names only the levers that can fix it: the minimum `--concurrency` when the schedule is concurrency-bound, and `--poll-interval` or a smaller alert set for single-request shapes concurrency cannot shorten. It never prescribes a single interval. +The error names only the levers that can fix it: the minimum `--concurrency` when the schedule is concurrency-bound, and a smaller alert set for single-request shapes concurrency cannot shorten. ## The startup pass and `ready_at` diff --git a/grafana-alertcheck/docs/how-alerts-are-evaluated.md b/grafana-alertcheck/docs/how-alerts-are-evaluated.md index 60c2c28f9..0a2042163 100644 --- a/grafana-alertcheck/docs/how-alerts-are-evaluated.md +++ b/grafana-alertcheck/docs/how-alerts-are-evaluated.md @@ -19,12 +19,13 @@ Grafana reports instance states in two vocabularies (`Alerting`/`Normal` at inst | `normal` | Healthy | | `firing` | The condition is true and `for` has elapsed | | `pending` | The condition is true, `for` has not elapsed | +| `recovering` | The condition has cleared but the rule's *keep firing for* has not elapsed | | `nodata` | The query returned no series (synthetic instance) | | `error` | The query failed (synthetic instance) | A rule's **rule-level** `state` and `health` are kept verbatim and only reported — they are never classified. The **instance** state is what the classifier reasons about. -A "bad" instance is one whose canonical state is in `--states` (default `firing`). `pending` and `nodata` are excluded by default. +A "bad" instance is one whose canonical state is in `--states` (default `firing,recovering`). `pending` and `nodata` are excluded by default. `recovering` is bad because the instance is still firing (Alertmanager keeps notifying) until its recovery period ends; `--states firing` opts out of tracking it. ## Verdict model @@ -35,7 +36,7 @@ For each instance the gate builds a timeline of bad spans over `[from, to]`, the | `healthy` | Good throughout, observed throughout | → 0 | | `new_failure` | Entered a bad state **inside** the window | → 1 | | `still_failing` | Bad at `from`, still bad at `to` | → 1 | -| `recovered` | Bad at `from`, cleared before `to`, stayed clear | → 0 | +| `recovered` | Bad at `from`, cleared and stayed clear (possibly during the recovery observation) | → 0 | | `unstable` | Cleared, then became bad again | → 1 | | `paused` | Paused **before** the window opened | counts against `--min-observed` unless `--allow-paused` | | `not_verified` | The window could not be observed: a gap, sustained `health=error`, a stale evaluation, or an absent rule | → 2 | @@ -91,8 +92,10 @@ A preexisting bad instance is deliberately **not** terminal: if it clears before Fail-fast is on by default and always preserves the failure: an early run can exit `1` or `2`, never `0`. The one difference from a full run is that an early exit may report `1` before an inability surfaces that would have made it `2`. `--fail-fast=false` disables the guard and always waits for the full window and its coverage proof. -## The drain wait and `transitionGrace` +## The drain wait, `transitionGrace`, and recovery observation A condition that arises just before `to` becomes `firing` only at the first evaluation after its `for` elapses. `transitionGrace` (derived from the watched rules' `for` values) extends the classification bound past `to` so such a surfacing condition is caught. After collection, a **drain wait** polls until each rule has evaluated through `to + transitionGrace` (bounded by `drainTimeout`); a rule that never does is `not_verified`. -Run time = `(to − from) + transitionGrace + drainTimeout`. This is printed at start. A requested window with a subsecond part is rounded up to the next whole second by extending `to`, so the plan never reads a window like `9m59.99445781s`. +The mirror case is recovery: an instance whose condition cleared just before `to` enters `recovering` and only resolves after its **keep firing for** elapses, which can be long after `to`. A first-seen `recovering` instance is always treated as preexisting — recovering is only reachable from `firing`, and every rule is sampled twice per evaluation interval, so an in-window fire cannot be missed between polls; its `activeAt` is the recovery onset, never the fire onset. The episode stays open until the instance reports `normal`. To decide that, `check` keeps observing the affected rules past `to + transitionGrace`, up to `activeAt + keep firing for` plus an evaluation/cadence margin. In recorder mode the recorder keeps polling so the extension evidence lands in the log; in single-step mode `check` polls the affected rules directly. Either way the extension is scoped to the recovering instances: a different instance going bad during the observation is outside the window and stays `healthy`. Each instance expires on its **own** deadline, and a clear past it — or one arriving after the rule paused or disappeared — does not resolve the episode: it stays `still_failing` (fail-closed). + +Run time = `(to − from) + transitionGrace + drainTimeout`, plus the recovery observation when one triggers; the extra deadline is printed when it starts. A requested window with a subsecond part is rounded up to the next whole second by extending `to`, so the plan never reads a window like `9m59.99445781s`. diff --git a/grafana-alertcheck/docs/reference/cli.md b/grafana-alertcheck/docs/reference/cli.md index 91ee6910f..cfbd90b01 100644 --- a/grafana-alertcheck/docs/reference/cli.md +++ b/grafana-alertcheck/docs/reference/cli.md @@ -27,7 +27,7 @@ grafana-alertcheck list ```bash grafana-alertcheck watch --out [--pidfile F] [--daemon-log F] \ (--alerts [--folder F] | --include-labels k=v,... [--exclude-labels k=v,...]) \ - [--exclude-alerts ] [--poll-interval D] [--concurrency N] [--until RFC3339] + [--exclude-alerts ] [--concurrency N] [--until RFC3339] ``` | Flag | Default | Meaning | @@ -40,7 +40,6 @@ grafana-alertcheck watch --out [--pidfile F] [--daemon-log F] \ | `--include-labels` | — | Comma-separated exact-match `key=value` pairs selecting rules by label (cannot be combined with `--alerts`) | | `--exclude-labels` | — | Comma-separated exact-match `key=value` pairs; a rule carrying any of them is dropped (requires `--include-labels`) | | `--exclude-alerts` | — | File of alert names, one per line, or `-` for stdin; subtracted from the selected set (works with `--alerts` and with labels) | -| `--poll-interval` | half the rule's interval | Override every rule's cadence (never clamped) | | `--concurrency` | `1` | Max concurrent requests to Grafana | | `--until` | run until signalled | Optional hard stop | @@ -83,7 +82,7 @@ grafana-alertcheck check [--in ] [--pidfile F] --from RFC3339 --to RFC3339 | `--include-labels` | — | Comma-separated exact-match `key=value` pairs selecting rules by label (cannot be combined with `--alerts`) | | `--exclude-labels` | — | Comma-separated exact-match `key=value` pairs; a rule carrying any of them is dropped (requires `--include-labels`) | | `--exclude-alerts` | — | File of alert names, one per line, or `-` for stdin; subtracted from the selected set (works with `--alerts` and with labels; refused **with** `--in`) | -| `--states` | `firing` | Comma-separated bad states: `firing,pending,nodata,error` | +| `--states` | `firing,recovering` | Comma-separated bad states: `firing,pending,recovering,nodata,error` | | `--preexisting` | `fail-unless-recovered` | `fail-unless-recovered` \| `fail` \| `ignore` | | `--min-observed` | every resolved rule | Minimum rules that must be observed | | `--allow-paused` | `false` | Don't count pre-window-paused rules against `--min-observed` | diff --git a/grafana-alertcheck/docs/reference/log-format.md b/grafana-alertcheck/docs/reference/log-format.md index 74694c8d7..7939bfbc1 100644 --- a/grafana-alertcheck/docs/reference/log-format.md +++ b/grafana-alertcheck/docs/reference/log-format.md @@ -52,7 +52,7 @@ The header must be line 1, appear once, and carry `schema_version` `1` (any othe - `url` and `rules` are the log's identity — `check` validates them against the current environment and a fresh ruler read. - `started_at` is when the recording opened; `ready_at` is when the first-observation pass completed and every watched, non-paused rule had been observed once. The pass is sequential, so `check` refuses a `from` before `ready_at` (a window opening inside the pass would rest on observations that do not exist). `ready_at` is absent on logs written before the field existed; `check` then falls back to `started_at`. - `is_paused` records the pause state at record start (the moment `paused` means). -- `poll_every_seconds` is the cadence the recording **actually used** (after any `--poll-interval` override). `check` derives `maxGap` from it, never from `interval_seconds`. +- `poll_every_seconds` is the cadence the recording used: always half the rule's `interval_seconds`. `check` derives `maxGap` from it, never by re-deriving from `interval_seconds`. - `for_seconds`, `interval_seconds`, `no_data_state`, `exec_err_state` are forensic only — `check` re-resolves definitions and never reads them back. ## Poll diff --git a/grafana-alertcheck/internal/gate/check.go b/grafana-alertcheck/internal/gate/check.go index ea07e3a88..eb6262687 100644 --- a/grafana-alertcheck/internal/gate/check.go +++ b/grafana-alertcheck/internal/gate/check.go @@ -67,12 +67,10 @@ type Config struct { Log string PidFile string - // There is deliberately NO PollEvery here, and `check` has no - // --poll-interval flag. In log mode the cadence comes from the header — - // the cadence the recording actually used — and a second authority would - // let an operator silently widen maxGap over evidence that was recorded at - // a different rate; in single-step mode the same process records and - // classifies, so the default cadence is the only cadence there is. + // There is deliberately NO cadence override anywhere. In log mode the + // cadence comes from the header — the cadence the recording actually used + // — and in single-step mode the same process records and classifies, so + // half the evaluation interval is the only cadence there is. Concurrency int Clock Clock @@ -229,13 +227,12 @@ func check(ctx context.Context, cfg Config, src Source) (Result, error) { // advisory: the authoritative header is re-read once collection ends and // the writer has exited. var ( - resolved []Definition - notes []string - earlyHdr Header - logHasHdr bool - rt map[string]RuleTimings - gt GlobalTimings - timingNote []string + resolved []Definition + notes []string + earlyHdr Header + logHasHdr bool + rt map[string]RuleTimings + gt GlobalTimings ) if cfg.Log != "" { earlyHdr, err = ReadLogHeader(cfg.Log) @@ -280,18 +277,14 @@ func check(ctx context.Context, cfg Config, src Source) (Result, error) { // ---- Derive the timings, print them, fit the request budget. ---------- if logHasHdr { // The header is the authority for the cadence actually recorded at; - // re-deriving it from defs would compare gaps recorded at an override - // cadence against thresholds computed from the default — fail-open in - // the faster-override direction. + // re-deriving it from defs would compare gaps against thresholds + // computed from a different cadence. rt, gt, err = DeriveTimingsFromLog(earlyHdr, resolved) if err != nil { return Result{}, fmt.Errorf("log identity: %w", err) } } else { - rt, gt, timingNote = DeriveTimings(resolved, 0) - for _, n := range timingNote { - fmt.Fprintf(cfg.Notes, "note: %s\n", n) - } + rt, gt = DeriveTimings(resolved) } fmt.Fprintln(cfg.Notes, StartupSummary(from, cfg.To, gt)) // MinObserved is printed with the plan, beside "planned run time", rather @@ -422,18 +415,34 @@ func check(ctx context.Context, cfg Config, src Source) (Result, error) { if err != nil { // Fail closed: in recorder mode the detached recorder is still running, // so reap it and keep the original error. - if logHasHdr { - if held, stopErr := stopRecorder(ctx, cfg); stopErr != nil { - err = errors.Join(err, stopErr) - } else if held != nil { - _ = held.Close() - } - } + err = reapRecorder(ctx, cfg, logHasHdr, err) // Nothing collected is classified; the count lets an operator tell a // run that failed at once from one that failed at minute nine. return Result{}, fmt.Errorf("collect evidence after %d poll(s): %w", len(collected.polls), err) } + // ---- Recovery observation, recorder mode. ----------------------------- + // The recorder keeps polling, so the extension evidence lands in the log + // ReadLog reads next and a re-classification sees the same evidence. The + // wait primes itself by reading the log before selecting the affected + // rules, so a Recovering poll written after collection's last read is not + // missed. + if logHasHdr && collected.term == nil && collected.sentinel == nil && + badStateSet(pol.States)[StateRecovering] { + + // --fail-fast=false never tailed the log, so create the tailer now. + if tail == nil { + tail, err = newLogTailer(cfg.Log) + if err != nil { + return Result{}, reapRecorder(ctx, cfg, true, err) + } + defer tail.Close() + } + if _, err := recoveryWait(ctx, cfg, resolved, rt, collected.polls, from, windowEnd, tailRecoverySource(tail), true); err != nil { + return Result{}, reapRecorder(ctx, cfg, true, err) + } + } + var ( polls []Poll sentinel *time.Time @@ -488,6 +497,20 @@ func check(ctx context.Context, cfg Config, src Source) (Result, error) { return result, err } + // ---- Recovery observation, single-step mode. -------------------------- + // Before the drain: a lagging rule can hold the drain past a recovering + // instance's deadline, and the drain keeps no instance evidence. + if !logHasHdr && badStateSet(pol.States)[StateRecovering] { + reducer := NewReducer() + reducer.seedFrom(polls) + recoveryPolls, err := recoveryWait(ctx, cfg, resolved, rt, polls, from, windowEnd, + directRecoverySource(src, reducer, cfg.Concurrency), false) + if err != nil { + return Result{}, err + } + polls = append(polls, recoveryPolls...) + } + // ---- The drain wait. -------------------------------------------------- // The last instance of the liveness check: did this rule evaluate through // the end of the window? It is I/O and it is deliberately NOT part of @@ -719,6 +742,22 @@ func stopRecorder(ctx context.Context, cfg Config) (*os.File, error) { }) } +// reapRecorder stops a still-running detached recorder after a failure so it is +// never left behind, joining any stop error. +func reapRecorder(ctx context.Context, cfg Config, logHasHdr bool, err error) error { + if !logHasHdr { + return err + } + held, stopErr := stopRecorder(ctx, cfg) + if stopErr != nil { + return errors.Join(err, stopErr) + } + if held != nil { + _ = held.Close() + } + return err +} + // drainVerdict is what the drain wait concluded about one rule it could not // clear. It carries the reason as well as the prose because the two outcomes // are genuinely different faults: drain_timeout means the rule is still there @@ -852,6 +891,238 @@ func drainWait(ctx context.Context, cfg Config, src Source, defs []Definition, p } } +// recoveryDeadline is the latest runner-domain time an instance that entered +// Recovering at activeAt can be expected to report its resolution: KeepFiringFor +// elapses at activeAt+KeepFiringFor, Grafana transitions on a later evaluation, +// and a poll reports that up to one cadence later. +func recoveryDeadline(activeAt time.Time, keepFiringFor time.Duration, t RuleTimings) time.Time { + return activeAt.Add(keepFiringFor + t.evalStaleAfter + t.pollEvery) +} + +// recoveringAtWindowEnd returns, per rule UID, the keys of instances whose most +// recent in-window observation is Recovering, each mapped to its recovery +// deadline. A rule with no such instance is absent. +func recoveringAtWindowEnd(polls []Poll, from, windowEnd time.Time, rt map[string]RuleTimings) map[string]map[string]time.Time { + type observed struct { + state State + activeAt time.Time + kff time.Duration + } + latest := make(map[string]map[string]observed) + for _, p := range polls { + if !pollInWindow(p, from, windowEnd) { + continue + } + byKey := latest[p.RuleUID] + if byKey == nil { + byKey = make(map[string]observed) + latest[p.RuleUID] = byKey + } + for _, inst := range p.Abnormal { + byKey[instanceKey(inst.Labels)] = observed{inst.State, runnerTime(p, inst.ActiveAt), p.KeepFiringFor()} + } + for _, key := range p.Cleared { + delete(byKey, key) + } + for _, key := range p.Vanished { + delete(byKey, key) + } + } + + out := make(map[string]map[string]time.Time) + for uid, byKey := range latest { + t, ok := rt[uid] + if !ok { + continue + } + for key, o := range byKey { + if o.state != StateRecovering { + continue + } + if out[uid] == nil { + out[uid] = make(map[string]time.Time) + } + out[uid][key] = recoveryDeadline(o.activeAt, o.kff, t) + } + } + return out +} + +// recoverySource returns the next batch of extension polls; sentinel is true +// when the recorder has finished (recorder mode only). +type recoverySource func(ctx context.Context, titles map[string]string, uids []string) (polls []Poll, sentinel bool, err error) + +// applyRecoveryPoll drops every key of poll.RuleUID that poll no longer reports +// Recovering: resolved, re-fired, paused, vanished, or absent all end the wait +// (an unresolved key then stays open and reads still_failing). +func applyRecoveryPoll(pending map[string]map[string]time.Time, poll Poll) { + keys, ok := pending[poll.RuleUID] + if !ok { + return + } + if poll.IsPaused || !poll.Found { + delete(pending, poll.RuleUID) + return + } + still := make(map[string]struct{}) + for _, inst := range poll.Abnormal { + if inst.State == StateRecovering { + still[instanceKey(inst.Labels)] = struct{}{} + } + } + for key := range keys { + if _, ok := still[key]; !ok { + delete(keys, key) + } + } + if len(keys) == 0 { + delete(pending, poll.RuleUID) + } +} + +// recoveryWait observes instances still Recovering at windowEnd until they +// resolve or their own deadline passes. It returns the extension polls; the +// recorder-mode caller discards them because ReadLog re-reads the log. +// prime reads the source once before selecting, so a Recovering poll written +// after collection stopped is not missed. +func recoveryWait(ctx context.Context, cfg Config, defs []Definition, rt map[string]RuleTimings, + polls []Poll, from, windowEnd time.Time, src recoverySource, prime bool) ([]Poll, error) { + + var out []Poll + if prime { + batch, sentinel, err := src(ctx, nil, nil) + if err != nil { + return nil, err + } + out = append(out, batch...) + polls = append(polls, batch...) + if sentinel { + return nil, nil // the recording finished; nothing left to observe + } + } + + pending := recoveringAtWindowEnd(polls, from, windowEnd, rt) + // Extension polls already in hand may resolve an episode before the first + // source read. + for _, p := range polls { + if pollAfterWindow(p, windowEnd) { + applyRecoveryPoll(pending, p) + } + } + if len(pending) == 0 { + return out, nil + } + + titles := make(map[string]string, len(pending)) + for _, d := range defs { + if keys, ok := pending[d.UID]; ok && len(keys) > 0 { + titles[d.UID] = d.Title + fmt.Fprintf(cfg.Notes, "recovery wait: rule %q has %d instance(s) still recovering\n", d.Title, len(keys)) + } + } + + for len(pending) > 0 { + now := cfg.Clock.Now() + // Expire keys individually: a key past its own deadline must not be + // resolved by a clear another instance kept the wait alive for. + for uid, keys := range pending { + for key, deadline := range keys { + if !now.Before(deadline) { + fmt.Fprintf(cfg.Notes, "recovery wait: rule %q: an instance did not resolve before %s; its episode stays open\n", + titles[uid], deadline.Format(time.RFC3339)) + delete(keys, key) + } + } + if len(keys) == 0 { + delete(pending, uid) + } + } + if len(pending) == 0 { + break + } + + uids := make([]string, 0, len(pending)) + for uid := range pending { + uids = append(uids, uid) + } + sort.Strings(uids) + + batch, sentinel, err := src(ctx, titles, uids) + if err != nil { + return out, err + } + relevant := make(map[string]struct{}, len(pending)) + for uid := range pending { + relevant[uid] = struct{}{} + } + for _, poll := range batch { + if _, ok := relevant[poll.RuleUID]; !ok { + continue + } + out = append(out, poll) + applyRecoveryPoll(pending, poll) + } + if sentinel || len(pending) == 0 { + break + } + + // Re-ask no faster than the tightest affected cadence, and never sleep + // past the soonest per-instance recovery deadline. + var wait time.Duration + soonest := false + for uid, keys := range pending { + for _, deadline := range keys { + if d := deadline.Sub(now); !soonest || d < wait { + wait, soonest = d, true + } + } + if every := rt[uid].pollEvery; every > 0 && (!soonest || every < wait) { + wait, soonest = every, true + } + } + select { + case <-ctx.Done(): + return out, ctx.Err() + case <-cfg.Clock.After(max(wait, 0)): + } + } + return out, nil +} + +// directRecoverySource polls the pending rules directly (single-step mode). +// The reducer is seeded from the polls already taken so the first extension +// poll's Cleared marker compares against the last recorded abnormal set. +func directRecoverySource(src Source, reducer *Reducer, concurrency int) recoverySource { + return func(ctx context.Context, titles map[string]string, uids []string) ([]Poll, bool, error) { + observed, err := observeAll(ctx, src, titles, uids, concurrency) + if err != nil { + return nil, false, fmt.Errorf("recovery wait: %w", err) + } + out := make([]Poll, 0, len(uids)) + for _, uid := range uids { + obs, ok := observed[uid] + if !ok { + continue + } + out = append(out, reducer.Reduce(uid, obs)) + } + return out, false, nil + } +} + +// tailRecoverySource reads the recorder's log (recorder mode). The recorder +// keeps polling, so the extension evidence is written to the log that ReadLog +// reads next. +func tailRecoverySource(tail *logTailer) recoverySource { + return func(context.Context, map[string]string, []string) ([]Poll, bool, error) { + polls, sentinel, err := tail.read() + if err != nil { + return nil, false, fmt.Errorf("recovery wait: %w", err) + } + return polls, sentinel != nil, nil + } +} + // anyPollEvaluatedThrough reports whether any recorded poll of a rule already // proves it evaluated through windowEnd. func anyPollEvaluatedThrough(polls []Poll, windowEnd time.Time) bool { diff --git a/grafana-alertcheck/internal/gate/check_test.go b/grafana-alertcheck/internal/gate/check_test.go index b53b12c3a..5c48aad89 100644 --- a/grafana-alertcheck/internal/gate/check_test.go +++ b/grafana-alertcheck/internal/gate/check_test.go @@ -798,7 +798,7 @@ func TestCheckSingleStepRefusesAScheduleThatDoesNotFit(t *testing.T) { _, err := check(context.Background(), cfg, src) require.Error(t, err, "the budget check to refuse the schedule") - for _, want := range []string{"raising --concurrency to at least 2", "raising poll-interval", "watching fewer alerts"} { + for _, want := range []string{"raising --concurrency to at least 2", "watching fewer alerts"} { require.Contains(t, err.Error(), want) } } @@ -1704,3 +1704,521 @@ func TestReadLogHeader(t *testing.T) { require.Contains(t, err.Error(), "schema version 99") }) } + +// --------------------------------------------------------------------------- +// Recovery observation ("keep firing for") +// --------------------------------------------------------------------------- + +// recoveryObservation is one state response for a rule with a recovery period. +func recoveryObservation(uid, title string, now time.Time, kff time.Duration, insts ...Instance) Observation { + totals := map[string]int{} + for _, i := range insts { + totals[string(i.State)]++ + } + return Observation{ + Rules: []StateRule{{ + UID: uid, Title: title, Folder: "F", Group: "G", + Interval: time.Minute, State: "inactive", Health: "ok", + LastEvaluation: now, KeepFiringFor: kff, + Instances: insts, Totals: totals, + }}, + GrafanaNow: now, Latency: 200 * time.Millisecond, + } +} + +func recoveringScenarioDef(uid, title string) Definition { + return Definition{ + UID: uid, Title: title, Folder: "F", Group: "G", + For: 3 * time.Minute, IntervalSeconds: 60, + NoDataState: "OK", ExecErrState: "OK", Kind: KindGrafanaManaged, + } +} + +// recoveringAtWindowEnd reports the instances still recovering at windowEnd and +// their deadline; an extension poll (after windowEnd) must not count. +func TestRecoveringAtWindowEnd(t *testing.T) { + from := testNow + to := from.Add(10 * time.Minute) + recoveryAt := to.Add(-time.Minute) + rt := map[string]RuleTimings{"r1": newRuleTimings(30*time.Second, 60)} + key := instanceKey(lbl("a")) + + polls := []Poll{ + abnormalPoll("r1", from, StateFiring, lbl("a"), from.Add(-time.Hour)), + recoveringPoll("r1", recoveryAt, lbl("a"), recoveryAt, 10*time.Minute), + recoveringPoll("r1", to.Add(time.Minute), lbl("b"), to.Add(time.Minute), 10*time.Minute), + } + got := recoveringAtWindowEnd(polls, from, to, rt) + require.Contains(t, got["r1"], key) + require.NotContains(t, got["r1"], instanceKey(lbl("b")), "an extension poll is not in-window evidence") + require.Equal(t, recoveryAt.Add(10*time.Minute+rt["r1"].evalStaleAfter+rt["r1"].pollEvery), got["r1"][key]) + + // A later in-window clear removes the instance: nothing to observe. + polls = append(polls, clearedPoll("r1", to, key)) + require.NotContains(t, recoveringAtWindowEnd(polls, from, to, rt), "r1") +} + +// The scenario: A is preexisting, starts recovering just before `to`, and only +// resolves after windowEnd; B goes pending at windowEnd and fires during A's +// recovery observation. A must end recovered and B healthy, because B's badness +// is outside the observation window. +func TestCheckRecoveringInstanceResolvesAndPostWindowOnsetIsIgnored(t *testing.T) { + clock := newVirtualClock(testNow) + from := clock.Now() + to := from.Add(5 * time.Minute) + windowEnd := to.Add(4 * time.Minute) // for=3m + interval=60s + recoveryAt := to.Add(-10 * time.Second) + resolveAt := recoveryAt.Add(10 * time.Minute) // first evaluation at/after KeepFiringFor + fireBAt := windowEnd.Add(3 * time.Minute) + + a := map[string]string{"instance": "a"} + b := map[string]string{"instance": "b"} + + src := newCheckSource(func(_ string, _ int) (Observation, error) { + now := clock.Now() + aState, aActive := StateRecovering, recoveryAt + switch { + case now.Before(recoveryAt): + aState, aActive = StateFiring, from.Add(-time.Hour) + case !now.Before(resolveAt): + aState, aActive = StateNormal, now + } + bState, bActive := StateNormal, now + switch { + case !now.Before(fireBAt): + bState, bActive = StateFiring, windowEnd + case !now.Before(windowEnd): + bState, bActive = StatePending, windowEnd + } + return recoveryObservation("r1", "Rule One", now, 10*time.Minute, + Instance{Labels: a, State: aState, ActiveAt: aActive}, + Instance{Labels: b, State: bState, ActiveAt: bActive}, + ), nil + }) + src.defs = []Definition{recoveringScenarioDef("r1", "Rule One")} + + cfg := Config{ + URL: "https://grafana.example.com", + Alerts: []string{"uid:r1"}, + From: from, + To: to, + Clock: clock, + Notes: &strings.Builder{}, + }.withDefaults() + + res, err := check(context.Background(), cfg, src) + require.NoError(t, err) + require.Empty(t, res.Violations, "B's post-window onset must not count") + require.Len(t, res.Verdicts, 1) + require.Equal(t, OutcomeRecovered, res.Verdicts[0].Outcome) + require.Contains(t, notesOf(cfg), "recovery wait") + require.Contains(t, notesOf(cfg), `rule "Rule One"`) +} + +// A rule with no recovering instance is not observed past windowEnd: only the +// affected rules pay the recovery wait. +func TestCheckRecoveryObservationPollsOnlyAffectedRules(t *testing.T) { + clock := newVirtualClock(testNow) + from := clock.Now() + to := from.Add(5 * time.Minute) + windowEnd := to.Add(4 * time.Minute) + recoveryAt := to.Add(-10 * time.Second) + resolveAt := recoveryAt.Add(10 * time.Minute) + fireBAt := windowEnd.Add(3 * time.Minute) + + src := newCheckSource(func(title string, _ int) (Observation, error) { + now := clock.Now() + switch title { + case "Alert A": + state, active := StateRecovering, recoveryAt + switch { + case now.Before(recoveryAt): + state, active = StateFiring, from.Add(-time.Hour) + case !now.Before(resolveAt): + state, active = StateNormal, now + } + return recoveryObservation("a", "Alert A", now, 10*time.Minute, + Instance{Labels: map[string]string{"instance": "a"}, State: state, ActiveAt: active}), nil + case "Alert B": + state, active := StateNormal, now + switch { + case !now.Before(fireBAt): + state, active = StateFiring, windowEnd + case !now.Before(windowEnd): + state, active = StatePending, windowEnd + } + return recoveryObservation("b", "Alert B", now, 0, + Instance{Labels: map[string]string{"instance": "b"}, State: state, ActiveAt: active}), nil + } + return Observation{}, fmt.Errorf("unexpected title %q", title) + }) + src.defs = []Definition{ + recoveringScenarioDef("a", "Alert A"), + recoveringScenarioDef("b", "Alert B"), + } + + cfg := Config{ + URL: "https://grafana.example.com", + Alerts: []string{"uid:a", "uid:b"}, + From: from, + To: to, + Clock: clock, + Notes: &strings.Builder{}, + }.withDefaults() + + res, err := check(context.Background(), cfg, src) + require.NoError(t, err) + require.Empty(t, res.Violations) + for _, v := range res.Verdicts { + switch v.RuleUID { + case "a": + require.Equal(t, OutcomeRecovered, v.Outcome) + case "b": + require.Equal(t, OutcomeHealthy, v.Outcome) + default: + t.Fatalf("unexpected verdict %q", v.RuleUID) + } + } + require.Greater(t, src.callCount("Alert A"), src.callCount("Alert B"), + "only the recovering rule is observed past windowEnd") +} + +// A recovering instance that never resolves inside its KeepFiringFor fails +// closed: the run terminates at the deadline and the episode stays open. +func TestCheckRecoveringNeverResolvesFailsClosed(t *testing.T) { + clock := newVirtualClock(testNow) + from := clock.Now() + to := from.Add(5 * time.Minute) + recoveryAt := to.Add(-10 * time.Second) + + src := newCheckSource(func(_ string, _ int) (Observation, error) { + now := clock.Now() + state, active := StateRecovering, recoveryAt + if now.Before(recoveryAt) { + state, active = StateFiring, from.Add(-time.Hour) + } + return recoveryObservation("r1", "Rule One", now, time.Minute, + Instance{Labels: map[string]string{"instance": "a"}, State: state, ActiveAt: active}), nil + }) + src.defs = []Definition{recoveringScenarioDef("r1", "Rule One")} + + cfg := Config{ + URL: "https://grafana.example.com", + Alerts: []string{"uid:r1"}, + From: from, + To: to, + Clock: clock, + Notes: &strings.Builder{}, + }.withDefaults() + + res, err := check(context.Background(), cfg, src) + require.NoError(t, err) + require.Len(t, res.Violations, 1) + require.Equal(t, OutcomeStillFailing, res.Violations[0].Outcome) + require.Equal(t, OutcomeStillFailing, res.Verdicts[0].Outcome) + require.Contains(t, notesOf(cfg), "did not resolve before") +} + +// Recovery runs before the drain: B's drain alone runs past A's deadline, and +// A must still resolve. +func TestCheckRecoveryRunsBeforeTheDrain(t *testing.T) { + clock := newVirtualClock(testNow) + from := clock.Now() + to := from.Add(5 * time.Minute) + windowEnd := to.Add(10 * time.Minute) // transitionGrace = B's 600s interval + recoveryAt := to.Add(-10 * time.Second) + kff := 12 * time.Minute + resolveAt := to.Add(11*time.Minute + 50*time.Second) + catchup := to.Add(15 * time.Minute) // after A's deadline + + src := newCheckSource(func(title string, _ int) (Observation, error) { + now := clock.Now() + switch title { + case "Alert A": + state, active := StateRecovering, recoveryAt + switch { + case now.Before(recoveryAt): + state, active = StateFiring, from.Add(-time.Hour) + case !now.Before(resolveAt): + state, active = StateNormal, now + } + return recoveryObservation("a", "Alert A", now, kff, + Instance{Labels: map[string]string{"instance": "a"}, State: state, ActiveAt: active}), nil + case "Alert B": + lastEval := now + if now.Before(catchup) { + lastEval = windowEnd.Add(-time.Minute) // behind, so the drain waits + if lastEval.After(now) { + lastEval = now + } + } + obs := recoveryObservation("b", "Alert B", now, 0, + Instance{Labels: map[string]string{"instance": "b"}, State: StateNormal, ActiveAt: now}) + obs.Rules[0].LastEvaluation = lastEval + return obs, nil + } + return Observation{}, fmt.Errorf("unexpected title %q", title) + }) + src.defs = []Definition{ + {UID: "a", Title: "Alert A", Folder: "F", Group: "G", IntervalSeconds: 60, NoDataState: "OK", ExecErrState: "OK", Kind: KindGrafanaManaged}, + {UID: "b", Title: "Alert B", Folder: "F", Group: "G", IntervalSeconds: 600, NoDataState: "OK", ExecErrState: "OK", Kind: KindGrafanaManaged}, + } + + cfg := Config{ + URL: "https://grafana.example.com", + Alerts: []string{"uid:a", "uid:b"}, + From: from, + To: to, + Clock: clock, + Notes: &strings.Builder{}, + }.withDefaults() + + res, err := check(context.Background(), cfg, src) + require.NoError(t, err) + require.Empty(t, res.Violations) + for _, v := range res.Verdicts { + switch v.RuleUID { + case "a": + require.Equal(t, OutcomeRecovered, v.Outcome, "A must resolve before its deadline, not expire during the drain") + case "b": + require.Equal(t, OutcomeHealthy, v.Outcome) + default: + t.Fatalf("unexpected verdict %q", v.RuleUID) + } + } + require.False(t, clock.Now().Before(catchup), "the drain still ran past A's deadline") +} + +// afterHookClock runs hook once, on the first wait, so a test can append to a +// log the recovery tail is reading. +type afterHookClock struct { + *virtualClock + hook func() +} + +func (c *afterHookClock) After(d time.Duration) <-chan time.Time { + if c.hook != nil { + h := c.hook + c.hook = nil + h() + } + return c.virtualClock.After(d) +} + +// The recovery tail reads the running recorder's log: the resolution poll is +// appended after the first read and must end the wait. +func TestRecoveryWaitTailSourceReadsTheRunningRecordersLog(t *testing.T) { + path := filepath.Join(t.TempDir(), "log.jsonl") + from := testNow + to := from.Add(5 * time.Minute) + windowEnd := to.Add(checkGrace) + recoveryAt := to.Add(-time.Minute) + resolveAt := windowEnd.Add(checkPollEvery) + key := instanceKey(lbl("a")) + + w, err := NewWriter(path, newFakeClock(resolveAt)) + require.NoError(t, err) + require.NoError(t, w.WriteHeader(Header{ + URL: "https://grafana.example.com", GrafanaVersion: "13.1.0", StartedAt: from.Add(-time.Minute), + Rules: []LoggedRule{{UID: checkUID, Title: checkTitle, IntervalSeconds: 60, PollEverySeconds: checkPollEvery.Seconds()}}, + })) + for at := from; !at.After(windowEnd); at = at.Add(checkPollEvery) { + state, active := StateFiring, from.Add(-time.Hour) + if !at.Before(recoveryAt) { + state, active = StateRecovering, recoveryAt + } + require.NoError(t, w.WritePoll(Poll{ + RuleUID: checkUID, GrafanaNow: at, Found: true, Health: "ok", LastEvaluation: at, + KeepFiringForMS: (10 * time.Minute).Milliseconds(), + Abnormal: []Instance{{Labels: lbl("a"), State: state, ActiveAt: active}}, + })) + } + + tail, err := newLogTailer(path) + require.NoError(t, err) + defer tail.Close() + base, sentinel, err := tail.read() + require.NoError(t, err) + require.Nil(t, sentinel) + + clock := &afterHookClock{virtualClock: newVirtualClock(windowEnd)} + clock.hook = func() { + require.NoError(t, w.WritePoll(Poll{ + RuleUID: checkUID, GrafanaNow: resolveAt, Found: true, Health: "ok", + LastEvaluation: resolveAt, Cleared: []string{key}, + })) + } + cfg := Config{URL: "https://grafana.example.com", Clock: clock, Notes: &strings.Builder{}}.withDefaults() + defs := []Definition{recoveringScenarioDef(checkUID, checkTitle)} + rt := map[string]RuleTimings{checkUID: newRuleTimings(checkPollEvery, 60)} + + out, err := recoveryWait(context.Background(), cfg, defs, rt, base, from, windowEnd, tailRecoverySource(tail), true) + require.NoError(t, err) + require.Len(t, out, 1) + require.Equal(t, resolveAt, out[0].GrafanaNow) + require.Contains(t, out[0].Cleared, key) + require.NoError(t, w.Stop()) +} + +// The prime read must see a Recovering poll written after collection's last +// read: without it the wait is skipped and the episode fails closed for a +// reason the recorder did observe. +func TestRecoveryWaitPrimeSeesAPollWrittenAfterCollection(t *testing.T) { + path := filepath.Join(t.TempDir(), "log.jsonl") + from := testNow + to := from.Add(5 * time.Minute) + windowEnd := to.Add(checkGrace) + recoveryAt := windowEnd // in-window, written after collection's last read + resolveAt := windowEnd.Add(checkPollEvery) + key := instanceKey(lbl("a")) + + w, err := NewWriter(path, newFakeClock(resolveAt)) + require.NoError(t, err) + require.NoError(t, w.WriteHeader(Header{ + URL: "https://grafana.example.com", GrafanaVersion: "13.1.0", StartedAt: from.Add(-time.Minute), + Rules: []LoggedRule{{UID: checkUID, Title: checkTitle, IntervalSeconds: 60, PollEverySeconds: checkPollEvery.Seconds()}}, + })) + // Collection's evidence: Alerting only, ending one cadence short of windowEnd. + for at := from; at.Before(windowEnd); at = at.Add(checkPollEvery) { + require.NoError(t, w.WritePoll(Poll{ + RuleUID: checkUID, GrafanaNow: at, Found: true, Health: "ok", LastEvaluation: at, + Abnormal: []Instance{{Labels: lbl("a"), State: StateFiring, ActiveAt: from.Add(-time.Hour)}}, + })) + } + + tail, err := newLogTailer(path) + require.NoError(t, err) + defer tail.Close() + base, sentinel, err := tail.read() + require.NoError(t, err) + require.Nil(t, sentinel) + + // The Recovering transition lands after collection's last read. + require.NoError(t, w.WritePoll(Poll{ + RuleUID: checkUID, GrafanaNow: recoveryAt, Found: true, Health: "ok", LastEvaluation: recoveryAt, + KeepFiringForMS: (10 * time.Minute).Milliseconds(), + Abnormal: []Instance{{Labels: lbl("a"), State: StateRecovering, ActiveAt: recoveryAt}}, + })) + + clock := &afterHookClock{virtualClock: newVirtualClock(windowEnd)} + clock.hook = func() { + require.NoError(t, w.WritePoll(Poll{ + RuleUID: checkUID, GrafanaNow: resolveAt, Found: true, Health: "ok", + LastEvaluation: resolveAt, Cleared: []string{key}, + })) + } + cfg := Config{URL: "https://grafana.example.com", Clock: clock, Notes: &strings.Builder{}}.withDefaults() + defs := []Definition{recoveringScenarioDef(checkUID, checkTitle)} + rt := map[string]RuleTimings{checkUID: newRuleTimings(checkPollEvery, 60)} + + out, err := recoveryWait(context.Background(), cfg, defs, rt, base, from, windowEnd, tailRecoverySource(tail), true) + require.NoError(t, err) + require.Contains(t, notesOf(cfg), "still recovering", "the primed Recovering poll must select the instance") + require.Len(t, out, 2, "the primed Recovering poll and the resolution") + require.Contains(t, out[1].Cleared, key) + require.NoError(t, w.Stop()) +} + +// Per-instance deadlines: an overdue instance is expired on its own deadline +// while another instance keeps the wait running for its own. +func TestRecoveryWaitExpiresInstancesIndividually(t *testing.T) { + from := testNow + to := from.Add(10 * time.Minute) + windowEnd := to + timings := newRuleTimings(30*time.Second, 60) + kff := time.Minute + aActive := to.Add(-5 * time.Minute) // deadline = to-1m30s: already past + bActive := to.Add(-time.Minute) // deadline = to+2m30s + keyA, keyB := instanceKey(lbl("a")), instanceKey(lbl("b")) + + base := []Poll{{ + RuleUID: checkUID, GrafanaNow: to, Found: true, Health: "ok", LastEvaluation: to, + KeepFiringForMS: kff.Milliseconds(), + Abnormal: []Instance{ + {Labels: lbl("a"), State: StateRecovering, ActiveAt: aActive}, + {Labels: lbl("b"), State: StateRecovering, ActiveAt: bActive}, + }, + }} + clock := newVirtualClock(to) + cfg := Config{URL: "https://grafana.example.com", Clock: clock, Notes: &strings.Builder{}}.withDefaults() + defs := []Definition{{UID: checkUID, Title: checkTitle, IntervalSeconds: 60}} + rt := map[string]RuleTimings{checkUID: timings} + + calls := 0 + src := func(context.Context, map[string]string, []string) ([]Poll, bool, error) { + calls++ + p := Poll{RuleUID: checkUID, GrafanaNow: clock.Now(), Found: true, Health: "ok", LastEvaluation: clock.Now()} + if calls == 1 { + // A's late clear arrives while B is still recovering. + p.Cleared = []string{keyA} + p.Abnormal = []Instance{{Labels: lbl("b"), State: StateRecovering, ActiveAt: bActive}} + } else { + p.Cleared = []string{keyA, keyB} + } + return []Poll{p}, false, nil + } + + out, err := recoveryWait(context.Background(), cfg, defs, rt, base, from, windowEnd, src, false) + require.NoError(t, err) + require.Contains(t, notesOf(cfg), "an instance did not resolve before", "A's own deadline must expire it") + require.GreaterOrEqual(t, calls, 2, "B must keep the wait running after A expires") + require.Len(t, out, 2) +} + +// A recording that already contains the recovery resolution classifies from +// the log alone: no live source is polled. +func TestCheckRecorderModeReadsRecoveryResolutionFromTheLog(t *testing.T) { + from := testNow + to := from.Add(5 * time.Minute) + windowEnd := to.Add(checkGrace) + recoveryAt := to.Add(-time.Minute) + resolveAt := windowEnd.Add(checkPollEvery) + key := instanceKey(lbl("a")) + + path := filepath.Join(t.TempDir(), "log.jsonl") + w, err := NewWriter(path, newFakeClock(resolveAt)) + require.NoError(t, err) + require.NoError(t, w.WriteHeader(Header{ + URL: "https://grafana.example.com", GrafanaVersion: "13.1.0", + StartedAt: from.Add(-time.Minute), ReadyAt: from.Add(-time.Minute), + Rules: []LoggedRule{{ + UID: checkUID, Title: checkTitle, Folder: "F", Group: "G", + IntervalSeconds: 60, NoDataState: "OK", ExecErrState: "OK", + PollEverySeconds: checkPollEvery.Seconds(), + }}, + })) + for at := from; !at.After(resolveAt); at = at.Add(checkPollEvery) { + p := Poll{ + RuleUID: checkUID, GrafanaNow: at, Found: true, State: "inactive", Health: "ok", + LastEvaluation: at, KeepFiringForMS: (10 * time.Minute).Milliseconds(), + } + switch { + case at.Before(recoveryAt): + p.State = "firing" + p.Abnormal = []Instance{{Labels: lbl("a"), State: StateFiring, ActiveAt: from.Add(-time.Hour)}} + case at.Before(resolveAt): + p.Abnormal = []Instance{{Labels: lbl("a"), State: StateRecovering, ActiveAt: recoveryAt}} + default: + p.Cleared = []string{key} + } + require.NoError(t, w.WritePoll(p)) + } + require.NoError(t, w.Stop()) + writePid(t, path+".pid", fmt.Sprintf("%d\n", deadPid(t))) + + clock := newVirtualClock(testNow) + cfg := recorderConfig(t, clock, path) + cfg.NoFailFast = true + // The log proves the evaluations and holds the resolution, so no live poll + // may happen. + src := newCheckSource(func(title string, _ int) (Observation, error) { + require.Fail(t, "unexpected live poll: the log is the evidence") + return Observation{}, errors.New("unexpected poll") + }) + + res, err := check(context.Background(), cfg, src) + require.NoError(t, err) + require.Empty(t, res.Violations) + require.Len(t, res.Verdicts, 1) + require.Equal(t, OutcomeRecovered, res.Verdicts[0].Outcome) +} diff --git a/grafana-alertcheck/internal/gate/classify.go b/grafana-alertcheck/internal/gate/classify.go index ab31edb3b..3f0872fed 100644 --- a/grafana-alertcheck/internal/gate/classify.go +++ b/grafana-alertcheck/internal/gate/classify.go @@ -160,16 +160,21 @@ type episode struct { // the translated ActiveAt against `from`, never by which poll happened to // report it first — a poll's own cadence is not evidence of when the condition // actually began. +// +// recoveryOpen marks an episode last seen Recovering: an extension poll may +// close it only until recoveryDeadline. type instanceTimeline struct { - labels map[string]string - preexisting bool - seen bool - badOpen bool - episodeStart time.Time - lastState State - lastHealth string - lastError string - episodes []episode + labels map[string]string + preexisting bool + seen bool + badOpen bool + recoveryOpen bool + recoveryDeadline time.Time + episodeStart time.Time + lastState State + lastHealth string + lastError string + episodes []episode } // runnerTime translates a Grafana-domain timestamp into the runner domain by @@ -183,9 +188,9 @@ func runnerTime(p Poll, grafanaDomain time.Time) time.Time { // [from, windowEnd] and reduces them to the rule's worst outcome, merged // BadFor, and the Violations the preexisting policy charges against the run. // PURE: no I/O, no clock reads; polls need not be pre-filtered to this rule. -func classifyRule(def Definition, polls []Poll, from, windowEnd time.Time, badStates map[State]bool, pol PreexistingPolicy) (Outcome, time.Duration, []Violation) { +// t bounds a recovering episode, identically to the recovery wait. +func classifyRule(def Definition, t RuleTimings, polls []Poll, from, windowEnd time.Time, badStates map[State]bool, pol PreexistingPolicy) (Outcome, time.Duration, []Violation) { rulePolls := pollsForRule(polls, def.UID) - inWindow := inWindowPolls(rulePolls, from, windowEnd) timelines := make(map[string]*instanceTimeline) order := make([]string, 0) @@ -223,6 +228,8 @@ func classifyRule(def Definition, polls []Poll, from, windowEnd time.Time, badSt } tl.episodes = append(tl.episodes, episode{start: tl.episodeStart, end: end, closedByRealClear: isReal}) tl.badOpen = false + tl.recoveryOpen = false // resolved; a later re-fire is out of scope + tl.recoveryDeadline = time.Time{} } // onsetOf resolves a fresh episode's start: the instance's own ActiveAt, // translated to the runner domain by this poll's skew, clamped to @@ -237,41 +244,89 @@ func classifyRule(def Definition, polls []Poll, from, windowEnd time.Time, badSt } return start } + // openRecovering: Recovering is only reachable from Alerting, so the fire + // predates this observation; ActiveAt is the recovery onset, not the fire. + openRecovering := func(tl *instanceTimeline) { + tl.preexisting = true + openEpisode(tl, from) + } + + // recoveryExpired: an extension clear past the deadline does not resolve. + recoveryExpired := func(tl *instanceTimeline, p Poll) bool { + return tl.recoveryDeadline.IsZero() || runnerTime(p, p.GrafanaNow).After(tl.recoveryDeadline) + } + + for _, p := range rulePolls { + if pollBeforeWindow(p, from) { + continue + } + extension := pollAfterWindow(p, windowEnd) + if extension && (p.IsPaused || !p.Found) { + // Not a resolution: retire eligibility, leave the episode open. + for _, tl := range timelines { + tl.recoveryOpen = false + } + continue + } - for _, p := range inWindow { byKey := make(map[string]Instance, len(p.Abnormal)) for _, inst := range p.Abnormal { byKey[instanceKey(inst.Labels)] = inst } for key, inst := range byKey { - tl := get(key, inst.Labels) + tl, known := timelines[key] + if extension && (!known || !tl.recoveryOpen) { + continue // post-window onset: not ours to classify + } + tl = get(key, inst.Labels) // backfills labels on an existing bare timeline bad := badStates[inst.State] switch { case !tl.seen: tl.seen = true if bad { - // Fail-closed: "preexisting" only when even the worst-case - // skew error places the onset at or before `from`; an onset - // that might be in-window must classify as a new episode. - activeAtRunner := runnerTime(p, inst.ActiveAt) - tl.preexisting = !activeAtRunner.Add(p.SkewBound()).After(from) - if tl.preexisting { - openEpisode(tl, from) + if inst.State == StateRecovering { + openRecovering(tl) } else { - openEpisode(tl, onsetOf(p, inst)) + // Fail-closed: "preexisting" only when even the + // worst-case skew error places the onset at or before + // `from`; an onset that might be in-window must + // classify as a new episode. + activeAtRunner := runnerTime(p, inst.ActiveAt) + tl.preexisting = !activeAtRunner.Add(p.SkewBound()).After(from) + if tl.preexisting { + openEpisode(tl, from) + } else { + openEpisode(tl, onsetOf(p, inst)) + } } } case bad && !tl.badOpen: + // Already observed in-window and not bad: a Recovering here + // means the fire happened in-window (a missed Alerting poll + // hides it), so it must never be preexisting. openEpisode(tl, onsetOf(p, inst)) case !bad && tl.badOpen: + if extension && recoveryExpired(tl, p) { + continue // a clear past the deadline does not resolve + } closeEpisode(tl, runnerTime(p, p.GrafanaNow), true) } + if bad && inst.State == StateRecovering { + tl.recoveryOpen = true + tl.recoveryDeadline = runnerTime(p, inst.ActiveAt).Add(p.KeepFiringFor() + t.evalStaleAfter + t.pollEvery) + } tl.lastState, tl.lastHealth, tl.lastError = inst.State, p.Health, p.LastError } for _, key := range p.Cleared { - tl := get(key, nil) + tl, known := timelines[key] + if extension && (!known || !tl.recoveryOpen) { + continue + } + if !known { + tl = get(key, nil) + } if !tl.seen { // Cleared on first mention: the transition happened pre-window, // with no in-window evidence it was ever bad. @@ -279,7 +334,9 @@ func classifyRule(def Definition, polls []Poll, from, windowEnd time.Time, badSt continue } if tl.badOpen { - closeEpisode(tl, runnerTime(p, p.GrafanaNow), true) + if !extension || !recoveryExpired(tl, p) { + closeEpisode(tl, runnerTime(p, p.GrafanaNow), true) + } } tl.lastHealth, tl.lastError = p.Health, p.LastError } @@ -287,7 +344,13 @@ func classifyRule(def Definition, polls []Poll, from, windowEnd time.Time, badSt // Vanished is a deliberate no-op: freeze badOpen/preexisting as-is, so a // vanish while bad stays bad (never reading as a recovery). for _, key := range p.Vanished { - tl := get(key, nil) + tl, known := timelines[key] + if extension && (!known || !tl.recoveryOpen) { + continue + } + if !known { + tl = get(key, nil) + } tl.seen = true tl.lastHealth = p.Health } @@ -445,13 +508,15 @@ func pollsForRule(polls []Poll, uid string) []Poll { return out } -// badStateSet turns Policy.States into a lookup set, defaulting to {firing} -// when the caller leaves States empty — decide applies the default itself so a -// test can pass a zero-value Policy and get the real default, rather than -// depending on the CLI to have filled it in. +// badStateSet turns Policy.States into a lookup set, defaulting to +// {firing, recovering} when the caller leaves States empty — decide applies +// the default itself so a test can pass a zero-value Policy and get the real +// default, rather than depending on the CLI to have filled it in. Recovering +// is bad by default because the instance is still firing until its +// KeepFiringFor elapses; an explicit --states firing opts out. func badStateSet(states []State) map[State]bool { if len(states) == 0 { - states = []State{StateFiring} + states = []State{StateFiring, StateRecovering} } set := make(map[State]bool, len(states)) for _, s := range states { @@ -573,7 +638,7 @@ func decide(h Header, polls []Poll, sentinel *time.Time, defs []Definition, EvalStaleAfter: t.evalStaleAfter, } - outcome, badFor, viols := classifyRule(def, polls, pol.From, windowEnd, badStates, pol.Preexisting) + outcome, badFor, viols := classifyRule(def, t, polls, pol.From, windowEnd, badStates, pol.Preexisting) if cov.Unobservable { outcome = OutcomeNotVerified anyUnobservable = true diff --git a/grafana-alertcheck/internal/gate/classify_test.go b/grafana-alertcheck/internal/gate/classify_test.go index 14cc95248..73ba50ff9 100644 --- a/grafana-alertcheck/internal/gate/classify_test.go +++ b/grafana-alertcheck/internal/gate/classify_test.go @@ -35,7 +35,17 @@ func quietPoll(uid string, at time.Time) Poll { return Poll{RuleUID: uid, GrafanaNow: at, Found: true, Health: "ok", LastEvaluation: at} } -var defaultBad = badStateSet(nil) // {firing} +var defaultBad = badStateSet(nil) // {firing, recovering} + +// classify is classifyRule with default timings (a 60s evaluation interval), +// for tests that do not exercise the recovery deadline margin directly. +func classify(def Definition, polls []Poll, from, to time.Time, badStates map[State]bool, pol PreexistingPolicy) (Outcome, time.Duration, []Violation) { + interval := def.IntervalSeconds + if interval == 0 { + interval = 60 + } + return classifyRule(def, newRuleTimings(defaultPollEvery(interval), interval), polls, from, to, badStates, pol) +} // pausedHeader builds the header decide reads `skipped` from: the pause state // as of record start. Definition.IsPaused is deliberately NOT that authority @@ -57,7 +67,7 @@ func TestClassifyRule_NoEvidenceIsClean(t *testing.T) { def := Definition{UID: "r1", Title: "R1"} polls := []Poll{quietPoll("r1", from), quietPoll("r1", to)} - outcome, badFor, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + outcome, badFor, viols := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) require.Equal(t, OutcomeHealthy, outcome) require.Zero(t, badFor) require.Empty(t, viols) @@ -74,7 +84,7 @@ func TestClassifyRule_NewOnsetInsideWindowIsNewlyBad(t *testing.T) { abnormalPoll("r1", onset, StateFiring, lbl("a"), onset), abnormalPoll("r1", to, StateFiring, lbl("a"), onset), } - outcome, badFor, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + outcome, badFor, viols := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) require.Equal(t, OutcomeNewFailure, outcome) require.Equal(t, to.Sub(onset), badFor) require.Len(t, viols, 1) @@ -96,7 +106,7 @@ func TestClassifyRule_NewOnsetThatClearsStillFails(t *testing.T) { clearedPoll("r1", clearAt, instanceKey(lbl("a"))), quietPoll("r1", to), } - outcome, _, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + outcome, _, viols := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) require.Equal(t, OutcomeNewFailure, outcome, "even though it cleared") require.Len(t, viols, 1) } @@ -115,7 +125,7 @@ func TestClassifyRule_PreexistingThatRecoversIsRecoveredAndNotAViolation(t *test clearedPoll("r1", clearAt, key), quietPoll("r1", to), } - outcome, badFor, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + outcome, badFor, viols := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) require.Equal(t, OutcomeRecovered, outcome) require.Equal(t, clearAt.Sub(from), badFor) require.Empty(t, viols, "default policy passes a recovered preexisting instance") @@ -136,7 +146,7 @@ func TestClassifyRule_LateRecoveryPassesRegardlessOfHowLateItIs(t *testing.T) { clearedPoll("r1", clearAt, key), quietPoll("r1", to), } - outcome, badFor, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + outcome, badFor, viols := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) require.Equal(t, OutcomeRecovered, outcome, "even 58 minutes into a 60-minute window") require.Equal(t, clearAt.Sub(from), badFor, "not a value clamped against a deadline") require.Empty(t, viols, "there is no deadline a preexisting recovery must beat") @@ -151,7 +161,7 @@ func TestClassifyRule_PreexistingStillBadAtWindowEndIsPersistentlyBad(t *testing abnormalPoll("r1", from, StateFiring, lbl("a"), from.Add(-time.Hour)), abnormalPoll("r1", to, StateFiring, lbl("a"), from.Add(-time.Hour)), } - outcome, badFor, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + outcome, badFor, viols := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) require.Equal(t, OutcomeStillFailing, outcome) require.Equal(t, to.Sub(from), badFor) require.Len(t, viols, 1) @@ -172,7 +182,7 @@ func TestClassifyRule_ClearThenBadAgainIsFlapping(t *testing.T) { abnormalPoll("r1", from.Add(5*time.Minute), StateFiring, lbl("a"), from.Add(5*time.Minute)), quietPoll("r1", to), } - outcome, _, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + outcome, _, viols := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) require.Equal(t, OutcomeUnstable, outcome) require.Len(t, viols, 1) require.Equal(t, OutcomeUnstable, viols[0].Outcome, "always a fail regardless of policy") @@ -206,7 +216,7 @@ func TestClassifyRule_FlappingAtEveryTimingOfTheSecondOnset(t *testing.T) { abnormalPoll("r1", tc.secondOnset, StateFiring, lbl("a"), tc.secondOnset), quietPoll("r1", to), } - outcome, _, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + outcome, _, viols := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) require.Equalf(t, OutcomeUnstable, outcome, "second onset at %s", tc.secondOnset) require.Len(t, viols, 1) require.Equal(t, OutcomeUnstable, viols[0].Outcome) @@ -227,7 +237,7 @@ func TestClassifyRule_VanishedWhileBadStaysPersistentlyBad(t *testing.T) { vanishedPoll("r1", from.Add(5*time.Minute), key), quietPoll("r1", to), } - outcome, badFor, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + outcome, badFor, viols := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) require.Equal(t, OutcomeStillFailing, outcome, "a vanish must never read as a recovery") require.Equal(t, to.Sub(from), badFor, "the freeze must hold the episode open to windowEnd") require.Len(t, viols, 1) @@ -246,7 +256,7 @@ func TestClassifyRule_VanishedWhileNeverBadIsUninteresting(t *testing.T) { vanishedPoll("r1", from.Add(5*time.Minute), key), quietPoll("r1", to), } - outcome, badFor, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + outcome, badFor, viols := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) require.Equal(t, OutcomeHealthy, outcome) require.Zero(t, badFor) require.Empty(t, viols) @@ -265,7 +275,7 @@ func TestClassifyRule_PreexistingPolicyFailFailsARecoveredInstance(t *testing.T) clearedPoll("r1", from.Add(2*time.Minute), key), quietPoll("r1", to), } - outcome, _, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFail) + outcome, _, viols := classify(def, polls, from, to, defaultBad, PreexistingFail) require.Equal(t, OutcomeRecovered, outcome, "the descriptive outcome does not change under policy=fail") require.Len(t, viols, 1) require.Equal(t, OutcomeRecovered, viols[0].Outcome, @@ -281,7 +291,7 @@ func TestClassifyRule_PreexistingPolicyIgnoreForgivesPersistentlyBad(t *testing. abnormalPoll("r1", from, StateFiring, lbl("a"), from.Add(-time.Hour)), abnormalPoll("r1", to, StateFiring, lbl("a"), from.Add(-time.Hour)), } - outcome, _, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingIgnore) + outcome, _, viols := classify(def, polls, from, to, defaultBad, PreexistingIgnore) require.Equal(t, OutcomeStillFailing, outcome, "the descriptive outcome does not change under policy=ignore") require.Empty(t, viols, "policy=ignore disregards a preexisting instance even if it never recovers") } @@ -297,7 +307,7 @@ func TestClassifyRule_PreexistingPolicyIgnoreStillFailsANewOnset(t *testing.T) { abnormalPoll("r1", onset, StateFiring, lbl("a"), onset), abnormalPoll("r1", to, StateFiring, lbl("a"), onset), } - outcome, _, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingIgnore) + outcome, _, viols := classify(def, polls, from, to, defaultBad, PreexistingIgnore) require.Equal(t, OutcomeNewFailure, outcome) require.Len(t, viols, 1, "ignore only forgives PREEXISTING badness") } @@ -323,7 +333,7 @@ func TestClassifyRule_WorstOfMultipleInstancesWins(t *testing.T) { Abnormal: []Instance{{Labels: lbl("b"), State: StateFiring, ActiveAt: from}}, }, } - outcome, _, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + outcome, _, viols := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) require.Equal(t, OutcomeStillFailing, outcome, "the worse of {recovered, still_failing}") require.Len(t, viols, 1) require.Equal(t, OutcomeStillFailing, viols[0].Outcome) @@ -758,7 +768,7 @@ func TestClassifyRule_OnsetBetweenFromAndFirstPollIsNewlyBadNotRecovered(t *test clearedPoll("r1", clearAt, key), quietPoll("r1", to), } - outcome, badFor, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + outcome, badFor, viols := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) require.Equal(t, OutcomeNewFailure, outcome, "the onset is after `from`, so it is not preexisting even though the FIRST in-window poll already observes it bad") require.Len(t, viols, 1) @@ -783,7 +793,7 @@ func TestClassifyRule_OnsetJustBeforeFromIsPreexisting(t *testing.T) { clearedPoll("r1", clearAt, key), quietPoll("r1", to), } - outcome, badFor, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + outcome, badFor, viols := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) require.Equal(t, OutcomeRecovered, outcome, "the onset is at/before `from`, genuinely preexisting") require.Empty(t, viols, "default policy passes a recovered preexisting instance") require.Equal(t, clearAt.Sub(from), badFor, @@ -814,7 +824,7 @@ func TestClassifyRule_SkewTranslatesActiveAtAcrossTheWindowBoundary(t *testing.T stillBad.GrafanaNow = to.Add(skew) stillBad.LastEvaluation = to.Add(skew) - outcome, badFor, _ := classifyRule(def, []Poll{poll, stillBad}, from, to, defaultBad, PreexistingFailUnlessRecovered) + outcome, badFor, _ := classify(def, []Poll{poll, stillBad}, from, to, defaultBad, PreexistingFailUnlessRecovered) require.Equal(t, OutcomeStillFailing, outcome, "a +90s skew must translate ActiveAt back to exactly `from`") require.Equal(t, to.Sub(from), badFor) @@ -837,7 +847,7 @@ func TestClassifyRule_LabelsSurviveWhenTimelineStartsFromAClearedMarker(t *testi abnormalPoll("r1", newOnset, StateFiring, lbl("a"), newOnset), quietPoll("r1", to), } - _, _, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + _, _, viols := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) require.Len(t, viols, 1) require.NotNil(t, viols[0].InstanceLabels) require.Equal(t, "a", viols[0].InstanceLabels["instance"], @@ -858,7 +868,7 @@ func TestClassifyRule_ViolationFieldsArePrecise(t *testing.T) { clearedPoll("r1", clearAt, instanceKey(lbl("a"))), quietPoll("r1", to), } - _, _, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + _, _, viols := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) require.Len(t, viols, 1) v := viols[0] require.True(t, v.FirstSeen.Equal(onset)) @@ -885,7 +895,7 @@ func TestClassifyRule_ClearedEventPastWindowEndClampsToWindowEnd(t *testing.T) { SkewBoundMS: bound.Milliseconds(), Cleared: []string{key}, }, } - outcome, badFor, _ := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + outcome, badFor, _ := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) require.Equal(t, OutcomeRecovered, outcome) require.Equal(t, to.Sub(from), badFor, "the episode end must clamp to windowEnd, not extend to the late Cleared event's raw time") @@ -915,7 +925,7 @@ func TestClassifyRule_OnsetJustPastWindowEndIsNewlyBadNotClean(t *testing.T) { Abnormal: []Instance{{Labels: lbl("a"), State: StateFiring, ActiveAt: windowEnd.Add(10 * time.Second)}}, } - outcome, badFor, viols := classifyRule(def, []Poll{poll}, from, windowEnd, defaultBad, PreexistingFailUnlessRecovered) + outcome, badFor, viols := classify(def, []Poll{poll}, from, windowEnd, defaultBad, PreexistingFailUnlessRecovered) require.Equal(t, OutcomeNewFailure, outcome, "an onset past windowEnd seen only via the skew bound must fail closed") require.Zero(t, badFor, "the zero-length episode must truncate to the window end") @@ -937,7 +947,7 @@ func TestClassifyRule_ClearAfterWindowEndIsPersistentlyBad(t *testing.T) { abnormalPoll("r1", from, StateFiring, lbl("a"), from.Add(-time.Hour)), clearedPoll("r1", to.Add(time.Hour), key), // far past `to`, not a boundary case } - outcome, badFor, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + outcome, badFor, viols := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) require.Equal(t, OutcomeStillFailing, outcome, "a clear outside the window must not read as a recovery") require.Equal(t, to.Sub(from), badFor) require.Len(t, viols, 1) @@ -967,7 +977,7 @@ func TestClassifyRule_CloseBeforeOpenClampsToZeroNotNegative(t *testing.T) { }, quietPoll("r1", to), } - outcome, badFor, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + outcome, badFor, viols := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) require.Equal(t, OutcomeNewFailure, outcome) require.GreaterOrEqual(t, badFor, time.Duration(0), "a non-negative duration even though the closing poll's translated time landed before the opening poll's") @@ -975,6 +985,250 @@ func TestClassifyRule_CloseBeforeOpenClampsToZeroNotNegative(t *testing.T) { require.Len(t, viols, 1) } +// --- Recovering (keep firing for) --- + +// recoveringPoll is abnormalPoll with the rule's recovery period attached. +func recoveringPoll(uid string, at time.Time, labels map[string]string, activeAt time.Time, kff time.Duration) Poll { + p := abnormalPoll(uid, at, StateRecovering, labels, activeAt) + p.KeepFiringForMS = kff.Milliseconds() + return p +} + +// A first-seen Recovering instance is never a new failure: Recovering is only +// reachable from Alerting, so the fire predates the observation and its +// ActiveAt (the recovery onset) must not start an in-window episode. +func TestClassifyRule_FirstSeenRecoveringIsPreexistingAndRecovers(t *testing.T) { + from := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) + to := from.Add(10 * time.Minute) + recoveryAt := from.Add(2 * time.Minute) + def := Definition{UID: "r1", Title: "R1"} + + polls := []Poll{ + recoveringPoll("r1", recoveryAt, lbl("a"), recoveryAt, 10*time.Minute), + clearedPoll("r1", to.Add(2*time.Minute), instanceKey(lbl("a"))), // extension poll + } + outcome, _, viols := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + require.Equal(t, OutcomeRecovered, outcome) + require.Empty(t, viols) +} + +// A Recovering instance first observed after an in-window non-bad observation +// is an in-window fire whose Alerting poll was missed: it must fail closed as a +// new failure, never read as preexisting (which would let it pass as recovered). +func TestClassifyRule_RecoveringAfterInWindowPendingIsNewFailure(t *testing.T) { + from := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) + to := from.Add(10 * time.Minute) + pendingAt := from.Add(2 * time.Minute) + recoveryAt := from.Add(3 * time.Minute) + def := Definition{UID: "r1", Title: "R1"} + + polls := []Poll{ + abnormalPoll("r1", pendingAt, StatePending, lbl("a"), pendingAt), + recoveringPoll("r1", recoveryAt, lbl("a"), recoveryAt, 10*time.Minute), + clearedPoll("r1", to.Add(2*time.Minute), instanceKey(lbl("a"))), // extension poll + } + outcome, _, viols := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + require.Equal(t, OutcomeNewFailure, outcome) + require.Len(t, viols, 1) + require.Equal(t, OutcomeNewFailure, viols[0].Outcome) +} + +// A recovering episode that never resolves inside its KeepFiringFor stays open +// and fails closed. +func TestClassifyRule_UnresolvedRecoveringIsStillFailing(t *testing.T) { + from := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) + to := from.Add(10 * time.Minute) + recoveryAt := to.Add(-time.Minute) + def := Definition{UID: "r1", Title: "R1"} + + polls := []Poll{ + recoveringPoll("r1", recoveryAt, lbl("a"), recoveryAt, 10*time.Minute), + } + outcome, _, viols := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + require.Equal(t, OutcomeStillFailing, outcome) + require.Len(t, viols, 1) + require.Equal(t, OutcomeStillFailing, viols[0].Outcome) +} + +// An extension poll may only close the recovering episode: a different +// instance going bad after windowEnd is outside the window and stays healthy. +func TestClassifyRule_ExtensionPollDoesNotOpenNewOnset(t *testing.T) { + from := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) + to := from.Add(10 * time.Minute) + recoveryAt := to.Add(-time.Minute) + def := Definition{UID: "r1", Title: "R1"} + + polls := []Poll{ + abnormalPoll("r1", from, StateFiring, lbl("a"), from.Add(-time.Hour)), // A preexisting + recoveringPoll("r1", recoveryAt, lbl("a"), recoveryAt, 10*time.Minute), + abnormalPoll("r1", to.Add(time.Minute), StateFiring, lbl("b"), to.Add(time.Minute)), // B post-window + clearedPoll("r1", to.Add(2*time.Minute), instanceKey(lbl("a"))), + } + outcome, _, viols := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + require.Equal(t, OutcomeRecovered, outcome, "B's post-window onset must not count") + require.Empty(t, viols) +} + +// A re-fire while the linger is still open is the same episode, not a second +// one: the episode stays open and reads still_failing. +func TestClassifyRule_RefireDuringRecoveryStaysOneEpisode(t *testing.T) { + from := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) + to := from.Add(10 * time.Minute) + recoveryAt := to.Add(-time.Minute) + def := Definition{UID: "r1", Title: "R1"} + + polls := []Poll{ + abnormalPoll("r1", from, StateFiring, lbl("a"), from.Add(-time.Hour)), + recoveringPoll("r1", recoveryAt, lbl("a"), recoveryAt, 10*time.Minute), + abnormalPoll("r1", to.Add(time.Minute), StateFiring, lbl("a"), to.Add(time.Minute)), + } + outcome, _, viols := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + require.Equal(t, OutcomeStillFailing, outcome) + require.Len(t, viols, 1) +} + +// A NEW onset that later enters Recovering still fails even though the +// extension observes its resolution. +func TestClassifyRule_NewOnsetRecoveringInExtensionStillNewFailure(t *testing.T) { + from := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) + to := from.Add(10 * time.Minute) + onset := from.Add(2 * time.Minute) + recoveryAt := to.Add(-time.Minute) + def := Definition{UID: "r1", Title: "R1"} + + polls := []Poll{ + abnormalPoll("r1", onset, StateFiring, lbl("a"), onset), + recoveringPoll("r1", recoveryAt, lbl("a"), recoveryAt, 10*time.Minute), + clearedPoll("r1", to.Add(2*time.Minute), instanceKey(lbl("a"))), + } + outcome, _, viols := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + require.Equal(t, OutcomeNewFailure, outcome) + require.Len(t, viols, 1) +} + +// Excluding recovering from --states opts out: it closes an open episode at +// the recovery onset, like any other non-bad state. +func TestClassifyRule_RecoveringExcludedClosesEpisode(t *testing.T) { + from := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) + to := from.Add(10 * time.Minute) + recoveryAt := from.Add(8 * time.Minute) + def := Definition{UID: "r1", Title: "R1"} + + polls := []Poll{ + abnormalPoll("r1", from, StateFiring, lbl("a"), from.Add(-time.Hour)), + recoveringPoll("r1", recoveryAt, lbl("a"), recoveryAt, 10*time.Minute), + } + outcome, badFor, viols := classify(def, polls, from, to, badStateSet([]State{StateFiring}), PreexistingFailUnlessRecovered) + require.Equal(t, OutcomeRecovered, outcome) + require.Equal(t, recoveryAt.Sub(from), badFor) + require.Empty(t, viols) +} + +// A paused extension poll must not resolve the episode: a pause is not a +// recovery, so the episode stays open and fails closed. +func TestClassifyRule_PausedExtensionPollDoesNotResolveRecovering(t *testing.T) { + from := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) + to := from.Add(10 * time.Minute) + recoveryAt := to.Add(-time.Minute) + def := Definition{UID: "r1", Title: "R1"} + + pausedClear := clearedPoll("r1", to.Add(2*time.Minute), instanceKey(lbl("a"))) + pausedClear.IsPaused = true + + polls := []Poll{ + abnormalPoll("r1", from, StateFiring, lbl("a"), from.Add(-time.Hour)), + recoveringPoll("r1", recoveryAt, lbl("a"), recoveryAt, 10*time.Minute), + pausedClear, + } + outcome, _, viols := classify(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + require.Equal(t, OutcomeStillFailing, outcome) + require.Len(t, viols, 1) +} + +// A clear observed after the instance's own recovery deadline must not resolve +// the episode: the deadline is the bound the wait gave up at. +func TestClassifyRule_ClearPastRecoveryDeadlineStaysStillFailing(t *testing.T) { + from := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) + to := from.Add(10 * time.Minute) + recoveryAt := to.Add(-time.Minute) + kff := time.Minute + timings := newRuleTimings(30*time.Second, 60) + def := Definition{UID: "r1", Title: "R1", IntervalSeconds: 60} + deadline := recoveryAt.Add(kff + timings.evalStaleAfter + timings.pollEvery) + + polls := []Poll{ + abnormalPoll("r1", from, StateFiring, lbl("a"), from.Add(-time.Hour)), + recoveringPoll("r1", recoveryAt, lbl("a"), recoveryAt, kff), + clearedPoll("r1", deadline.Add(time.Minute), instanceKey(lbl("a"))), + } + outcome, _, viols := classifyRule(def, timings, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + require.Equal(t, OutcomeStillFailing, outcome) + require.Len(t, viols, 1) +} + +// The same clear one minute inside the deadline is a genuine recovery. +func TestClassifyRule_ClearWithinRecoveryDeadlineIsRecovered(t *testing.T) { + from := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) + to := from.Add(10 * time.Minute) + recoveryAt := to.Add(-time.Minute) + kff := time.Minute + timings := newRuleTimings(30*time.Second, 60) + def := Definition{UID: "r1", Title: "R1", IntervalSeconds: 60} + deadline := recoveryAt.Add(kff + timings.evalStaleAfter + timings.pollEvery) + + polls := []Poll{ + abnormalPoll("r1", from, StateFiring, lbl("a"), from.Add(-time.Hour)), + recoveringPoll("r1", recoveryAt, lbl("a"), recoveryAt, kff), + clearedPoll("r1", deadline.Add(-time.Minute), instanceKey(lbl("a"))), + } + outcome, _, viols := classifyRule(def, timings, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + require.Equal(t, OutcomeRecovered, outcome) + require.Empty(t, viols) +} + +// A pause during the recovery observation retires eligibility: a later clear, +// even an unpaused one, must not resolve the episode. +func TestClassifyRule_PauseDuringRecoveryRetiresEligibility(t *testing.T) { + from := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) + to := from.Add(10 * time.Minute) + recoveryAt := to.Add(-time.Minute) + def := Definition{UID: "r1", Title: "R1"} + paused := Poll{ + RuleUID: "r1", GrafanaNow: to.Add(time.Minute), Found: true, Health: "ok", + LastEvaluation: to.Add(time.Minute), IsPaused: true, + } + + polls := []Poll{ + abnormalPoll("r1", from, StateFiring, lbl("a"), from.Add(-time.Hour)), + recoveringPoll("r1", recoveryAt, lbl("a"), recoveryAt, 10*time.Minute), + paused, + clearedPoll("r1", to.Add(2*time.Minute), instanceKey(lbl("a"))), + } + outcome, _, viols := classifyRule(def, newRuleTimings(30*time.Second, 60), polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + require.Equal(t, OutcomeStillFailing, outcome) + require.Len(t, viols, 1) +} + +// An absent rule during the recovery observation retires eligibility the same +// way a pause does. +func TestClassifyRule_AbsentRuleDuringRecoveryRetiresEligibility(t *testing.T) { + from := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) + to := from.Add(10 * time.Minute) + recoveryAt := to.Add(-time.Minute) + def := Definition{UID: "r1", Title: "R1"} + absent := Poll{RuleUID: "r1", GrafanaNow: to.Add(time.Minute), Found: false} + + polls := []Poll{ + abnormalPoll("r1", from, StateFiring, lbl("a"), from.Add(-time.Hour)), + recoveringPoll("r1", recoveryAt, lbl("a"), recoveryAt, 10*time.Minute), + absent, + clearedPoll("r1", to.Add(2*time.Minute), instanceKey(lbl("a"))), + } + outcome, _, viols := classifyRule(def, newRuleTimings(30*time.Second, 60), polls, from, to, defaultBad, PreexistingFailUnlessRecovered) + require.Equal(t, OutcomeStillFailing, outcome) + require.Len(t, viols, 1) +} + // --- mergeDurations --- func TestMergeDurations_OverlappingEpisodesCountOnce(t *testing.T) { diff --git a/grafana-alertcheck/internal/gate/coverage.go b/grafana-alertcheck/internal/gate/coverage.go index 977c45c9e..afdf00d36 100644 --- a/grafana-alertcheck/internal/gate/coverage.go +++ b/grafana-alertcheck/internal/gate/coverage.go @@ -242,16 +242,31 @@ func proveCoverage(h Header, polls []Poll, sentinel *time.Time, t RuleTimings, d func inWindowPolls(polls []Poll, from, windowEnd time.Time) []Poll { var out []Poll for _, p := range polls { - bound := p.SkewBound() - runner := runnerTime(p, p.GrafanaNow) - if runner.Before(from.Add(-bound)) || runner.After(windowEnd.Add(bound)) { - continue + if pollInWindow(p, from, windowEnd) { + out = append(out, p) } - out = append(out, p) } return out } +// pollInWindow is the one membership test for [from, windowEnd], with the +// boundary widened by the poll's own skew bound. +func pollInWindow(p Poll, from, windowEnd time.Time) bool { + return !pollBeforeWindow(p, from) && !pollAfterWindow(p, windowEnd) +} + +// pollBeforeWindow reports a poll whose translated time is before `from` +// (widened by its bound): never classified. +func pollBeforeWindow(p Poll, from time.Time) bool { + return runnerTime(p, p.GrafanaNow).Before(from.Add(-p.SkewBound())) +} + +// pollAfterWindow reports a poll whose translated time is after windowEnd +// (widened by its bound): an extension poll. +func pollAfterWindow(p Poll, windowEnd time.Time) bool { + return runnerTime(p, p.GrafanaNow).After(windowEnd.Add(p.SkewBound())) +} + // ruleHeartbeatGap finds the largest unobserved span inside [from, windowEnd], // including the two boundary segments — which is why data at both ends with a // hole in the middle still fails. in must be filtered (inWindowPolls) and diff --git a/grafana-alertcheck/internal/gate/log.go b/grafana-alertcheck/internal/gate/log.go index 68854e359..b0ef8c1a4 100644 --- a/grafana-alertcheck/internal/gate/log.go +++ b/grafana-alertcheck/internal/gate/log.go @@ -121,6 +121,9 @@ type Poll struct { SkewMS int64 `json:"skew_ms"` SkewBoundMS int64 `json:"skew_bound_ms"` LatencyMS int64 `json:"latency_ms"` + // KeepFiringForMS is the rule's recovery period as reported by this + // response; 0/absent means no instance can be Recovering. + KeepFiringForMS int64 `json:"keep_firing_for_ms,omitempty"` // Found false means an authoritative 2xx in which this rule was absent — // never a transport failure, which the transport retries and never turns // into a Poll. The coverage proof turns it into unobservable. @@ -164,6 +167,11 @@ func (p Poll) SkewBound() time.Duration { return time.Duration(p.SkewBoundMS) * // Latency is the wall time this poll's request took, feeding the budget check. func (p Poll) Latency() time.Duration { return time.Duration(p.LatencyMS) * time.Millisecond } +// KeepFiringFor is the rule's recovery period as reported by this poll. +func (p Poll) KeepFiringFor() time.Duration { + return time.Duration(p.KeepFiringForMS) * time.Millisecond +} + // Reducer turns each Observation into the single Poll record that goes into // the log. It holds the previous poll's abnormal instance keys per rule, which // is all the state the transition markers need. @@ -219,6 +227,7 @@ func (r *Reducer) Reduce(uid string, obs Observation) Poll { p.LastEvaluation = rule.LastEvaluation p.IsPaused = rule.IsPaused p.Histogram = rule.Totals + p.KeepFiringForMS = rule.KeepFiringFor.Milliseconds() // present indexes every instance in THIS response, normal ones included — // the markers below must resolve a departed key against the same response, diff --git a/grafana-alertcheck/internal/gate/parse_state.go b/grafana-alertcheck/internal/gate/parse_state.go index 7c35c7e7d..e4e6e98b5 100644 --- a/grafana-alertcheck/internal/gate/parse_state.go +++ b/grafana-alertcheck/internal/gate/parse_state.go @@ -13,11 +13,12 @@ import ( type State string const ( - StateNormal State = "normal" - StateFiring State = "firing" - StatePending State = "pending" - StateNodata State = "nodata" - StateError State = "error" + StateNormal State = "normal" + StateFiring State = "firing" + StatePending State = "pending" + StateNodata State = "nodata" + StateError State = "error" + StateRecovering State = "recovering" ) // Instance is one entry of a rule's alerts[]. State is always canonical; Reason @@ -41,6 +42,9 @@ type Instance struct { type StateRule struct { UID, Title, Folder, Group string Interval time.Duration + // KeepFiringFor is the rule's recovery period in seconds; 0/absent means + // no instance can be Recovering. Carried onto the Poll. + KeepFiringFor time.Duration // State and Health are raw, lowercase, and reporting-only — never // classified. State in particular is never normalized. State, Health string @@ -159,6 +163,13 @@ func parseStateRule(raw json.RawMessage, folder, group string, interval time.Dur return StateRule{}, fmt.Errorf("rule %q: %w", uid, err) } + // Optional; absent means 0 (no recovery linger). + var keepFiringForSeconds float64 + if err := opt(m, "keepFiringFor", &keepFiringForSeconds); err != nil { + return StateRule{}, fmt.Errorf("rule %q: %w", uid, err) + } + r.KeepFiringFor = time.Duration(keepFiringForSeconds * float64(time.Second)) + if err := opt(m, "totals", &r.Totals); err != nil { return StateRule{}, fmt.Errorf("rule %q: %w", uid, err) } @@ -222,15 +233,17 @@ func parseInstance(raw json.RawMessage) (Instance, error) { return inst, nil } -// baseInstanceStates is the strict 5-value allowlist for the base of an -// instance state. Anything else — including an unrecognized base inside a -// "Base (Reason)" composite — is a parse error. +// baseInstanceStates is the strict allowlist for the base of an instance +// state; anything else is a parse error. Recovering is Grafana's "keep firing +// for" state — condition cleared, linger not elapsed, only reachable from +// Alerting. var baseInstanceStates = map[string]State{ - "Normal": StateNormal, - "Alerting": StateFiring, - "Pending": StatePending, - "NoData": StateNodata, - "Error": StateError, + "Normal": StateNormal, + "Alerting": StateFiring, + "Pending": StatePending, + "NoData": StateNodata, + "Error": StateError, + "Recovering": StateRecovering, } // normalizeInstanceState normalizes an instance-level state string into its diff --git a/grafana-alertcheck/internal/gate/parse_state_test.go b/grafana-alertcheck/internal/gate/parse_state_test.go index f55c10185..a37591ddb 100644 --- a/grafana-alertcheck/internal/gate/parse_state_test.go +++ b/grafana-alertcheck/internal/gate/parse_state_test.go @@ -186,6 +186,8 @@ func TestParseNormalizeInstanceState(t *testing.T) { {"Pending", StatePending, "", false}, {"NoData", StateNodata, "", false}, {"Error", StateError, "", false}, + {"Recovering", StateRecovering, "", false}, + {"Recovering (NoData)", StateRecovering, "NoData", false}, {"Normal (NoData)", StateNormal, "NoData", false}, {"Normal (Error)", StateNormal, "Error", false}, {"Normal (MissingSeries)", StateNormal, "MissingSeries", false}, @@ -242,15 +244,17 @@ func TestParseState_KeepFiringForIsOptional(t *testing.T) { tests := []struct { name string extra string + want time.Duration }{ - {"present", `,"keepFiringFor":300`}, - {"absent", ""}, + {"present", `,"keepFiringFor":300`, 300 * time.Second}, + {"absent", "", 0}, } for _, tc := range tests { t.Run(tc.name, func(t *testing.T) { rules, err := ParseState(minimalStateBody(tc.extra)) require.NoError(t, err) require.Len(t, rules, 1) + require.Equal(t, tc.want, rules[0].KeepFiringFor) }) } } diff --git a/grafana-alertcheck/internal/gate/schedule.go b/grafana-alertcheck/internal/gate/schedule.go index 65bb7d419..a5a745104 100644 --- a/grafana-alertcheck/internal/gate/schedule.go +++ b/grafana-alertcheck/internal/gate/schedule.go @@ -69,29 +69,20 @@ func defaultPollEvery(intervalSeconds int) time.Duration { } // DeriveTimings computes every resolved rule's ruleTimings (keyed by UID) plus -// the shared globalTimings. A non-zero override is used verbatim for every rule -// and never clamped to the default — an override above intervalSeconds/2 widens -// maxGap and is reported as a note, not corrected. -func DeriveTimings(defs []Definition, override time.Duration) (rules map[string]RuleTimings, global GlobalTimings, notes []string) { +// the shared globalTimings. Every rule polls at half its own evaluation +// interval: the cadence is fixed so no state that lasts a full evaluation +// interval can fall between polls, and there is deliberately no override that +// could widen maxGap and hide one. +func DeriveTimings(defs []Definition) (rules map[string]RuleTimings, global GlobalTimings) { rules = make(map[string]RuleTimings, len(defs)) for _, d := range defs { - def := defaultPollEvery(d.IntervalSeconds) - pollEvery := def - if override > 0 { - pollEvery = override - if override > def { - notes = append(notes, fmt.Sprintf( - "rule %s: --poll-interval %s exceeds half its %ds evaluation interval (%s); maxGap widens accordingly", - d.Title, override, d.IntervalSeconds, def)) - } - } - rt := newRuleTimings(pollEvery, d.IntervalSeconds) + rt := newRuleTimings(defaultPollEvery(d.IntervalSeconds), d.IntervalSeconds) rt.title = d.Title rules[d.UID] = rt } // In this mode defs ARE the start-of-step snapshot, so they answer what was // paused at the window open; only the log-mode counterpart uses the header. - return rules, deriveGlobalTimings(defs, pausedSet(defs)), notes + return rules, deriveGlobalTimings(defs, pausedSet(defs)) } // pausedSet is Header.pausedAtStart's counterpart for a set of definitions @@ -316,8 +307,8 @@ func (s *Scheduler) earliestDue() (time.Time, bool) { // - the burst bound — the slowest measured request is slower than the fleet's // tightest cadence, which can open a mid-run gap beyond that rule's maxGap. // -// The message names only the three operator controls: concurrency, -// poll-interval, and the alert list. +// The message names only the two operator controls: concurrency and the alert +// list. func CheckBudget(t map[string]RuleTimings, measured map[string]time.Duration, concurrency int) error { if len(t) == 0 { return nil @@ -389,10 +380,10 @@ func CheckBudget(t map[string]RuleTimings, measured map[string]time.Duration, co // Utilization is the one condition concurrency fixes, so name the exact // value; the other two are single-request shapes no concurrency shortens. if utilization > float64(concurrency) { - fmt.Fprintf(&b, "fix by: raising --concurrency to at least %d (currently %d), raising poll-interval, or watching fewer alerts", + fmt.Fprintf(&b, "fix by: raising --concurrency to at least %d (currently %d), or watching fewer alerts", int(math.Ceil(utilization)), concurrency) } else { - b.WriteString("fix by: raising poll-interval or watching fewer alerts") + b.WriteString("fix by: watching fewer alerts") } return fmt.Errorf("%s", b.String()) } diff --git a/grafana-alertcheck/internal/gate/schedule_test.go b/grafana-alertcheck/internal/gate/schedule_test.go index b4ae02669..53c71dde7 100644 --- a/grafana-alertcheck/internal/gate/schedule_test.go +++ b/grafana-alertcheck/internal/gate/schedule_test.go @@ -11,8 +11,7 @@ func TestDeriveTimings_Default(t *testing.T) { defs := []Definition{ {UID: "r1", Title: "R1", IntervalSeconds: 60}, } - rules, _, notes := DeriveTimings(defs, 0) - require.Empty(t, notes) + rules, _ := DeriveTimings(defs) rt := rules["r1"] require.Equal(t, 30*time.Second, rt.pollEvery) require.Equal(t, 60*time.Second, rt.maxGap) @@ -20,30 +19,12 @@ func TestDeriveTimings_Default(t *testing.T) { require.Equal(t, 120*time.Second, rt.evalStaleAfter) } -func TestDeriveTimings_OverrideVerbatimNoClamp(t *testing.T) { - defs := []Definition{ - {UID: "r1", Title: "R1", IntervalSeconds: 10}, // default pollEvery = 5s - } - rules, _, notes := DeriveTimings(defs, 20*time.Second) - rt := rules["r1"] - require.Equal(t, 20*time.Second, rt.pollEvery, "the override verbatim (20s), never clamped down to the 5s default") - require.Equal(t, 40*time.Second, rt.maxGap) - require.Len(t, notes, 1) - require.Contains(t, notes[0], "R1") -} - -func TestDeriveTimings_OverrideBelowDefaultNoNote(t *testing.T) { - defs := []Definition{{UID: "r1", Title: "R1", IntervalSeconds: 60}} // default pollEvery = 30s - _, _, notes := DeriveTimings(defs, 5*time.Second) - require.Empty(t, notes) -} - func TestDeriveTimings_TransitionGraceExcludesSkippedRule(t *testing.T) { defs := []Definition{ {UID: "r1", Title: "Tight", IntervalSeconds: 60, For: time.Minute}, {UID: "r2", Title: "PausedLongFor", IntervalSeconds: 60, For: time.Hour, IsPaused: true}, } - _, global, _ := DeriveTimings(defs, 0) + _, global := DeriveTimings(defs) want := time.Minute + 60*time.Second // r1's for+interval; r2 (skipped) must not win despite its huge `for` require.Equal(t, want, global.transitionGrace) require.Contains(t, global.graceSource, "Tight") @@ -58,8 +39,7 @@ func TestDeriveTimings_TransitionGraceExcludesSkippedRule(t *testing.T) { // w unit — testdata/README.md) through ParseDefinitions and DeriveTimings. func TestDeriveTimings_RealForOneWeekRuleSetsTransitionGrace(t *testing.T) { defs := rulerDefs(t) - _, global, notes := DeriveTimings(defs, 0) - require.Empty(t, notes, "no --poll-interval override is given, so no override note should fire") + _, global := DeriveTimings(defs) want := 7*24*time.Hour + 60*time.Second // rule0000010: for=1w, intervalSeconds=60 require.Equal(t, want, global.transitionGrace) @@ -68,7 +48,7 @@ func TestDeriveTimings_RealForOneWeekRuleSetsTransitionGrace(t *testing.T) { func TestDeriveTimings_TransitionGraceZeroWhenAllSkipped(t *testing.T) { defs := []Definition{{UID: "r1", Title: "R1", IntervalSeconds: 60, For: time.Hour, IsPaused: true}} - _, global, _ := DeriveTimings(defs, 0) + _, global := DeriveTimings(defs) require.Zero(t, global.transitionGrace) } @@ -123,7 +103,7 @@ func TestDeriveTimingsFromLog_TransitionGraceFollowsTheHeaderNotTheDefinition(t func TestDeriveTimings_DrainTimeoutIncludesPaused(t *testing.T) { defs := []Definition{{UID: "r1", Title: "R1", IntervalSeconds: 10}} - _, global, _ := DeriveTimings(defs, 0) + _, global := DeriveTimings(defs) require.Equal(t, minDrainTimeout, global.drainTimeout) } @@ -132,14 +112,14 @@ func TestDeriveTimings_DrainTimeoutFloor(t *testing.T) { {UID: "r1", Title: "Tight", IntervalSeconds: 60, For: time.Minute}, {UID: "r2", Title: "PausedLongFor", IntervalSeconds: 180, For: time.Hour, IsPaused: true}, } - _, global, _ := DeriveTimings(defs, 0) + _, global := DeriveTimings(defs) // double the longest interval (2 * 180s) should be the drain timeout require.Equal(t, 2*180*time.Second, global.drainTimeout) } func TestDeriveTimings_DrainTimeoutAboveFloor(t *testing.T) { defs := []Definition{{UID: "r1", Title: "R1", IntervalSeconds: 300}} // 2x300s = 600s > 2m floor - _, global, _ := DeriveTimings(defs, 0) + _, global := DeriveTimings(defs) require.Equal(t, 600*time.Second, global.drainTimeout) } @@ -418,11 +398,10 @@ func TestCheckBudget_NonPositivePollIntervalIsAnError(t *testing.T) { } // assertBudgetMessage checks the message contents: a measured duration is -// present, and all three controls are named — never a single suggested -// interval. +// present, and the remaining levers are named. func assertBudgetMessage(t *testing.T, msg string) { t.Helper() - for _, want := range []string{"measured", "concurrency", "poll-interval", "fewer"} { + for _, want := range []string{"measured", "concurrency", "fewer"} { require.Contains(t, msg, want) } } diff --git a/grafana-alertcheck/internal/gate/terminal.go b/grafana-alertcheck/internal/gate/terminal.go index c8b4197bb..863ad28cb 100644 --- a/grafana-alertcheck/internal/gate/terminal.go +++ b/grafana-alertcheck/internal/gate/terminal.go @@ -64,7 +64,7 @@ func terminalVerdict(h Header, polls []Poll, defs []Definition, rt map[string]Ru }, true } - outcome, _, _ := classifyRule(def, polls, from, at, badStates, pol.Preexisting) + outcome, _, _ := classifyRule(def, rt[def.UID], polls, from, at, badStates, pol.Preexisting) if outcome == OutcomeNewFailure || outcome == OutcomeUnstable { if violation == nil { v := Termination{ diff --git a/grafana-alertcheck/internal/gate/watch.go b/grafana-alertcheck/internal/gate/watch.go index 98e29145f..e907cb646 100644 --- a/grafana-alertcheck/internal/gate/watch.go +++ b/grafana-alertcheck/internal/gate/watch.go @@ -70,13 +70,6 @@ type WatchConfig struct { // collection loop ends. Until time.Time - // PollEvery is the --poll-interval override, used verbatim for every rule - // and never clamped. Zero means each rule polls at half its own evaluation - // interval. Whatever this resolves to is written into the header as the - // cadence actually used, and that header value — never a re-derivation from - // the definitions — is what check derives maxGap from. - PollEvery time.Duration - Concurrency int Clock Clock @@ -306,10 +299,7 @@ func prepareWatch(ctx context.Context, cfg WatchConfig, src Source) (*preparedWa printLabelSelection(cfg.Notes, resolved, cfg.IncludeLabels, cfg.ExcludeLabels, len(cfg.ExcludeAlerts)) } - rt, _, timingNotes := DeriveTimings(resolved, cfg.PollEvery) - for _, n := range timingNotes { - fmt.Fprintf(cfg.Notes, "note: %s\n", n) - } + rt, _ := DeriveTimings(resolved) for _, d := range resolved { // A cadence of zero would make the child spin: every rule is due the // instant it was marked. It also cannot be written into the header, diff --git a/grafana-alertcheck/internal/gate/watch_daemon_test.go b/grafana-alertcheck/internal/gate/watch_daemon_test.go index 04c1276e9..58f16a4fd 100644 --- a/grafana-alertcheck/internal/gate/watch_daemon_test.go +++ b/grafana-alertcheck/internal/gate/watch_daemon_test.go @@ -131,7 +131,7 @@ const testBearerToken = "test-token" // the reducer's select-by-UID finds it. func grafanaTestServer(t *testing.T) *httptest.Server { t.Helper() - ruler := readFixture(t, "ruler_rules.json") + ruler := patchedRulerBody(t) state := patchedStateBody(t) srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { @@ -162,6 +162,35 @@ func grafanaTestServer(t *testing.T) *httptest.Server { return srv } +// patchedRulerBody serves the ruler fixture with the watched rule's evaluation +// interval shortened to 2s, so the detached recorder appends polls fast enough +// for a real-process test. The cadence is half that interval — the recorder has +// no override any more. +func patchedRulerBody(t *testing.T) []byte { + t.Helper() + var body map[string][]map[string]any + require.NoError(t, json.Unmarshal(readFixture(t, "ruler_rules.json"), &body)) + found := false + for _, groups := range body { + for _, group := range groups { + rules, _ := group["rules"].([]any) + for _, r := range rules { + rule, _ := r.(map[string]any) + ga, _ := rule["grafana_alert"].(map[string]any) + if ga["uid"] != watchActiveUID { + continue + } + ga["intervalSeconds"] = 2 + found = true + } + } + } + require.True(t, found, "ruler fixture has no rule %s", watchActiveUID) + b, err := json.Marshal(body) + require.NoError(t, err) + return b +} + func patchedStateBody(t *testing.T) []byte { t.Helper() var body map[string]any @@ -208,8 +237,8 @@ func waitFor(t *testing.T, what string, timeout time.Duration, cond func() bool) // It asserts the four things only a real spawn can show — the pidfile points // at a live process, that process is in its own session (setsid, not a bare // `&`), it keeps appending after Watch returned, and SIGTERM makes it finish -// the log in the stop order — with a 2s --poll-interval, which also exercises -// the unclamped-override path. +// the log in the stop order. The served rule evaluates every 2s, so its 1s +// cadence appends polls fast enough for a real-time test. func TestWatchSpawnsADetachedRecorder(t *testing.T) { srv := grafanaTestServer(t) t.Setenv("GRAFANA_URL", srv.URL) @@ -222,7 +251,6 @@ func TestWatchSpawnsADetachedRecorder(t *testing.T) { Token: testBearerToken, Alerts: []string{"uid:" + watchActiveUID}, Out: out, - PollEvery: 2 * time.Second, Concurrency: 2, Notes: ¬es, } @@ -264,7 +292,7 @@ func TestWatchSpawnsADetachedRecorder(t *testing.T) { require.Equal(t, srv.URL, header.URL) require.Equal(t, "13.1.0", header.GrafanaVersion) require.Len(t, header.Rules, 1) - require.Equal(t, float64(2), header.Rules[0].PollEverySeconds) + require.Equal(t, float64(1), header.Rules[0].PollEverySeconds, "half the served 2s evaluation interval") for i, p := range polls { require.Equalf(t, watchActiveUID, p.RuleUID, "poll %d", i) require.Truef(t, p.Found, "poll %d", i) diff --git a/grafana-alertcheck/internal/gate/watch_process.go b/grafana-alertcheck/internal/gate/watch_process.go index 169a4f65a..b730b40a3 100644 --- a/grafana-alertcheck/internal/gate/watch_process.go +++ b/grafana-alertcheck/internal/gate/watch_process.go @@ -90,8 +90,8 @@ func spawnChild(cfg WatchConfig) (detachedChild, error) { // optional hard stop and the concurrency limit. // // Notably absent: --pidfile (the parent writes it, so check can find the pid -// the instant Watch returns), --alerts, --folder, --poll-interval, and -// anything derived from them. The CLI's `watch` FlagSet needs one flag of its +// the instant Watch returns), --alerts, --folder, and anything derived from +// them. The CLI's `watch` FlagSet needs one flag of its // own for this path — ReadyFDFlag — and dispatches to RunDaemonChild when it // sees DaemonChildFlag. func childArgs(cfg WatchConfig) []string { diff --git a/grafana-alertcheck/internal/gate/watch_test.go b/grafana-alertcheck/internal/gate/watch_test.go index 576745da2..34d9c7eb4 100644 --- a/grafana-alertcheck/internal/gate/watch_test.go +++ b/grafana-alertcheck/internal/gate/watch_test.go @@ -426,25 +426,6 @@ func TestPrepareWatchDoesNotWaitForPausedRules(t *testing.T) { require.Contains(t, notes.String(), "paused") } -// One authority for the cadence, from the writing side: whatever -// --poll-interval resolves to is what the header records, because that is the -// only value check may derive maxGap from. -func TestPrepareWatchHeaderRecordsTheOverriddenCadence(t *testing.T) { - var notes strings.Builder - cfg := watchTestConfig(t, ¬es, "uid:"+watchActiveUID) - cfg.PollEvery = 120 * time.Second // the rule evaluates every 60s - src := watchTestSource(t, liveObservation(testNow)) - - prep, err := prepareWatch(context.Background(), cfg, src) - require.NoError(t, err) - defer prep.writer.Close() - - require.Equal(t, float64(120), prep.header.Rules[0].PollEverySeconds, - "the override, used verbatim and never clamped") - require.Equal(t, 240*time.Second, prep.timings[watchActiveUID].maxGap) - require.Contains(t, notes.String(), "--poll-interval") -} - // ReadyAt must be stamped after the first-observation pass, never before it: // that is what lets check refuse a `from` inside the pass. StartedAt stays the // record start, captured before the pass. @@ -722,7 +703,6 @@ func TestChildArgsCarryNoSecretsAndNoRuleSet(t *testing.T) { Out: "/tmp/log.jsonl", PidFile: "/tmp/log.jsonl.pid", Until: testNow.Add(time.Hour), - PollEvery: 17 * time.Second, Concurrency: 3, } args := childArgs(cfg)