diff --git a/grafana-alertcheck/.changeset/v0.1.2.md b/grafana-alertcheck/.changeset/v0.1.2.md new file mode 100644 index 000000000..2ba1c19a8 --- /dev/null +++ b/grafana-alertcheck/.changeset/v0.1.2.md @@ -0,0 +1,2 @@ +- Add the `stop` subcommand: reap a detached recorder after a failed work step. It is idempotent, so it is safe as an `if: always()` step. +- `check` now exits early (fail-fast) on a condition that cannot become a pass — a post-`from` bad onset or an inability — instead of always waiting for `to + transitionGrace + drainTimeout`. An early exit is never a pass; pass `--no-fail-fast` to always wait for the full window and its coverage proof. diff --git a/grafana-alertcheck/.changeset/v0.1.3.md b/grafana-alertcheck/.changeset/v0.1.3.md new file mode 100644 index 000000000..66b3aa5fc --- /dev/null +++ b/grafana-alertcheck/.changeset/v0.1.3.md @@ -0,0 +1,2 @@ +- The `check` output now speaks plain words. Outcome values are renamed: `clean` → `healthy`, `newly_bad` → `new_failure`, `persistently_bad` → `still_failing`, `flapping` → `unstable`, `skipped` → `paused`, `unobservable` → `not_verified`, and the early-exit `terminated_early.kind` follows the same rename. A `--min-observed` deficit that no rule explains is now reported as `not_counted` instead of being blamed on a paused rule. This changes the `--output json` vocabulary and anything downstream of it, including the action's `outcomes` output. +- The human table dropped its internal column names. `RESULTS` is now `ALERT`/`VERDICT`/`BROKEN FOR`/`CHECKED EVERY`/`WINDOW COVERED`/`DETAILS`; `VIOLATIONS` uses `GRAFANA STATE`/`GRAFANA HEALTH` and a single-word `INSTANCES` column (the previous `INSTANCE COUNT` header read as two columns, one of them empty); `THRESHOLDS` became `LIMITS USED`, with limits named in plain words and explained by a legend under the table. The footer now spells out the extra observation time, the evaluation wait and the clock difference from Grafana. diff --git a/grafana-alertcheck/README.md b/grafana-alertcheck/README.md index 96aa0e999..6dd5e038e 100644 --- a/grafana-alertcheck/README.md +++ b/grafana-alertcheck/README.md @@ -9,7 +9,7 @@ watch → your work → check `watch` starts a background recorder that polls each named alert into a JSONL log. After the work emits a `from`/`to` pair, `check` proves continuous coverage of that window, classifies each alert's state -timeline, and exits `0`, `1`, or `2`. +timeline, and exits `0`, `1`, or `2`. If the work fails first, `stop` reaps the recorder. It **fails closed**: if it cannot get an answer, it stops the release — never a pass on an unproven window. diff --git a/grafana-alertcheck/cmd/check.go b/grafana-alertcheck/cmd/check.go index e95fce376..ab063562a 100644 --- a/grafana-alertcheck/cmd/check.go +++ b/grafana-alertcheck/cmd/check.go @@ -16,7 +16,7 @@ import ( const checkUsage = "usage: grafana-alertcheck check [--in ] [--pidfile F] --from RFC3339 --to RFC3339 " + "[--alerts ...] [--folder F] [--states ...] [--preexisting ...] [--min-observed N] [--allow-paused] " + - "[--nodata-is-unobservable] [--concurrency N] [--output json]" + "[--nodata-is-unobservable] [--no-fail-fast] [--concurrency N] [--output json]" // runCheck is the classify step's CLI surface: parse flags into a gate.Config, // run gate.Check, and translate its (Result, error) into output and an exit @@ -38,6 +38,7 @@ func runCheck(args []string, stdin io.Reader, stdout, stderr io.Writer) int { 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") nodataIsUnobservable := fs.Bool("nodata-is-unobservable", false, "treat a sustained health=nodata as unobservable rather than a note") + noFailFast := fs.Bool("no-fail-fast", false, "collect to to+transitionGrace even after a certain failure, for a full-window coverage proof instead of the fastest feedback") output := fs.String("output", "", `"json" writes the machine-readable Result to stdout in addition to the table; default is the table alone`) if err := fs.Parse(args); err != nil { @@ -86,6 +87,7 @@ func runCheck(args []string, stdin io.Reader, stdout, stderr io.Writer) int { MinObserved: *minObserved, AllowPaused: *allowPaused, NodataIsUnobservable: *nodataIsUnobservable, + NoFailFast: *noFailFast, Log: *in, PidFile: *pidfile, Concurrency: *common.concurrency, diff --git a/grafana-alertcheck/cmd/main.go b/grafana-alertcheck/cmd/main.go index 7ab5d6397..2cc2e5793 100644 --- a/grafana-alertcheck/cmd/main.go +++ b/grafana-alertcheck/cmd/main.go @@ -12,7 +12,7 @@ func main() { os.Exit(run(os.Args[1:], os.Stdout, os.Stderr)) } -const usage = "usage: grafana-alertcheck " +const usage = "usage: grafana-alertcheck " // run is the whole of main's testable surface: parse the subcommand, dispatch, // return the process exit code. Exit codes below 2 (pass/violations) belong to @@ -37,6 +37,8 @@ func run(args []string, stdout, stderr io.Writer) int { return runWatch(args[1:], os.Stdin, stdout, stderr) case "check": return runCheck(args[1:], os.Stdin, stdout, stderr) + case "stop": + return runStop(args[1:], os.Stdin, stdout, stderr) case "-h", "-help", "--help": fmt.Fprintln(stdout, usage) return 0 diff --git a/grafana-alertcheck/cmd/stop.go b/grafana-alertcheck/cmd/stop.go new file mode 100644 index 000000000..be282b483 --- /dev/null +++ b/grafana-alertcheck/cmd/stop.go @@ -0,0 +1,63 @@ +package main + +import ( + "context" + "errors" + "flag" + "fmt" + "io" + "os/signal" + "syscall" + + "github.com/smartcontractkit/chainlink-testing-framework/grafana-alertcheck/internal/gate" +) + +const stopUsage = "usage: grafana-alertcheck stop --out [--pidfile F]" + +// runStop reaps a detached recorder without reading or classifying its log — +// what an `if: always()` step calls when the work failed. It is idempotent, so +// it is a no-op after check has already stopped the recorder. +func runStop(args []string, _ io.Reader, _, stderr io.Writer) int { + fs := flag.NewFlagSet("stop", flag.ContinueOnError) + fs.SetOutput(stderr) + fs.Usage = func() { fmt.Fprintln(stderr, stopUsage) } + + out := fs.String("out", "", "JSONL log path whose recorder to stop") + pidfile := fs.String("pidfile", "", "pidfile path (default .pid)") + + if err := fs.Parse(args); err != nil { + if errors.Is(err, flag.ErrHelp) { + return 0 + } + return 2 + } + if fs.NArg() != 0 { + fmt.Fprintf(stderr, "stop: unexpected arguments %v\n", fs.Args()) + return 2 + } + if *out == "" { + fmt.Fprintln(stderr, "stop: --out is required") + return 2 + } + + // SIGINT/SIGTERM cancel the wait cleanly; an interrupted cleanup is a + // could-not-complete, never a silent success. + ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) + defer cancel() + + held, err := gate.StopRecorder(ctx, gate.StopConfig{ + Log: *out, + PidFile: *pidfile, + Clock: gate.SystemClock{}, + Notes: newNoteStyler(stderr), + Cleanup: true, + }) + if held != nil { + _ = held.Close() + } + if err != nil { + fmt.Fprintln(stderr, err) + return 2 + } + return 0 +} diff --git a/grafana-alertcheck/cmd/stop_test.go b/grafana-alertcheck/cmd/stop_test.go new file mode 100644 index 000000000..6021730f8 --- /dev/null +++ b/grafana-alertcheck/cmd/stop_test.go @@ -0,0 +1,28 @@ +package main + +import ( + "bytes" + "os" + "path/filepath" + "testing" + + "github.com/stretchr/testify/require" +) + +func TestRunStop_RequiresOut(t *testing.T) { + var stdout, stderr bytes.Buffer + code := run([]string{"stop"}, &stdout, &stderr) + require.Equal(t, 2, code) + require.Contains(t, stderr.String(), "--out is required") +} + +// Nothing recorded: an `if: always()` stop still exits 0. +func TestRunStop_NoPidfileIsANoOp(t *testing.T) { + out := filepath.Join(t.TempDir(), "log.jsonl") + require.NoError(t, os.WriteFile(out, nil, 0o644)) // nolint:gosec // test-only temp file + + var stdout, stderr bytes.Buffer + code := run([]string{"stop", "--out", out}, &stdout, &stderr) + require.Equal(t, 0, code) + require.Contains(t, stderr.String(), "nothing to stop") +} diff --git a/grafana-alertcheck/cmd/table.go b/grafana-alertcheck/cmd/table.go index 26a8944de..085444b88 100644 --- a/grafana-alertcheck/cmd/table.go +++ b/grafana-alertcheck/cmd/table.go @@ -11,22 +11,29 @@ import ( "github.com/smartcontractkit/chainlink-testing-framework/grafana-alertcheck/internal/gate" ) +// limitsLegend explains each LIMITS USED column in one plain sentence, so the +// table needs no documentation lookup. +const limitsLegend = ` max gap without check — the longest gap between two checks we accept before we say the alert was not watched. + query failing for — how long Grafana may keep failing to run the alert's query before we stop trusting its state. + no evaluation for — how long Grafana may go without evaluating the alert before we stop trusting its state.` + // renderTable is the human table. It always writes to the writer it is given, // which the caller (runCheck) always points at stderr — stdout is reserved for // the machine-readable --output json. // -// Three titled tables, in order (the name column is RULE in all of them — one -// row is one resolved alert rule, never a firing instance): +// Three titled tables, in order (the ALERT column is one resolved alert rule, +// never a firing instance): +// +// 1. RESULTS, one line per rule: verdict, time broken, check cadence, +// whether the window was observed, and any notes; +// 2. VIOLATIONS, one line per distinct violation, with the raw Grafana state +// and health and the number of instances it stands for; +// 3. LIMITS USED, the coverage thresholds that answer "why" on exit 2. Each +// column is named in plain words and explained by the legend below it, so +// the table needs no documentation lookup. // -// 1. RESULTS, one line per rule: outcome, BadFor, pollEvery, proved-or-not -// with the largest gap; -// 2. VIOLATIONS, one line per distinct violation. -// 3. THRESHOLDS, the numbers that answer "why" on exit 2: each non-skipped -// rule's maxGap/healthGrace/evalStaleAfter, followed by the global -// transitionGrace and drainTimeout, and the largest measured clock skew -// alongside its own error bound (RTT/2) — SkewHardLimit is a separate, -// fixed input threshold and is reported next to it, never as if it were -// that bound. +// The global footer then reports the extra observation time, the drain limit +// and the largest measured clock difference, also in plain words. func renderTable(w io.Writer, res gate.Result) error { alertOf := make(map[string]string, len(res.Verdicts)) for _, v := range res.Verdicts { @@ -40,9 +47,19 @@ func renderTable(w io.Writer, res gate.Result) error { // wall of progress text. fmt.Fprintln(w) + if te := res.TerminatedEarly; te != nil { + detail := string(te.Reason) + if detail == "" { + detail = string(te.Outcome) + } + fmt.Fprintf(w, "EARLY EXIT: %s %q at %s (%s); the window [%s, %s] was not fully observed\n\n", + te.Kind, te.Alert, te.At.Format(time.RFC3339), detail, + res.From.Format(time.RFC3339), res.To.Format(time.RFC3339)) + } + fmt.Fprintln(w, "RESULTS") tw := tabwriter.NewWriter(w, 0, 4, 2, ' ', 0) - fmt.Fprintln(tw, "RULE\tOUTCOME\tBADFOR\tPOLLEVERY\tPROVED\tNOTE") + fmt.Fprintln(tw, "ALERT\tVERDICT\tBROKEN FOR\tCHECKED EVERY\tWINDOW COVERED\tDETAILS") for _, v := range sortedVerdicts(res.Verdicts) { fmt.Fprintf(tw, "%s\t%s\t%s\t%s\t%s\t%s\n", v.Alert, v.Outcome, v.BadFor.Round(time.Second), v.PollEvery.Round(time.Second), @@ -55,7 +72,7 @@ func renderTable(w io.Writer, res gate.Result) error { if len(res.Violations) > 0 { fmt.Fprintln(w, "\nVIOLATIONS") vtw := tabwriter.NewWriter(w, 0, 4, 2, ' ', 0) - fmt.Fprintln(vtw, "RULE\tOUTCOME\tSTATE\tHEALTH\tINSTANCE COUNT\tNOTE") + fmt.Fprintln(vtw, "ALERT\tVERDICT\tGRAFANA STATE\tGRAFANA HEALTH\tINSTANCES\tDETAILS") for _, g := range groupedViolations(res.Violations) { fmt.Fprintf(vtw, "%s\t%s\t%s\t%s\t%s\t%s\n", alertLabel(g.v, alertOf), g.v.Outcome, g.v.State, g.v.Health, instanceCount(g), g.v.Note) } @@ -64,14 +81,13 @@ func renderTable(w io.Writer, res gate.Result) error { } } - // The per-rule thresholds answer "why" on exit 2: a table, not the prose - // "rule NAME: maxGap=... healthGrace=... evalStaleAfter=..." that repeated - // the rule name a fourth time. It is separated from the result above by a - // blank line. + // The per-rule limits answer "why" on exit 2. The columns are spelled out + // and explained by limitsLegend right below, so an operator does not have + // to look anything up. fmt.Fprintln(w) - fmt.Fprintln(w, "THRESHOLDS") + fmt.Fprintln(w, "LIMITS USED") ttw := tabwriter.NewWriter(w, 0, 4, 2, ' ', 0) - fmt.Fprintln(ttw, "RULE\tMAXGAP\tHEALTHGRACE\tEVALSTALEAFTER") + fmt.Fprintln(ttw, "ALERT\tMAX GAP WITHOUT CHECK\tQUERY FAILING FOR\tNO EVALUATION FOR") for _, uid := range sortedThresholdUIDs(res.Thresholds, alertOf) { t := res.Thresholds[uid] fmt.Fprintf(ttw, "%s\t%s\t%s\t%s\n", @@ -80,11 +96,17 @@ func renderTable(w io.Writer, res gate.Result) error { if err := ttw.Flush(); err != nil { return fmt.Errorf("render table: %w", err) } + fmt.Fprintln(w, limitsLegend) fmt.Fprintln(w) - fmt.Fprintf(w, "global: transitionGrace=%s (source: %s) drainTimeout=%s\n", - res.Global.TransitionGrace, res.Global.GraceSource, res.Global.DrainTimeout) - fmt.Fprintf(w, "largest measured clock skew: %s (bound ±%s, hard limit %s), grafana %s\n", + if res.Global.TransitionGrace > 0 { + fmt.Fprintf(w, "extra watching after your window: +%s — so an alert that only starts firing at the end is still caught (slowest: %s)\n", + res.Global.TransitionGrace, res.Global.GraceSource) + } else { + fmt.Fprintln(w, "extra watching after your window: none") + } + fmt.Fprintf(w, "max wait for all alerts to finish evaluating: %s\n", res.Global.DrainTimeout) + fmt.Fprintf(w, "clock difference from Grafana: %s, accurate to ±%s (checks fail above %s); Grafana %s\n", res.ClockSkew.Round(time.Millisecond), res.ClockSkewBound.Round(time.Millisecond), gate.SkewHardLimit, res.GrafanaVersion) // The verdict — the single number a terminal operator reads last — sits on @@ -94,8 +116,8 @@ func renderTable(w io.Writer, res gate.Result) error { return nil } -// violationsLabel colours the "violations: N" prefix of the footer: green for a -// clean run, red otherwise. The rest of the line is written uncoloured. +// violationsLabel colours the "violations: N" prefix of the footer: green when +// there are none, red otherwise. The rest of the line is written uncoloured. func violationsLabel(n int, enabled bool) string { s := fmt.Sprintf("violations: %d", n) if !enabled { @@ -107,10 +129,10 @@ func violationsLabel(n int, enabled bool) string { return ansiRed + s + ansiReset } -// provedLabel is the table's PROVED column: "yes" for a clean coverage -// proof, "no" with the reason and largest gap for an unobservable rule, and -// "-" for a rule decide never asked proveCoverage about at all (skipped — -// paused before the window opened). +// provedLabel is the table's WINDOW COVERED column: "yes" for a fully +// observed window, "no" with the reason and largest gap for a not-verified +// rule, and "-" for a rule decide never asked proveCoverage about at all +// (paused before the window opened). func provedLabel(cov gate.CoverageResult) string { if cov.Reason == "" && !cov.Unobservable && !cov.Proved { return "-" @@ -182,8 +204,11 @@ func sameRendered(a, b gate.Violation) bool { return violationSignature(a) == violationSignature(b) } +// instanceCount is the INSTANCES column: how many alert instances one grouped +// violation row stands for. Paused and not-counted rows stand for no instance +// at all, so they render "-". func instanceCount(g violationGroup) string { - if g.v.Outcome == gate.OutcomeSkipped { + if g.v.Outcome == gate.OutcomePaused || g.v.Outcome == gate.OutcomeNotCounted { return "-" } return strconv.Itoa(g.n) diff --git a/grafana-alertcheck/cmd/table_test.go b/grafana-alertcheck/cmd/table_test.go index b050f0d30..acf987423 100644 --- a/grafana-alertcheck/cmd/table_test.go +++ b/grafana-alertcheck/cmd/table_test.go @@ -11,8 +11,8 @@ import ( ) // The golden table test: a fixed Result renders a deterministic, ordered rule -// table, a violations section and a footer carrying the per-rule and global -// thresholds plus the skew and its bound — with no live Check involved. +// table, a violations section and a footer carrying the per-rule limits and +// globals in plain words — with no live Check involved. func TestRenderTable(t *testing.T) { gapAt := time.Date(2026, 1, 1, 12, 0, 0, 0, time.UTC) res := gate.Result{ @@ -20,15 +20,15 @@ func TestRenderTable(t *testing.T) { ClockSkew: 1500 * time.Millisecond, ClockSkewBound: 250 * time.Millisecond, Verdicts: []gate.RuleVerdict{ - {Alert: "Zebra Alert", RuleUID: "uid-z", Outcome: gate.OutcomeClean, PollEvery: 30 * time.Second}, - {Alert: "Ape Alert", RuleUID: "uid-a", Outcome: gate.OutcomeUnobservable, + {Alert: "Zebra Alert", RuleUID: "uid-z", Outcome: gate.OutcomeHealthy, PollEvery: 30 * time.Second}, + {Alert: "Ape Alert", RuleUID: "uid-a", Outcome: gate.OutcomeNotVerified, PollEvery: 30 * time.Second, Note: "gap of 5m0s starting at 2026-01-01T12:00:00Z exceeds maxGap 1m0s"}, - {Alert: "Paused Alert", RuleUID: "uid-p", Outcome: gate.OutcomeSkipped, + {Alert: "Paused Alert", RuleUID: "uid-p", Outcome: gate.OutcomePaused, Note: "paused before the window opened; counts against --min-observed unless --allow-paused is set"}, }, Violations: []gate.Violation{ - {Alert: "Ape Alert", RuleUID: "uid-a", Outcome: gate.OutcomeUnobservable, State: gate.StateFiring, Health: "error", Note: "unobservable"}, - {Alert: "Paused Alert", RuleUID: "uid-p", Outcome: gate.OutcomeSkipped, + {Alert: "Ape Alert", RuleUID: "uid-a", Outcome: gate.OutcomeNotVerified, State: gate.StateFiring, Health: "error", Note: "not verified"}, + {Alert: "Paused Alert", RuleUID: "uid-p", Outcome: gate.OutcomePaused, Note: "paused before the window opened; counts against --min-observed unless --allow-paused is set"}, }, Coverage: map[string]gate.CoverageResult{ @@ -50,42 +50,49 @@ func TestRenderTable(t *testing.T) { require.NoError(t, renderTable(&buf, res)) out := buf.String() - // Rule table: Ape sorts before Zebra sorts before... Paused is skipped and - // carries no coverage entry, so it renders "-" for PROVED. + // The rule table: Ape sorts before Zebra sorts before... Paused carries no + // coverage entry, so it renders "-" for WINDOW COVERED. + require.Contains(t, out, "ALERT") require.Contains(t, out, "Ape Alert") - require.Contains(t, out, "unobservable") + require.Contains(t, out, "not_verified") require.Contains(t, out, "heartbeat_gap") require.Contains(t, out, "largest gap 5m0s") require.Contains(t, out, "Zebra Alert") - require.Contains(t, out, "clean") + require.Contains(t, out, "healthy") + require.Contains(t, out, "WINDOW COVERED") // The violations section must show up even without --output json, and must - // carry the --allow-paused hint text verbatim. + // carry the --allow-paused hint text verbatim. INSTANCES is a single word + // so the count under it cannot read as a second, empty column. require.Contains(t, out, "VIOLATIONS") require.Contains(t, out, "--allow-paused") - require.Contains(t, out, "STATE") - require.Contains(t, out, "HEALTH") - require.Contains(t, out, "INSTANCE COUNT") + require.Contains(t, out, "GRAFANA STATE") + require.Contains(t, out, "GRAFANA HEALTH") + require.Contains(t, out, "INSTANCES") require.Contains(t, out, string(gate.StateFiring)) require.Contains(t, out, "error") - // The footer: per-rule thresholds are a table (RULE/MAXGAP/HEALTHGRACE/ - // EVALSTALEAFTER) rather than prose, followed by the global thresholds and - // the violations count with the skew and its own bound rather than the - // fixed hard limit. - require.Contains(t, out, "MAXGAP") - require.Contains(t, out, "HEALTHGRACE") - require.Contains(t, out, "EVALSTALEAFTER") - require.Contains(t, out, "global: transitionGrace=5m0s (source: Ape Alert (for=5m)) drainTimeout=2m0s") - require.Contains(t, out, "largest measured clock skew: 1.5s (bound ±250ms, hard limit 1m0s)") + // The limits table names each threshold in plain words and explains it + // right below, so an operator does not have to consult the docs. + require.Contains(t, out, "LIMITS USED") + require.Contains(t, out, "MAX GAP WITHOUT CHECK") + require.Contains(t, out, "QUERY FAILING FOR") + require.Contains(t, out, "NO EVALUATION FOR") + require.Contains(t, out, "the longest gap between two checks") + require.Contains(t, out, "without evaluating the alert") + + // The global footer in plain words. + require.Contains(t, out, "extra watching after your window: +5m0s") + require.Contains(t, out, "slowest: Ape Alert (for=5m)") + require.Contains(t, out, "max wait for all alerts to finish evaluating: 2m0s") + require.Contains(t, out, "clock difference from Grafana: 1.5s, accurate to ±250ms (checks fail above 1m0s); Grafana 13.1.0") require.Contains(t, out, "violations: 2") - require.Contains(t, out, "13.1.0") } // The "-" case: a rule decide never asked proveCoverage about (paused before // the window opened) has an empty CoverageResult and must not be reported as -// either proved or unobservable. -func TestProvedLabel_Skipped(t *testing.T) { +// either covered or not verified. +func TestProvedLabel_Paused(t *testing.T) { require.Equal(t, "-", provedLabel(gate.CoverageResult{})) } @@ -93,12 +100,12 @@ func TestProvedLabel_Skipped(t *testing.T) { // rendered signature, each with a count. func TestGroupedViolations(t *testing.T) { in := []gate.Violation{ - {Alert: "OCR2 Consensus failure", RuleUID: "uid-o", Outcome: gate.OutcomePersistentlyBad, State: gate.StateFiring, Health: "ok"}, - {Alert: "OCR2 Consensus failure", RuleUID: "uid-o", Outcome: gate.OutcomePersistentlyBad, State: gate.StateFiring, Health: "ok"}, - {Alert: "OCR2 Consensus failure", RuleUID: "uid-o", Outcome: gate.OutcomePersistentlyBad, State: gate.StateFiring, Health: "ok"}, - {Alert: "OCR2 Consensus failure", RuleUID: "uid-o", Outcome: gate.OutcomePersistentlyBad, State: gate.StateFiring, Health: "error"}, - {Alert: "Other Alert", RuleUID: "uid-p", Outcome: gate.OutcomeNewlyBad, State: gate.StateFiring, Health: "ok", Note: "x"}, - {Alert: "Other Alert", RuleUID: "uid-p", Outcome: gate.OutcomeNewlyBad, State: gate.StateFiring, Health: "ok", Note: "x"}, + {Alert: "OCR2 Consensus failure", RuleUID: "uid-o", Outcome: gate.OutcomeStillFailing, State: gate.StateFiring, Health: "ok"}, + {Alert: "OCR2 Consensus failure", RuleUID: "uid-o", Outcome: gate.OutcomeStillFailing, State: gate.StateFiring, Health: "ok"}, + {Alert: "OCR2 Consensus failure", RuleUID: "uid-o", Outcome: gate.OutcomeStillFailing, State: gate.StateFiring, Health: "ok"}, + {Alert: "OCR2 Consensus failure", RuleUID: "uid-o", Outcome: gate.OutcomeStillFailing, State: gate.StateFiring, Health: "error"}, + {Alert: "Other Alert", RuleUID: "uid-p", Outcome: gate.OutcomeNewFailure, State: gate.StateFiring, Health: "ok", Note: "x"}, + {Alert: "Other Alert", RuleUID: "uid-p", Outcome: gate.OutcomeNewFailure, State: gate.StateFiring, Health: "ok", Note: "x"}, } got := groupedViolations(in) @@ -108,7 +115,7 @@ func TestGroupedViolations(t *testing.T) { for _, g := range got { counts[g.v.Health+"|"+string(g.v.Outcome)] = g.n } - require.Equal(t, 3, counts["ok|"+string(gate.OutcomePersistentlyBad)]) - require.Equal(t, 1, counts["error|"+string(gate.OutcomePersistentlyBad)]) - require.Equal(t, 2, counts["ok|"+string(gate.OutcomeNewlyBad)]) + require.Equal(t, 3, counts["ok|"+string(gate.OutcomeStillFailing)]) + require.Equal(t, 1, counts["error|"+string(gate.OutcomeStillFailing)]) + require.Equal(t, 2, counts["ok|"+string(gate.OutcomeNewFailure)]) } diff --git a/grafana-alertcheck/docs/architecture.md b/grafana-alertcheck/docs/architecture.md index 455053163..fa2971178 100644 --- a/grafana-alertcheck/docs/architecture.md +++ b/grafana-alertcheck/docs/architecture.md @@ -15,11 +15,11 @@ This page documents the invariants and seams a maintainer must not break. It exi The gate must fail if it cannot get an answer. Every rule below is a specific instance of that: - **An error is never a pass.** A pass is exactly `len(Violations) == 0 && err == nil`. Every error path leaves `err` non-nil, and the CLI maps that to exit `2` unconditionally. -- **Inability beats violation.** Any `unobservable` rule is exit `2`, even alongside a real violation found first. +- **Inability beats violation.** Any `not_verified` rule is exit `2`, even alongside a real violation found first. - **Absent never means normal.** An instance that leaves the bad set is looked up in the *same* response: present as `normal` → cleared; absent (or `MissingSeries`) → vanished (a discontinuity, not a recovery). - **Staleness is absolute.** `grafana_now − lastEvaluation` is compared against a threshold, never "did it increase since the last poll" — a delta check reports stale on ~half the polls of a healthy rule (we poll at half of `intervalSeconds` of each rule). - **`grafana_now` is the response `Date` header.** Never the runner clock, in any comparison against a Grafana timestamp. -- **No early exit.** `check` collects to `to + transitionGrace` before classifying once. +- **An early exit can never be a pass.** `check` may stop collecting before `to + transitionGrace` (fail-fast), but only on a *monotone* terminal verdict: an inability that has already happened, or a post-`from` bad onset (which the full classifier would call `new_failure`/`unstable`). The one outcome that forgives an observed bad state, `recovered`, is reserved for bad-at-`from`, so a preexisting condition is never terminal. `--no-fail-fast` removes the guard entirely. - **No replay.** No run-id key, no artifact download, no state between attempts. A retry is a new piece of work and observation. ## The pure-function seam @@ -36,6 +36,15 @@ HTTP ──> Source ──> []StateRule ──> reduce ──> []Poll ──> pr - `Check`/`Watch` are I/O shells: HTTP, signals, the pidfile, file reads, the countdown print. The only test doubles needed are the `Source` and `Clock` interfaces. - `Policy` is the narrowed view of `Config` that reaches the pure layer — classification knobs and the window, no URL and no token. The token must never cross that line, which is the cheapest guarantee it never lands in an error string or a result. +## Early exit (fail-fast) + +By default `check` stops as soon as it knows the run cannot pass, rather than holding the runner for a window whose verdict is already decided. The guard is pure (`terminalVerdict`): it evaluates `proveCoverage` and `classifyRule` over the sub-window observed so far, with a synthetic sentinel, and only ever fires on a verdict that cannot become a pass. + +- In single-step mode `check` evaluates the guard after each in-process poll batch. +- In recorder mode the evidence lives in another process, so `check` **tails the recorder's log** while it waits, consuming complete newline-terminated records only. This is the one place the project reads a log a writer can still append to, and only because fail-fast wants an early answer, not the authoritative one — the strict, whole-file `ReadLog` still runs after the writer exits, and only its result is classified. + +On a terminal verdict the run is classified over `[from, At]` by the **same `decide`**, with the policy window clamped to `At` and the transition grace zeroed; `decide` is not forked. The requested window and the real thresholds are restored on the `Result`, which carries a `TerminatedEarly` marker. The drain wait is skipped — its question no longer applies. `--no-fail-fast` leaves the guard unset and restores the full-window behavior exactly. + ## Strict parsing as the version guard Both API responses are parsed strictly: a missing or unparseable **required** field (`health`, `state`, `lastEvaluation`, `interval`) is an error, never a zero value. Optional keys (`alerts`, `totals`, `labels`, `keepFiringFor`) are absent-tolerant, and unknown keys are ignored — so Grafana can add fields without breaking the parser, but removing one fails loudly. @@ -60,6 +69,8 @@ On a clean stop (SIGTERM/SIGINT/`--until`) the child finishes the in-flight writ `check` signals via the pidfile, waits for the **lock** to release (never the pid), and only then reads the log once. Reading while a writer can still append can only produce a shorter window than was recorded. +`stop` uses the same pidfile-plus-flock protocol, but as a cleanup operation rather than evidence-gathering: it SIGKILLs a recorder that ignores SIGTERM, removes the pidfile, and treats a missing pidfile as "nothing to stop". `check` inherits none of that — a writer that will not exit is a could-not-check, never a silent kill. + ## The log is the source of truth `watch` records raw evidence, so nothing trusts a state that could become unreachable. Two consequences a maintainer must preserve: diff --git a/grafana-alertcheck/docs/how-alerts-are-evaluated.md b/grafana-alertcheck/docs/how-alerts-are-evaluated.md index 9e09d4cbc..58ac89c14 100644 --- a/grafana-alertcheck/docs/how-alerts-are-evaluated.md +++ b/grafana-alertcheck/docs/how-alerts-are-evaluated.md @@ -32,13 +32,15 @@ For each instance the gate builds a timeline of bad spans over `[from, to]`, the | Outcome | Shape | Exit | | ------- | ----- | ---- | -| `clean` | Good throughout, observed throughout | → 0 | -| `newly_bad` | Entered a bad state **inside** the window | → 1 | -| `persistently_bad` | Bad at `from`, still bad at `to` | → 1 | +| `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 | -| `flapping` | Cleared, then became bad again | → 1 | -| `skipped` | Paused **before** the window opened | reported, not observable | -| `unobservable` | Coverage gap / sustained `health=error` / stale / absent | → 2 | +| `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 | + +A `--min-observed` deficit that no rule explains is reported as `not_counted`: it is not a verdict on any alert. `recovered` has **no deadline** — an alert that clears at minute 58 of a 60-minute window still passes. The total bad time is reported as `BadFor`. @@ -57,11 +59,11 @@ When an instance leaves the bad set, the gate looks it up **in the same response - Present as `normal` → `cleared` (a real recovery). - Absent, or present as `normal (MissingSeries)` → `vanished` (a discontinuity, **not** a recovery). -A vanished instance that was bad stays `persistently_bad`. A metric that stops being emitted is not evidence of health — this is deliberate and can surprise users whose fix is to remove a metric rather than drive it to a good value. +A vanished instance that was bad stays `still_failing`. A metric that stops being emitted is not evidence of health — this is deliberate and can surprise users whose fix is to remove a metric rather than drive it to a good value. ## Coverage proof -Before classifying, `check` must **prove** continuous coverage of `[from, to]` for each alert. Nine checks run; any failure makes the rule `unobservable`: +Before classifying, `check` must **prove** continuous coverage of `[from, to]` for each alert. Nine checks run; any failure makes the rule `not_verified`: 1. **Sentinel** — a clean recorder stop, timestamped at or after `to + transitionGrace`. A recorder that died mid-window looks exactly like a coverage gap and is one. 2. **`from` bounds** — `from` earlier than the recording start is unprovable. @@ -69,17 +71,28 @@ Before classifying, `check` must **prove** continuous coverage of `[from, to]` f 4. **`health=error`** — a contiguous run longer than `healthGrace` consumes coverage; a short blip is a note. 5. **`health=nodata`** — a note, never fatal (unless `--nodata-is-unobservable`). 6. **Liveness** — `grafana_now − lastEvaluation` must not exceed `evalStaleAfter`. This is an **absolute** check, never a "did it increase since the last poll" delta. -7. **In-window pause** — a poll reporting `isPaused` mid-window is `unobservable` (the primary pause detector). +7. **In-window pause** — a poll reporting `isPaused` mid-window is `not_verified` (the primary pause detector). 8. **Rule absent** — an authoritative `2xx` with no matching rule. 9. **`KeepLast`** — a note naming a stale-state blind spot. ## Health: `error` vs `nodata` -- `health=error` means the query **failed** — a malfunction. Sustained past `healthGrace`, it makes the rule `unobservable`. +- `health=error` means the query **failed** — a malfunction. Sustained past `healthGrace`, it makes the rule `not_verified`. - `health=nodata` means the query **ran and returned no series** — indistinguishable from a quiet system. It is not fatal by default; most of a fleet runs `no_data_state: OK`. +## Early exit (fail-fast) + +`check` does not have to wait for the whole window to know the run has failed. As soon as it observes a condition that cannot become a pass, it stops and classifies the sub-window it did see: + +- a **post-`from` bad onset** — the full classifier would call it `new_failure` (or `unstable`), which fails whether or not it later clears; or +- an **inability** — a heartbeat gap, a sustained `health=error` run, a stale evaluation, an in-window pause, or an absent rule. + +A preexisting bad instance is deliberately **not** terminal: if it clears before `to` the full run would call it `recovered`, which passes. + +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`. `--no-fail-fast` disables the guard and always waits for the full window and its coverage proof. + ## The drain wait and `transitionGrace` -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 `unobservable`. +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, and the grace is warned about when it exceeds a quarter of the window — the window may be too short for the alert's `for`. diff --git a/grafana-alertcheck/docs/index.md b/grafana-alertcheck/docs/index.md index e8ec76928..f9d8cac0a 100644 --- a/grafana-alertcheck/docs/index.md +++ b/grafana-alertcheck/docs/index.md @@ -50,6 +50,8 @@ grafana-alertcheck check --in /tmp/run.jsonl --from "$deployed_at" --to "$finish `watch` returns only after the recorder has observed every named, non-paused alert once and reported ready — so auth, name-resolution, and parse failures surface **before** your deploy runs. +If your work fails before `check` runs and the alert verdict no longer matters, reap the recorder with `grafana-alertcheck stop --out /tmp/run.jsonl`. It is idempotent, so it is safe as an `if: always()` step: after `check` has already stopped the recorder it is a no-op. + ## Quickstart — single-step mode Skip the recorder and observe the window inline, from inside `check` itself: @@ -78,7 +80,7 @@ An error is never a pass: `2` wins over any violation found alongside it. - `recovered` has **no deadline** — a bad-at-`from` alert that clears by `to` passes; set `--preexisting fail` to forbid it. - If you retry the check, then the work also needs to be retried - there is no way to check the past. - `watch` and `check` must run in **one job, one runner, one filesystem** — nothing persists across jobs or attempts. -- The gate **never exits early** — a violation at minute 2 still holds the runner to `to + transitionGrace + drainTimeout`; size the job timeout to the planned run time the gate prints at start. +- The gate **stops early on a certain failure** — as soon as a post-`from` bad onset or an inability is observed, `check` returns instead of holding the runner to `to + transitionGrace + drainTimeout`. This can never turn into a pass, but it can report exit `1` where a full run would have reported exit `2` (inability beats violation only when the inability is observed). Pass `--no-fail-fast` to always wait for the full window and its coverage proof; size the job timeout to the planned run time the gate prints at start either way. ## More diff --git a/grafana-alertcheck/docs/reference/cli.md b/grafana-alertcheck/docs/reference/cli.md index bd4c95718..312ea4860 100644 --- a/grafana-alertcheck/docs/reference/cli.md +++ b/grafana-alertcheck/docs/reference/cli.md @@ -9,7 +9,7 @@ description: "Full reference for the grafana-alertcheck CLI: watch, check, list, # CLI reference ``` -grafana-alertcheck +grafana-alertcheck ``` Connection details are always from the environment: `GRAFANA_URL` and `GRAFANA_TOKEN`. The token is never a flag and never logged. @@ -42,12 +42,29 @@ grafana-alertcheck watch --out [--pidfile F] [--daemon-log F] \ `watch` writes the header, observes every non-paused rule once, checks the budget, then detaches a background recorder and returns. Recording is **unfiltered** — there is no `--states` here, so the same log can be re-classified later under different `--states` without re-recording. +## `stop` — reap the recorder + +```bash +grafana-alertcheck stop --out [--pidfile F] +``` + +| Flag | Default | Meaning | +| ---- | ------- | ------- | +| `--out` | — | JSONL log path whose recorder to stop (required) | +| `--pidfile` | `.pid` | Pidfile of the recorder to stop | + +Stops a detached recorder that is still running, or confirms it has already finished. It reads the pidfile, then asks the log's **flock** whether a writer exists right now (the lock is authoritative; a pid can be reused): if a writer is alive it is sent `SIGTERM`, and if it ignores that it is killed, then the pidfile is removed. + +It is **idempotent** — after `check` has already stopped the recorder, or after a previous `stop`, it reports that there is nothing to stop and exits `0`. That is what lets an `if: always()` step call it on both the success and failure paths. + +Use it when the work failed and the alert verdict no longer matters, but the recorder must still be reaped: the recorder is detached in its own session, so neither `check` nor the runner's cleanup will stop it, and it would keep polling Grafana until its window elapsed. + ## `check` — classify ```bash grafana-alertcheck check [--in ] [--pidfile F] --from RFC3339 --to RFC3339 \ [--alerts ...] [--folder F] [--states ...] [--preexisting ...] [--min-observed N] \ - [--allow-paused] [--nodata-is-unobservable] [--concurrency N] [--output json] + [--allow-paused] [--nodata-is-unobservable] [--no-fail-fast] [--concurrency N] [--output json] ``` | Flag | Default | Meaning | @@ -62,9 +79,12 @@ grafana-alertcheck check [--in ] [--pidfile F] --from RFC3339 --to RFC3339 | `--min-observed` | every resolved rule | Minimum rules that must be observed | | `--allow-paused` | `false` | Don't count pre-window-paused rules against `--min-observed` | | `--nodata-is-unobservable` | `false` | Treat sustained `health=nodata` as unobservable | +| `--no-fail-fast` | `false` | Wait for the full window even after a certain failure | | `--concurrency` | `1` | Max concurrent requests | | `--output` | `table` | `json` also writes the machine-readable result to stdout | +By default `check` **exits early** on a failure that cannot become a pass: a post-`from` bad onset, or an inability (a heartbeat gap, a sustained `health=error`, a stale evaluation, an in-window pause, an absent rule). This is a latency optimization, not a weaker gate — it never exits `0` early. The one observable difference is that an early exit can report `1` where a full run would have discovered an inability later and reported `2`. `--no-fail-fast` always waits for `to + transitionGrace` and the full coverage proof; the `Result` then carries no `terminated_early` marker. With early exit the JSON result includes `terminated_early` naming the rule, kind, reason and time. + `--from` and `--to` are RFC3339 with an explicit offset and must come from your work — `from` from the deploy step, `to` from the step that finishes. In recorder mode an absent `--from` is a hard error; in single-step mode it falls back (with a warning) to the start of the step. ## Naming alerts @@ -82,7 +102,7 @@ Datasource-managed and recording rules are refused with a specific error. A name ## Output and exit codes -The human table goes to **stderr**: `RESULTS` (one row per rule), `VIOLATIONS` (one per distinct rule/outcome/state/health/note signature, with a `COUNT` of the instances it stands for — instance identity is only in the JSON), and `THRESHOLDS` (each rule's `maxGap`/`healthGrace`/`evalStaleAfter` plus global `transitionGrace`/`drainTimeout` and the largest measured clock skew). `--output json` writes the result to stdout. +The human table goes to **stderr**: `RESULTS` (one row per rule, with the verdict, time broken, check cadence and whether the window was observed), `VIOLATIONS` (one per distinct rule/verdict/state/health/note signature, with an `INSTANCES` count of the instances it stands for — instance identity is only in the JSON), and `LIMITS USED` (each rule's observation limits in plain words, explained by a legend under the table, plus the extra observation time, the evaluation wait and the largest measured clock difference). The JSON outcome values are `healthy`, `new_failure`, `still_failing`, `recovered`, `unstable`, `paused`, `not_verified` and the synthetic `not_counted`. `--output json` writes the result to stdout. | Code | Meaning | | ---- | ------- | diff --git a/grafana-alertcheck/docs/reference/log-format.md b/grafana-alertcheck/docs/reference/log-format.md index 9dc27c29c..0bf1830af 100644 --- a/grafana-alertcheck/docs/reference/log-format.md +++ b/grafana-alertcheck/docs/reference/log-format.md @@ -49,7 +49,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. -- `is_paused` records the pause state at record start (the moment `skipped` means). +- `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`. - `for_seconds`, `interval_seconds`, `no_data_state`, `exec_err_state` are forensic only — `check` re-resolves definitions and never reads them back. @@ -95,4 +95,4 @@ Instance keys are the JSON encoding of the labels map (with stable key order), s { "type": "stopped", "at": "2026-09-07T10:10:30Z" } ``` -`at` is the recorder's own stop time. `check` compares it against `to + transitionGrace`; absent or earlier is `unobservable` — never a pass. +`at` is the recorder's own stop time. `check` compares it against `to + transitionGrace`; absent or earlier is `not_verified` — never a pass. diff --git a/grafana-alertcheck/internal/gate/check.go b/grafana-alertcheck/internal/gate/check.go index cc83fb725..9faa8f1f1 100644 --- a/grafana-alertcheck/internal/gate/check.go +++ b/grafana-alertcheck/internal/gate/check.go @@ -21,16 +21,6 @@ import ( // for; a silent wait is indistinguishable from a hung process. const countdownEvery = 30 * time.Second -// recorderStopTimeout bounds the wait for the recorder's exit after SIGTERM. -// Everything after the signal is local (finish the in-flight write, sentinel, -// fsync), so this is loose; it stays a hard error because a log a writer still -// holds cannot be read. -const recorderStopTimeout = 30 * time.Second - -// recorderStopPoll is how often the wait re-checks the lock. With no wait(2) -// on a detached session leader, its exit is observable only by polling. -const recorderStopPoll = 100 * time.Millisecond - // Config is check's whole input. It is the CLI's view of a run, and it is // deliberately wider than Policy: Policy is the narrowed, pure-layer subset // that reaches decide (classify.go), and the token is the field that must @@ -53,6 +43,11 @@ type Config struct { AllowPaused bool NodataIsUnobservable bool + // NoFailFast disables early exit, so the loop always runs to + // to+transitionGrace for a full-window coverage proof. The CLI sets it + // with --no-fail-fast; the gate never exits 0 early either way. + NoFailFast bool + // From is the moment the deploy finished and To is the end of the work. // They are different moments and both come from the work. In recorder mode // an absent From is a hard error; in single-step mode it falls back to the @@ -337,22 +332,66 @@ func check(ctx context.Context, cfg Config, src Source) (Result, error) { } } + // The fail-fast guard needs a header now; the authoritative ReadLog + // overwrites this advisory one after the writer exits. + if logHasHdr { + header = earlyHdr + } + + // Fixed once, before collection, so the guard and decide share one policy. + pol := Policy{ + States: cfg.States, + Preexisting: cfg.Preexisting, + MinObserved: minObserved, + AllowPaused: cfg.AllowPaused, + NodataIsUnobservable: cfg.NodataIsUnobservable, + From: from, + To: cfg.To, + } + // ---- Collect the evidence. -------------------------------------------- - // Collect ONLY. No classification happens here and there is no early exit, - // even once a violation is certain: the loop always runs to - // to + transitionGrace, which is what makes "did the early exit lose the - // coverage proof?" a question that cannot be asked. + // The loop never classifies; its only access to the pure layer is the + // fail-fast guard, which ends the run only on a condition that cannot + // become a pass. With --no-fail-fast there is no guard. windowEnd := cfg.To.Add(gt.transitionGrace) - var poller *livePoller + var ( + poller *livePoller + tail *logTailer + guard terminalCheck + ) if !logHasHdr { poller = newLivePoller(src, reducer, activeRules(resolved), rt, cfg.Concurrency, cfg.Clock.Now()) } - collected, err := collectUntil(ctx, cfg, windowEnd, poller) + failFast := !cfg.NoFailFast + if failFast { + guard = func(polls []Poll, at time.Time) (Termination, bool) { + return terminalVerdict(header, polls, resolved, rt, pol, from, at) + } + if logHasHdr { + // In recorder mode check otherwise only sleeps, so the guard reads + // the recorder's log as it lands. ReadLog is still the evidence. + tail, err = newLogTailer(cfg.Log) + if err != nil { + return Result{}, err + } + defer tail.Close() + } + } + collected, err := collectUntil(ctx, cfg, windowEnd, initial, poller, tail, guard) 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() + } + } // 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), err) + return Result{}, fmt.Errorf("collect evidence after %d poll(s): %w", len(collected.polls), err) } var ( @@ -387,15 +426,26 @@ func check(ctx context.Context, cfg Config, src Source) (Result, error) { return Result{}, fmt.Errorf("log identity: %w", err) } } else { - polls = make([]Poll, 0, len(initial)+len(collected)) - polls = append(polls, initial...) - polls = append(polls, collected...) - // The shell stamps the sentinel itself, when the collection loop - // exits: by construction that is at or after to + transitionGrace, so - // the sentinel check passes for the same reason a clean recorder stop - // does, and for no other. - stoppedAt := cfg.Clock.Now() - sentinel = &stoppedAt + // Seeded with the measurement pass, so this is initial + the loop's polls. + polls = collected.polls + if collected.term == nil { + // The loop ran to to+transitionGrace, so a sentinel stamped now + // proves it. A fail-fast exit has no such claim; earlyResult + // supplies its own. + stoppedAt := cfg.Clock.Now() + sentinel = &stoppedAt + } + } + + // ---- Fail-fast exit. -------------------------------------------------- + // The window is not proved, so earlyResult classifies the observed + // sub-window, preserves the non-zero verdict, and skips the drain wait. + if collected.term != nil { + term := *collected.term + fmt.Fprintf(cfg.Notes, "note: fail-fast: %s (%s at %s); stopping before the window closed\n", + term.Kind, term.Alert, term.At.Format(time.RFC3339)) + result, err := earlyResult(header, polls, resolved, rt, gt, pol, term) + return result, err } // ---- The drain wait. -------------------------------------------------- @@ -409,15 +459,6 @@ func check(ctx context.Context, cfg Config, src Source) (Result, error) { } // ---- Classify. -------------------------------------------------------- - pol := Policy{ - States: cfg.States, - Preexisting: cfg.Preexisting, - MinObserved: minObserved, - AllowPaused: cfg.AllowPaused, - NodataIsUnobservable: cfg.NodataIsUnobservable, - From: from, - To: cfg.To, - } result, decideErr := decide(header, polls, sentinel, resolved, rt, gt, pol) result, drainErr := mergeDrainTimeouts(result, drained) @@ -531,22 +572,38 @@ func (p *livePoller) poll(ctx context.Context, uids []string) ([]Poll, error) { return out, obsErr } -// collectUntil is the collection loop, shared by both modes. With a poller it -// polls each rule on its own cadence; with nil it only waits, because in -// recorder mode the evidence is being written by another process. Both print -// the same countdown, because both are the same silence to an operator -// watching a job. -// -// It never classifies and never exits early. -func collectUntil(ctx context.Context, cfg Config, deadline time.Time, p *livePoller) ([]Poll, error) { - var ( - polls []Poll - lastPrint time.Time - ) +// tailPollInterval bounds how long a fail-fast condition can sit unnoticed in +// recorder mode. The read is a cheap seek-and-append from the last offset. +const tailPollInterval = time.Second + +// terminalCheck is the fail-fast guard: polls observed so far plus runner time +// in, a terminal verdict out. nil disables fail-fast. +type terminalCheck func(polls []Poll, at time.Time) (Termination, bool) + +// collection is what collectUntil observed. polls is the evidence in +// single-step mode and the tailed subset in recorder mode, where ReadLog +// re-reads the authoritative set after the writer exits. term is set only when +// the guard ended the loop early. +type collection struct { + polls []Poll + sentinel *time.Time + term *Termination +} + +// collectUntil is the collection loop. With a poller it polls each rule on its +// own cadence; in recorder mode it tails the log the recorder is writing. A +// guard is evaluated after each observation and ends the loop early on a +// condition that cannot become a pass. +func collectUntil(ctx context.Context, cfg Config, deadline time.Time, seed []Poll, + p *livePoller, tail *logTailer, guard terminalCheck) (collection, error) { + + out := collection{polls: seed} + var lastPrint time.Time + for { now := cfg.Clock.Now() if !now.Before(deadline) { - return polls, nil + return out, nil } if lastPrint.IsZero() || now.Sub(lastPrint) >= countdownEvery { fmt.Fprintf(cfg.Notes, "collecting: %s until the window closes at %s\n", @@ -563,111 +620,56 @@ func collectUntil(ctx context.Context, cfg Config, deadline time.Time, p *livePo due := p.sched.Due(now) for _, uid := range due { if merr := p.sched.Mark(uid, now); merr != nil { - return polls, fmt.Errorf("mark %s: %w", uid, merr) + return out, fmt.Errorf("mark %s: %w", uid, merr) } } if len(due) > 0 { batch, err := p.poll(ctx, due) - polls = append(polls, batch...) + out.polls = append(out.polls, batch...) if err != nil { - return polls, err + return out, err } } if next, ok := p.sched.earliestDue(); ok { wait = min(wait, next.Sub(cfg.Clock.Now())) } } + if tail != nil { + batch, sentinel, err := tail.read() + if err != nil { + return out, err + } + out.polls = append(out.polls, batch...) + if sentinel != nil { + out.sentinel = sentinel + } + wait = min(wait, tailPollInterval) + } + + if guard != nil { + if term, ok := guard(out.polls, cfg.Clock.Now()); ok { + out.term = &term + return out, nil + } + } select { case <-ctx.Done(): - return polls, ctx.Err() + return out, ctx.Err() case <-cfg.Clock.After(max(wait, 0)): } } } -// stopRecorder signals the recorder and waits for it to go; the log may not be -// read until the writer has provably gone, so every failure is a hard error. -// It returns the log held under an exclusive flock, which the caller must keep -// open across ReadLog — the lock is the proof that no writer exists. -// -// Two authorities, only one of which is evidence: -// -// - the PIDFILE says whether a recording ever started (it is written only -// after the child reports ready, and removed on failure). -// - the FLOCK says whether a writer exists right now. A pidfile can go stale -// — nothing removes it on a clean --until stop, so it may name a pid -// somebody else now owns — but the kernel drops a flock when the holder -// exits, so the lock is always authoritative. -// -// So: read the pidfile to learn a recording happened, ask the lock whether it -// is still running, and signal only if it is. +// stopRecorder is check's use of StopRecorder: no cleanup semantics, and the +// log it returns must stay held across ReadLog. func stopRecorder(ctx context.Context, cfg Config) (*os.File, error) { - pid, err := ReadPidFile(cfg.PidFile) - if err != nil { - return nil, fmt.Errorf("cannot stop the recorder: %w; a pidfile is written only once a recorder reports that it is running, so an unreadable one means the recording never started", err) - } - - log, err := os.Open(cfg.Log) - if err != nil { - return nil, fmt.Errorf("open %s to check for a writer: %w", cfg.Log, err) - } - - held, err := tryLockExclusive(log) - if err != nil { - log.Close() - return nil, err - } - if held { - // No writer. Send no signal, whatever the pidfile says — the pid may - // belong to somebody else entirely by now. Which of --until, a clean - // stop and a death ended the recording is the sentinel's question, - // answered by the coverage proof over the log this unblocks. - fmt.Fprintf(cfg.Notes, "note: no writer holds %s; the recorder has already finished\n", cfg.Log) - return log, nil - } - - // The lock is held, so a writer is alive and the pidfile's pid cannot be - // stale — the recorder that took the lock is the one the parent recorded. - gone, err := signalRecorder(pid) - if err != nil { - log.Close() - return nil, err - } - if gone { - // A live writer holds the log and the pidfile names a process that - // does not exist. That is a broken contract, not a case to reason - // around: signalling the real holder would mean guessing who it is. - log.Close() - return nil, fmt.Errorf("a writer holds %s but pidfile %s names pid %d, which does not exist: the pidfile does not name the process that holds the log", - cfg.Log, cfg.PidFile, pid) - } - - // Wait on the LOCK, not on the pid: its release is the kernel-guaranteed - // writer-is-gone event, and it carries no pid-reuse hazard. - deadline := cfg.Clock.Now().Add(recorderStopTimeout) - for { - select { - case <-ctx.Done(): - log.Close() - return nil, ctx.Err() - case <-cfg.Clock.After(recorderStopPoll): - } - - held, err := tryLockExclusive(log) - if err != nil { - log.Close() - return nil, err - } - if held { - return log, nil - } - if !cfg.Clock.Now().Before(deadline) { - log.Close() - return nil, fmt.Errorf("recorder pid %d still holds %s %s after SIGTERM; refusing to read a log a writer can still append to", - pid, cfg.Log, recorderStopTimeout) - } - } + return StopRecorder(ctx, StopConfig{ + Log: cfg.Log, + PidFile: cfg.PidFile, + Clock: cfg.Clock, + Notes: cfg.Notes, + }) } // drainVerdict is what the drain wait concluded about one rule it could not @@ -683,13 +685,13 @@ type drainVerdict struct { // drainWait is the final liveness check: did each rule evaluate through the // end of the window? A rule that cannot answer within drainTimeout is -// unobservable, never a pass. It returns one verdict per rule it could not +// not_verified, never a pass. It returns one verdict per rule it could not // clear (keyed by UID); an error only for a hard failure of the wait itself. // // Two kinds of rule are excluded up front because draining them could not // change a verdict: a rule the HEADER says was paused at the window open (the // header, not the late-resolved definitions — see Header.pausedAtStart), and a -// rule whose last poll says Found == false (already unobservable via rule_absent). +// rule whose last poll says Found == false (already not_verified via rule_absent). func drainWait(ctx context.Context, cfg Config, src Source, defs []Definition, pausedAtStart map[string]bool, rt map[string]RuleTimings, polls []Poll, windowEnd time.Time, timeout time.Duration) (map[string]drainVerdict, error) { @@ -860,17 +862,17 @@ func mergeDrainTimeouts(res Result, drained map[string]drainVerdict) (Result, er cov.Notes = append(cov.Notes, verdict.note) res.Coverage[uid] = cov - if res.Verdicts[i].Outcome != OutcomeUnobservable { + if res.Verdicts[i].Outcome != OutcomeNotVerified { names = append(names, fmt.Sprintf("%s (%s)", res.Verdicts[i].Alert, verdict.reason)) } - res.Verdicts[i].Outcome = OutcomeUnobservable + res.Verdicts[i].Outcome = OutcomeNotVerified res.Verdicts[i].Note = strings.Join(cov.Notes, "; ") } if len(names) == 0 { - // Every drained rule was already unobservable for an earlier reason, + // Every drained rule was already not_verified for an earlier reason, // so decide's own error already stops the run. Adding a second error // saying the same thing would only make the message longer. return res, nil } - return res, fmt.Errorf("gate: %d rule(s) unobservable at the drain wait: %s", len(names), strings.Join(names, "; ")) + return res, fmt.Errorf("gate: %d rule(s) not verified at the drain wait: %s", len(names), strings.Join(names, "; ")) } diff --git a/grafana-alertcheck/internal/gate/check_process.go b/grafana-alertcheck/internal/gate/check_process.go index 86dd73148..9deca1993 100644 --- a/grafana-alertcheck/internal/gate/check_process.go +++ b/grafana-alertcheck/internal/gate/check_process.go @@ -27,3 +27,12 @@ func signalRecorder(pid int) (gone bool, err error) { return false, fmt.Errorf("signal recorder pid %d: %w", pid, err) } } + +// killRecorder is the cleanup-only counterpart to signalRecorder: SIGKILL to a +// recorder that ignored SIGTERM. An already-gone pid is not an error. +func killRecorder(pid int) error { + if err := syscall.Kill(pid, syscall.SIGKILL); err != nil && !errors.Is(err, syscall.ESRCH) { + return fmt.Errorf("kill recorder pid %d: %w", pid, err) + } + return nil +} diff --git a/grafana-alertcheck/internal/gate/check_test.go b/grafana-alertcheck/internal/gate/check_test.go index 0fdc6dffa..2dc6cb1a6 100644 --- a/grafana-alertcheck/internal/gate/check_test.go +++ b/grafana-alertcheck/internal/gate/check_test.go @@ -268,7 +268,7 @@ func TestCheckSingleStepCleanWindowPasses(t *testing.T) { // A pass is exactly this shape. require.Empty(t, res.Violations) require.Len(t, res.Verdicts, 1) - require.Equal(t, OutcomeClean, res.Verdicts[0].Outcome) + require.Equal(t, OutcomeHealthy, res.Verdicts[0].Outcome) cov := res.Coverage[checkUID] require.True(t, cov.Proved) require.False(t, cov.Unobservable) @@ -340,7 +340,7 @@ func TestCheckSingleStepContinuousHealthErrorIsUnobservable(t *testing.T) { res, err := check(context.Background(), cfg, src) require.Error(t, err, "continuous health=error must be unobservable") require.Len(t, res.Verdicts, 1) - require.Equal(t, OutcomeUnobservable, res.Verdicts[0].Outcome) + require.Equal(t, OutcomeNotVerified, res.Verdicts[0].Outcome) require.Equal(t, ReasonHealthError, res.Coverage[def.UID].Reason) } @@ -361,15 +361,13 @@ func TestCheckSingleStepFiringInstanceReportsWithoutExitingEarly(t *testing.T) { res, err := check(context.Background(), cfg, src) require.NoError(t, err, "a violation is exit 1, not an error") require.Len(t, res.Violations, 1) - require.Equal(t, OutcomePersistentlyBad, res.Violations[0].Outcome) + require.Equal(t, OutcomeStillFailing, res.Violations[0].Outcome) require.False(t, clock.Now().Before(cfg.To.Add(checkGrace)), "exited early; collection must run to to+grace") } -// A newly_bad instance at from+30s gives exit 1, but ONLY after -// to+transitionGrace. The test above covers a rule already bad before the -// window opened (persistently_bad); this covers a fresh onset just inside the -// window, which must not release the runner the instant it is first observed. -func TestCheckSingleStepNewOnsetDoesNotExitEarly(t *testing.T) { +// A fresh onset ends the run early, well before to+transitionGrace, and still +// reports the exit-1 shape with the requested window preserved. +func TestCheckSingleStepNewOnsetExitsEarlyByDefault(t *testing.T) { clock := newVirtualClock(testNow) cfg := baseConfig(t, clock) onset := testNow.Add(30 * time.Second) @@ -386,12 +384,68 @@ func TestCheckSingleStepNewOnsetDoesNotExitEarly(t *testing.T) { return healthyObservation(now, firing), nil }) + res, err := check(context.Background(), cfg, src) + require.NoError(t, err, "a violation is exit 1, not an error") + require.Len(t, res.Violations, 1) + require.Equal(t, OutcomeNewFailure, res.Violations[0].Outcome) + require.True(t, clock.Now().Before(cfg.To.Add(checkGrace)), + "the runner was released only after to+grace; fail-fast did not fire") + require.NotNil(t, res.TerminatedEarly) + require.Equal(t, TerminationViolation, res.TerminatedEarly.Kind) + require.Equal(t, OutcomeNewFailure, res.TerminatedEarly.Outcome) + require.True(t, res.To.Equal(cfg.To), "the requested window is still reported") + require.Contains(t, notesOf(cfg), "fail-fast") +} + +// A preexisting bad instance can still become `recovered` (a pass), so it is +// never terminal even with fail-fast on. +func TestCheckSingleStepPreexistingBadDoesNotExitEarly(t *testing.T) { + clock := newVirtualClock(testNow) + cfg := baseConfig(t, clock) + firing := Instance{ + Labels: map[string]string{"alertname": checkTitle, "instance": "a"}, + State: StateFiring, + ActiveAt: testNow.Add(-10 * time.Minute), // bad before the window opened + } + src := newCheckSource(func(_ string, _ int) (Observation, error) { + return healthyObservation(clock.Now(), firing), nil + }) + + 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.Nil(t, res.TerminatedEarly, "a preexisting condition can still recover, so it is not terminal") + require.False(t, clock.Now().Before(cfg.To.Add(checkGrace)), + "exited early; collection must run to to+grace for a preexisting bad instance") +} + +// --no-fail-fast runs the same fresh onset to to+transitionGrace. +func TestCheckSingleStepNewOnsetNoFailFastRunsToTheEnd(t *testing.T) { + clock := newVirtualClock(testNow) + cfg := baseConfig(t, clock) + cfg.NoFailFast = true + onset := testNow.Add(30 * time.Second) + src := newCheckSource(func(_ string, _ int) (Observation, error) { + now := clock.Now() + if now.Before(onset) { + return healthyObservation(now), nil + } + firing := Instance{ + Labels: map[string]string{"alertname": checkTitle, "instance": "a"}, + State: StateFiring, + ActiveAt: onset, + } + return healthyObservation(now, firing), nil + }) + res, err := check(context.Background(), cfg, src) require.NoError(t, err) require.Len(t, res.Violations, 1) - require.Equal(t, OutcomeNewlyBad, res.Violations[0].Outcome) + require.Equal(t, OutcomeNewFailure, res.Violations[0].Outcome) + require.Nil(t, res.TerminatedEarly) require.False(t, clock.Now().Before(cfg.To.Add(checkGrace)), - "exited early; collection must run to to+grace even for a fresh onset at from+30s") + "with --no-fail-fast the loop must run to to+grace") } // An ABSENT `from` in single-step mode (as opposed to recorder mode, which @@ -596,13 +650,82 @@ func TestCheckRecorderModeCleanWindowPasses(t *testing.T) { require.NoError(t, err) require.Empty(t, res.Violations) require.Len(t, res.Verdicts, 1) - require.Equal(t, OutcomeClean, res.Verdicts[0].Outcome) + require.Equal(t, OutcomeHealthy, res.Verdicts[0].Outcome) require.Equal(t, "13.1.0", res.GrafanaVersion) // The collection loop still waited out to+transitionGrace even though the // recorder had already finished. require.False(t, clock.Now().Before(windowEnd)) } +// recordedLogWithOnset records a full window where the rule fires from onset on. +func recordedLogWithOnset(t *testing.T, dir string, onset time.Time) string { + t.Helper() + windowEnd := testNow.Add(5*time.Minute + checkGrace) + path := filepath.Join(dir, "log.jsonl") + w, err := NewWriter(path, newFakeClock(windowEnd.Add(30*time.Second))) + require.NoError(t, err) + require.NoError(t, w.WriteHeader(Header{ + URL: "https://grafana.example.com", GrafanaVersion: "13.1.0", StartedAt: testNow.Add(-time.Minute), + Rules: []LoggedRule{{ + UID: checkUID, Title: checkTitle, Folder: "F", Group: "G", + IntervalSeconds: 60, NoDataState: "OK", ExecErrState: "OK", + PollEverySeconds: checkPollEvery.Seconds(), + }}, + })) + for at := testNow.Add(-time.Minute); !at.After(windowEnd.Add(30 * time.Second)); at = at.Add(checkPollEvery) { + p := Poll{RuleUID: checkUID, GrafanaNow: at, Found: true, State: "inactive", Health: "ok", LastEvaluation: at} + if !at.Before(onset) { + p.State = "firing" + p.Abnormal = []Instance{{ + Labels: map[string]string{"alertname": checkTitle, "instance": "a"}, + State: StateFiring, + ActiveAt: onset, + }} + } + require.NoError(t, w.WritePoll(p)) + } + require.NoError(t, w.Stop()) + return path +} + +// Recorder-mode fail-fast reads the recorder's log as it lands, ending the wait +// on an onset inside the window. No live source is polled. +func TestCheckRecorderModeExitsEarlyOnANewOnset(t *testing.T) { + onset := testNow.Add(time.Minute) + logPath := recordedLogWithOnset(t, t.TempDir(), onset) + writePid(t, logPath+".pid", fmt.Sprintf("%d\n", deadPid(t))) + + clock := newVirtualClock(testNow) + cfg := recorderConfig(t, clock, logPath) + res, err := check(context.Background(), cfg, newCheckSource(nil)) + require.NoError(t, err) + require.NotNil(t, res.TerminatedEarly) + require.Equal(t, TerminationViolation, res.TerminatedEarly.Kind) + require.Equal(t, OutcomeNewFailure, res.TerminatedEarly.Outcome) + require.Len(t, res.Violations, 1) + require.True(t, clock.Now().Before(cfg.To.Add(checkGrace)), + "the runner was held to the window; fail-fast did not fire") + require.True(t, res.To.Equal(cfg.To), "the requested window is still reported") +} + +// --no-fail-fast over the same recording runs to to+transitionGrace. +func TestCheckRecorderModeNoFailFastRunsToTheEnd(t *testing.T) { + onset := testNow.Add(time.Minute) + logPath := recordedLogWithOnset(t, t.TempDir(), onset) + writePid(t, logPath+".pid", fmt.Sprintf("%d\n", deadPid(t))) + + clock := newVirtualClock(testNow) + cfg := recorderConfig(t, clock, logPath) + cfg.NoFailFast = true + res, err := check(context.Background(), cfg, newCheckSource(nil)) + require.NoError(t, err) + require.Nil(t, res.TerminatedEarly) + require.Len(t, res.Violations, 1) + require.Equal(t, OutcomeNewFailure, res.Violations[0].Outcome) + require.False(t, clock.Now().Before(cfg.To.Add(checkGrace)), + "with --no-fail-fast the loop must run to to+grace") +} + // The identity of the log is not correct. The check runs against the header // read EARLY, so it fails before the window's wait rather than after it. func TestCheckFailClosedOnWrongLogIdentity(t *testing.T) { @@ -680,7 +803,7 @@ func TestCheckRecorderModeFromSameSecondAsStartedAtPasses(t *testing.T) { require.NoError(t, err) require.Empty(t, res.Violations) require.Len(t, res.Verdicts, 1) - require.Equal(t, OutcomeClean, res.Verdicts[0].Outcome) + require.Equal(t, OutcomeHealthy, res.Verdicts[0].Outcome) } // The coverage proof failed: a hole in the middle of the recording is not @@ -712,7 +835,7 @@ func TestCheckFailClosedOnCoverageGap(t *testing.T) { res, err := check(context.Background(), cfg, newCheckSource(nil)) require.Error(t, err, "the coverage gap to fail closed") require.Equal(t, ReasonHeartbeatGap, res.Coverage[checkUID].Reason) - require.Equal(t, OutcomeUnobservable, res.Verdicts[0].Outcome) + require.Equal(t, OutcomeNotVerified, res.Verdicts[0].Outcome) } // An episode fully between the deploy and the start of the check: recorder @@ -745,7 +868,7 @@ func TestCheckRecorderModeFindsAGapImmediatelyAfterTheDeploy(t *testing.T) { res, err := check(context.Background(), cfg, newCheckSource(nil)) require.Error(t, err, "a hole right after the deploy hides whatever happened there") require.Len(t, res.Verdicts, 1) - require.Equal(t, OutcomeUnobservable, res.Verdicts[0].Outcome, "never clean") + require.Equal(t, OutcomeNotVerified, res.Verdicts[0].Outcome, "never clean") } // The drain limit passed. The recording itself is clean, so this isolates the @@ -774,7 +897,7 @@ func TestCheckFailClosedOnDrainTimeout(t *testing.T) { res, err := check(context.Background(), cfg, src) require.Error(t, err, "the drain limit to fail closed") require.Equal(t, ReasonDrainTimeout, res.Coverage[checkUID].Reason) - require.Equal(t, OutcomeUnobservable, res.Verdicts[0].Outcome) + require.Equal(t, OutcomeNotVerified, res.Verdicts[0].Outcome) require.Contains(t, res.Verdicts[0].Note, "drain limit") require.GreaterOrEqual(t, clock.Now().Sub(windowEnd), checkDrainLimit, "the rule never evaluates through the window, so the drain wait must run its full limit") @@ -873,10 +996,10 @@ func pausedAfterWindowCheck(t *testing.T, allowPaused bool) (Result, Config, err func TestCheckPausingARuleAfterTheWindowDoesNotMakeItSkipped(t *testing.T) { res, _, err := pausedAfterWindowCheck(t, false) require.NoError(t, err) - require.Equal(t, OutcomeNewlyBad, res.Verdicts[0].Outcome, + require.Equal(t, OutcomeNewFailure, res.Verdicts[0].Outcome, "the rule was active for the whole window and fired inside it") require.Len(t, res.Violations, 1) - require.Equal(t, OutcomeNewlyBad, res.Violations[0].Outcome) + require.Equal(t, OutcomeNewFailure, res.Violations[0].Outcome) require.NotContains(t, res.Verdicts[0].Note, "paused before the window opened") } @@ -924,7 +1047,7 @@ func TestCheckHeaderPausedRuleStaysSkipped(t *testing.T) { res, err := run(false) require.NoError(t, err, "a skipped rule is a known condition, not an inability") - require.Equal(t, OutcomeSkipped, res.Verdicts[0].Outcome) + require.Equal(t, OutcomePaused, res.Verdicts[0].Outcome) _, ok := res.Coverage[checkUID] require.False(t, ok, "a skipped rule has no coverage to prove") require.Len(t, res.Violations, 1, "the MinObserved shortfall") @@ -1091,7 +1214,7 @@ func TestCheckDeadPidWithNoSentinelIsUnobservable(t *testing.T) { res, err := check(context.Background(), cfg, newCheckSource(nil)) require.Error(t, err, "no sentinel means the recorder never proved it ran to the end") require.Len(t, res.Verdicts, 1) - require.Equal(t, OutcomeUnobservable, res.Verdicts[0].Outcome) + require.Equal(t, OutcomeNotVerified, res.Verdicts[0].Outcome) } // An incomplete last line gives exit 2. log_test.go's TestReadLogRejectsBadLogs @@ -1209,8 +1332,8 @@ func TestMergeDrainTimeoutsNamesEveryUnobservableRule(t *testing.T) { "b": {Unobservable: true, Reason: ReasonHeartbeatGap, Notes: []string{"rule \"B\": gap"}}, }, Verdicts: []RuleVerdict{ - {Alert: "A", RuleUID: "a", Outcome: OutcomeClean}, - {Alert: "B", RuleUID: "b", Outcome: OutcomeUnobservable}, + {Alert: "A", RuleUID: "a", Outcome: OutcomeHealthy}, + {Alert: "B", RuleUID: "b", Outcome: OutcomeNotVerified}, }, } @@ -1218,9 +1341,9 @@ func TestMergeDrainTimeoutsNamesEveryUnobservableRule(t *testing.T) { "a": {reason: ReasonDrainTimeout, note: "rule \"A\": did not evaluate through the end within the drain limit"}, "b": {reason: ReasonDrainTimeout, note: "rule \"B\": did not evaluate through the end within the drain limit"}, }) - require.Error(t, err, "naming the newly unobservable rule") - require.Contains(t, err.Error(), "unobservable at the drain wait") - // Only A is newly unobservable; B was already, so naming it twice would + require.Error(t, err, "naming the newly not_verified rule") + require.Contains(t, err.Error(), "not verified at the drain wait") + // Only A is newly not_verified; B was already, so naming it twice would // only lengthen the message. require.Contains(t, err.Error(), "A ("+string(ReasonDrainTimeout)+")") require.NotContains(t, err.Error(), "B (") @@ -1228,7 +1351,7 @@ func TestMergeDrainTimeoutsNamesEveryUnobservableRule(t *testing.T) { // B keeps the reason the coverage proof gave it — the FIRST reason wins, // as it does inside proveCoverage. require.Equal(t, ReasonHeartbeatGap, merged.Coverage["b"].Reason) - require.Equal(t, OutcomeUnobservable, merged.Verdicts[0].Outcome) + require.Equal(t, OutcomeNotVerified, merged.Verdicts[0].Outcome) } // ReadLogHeader is the one read of a log a writer may still hold, so its diff --git a/grafana-alertcheck/internal/gate/classify.go b/grafana-alertcheck/internal/gate/classify.go index da4ec0002..ab31edb3b 100644 --- a/grafana-alertcheck/internal/gate/classify.go +++ b/grafana-alertcheck/internal/gate/classify.go @@ -19,26 +19,30 @@ const ReasonNodata UnobservableReason = "nodata" type Outcome string const ( - OutcomeClean Outcome = "clean" - OutcomeNewlyBad Outcome = "newly_bad" - OutcomeRecovered Outcome = "recovered" - OutcomePersistentlyBad Outcome = "persistently_bad" - OutcomeFlapping Outcome = "flapping" - OutcomeSkipped Outcome = "skipped" - OutcomeUnobservable Outcome = "unobservable" + OutcomeHealthy Outcome = "healthy" + OutcomeNewFailure Outcome = "new_failure" + OutcomeRecovered Outcome = "recovered" + OutcomeStillFailing Outcome = "still_failing" + OutcomeUnstable Outcome = "unstable" + OutcomePaused Outcome = "paused" + OutcomeNotVerified Outcome = "not_verified" + // OutcomeNotCounted is synthetic: decide uses it for a --min-observed + // deficit row that no resolved rule explains. It is not a verdict on an + // alert, so it is deliberately not named after one. + OutcomeNotCounted Outcome = "not_counted" ) // PreexistingPolicy governs only the ONE ambiguous case in the outcome table: -// an instance that was already bad when the window opened. A newly_bad or -// flapping instance is a fail under every policy, so this type only ever -// changes how `recovered` and `persistently_bad` are judged (isViolation +// an instance that was already bad when the window opened. A new_failure or +// unstable instance is a fail under every policy, so this type only ever +// changes how `recovered` and `still_failing` are judged (isViolation // below). type PreexistingPolicy string const ( // PreexistingFailUnlessRecovered is the default: a preexisting instance // that clears and stays clear is a pass (`recovered`); one that never - // clears is still a fail (`persistently_bad`). + // clears is still a fail (`still_failing`). PreexistingFailUnlessRecovered PreexistingPolicy = "fail-unless-recovered" // PreexistingFail makes ANY preexisting instance a fail, even one that // recovers — for a user who wants no benefit of the doubt for a @@ -46,7 +50,7 @@ const ( PreexistingFail PreexistingPolicy = "fail" // PreexistingIgnore disregards a preexisting instance entirely, whether // it recovers or stays bad for the whole window: only a genuinely NEW - // bad episode (newly_bad or flapping) can fail the rule. + // bad episode (new_failure or unstable) can fail the rule. PreexistingIgnore PreexistingPolicy = "ignore" ) @@ -135,6 +139,10 @@ type Result struct { Global GlobalThresholds Verdicts []RuleVerdict Violations []Violation + // TerminatedEarly is set only when fail-fast stopped before the window + // closed; the coverage proof is then over [from, At]. To and Global below + // still report the requested values. Published JSON output. + TerminatedEarly *Termination `json:"terminated_early,omitempty"` } // episode is one contiguous, policy-bad span of one instance's timeline, @@ -290,7 +298,7 @@ func classifyRule(def Definition, polls []Poll, from, windowEnd time.Time, badSt slices.Sort(order) var ( - outcome = OutcomeClean + outcome = OutcomeHealthy badFor []episode viols []Violation ) @@ -307,17 +315,17 @@ func classifyRule(def Definition, polls []Poll, from, windowEnd time.Time, badSt var instOutcome Outcome switch { case len(tl.episodes) > 1: - instOutcome = OutcomeFlapping + instOutcome = OutcomeUnstable case tl.preexisting: if tl.episodes[0].closedByRealClear { instOutcome = OutcomeRecovered } else { - instOutcome = OutcomePersistentlyBad + instOutcome = OutcomeStillFailing } default: // A genuinely new onset fails whether or not it clears in-window; // only a preexisting condition earns `recovered`. - instOutcome = OutcomeNewlyBad + instOutcome = OutcomeNewFailure } if outcomeRank(instOutcome) > outcomeRank(outcome) { @@ -349,17 +357,17 @@ func classifyRule(def Definition, polls []Poll, from, windowEnd time.Time, badSt } // isViolation decides whether one instance's outcome counts against the run, -// once the preexisting policy is applied. newly_bad and flapping always do: +// once the preexisting policy is applied. new_failure and unstable always do: // both contain a genuinely new bad episode, so no policy forgives them. -// recovered and persistently_bad are, by classifyRule's construction, -// ALWAYS preexisting (a non-preexisting single episode is newly_bad instead, +// recovered and still_failing are, by classifyRule's construction, +// ALWAYS preexisting (a non-preexisting single episode is new_failure instead, // regardless of whether it clears) — so these are the only two policy can // change, and isViolation needs no separate preexisting flag to know that. func isViolation(o Outcome, pol PreexistingPolicy) bool { switch o { - case OutcomeNewlyBad, OutcomeFlapping: + case OutcomeNewFailure, OutcomeUnstable: return true - case OutcomePersistentlyBad: + case OutcomeStillFailing: return pol != PreexistingIgnore case OutcomeRecovered: return pol == PreexistingFail @@ -371,25 +379,25 @@ func isViolation(o Outcome, pol PreexistingPolicy) bool { // outcomeRank orders outcomes for classifyRule's worst-of reduction across a // rule's instances: // -// unobservable > {flapping, persistently_bad, newly_bad} > recovered > -// skipped > clean +// not_verified > {unstable, still_failing, new_failure} > recovered > +// paused > healthy // -// with unobservable and skipped applied outside this function (decide owns -// both: unobservable from CoverageResult, skipped from the log header). The +// with not_verified and paused applied outside this function (decide owns +// both: not_verified from CoverageResult, paused from the log header). The // three fail values are not ranked against each other by anything that reads // this, so their relative order here is an arbitrary but fixed tie-break, not // a claim that one is worse than another. func outcomeRank(o Outcome) int { switch o { - case OutcomeFlapping: + case OutcomeUnstable: return 4 - case OutcomePersistentlyBad: + case OutcomeStillFailing: return 3 - case OutcomeNewlyBad: + case OutcomeNewFailure: return 2 case OutcomeRecovered: return 1 - default: // OutcomeClean + default: // OutcomeHealthy return 0 } } @@ -452,10 +460,34 @@ func badStateSet(states []State) map[State]bool { return set } +// applyNodataPolicy promotes a sustained health=nodata run to unobservable +// when the caller has checked Policy.NodataIsUnobservable. proveCoverage never +// escalates nodata itself (most fleets map no-data to OK), so both the +// end-of-run decide and the fail-fast guard must apply this identically — one +// helper, so the two can never drift. +func applyNodataPolicy(def Definition, polls []Poll, cov *CoverageResult, t RuleTimings, from, windowEnd time.Time) { + if cov.Unobservable { + return + } + inWindow := inWindowPolls(pollsForRule(polls, def.UID), from, windowEnd) + runLen, sawAny := longestHealthRun(inWindow, "nodata") + if !sawAny || runLen <= t.healthGrace { + return + } + cov.Unobservable = true + cov.Proved = false + if cov.Reason == "" { + cov.Reason = ReasonNodata + } + cov.Notes = append(cov.Notes, fmt.Sprintf( + "rule %q: health=nodata for %s exceeds healthGrace %s and --nodata-is-unobservable is set", + def.Title, runLen, t.healthGrace)) +} + // decide is the pure seam between the collected evidence and the CLI's exit // code, and carries nearly the whole test suite because of it. It combines // proveCoverage's nine checks with classifyRule's timelines under one Policy, -// and owns the inability-beats-violation rule: any unobservable rule makes +// and owns the inability-beats-violation rule: any not-verified rule makes // decide return a non-nil error, which the CLI maps to exit 2 unconditionally // — never to 0 or 1, and never suppressed by a real violation found alongside // it. @@ -504,13 +536,13 @@ func decide(h Header, polls []Poll, sentinel *time.Time, defs []Definition, windowEnd := pol.To.Add(gt.transitionGrace) var ( - skippedRules []Definition + pausedRules []Definition watchedCount int anyUnobservable bool unobservableNames []string ) - // `skipped` is decided from the header, never from defs: defs are resolved + // `paused` is decided from the header, never from defs: defs are resolved // after the window closed, so Definition.IsPaused describes the present, // while Header.pausedAtStart describes the window open — the only moment // "paused before the window opened" can mean. @@ -518,9 +550,9 @@ func decide(h Header, polls []Poll, sentinel *time.Time, defs []Definition, for _, def := range defs { if pausedAtStart[def.UID] { - skippedRules = append(skippedRules, def) + pausedRules = append(pausedRules, def) result.Verdicts = append(result.Verdicts, RuleVerdict{ - Alert: def.Title, RuleUID: def.UID, Outcome: OutcomeSkipped, + Alert: def.Title, RuleUID: def.UID, Outcome: OutcomePaused, PollEvery: rt[def.UID].pollEvery, Note: "paused before the window opened", }) @@ -531,18 +563,8 @@ func decide(h Header, polls []Poll, sentinel *time.Time, defs []Definition, t := rt[def.UID] cov := proveCoverage(h, polls, sentinel, t, def, pol.From, pol.To, gt.transitionGrace) - if pol.NodataIsUnobservable && !cov.Unobservable { - inWindow := inWindowPolls(pollsForRule(polls, def.UID), pol.From, windowEnd) - if runLen, sawAny := longestHealthRun(inWindow, "nodata"); sawAny && runLen > t.healthGrace { - cov.Unobservable = true - cov.Proved = false - if cov.Reason == "" { - cov.Reason = ReasonNodata - } - cov.Notes = append(cov.Notes, fmt.Sprintf( - "rule %q: health=nodata for %s exceeds healthGrace %s and --nodata-is-unobservable is set", - def.Title, runLen, t.healthGrace)) - } + if pol.NodataIsUnobservable { + applyNodataPolicy(def, polls, &cov, t, pol.From, windowEnd) } result.Coverage[def.UID] = cov result.Thresholds[def.UID] = RuleThresholds{ @@ -553,7 +575,7 @@ func decide(h Header, polls []Poll, sentinel *time.Time, defs []Definition, outcome, badFor, viols := classifyRule(def, polls, pol.From, windowEnd, badStates, pol.Preexisting) if cov.Unobservable { - outcome = OutcomeUnobservable + outcome = OutcomeNotVerified anyUnobservable = true unobservableNames = append(unobservableNames, fmt.Sprintf("%s (%s)", def.Title, cov.Reason)) } @@ -570,9 +592,9 @@ func decide(h Header, polls []Poll, sentinel *time.Time, defs []Definition, counted := watchedCount var attributable []Definition if pol.AllowPaused { - counted += len(skippedRules) + counted += len(pausedRules) } else { - attributable = skippedRules + attributable = pausedRules } if shortfall := minObserved - counted; shortfall > 0 { attributed := 0 @@ -585,7 +607,7 @@ func decide(h Header, polls []Poll, sentinel *time.Time, defs []Definition, // prints Note verbatim rather than re-deriving the hint, so the // exact wording here is what an operator reads. result.Violations = append(result.Violations, Violation{ - Alert: def.Title, RuleUID: def.UID, Outcome: OutcomeSkipped, + Alert: def.Title, RuleUID: def.UID, Outcome: OutcomePaused, Note: "paused before the window opened; counts against --min-observed unless --allow-paused is set", }) attributed++ @@ -597,14 +619,14 @@ func decide(h Header, polls []Poll, sentinel *time.Time, defs []Definition, // rule state read from a real poll, and this Violation never // touched one. result.Violations = append(result.Violations, Violation{ - Outcome: OutcomeSkipped, + Outcome: OutcomeNotCounted, Note: fmt.Sprintf("min-observed %d exceeds the %d rule(s) counted as observed", minObserved, counted), }) } } if anyUnobservable { - return result, fmt.Errorf("gate: %d rule(s) unobservable: %s", len(unobservableNames), strings.Join(unobservableNames, "; ")) + return result, fmt.Errorf("gate: %d rule(s) not verified: %s", len(unobservableNames), strings.Join(unobservableNames, "; ")) } return result, nil } diff --git a/grafana-alertcheck/internal/gate/classify_test.go b/grafana-alertcheck/internal/gate/classify_test.go index 0e6b16c8e..14cc95248 100644 --- a/grafana-alertcheck/internal/gate/classify_test.go +++ b/grafana-alertcheck/internal/gate/classify_test.go @@ -1,6 +1,7 @@ package gate import ( + "encoding/json" "testing" "time" @@ -48,7 +49,7 @@ func pausedHeader(startedAt time.Time, pausedUIDs ...string) Header { return h } -// --- clean / newly_bad --- +// --- healthy / new_failure --- func TestClassifyRule_NoEvidenceIsClean(t *testing.T) { from := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) @@ -57,7 +58,7 @@ func TestClassifyRule_NoEvidenceIsClean(t *testing.T) { polls := []Poll{quietPoll("r1", from), quietPoll("r1", to)} outcome, badFor, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) - require.Equal(t, OutcomeClean, outcome) + require.Equal(t, OutcomeHealthy, outcome) require.Zero(t, badFor) require.Empty(t, viols) } @@ -74,10 +75,10 @@ func TestClassifyRule_NewOnsetInsideWindowIsNewlyBad(t *testing.T) { abnormalPoll("r1", to, StateFiring, lbl("a"), onset), } outcome, badFor, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) - require.Equal(t, OutcomeNewlyBad, outcome) + require.Equal(t, OutcomeNewFailure, outcome) require.Equal(t, to.Sub(onset), badFor) require.Len(t, viols, 1) - require.Equal(t, OutcomeNewlyBad, viols[0].Outcome) + require.Equal(t, OutcomeNewFailure, viols[0].Outcome) } // A genuinely new bad episode fails even if it clears again before the window @@ -96,11 +97,11 @@ func TestClassifyRule_NewOnsetThatClearsStillFails(t *testing.T) { quietPoll("r1", to), } outcome, _, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) - require.Equal(t, OutcomeNewlyBad, outcome, "even though it cleared") + require.Equal(t, OutcomeNewFailure, outcome, "even though it cleared") require.Len(t, viols, 1) } -// --- recovered / persistently_bad (preexisting) --- +// --- recovered / still_failing (preexisting) --- func TestClassifyRule_PreexistingThatRecoversIsRecoveredAndNotAViolation(t *testing.T) { from := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) @@ -151,13 +152,13 @@ func TestClassifyRule_PreexistingStillBadAtWindowEndIsPersistentlyBad(t *testing abnormalPoll("r1", to, StateFiring, lbl("a"), from.Add(-time.Hour)), } outcome, badFor, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) - require.Equal(t, OutcomePersistentlyBad, outcome) + require.Equal(t, OutcomeStillFailing, outcome) require.Equal(t, to.Sub(from), badFor) require.Len(t, viols, 1) - require.Equal(t, OutcomePersistentlyBad, viols[0].Outcome) + require.Equal(t, OutcomeStillFailing, viols[0].Outcome) } -// --- flapping --- +// --- unstable --- func TestClassifyRule_ClearThenBadAgainIsFlapping(t *testing.T) { from := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) @@ -172,12 +173,12 @@ func TestClassifyRule_ClearThenBadAgainIsFlapping(t *testing.T) { quietPoll("r1", to), } outcome, _, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) - require.Equal(t, OutcomeFlapping, outcome) + require.Equal(t, OutcomeUnstable, outcome) require.Len(t, viols, 1) - require.Equal(t, OutcomeFlapping, viols[0].Outcome, "always a fail regardless of policy") + require.Equal(t, OutcomeUnstable, viols[0].Outcome, "always a fail regardless of policy") } -// A clear and then a second bad state gives flapping, wherever the second bad +// A clear and then a second bad state gives unstable, wherever the second bad // state lands. A table over where the second onset falls — immediately after // the clear, mid-window, and right at the // last instant before windowEnd — closes the boundary this single fixed @@ -206,9 +207,9 @@ func TestClassifyRule_FlappingAtEveryTimingOfTheSecondOnset(t *testing.T) { quietPoll("r1", to), } outcome, _, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) - require.Equalf(t, OutcomeFlapping, outcome, "second onset at %s", tc.secondOnset) + require.Equalf(t, OutcomeUnstable, outcome, "second onset at %s", tc.secondOnset) require.Len(t, viols, 1) - require.Equal(t, OutcomeFlapping, viols[0].Outcome) + require.Equal(t, OutcomeUnstable, viols[0].Outcome) }) } } @@ -227,7 +228,7 @@ func TestClassifyRule_VanishedWhileBadStaysPersistentlyBad(t *testing.T) { quietPoll("r1", to), } outcome, badFor, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) - require.Equal(t, OutcomePersistentlyBad, outcome, "a vanish must never read as a recovery") + 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 +247,7 @@ func TestClassifyRule_VanishedWhileNeverBadIsUninteresting(t *testing.T) { quietPoll("r1", to), } outcome, badFor, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) - require.Equal(t, OutcomeClean, outcome) + require.Equal(t, OutcomeHealthy, outcome) require.Zero(t, badFor) require.Empty(t, viols) } @@ -281,7 +282,7 @@ func TestClassifyRule_PreexistingPolicyIgnoreForgivesPersistentlyBad(t *testing. abnormalPoll("r1", to, StateFiring, lbl("a"), from.Add(-time.Hour)), } outcome, _, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingIgnore) - require.Equal(t, OutcomePersistentlyBad, outcome, "the descriptive outcome does not change under policy=ignore") + 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 +298,7 @@ func TestClassifyRule_PreexistingPolicyIgnoreStillFailsANewOnset(t *testing.T) { abnormalPoll("r1", to, StateFiring, lbl("a"), onset), } outcome, _, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingIgnore) - require.Equal(t, OutcomeNewlyBad, outcome) + require.Equal(t, OutcomeNewFailure, outcome) require.Len(t, viols, 1, "ignore only forgives PREEXISTING badness") } @@ -323,9 +324,9 @@ func TestClassifyRule_WorstOfMultipleInstancesWins(t *testing.T) { }, } outcome, _, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) - require.Equal(t, OutcomePersistentlyBad, outcome, "the worse of {recovered, persistently_bad}") + require.Equal(t, OutcomeStillFailing, outcome, "the worse of {recovered, still_failing}") require.Len(t, viols, 1) - require.Equal(t, OutcomePersistentlyBad, viols[0].Outcome) + require.Equal(t, OutcomeStillFailing, viols[0].Outcome) } // --- decide(): skipped rules, unobservable, MinObserved, exit mapping --- @@ -344,9 +345,9 @@ func TestDecide_SkippedRuleNeverReachesProveCoverage(t *testing.T) { // The HEADER is what says paused — decide reads skipped from there, not // from def.IsPaused, which is a post-window reading (Header.pausedAtStart). res, err := decide(pausedHeader(from.Add(-time.Hour), "r1"), nil, nil, defs, rt, gt, pol) - require.NoError(t, err, "a rule paused before the window is skipped, not unobservable") + require.NoError(t, err, "a rule paused before the window is paused, not not_verified") require.Len(t, res.Verdicts, 1) - require.Equal(t, OutcomeSkipped, res.Verdicts[0].Outcome) + require.Equal(t, OutcomePaused, res.Verdicts[0].Outcome) _, ok := res.Coverage["r1"] require.False(t, ok, "a skipped rule has no coverage to prove") } @@ -365,11 +366,11 @@ func TestDecide_UnobservableRuleAlwaysReturnsAnError(t *testing.T) { res, err := decide(Header{StartedAt: from.Add(-time.Hour)}, nil, nil, defs, rt, gt, pol) require.Error(t, err, "an unobservable rule must always fail the run") require.Len(t, res.Verdicts, 1) - require.Equal(t, OutcomeUnobservable, res.Verdicts[0].Outcome) + require.Equal(t, OutcomeNotVerified, res.Verdicts[0].Outcome) } // Any unobservable rule means exit 2, with no exception — even alongside a -// real newly_bad. +// real new_failure. func TestDecide_UnobservableWinsEvenAlongsideARealViolation(t *testing.T) { from := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) to := from.Add(10 * time.Minute) @@ -407,11 +408,11 @@ func TestDecide_UnobservableWinsEvenAlongsideARealViolation(t *testing.T) { gotBad = v.Outcome } } - require.Equal(t, OutcomeUnobservable, gotBroken) - require.Equal(t, OutcomeNewlyBad, gotBad, + require.Equal(t, OutcomeNotVerified, gotBroken) + require.Equal(t, OutcomeNewFailure, gotBad, "classification still runs and is still visible in Verdicts") require.NotEmpty(t, res.Violations, - "the newly_bad instance still reported even though the run fails on the unobservable rule") + "the new_failure instance still reported even though the run fails on the not_verified rule") } // A clean verdict with a coverage gap must never give exit 0, and recovered @@ -429,7 +430,7 @@ func TestDecide_UnobservableRuleWinsOverEveryFavorableOutcome(t *testing.T) { wantOutcome Outcome }{ { - name: "clean", + name: "healthy", goodPolls: func() []Poll { var polls []Poll for ts := from; !ts.After(to); ts = ts.Add(30 * time.Second) { @@ -437,7 +438,7 @@ func TestDecide_UnobservableRuleWinsOverEveryFavorableOutcome(t *testing.T) { } return polls }(), - wantOutcome: OutcomeClean, + wantOutcome: OutcomeHealthy, }, { // Dense 30s-spaced polls throughout, so "good"'s own coverage @@ -466,9 +467,9 @@ func TestDecide_UnobservableRuleWinsOverEveryFavorableOutcome(t *testing.T) { wantOutcome: OutcomeRecovered, }, { - name: "skipped", + name: "paused", pausedAtStart: true, - wantOutcome: OutcomeSkipped, + wantOutcome: OutcomePaused, }, } @@ -502,7 +503,7 @@ func TestDecide_UnobservableRuleWinsOverEveryFavorableOutcome(t *testing.T) { } } require.Equal(t, tc.wantOutcome, gotGood) - require.Equal(t, OutcomeUnobservable, gotBroken) + require.Equal(t, OutcomeNotVerified, gotBroken) }) } } @@ -540,7 +541,7 @@ func TestDecide_RecoveredOutcomeOverriddenByItsOwnCoverageGap(t *testing.T) { res, err := decide(Header{StartedAt: from.Add(-time.Hour)}, polls, &sentinel, defs, rt, gt, pol) require.Error(t, err, "r1's own coverage gap must fail the run even though it recovered") require.Len(t, res.Verdicts, 1) - require.Equal(t, OutcomeUnobservable, res.Verdicts[0].Outcome, "never recovered") + require.Equal(t, OutcomeNotVerified, res.Verdicts[0].Outcome, "never recovered") require.False(t, res.Coverage["r1"].Proved) } @@ -562,7 +563,7 @@ func TestDecide_CleanWindowIsAPass(t *testing.T) { res, err := decide(Header{StartedAt: from.Add(-time.Hour)}, polls, &sentinel, defs, rt, gt, pol) require.NoError(t, err) require.Empty(t, res.Violations, "a pass is exactly len(Violations)==0 && err==nil") - require.Equal(t, OutcomeClean, res.Verdicts[0].Outcome) + require.Equal(t, OutcomeHealthy, res.Verdicts[0].Outcome) } // A pause and then an unpause inside the window, with an episode that would @@ -604,7 +605,7 @@ func TestDecide_PauseThenUnpauseWithHiddenEpisodeGivesUnobservableNotClean(t *te res, err := decide(Header{StartedAt: from.Add(-time.Hour)}, polls, &sentinel, defs, rt, gt, pol) require.Error(t, err, "the pause-then-unpause blind interval must fail closed") require.Len(t, res.Verdicts, 1) - require.Equal(t, OutcomeUnobservable, res.Verdicts[0].Outcome, "never clean") + require.Equal(t, OutcomeNotVerified, res.Verdicts[0].Outcome, "never clean") require.False(t, res.Coverage["r1"].Proved) } @@ -632,10 +633,10 @@ func TestDecide_SkippedOnlyShortfallProducesAViolationWithoutAnError(t *testing. sentinel := to res, err := decide(pausedHeader(from.Add(-time.Hour), "paused"), polls, &sentinel, defs, rt, gt, pol) - require.NoError(t, err, "a shortfall caused only by a skipped rule is exit 1, not exit 2") + require.NoError(t, err, "a shortfall caused only by a paused rule is exit 1, not exit 2") require.Len(t, res.Violations, 1) v := res.Violations[0] - require.Equal(t, OutcomeSkipped, v.Outcome) + require.Equal(t, OutcomePaused, v.Outcome) require.Equal(t, "paused", v.RuleUID) require.Equal(t, "Paused", v.Alert) require.NotEmpty(t, v.Note, @@ -665,7 +666,8 @@ func TestDecide_ExplicitMinObservedShortfallWithNoPausedRuleStillProducesAViolat require.NoError(t, err, "an unmet MinObserved is exit 1, never exit 2") require.Len(t, res.Violations, 2, "the shortfall (3-1=2) must surface directly rather than pass silently") for _, v := range res.Violations { - require.Equal(t, OutcomeSkipped, v.Outcome) + require.Equal(t, OutcomeNotCounted, v.Outcome, + "no resolved rule explains the deficit, so it must not be blamed on a paused one") } } @@ -740,7 +742,7 @@ func TestDecide_NodataIsANoteByDefault(t *testing.T) { // An instance whose true onset (ActiveAt) falls strictly inside the window — // even though the first poll that happens to observe it already shows it bad — -// must never be treated as preexisting. If it then clears, that is newly_bad +// must never be treated as preexisting. If it then clears, that is new_failure // (exit 1), not recovered (exit 0). func TestClassifyRule_OnsetBetweenFromAndFirstPollIsNewlyBadNotRecovered(t *testing.T) { from := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) @@ -757,10 +759,10 @@ func TestClassifyRule_OnsetBetweenFromAndFirstPollIsNewlyBadNotRecovered(t *test quietPoll("r1", to), } outcome, badFor, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) - require.Equal(t, OutcomeNewlyBad, outcome, + 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) - require.Equal(t, OutcomeNewlyBad, viols[0].Outcome) + require.Equal(t, OutcomeNewFailure, viols[0].Outcome) require.Equal(t, clearAt.Sub(onset), badFor, "BadFor must count from the true onset, not from `from`") } @@ -799,7 +801,7 @@ func TestClassifyRule_SkewTranslatesActiveAtAcrossTheWindowBoundary(t *testing.T // Grafana's clock reads 90s ahead of the runner's (skew = +90s). The // poll's raw GrafanaNow/ActiveAt both sit 90s past `from` in Grafana's // domain, but translate to exactly `from` in the runner domain — genuinely - // preexisting once translated, and wrongly "newly_bad" if the skew is + // preexisting once translated, and wrongly "new_failure" if the skew is // ignored. skew := 90 * time.Second rawActiveAt := from.Add(skew) @@ -813,7 +815,7 @@ func TestClassifyRule_SkewTranslatesActiveAtAcrossTheWindowBoundary(t *testing.T stillBad.LastEvaluation = to.Add(skew) outcome, badFor, _ := classifyRule(def, []Poll{poll, stillBad}, from, to, defaultBad, PreexistingFailUnlessRecovered) - require.Equal(t, OutcomePersistentlyBad, outcome, + require.Equal(t, OutcomeStillFailing, outcome, "a +90s skew must translate ActiveAt back to exactly `from`") require.Equal(t, to.Sub(from), badFor) } @@ -893,7 +895,7 @@ func TestClassifyRule_ClearedEventPastWindowEndClampsToWindowEnd(t *testing.T) { // reading of the upper boundary: an instance whose runner-domain onset lands // only slightly past windowEnd (to + transitionGrace) is reachable at all only // because inWindowPolls widens the boundary outward by the skew bound, so the -// gate cannot PROVE it belongs to the next window. It is charged as newly_bad — +// gate cannot PROVE it belongs to the next window. It is charged as new_failure — // with BadFor truncated to zero — rather than silently forgiven as clean. func TestClassifyRule_OnsetJustPastWindowEndIsNewlyBadNotClean(t *testing.T) { from := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) @@ -914,13 +916,13 @@ func TestClassifyRule_OnsetJustPastWindowEndIsNewlyBadNotClean(t *testing.T) { } outcome, badFor, viols := classifyRule(def, []Poll{poll}, from, windowEnd, defaultBad, PreexistingFailUnlessRecovered) - require.Equal(t, OutcomeNewlyBad, outcome, + 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") require.Len(t, viols, 1) } -// A clear after `to` gives persistently_bad. classifyRule filters +// A clear after `to` gives still_failing. classifyRule filters // its input to [from, windowEnd] itself (inWindowPolls), so a Cleared event // GENUINELY past windowEnd — well beyond any skew bound, unlike the clamp // case above — never reaches the timeline at all: the instance is still bad @@ -936,10 +938,10 @@ func TestClassifyRule_ClearAfterWindowEndIsPersistentlyBad(t *testing.T) { clearedPoll("r1", to.Add(time.Hour), key), // far past `to`, not a boundary case } outcome, badFor, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) - require.Equal(t, OutcomePersistentlyBad, outcome, "a clear outside the window must not read as a recovery") + 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) - require.Equal(t, OutcomePersistentlyBad, viols[0].Outcome) + require.Equal(t, OutcomeStillFailing, viols[0].Outcome) } // TestClassifyRule_CloseBeforeOpenClampsToZeroNotNegative pins the @@ -966,7 +968,7 @@ func TestClassifyRule_CloseBeforeOpenClampsToZeroNotNegative(t *testing.T) { quietPoll("r1", to), } outcome, badFor, viols := classifyRule(def, polls, from, to, defaultBad, PreexistingFailUnlessRecovered) - require.Equal(t, OutcomeNewlyBad, outcome) + 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") require.Zero(t, badFor, "the clamp collapses the inverted span to a zero-length episode") @@ -990,3 +992,30 @@ func TestMergeDurations_OverlappingEpisodesCountOnce(t *testing.T) { func TestMergeDurations_Empty(t *testing.T) { require.Zero(t, mergeDurations(nil)) } + +// --- published vocabulary --- + +// Outcome strings are published JSON, but every other test compares against +// the constants, so a typo in a constant's literal would pass them all. Pin +// the literals themselves. +func TestOutcomeJSONVocabulary(t *testing.T) { + want := map[Outcome]string{ + OutcomeHealthy: "healthy", + OutcomeNewFailure: "new_failure", + OutcomeRecovered: "recovered", + OutcomeStillFailing: "still_failing", + OutcomeUnstable: "unstable", + OutcomePaused: "paused", + OutcomeNotVerified: "not_verified", + OutcomeNotCounted: "not_counted", + } + require.Len(t, want, 8, "every Outcome constant must be pinned here") + + for outcome, literal := range want { + t.Run(string(outcome), func(t *testing.T) { + raw, err := json.Marshal(outcome) + require.NoError(t, err) + require.Equal(t, `"`+literal+`"`, string(raw)) + }) + } +} diff --git a/grafana-alertcheck/internal/gate/coverage_test.go b/grafana-alertcheck/internal/gate/coverage_test.go index 7126372ed..226b46243 100644 --- a/grafana-alertcheck/internal/gate/coverage_test.go +++ b/grafana-alertcheck/internal/gate/coverage_test.go @@ -79,7 +79,7 @@ func TestProveCoverage_SentinelBeforeGraceIsUnobservable(t *testing.T) { dres, err := decide(Header{StartedAt: from.Add(-time.Hour)}, nil, &sentinel, defs, drt, gt, pol) require.Error(t, err, "a sentinel short of to+grace must fail the run") require.Len(t, dres.Verdicts, 1) - require.Equal(t, OutcomeUnobservable, dres.Verdicts[0].Outcome) + require.Equal(t, OutcomeNotVerified, dres.Verdicts[0].Outcome) } func TestProveCoverage_SentinelExactlyAtGraceIsFine(t *testing.T) { @@ -123,7 +123,7 @@ func TestProveCoverage_FromBeforeRecordIsUnobservable(t *testing.T) { dres, err := decide(Header{StartedAt: started}, nil, &sentinel, defs, drt, gt, pol) require.Error(t, err, "`from` before the recording started must fail the run") require.Len(t, dres.Verdicts, 1) - require.Equal(t, OutcomeUnobservable, dres.Verdicts[0].Outcome) + require.Equal(t, OutcomeNotVerified, dres.Verdicts[0].Outcome) } // The from-bounds check compares at whole-second granularity: a whole-second @@ -688,15 +688,15 @@ func TestProveCoverage_MultipleFailuresReasonIsFirstButAllNoted(t *testing.T) { "a later failure must still be recorded, not swallowed once Reason is already set") } -// --- Skipped rules --- +// --- Paused rules --- // A known limit of this function's contract, not a bug in it: a rule paused // BEFORE the window opened is never scheduled or polled (watch.go), so it // reaches proveCoverage with zero polls at all. proveCoverage has no notion of -// "skipped" — that classification belongs to the definitions +// "paused" — that classification belongs to the definitions // (LoggedRule.IsPaused / Definition.IsPaused), never to the polls — so it // reports the whole window as one big heartbeat_gap instead. decide is what -// reads skipped status from the header and never calls this function for such +// reads paused status from the header and never calls this function for such // a rule; this pins the behavior it relies on not reaching. func TestProveCoverage_SkippedRuleWithZeroPollsPinnedAsHeartbeatGap(t *testing.T) { from := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) diff --git a/grafana-alertcheck/internal/gate/flock.go b/grafana-alertcheck/internal/gate/flock.go index 996844ca0..3333cff22 100644 --- a/grafana-alertcheck/internal/gate/flock.go +++ b/grafana-alertcheck/internal/gate/flock.go @@ -24,15 +24,8 @@ func isLockContention(err error) bool { return errors.Is(err, syscall.EWOULDBLOCK) || errors.Is(err, syscall.EAGAIN) } -// tryLockExclusive is the same call read as a question rather than as a -// demand: held is false when another process holds the lock, and err is -// non-nil only for a failure that is not contention. -// -// check needs that distinction where NewWriter does not. NewWriter is entitled -// to treat any refusal as "another writer has it", because it wants the lock; -// check only wants to know whether a writer EXISTS. The lock answers -// that directly, where a pid can only infer it — the kernel releases a flock -// when the holder exits, crash included, and pids get reused. +// tryLockExclusive is lockExclusive read as a question: held is false when +// another process holds the lock; err is non-nil only for a real failure. func tryLockExclusive(f *os.File) (held bool, err error) { switch err := syscall.Flock(int(f.Fd()), syscall.LOCK_EX|syscall.LOCK_NB); { case err == nil: @@ -43,3 +36,18 @@ func tryLockExclusive(f *os.File) (held bool, err error) { return false, fmt.Errorf("flock %s: %w", f.Name(), err) } } + +// tryLockShared is the reader side of the same question: held is false when an +// EXCLUSIVE holder exists. The recorder holds the log exclusively; readers and +// cleanup hold it shared, so shared contention always means the recorder, never +// another reader. +func tryLockShared(f *os.File) (held bool, err error) { + switch err := syscall.Flock(int(f.Fd()), syscall.LOCK_SH|syscall.LOCK_NB); { + case err == nil: + return true, nil + case isLockContention(err): + return false, nil + default: + return false, fmt.Errorf("flock %s: %w", f.Name(), err) + } +} diff --git a/grafana-alertcheck/internal/gate/logtail.go b/grafana-alertcheck/internal/gate/logtail.go new file mode 100644 index 000000000..cb713a780 --- /dev/null +++ b/grafana-alertcheck/internal/gate/logtail.go @@ -0,0 +1,104 @@ +package gate + +import ( + "bytes" + "encoding/json" + "errors" + "fmt" + "io" + "os" + "time" +) + +// tailReadChunk bounds one read of the growing log. +const tailReadChunk = 64 * 1024 + +// logTailer incrementally reads a recorder's log while the recorder may still +// be appending to it, for the fail-fast guard only. +// +// It consumes complete, newline-terminated records only — a torn write stays +// buffered until its newline arrives — and never feeds the final +// classification. Once the run stops, check still reads the log once through +// the strict whole-file ReadLog. +type logTailer struct { + f *os.File + offset int64 // file position up to which the file has been read + buf []byte // bytes read but not yet forming a complete line +} + +func newLogTailer(path string) (*logTailer, error) { + f, err := os.Open(path) + if err != nil { + return nil, fmt.Errorf("tail log %s: %w", path, err) + } + return &logTailer{f: f}, nil +} + +func (t *logTailer) Close() error { return t.f.Close() } + +// read returns polls appended since the previous call, plus the sentinel when +// the recorder has finished. +func (t *logTailer) read() ([]Poll, *time.Time, error) { + if _, err := t.f.Seek(t.offset, io.SeekStart); err != nil { + return nil, nil, fmt.Errorf("tail log: seek: %w", err) + } + chunk := make([]byte, tailReadChunk) + for { + n, err := t.f.Read(chunk) + if n > 0 { + t.buf = append(t.buf, chunk[:n]...) + t.offset += int64(n) + } + if err != nil { + if errors.Is(err, io.EOF) { + break + } + return nil, nil, fmt.Errorf("tail log: read: %w", err) + } + if n == 0 { + break + } + } + + var ( + polls []Poll + sentinel *time.Time + ) + for { + i := bytes.IndexByte(t.buf, '\n') + if i < 0 { + break + } + line := t.buf[:i] + t.buf = t.buf[i+1:] + if len(bytes.TrimSpace(line)) == 0 { + continue + } + var probe struct { + Type RecordType `json:"type"` + } + if err := json.Unmarshal(line, &probe); err != nil { + return nil, nil, fmt.Errorf("tail log: unparseable complete record: %w", err) + } + switch probe.Type { + case RecordHeader: + // Already read authoritatively by ReadLogHeader; nothing to do. + case RecordPoll: + var rec pollRecord + if err := json.Unmarshal(line, &rec); err != nil { + return nil, nil, fmt.Errorf("tail log: unparseable poll: %w", err) + } + polls = append(polls, rec.Poll) + case RecordStopped: + var rec stoppedRecord + if err := json.Unmarshal(line, &rec); err != nil { + return nil, nil, fmt.Errorf("tail log: unparseable sentinel: %w", err) + } + at := rec.At + sentinel = &at + default: + return nil, nil, fmt.Errorf("tail log: unknown record type %q", probe.Type) + } + } + return polls, sentinel, nil +} diff --git a/grafana-alertcheck/internal/gate/logtail_test.go b/grafana-alertcheck/internal/gate/logtail_test.go new file mode 100644 index 000000000..ac019a1e6 --- /dev/null +++ b/grafana-alertcheck/internal/gate/logtail_test.go @@ -0,0 +1,89 @@ +package gate + +import ( + "encoding/json" + "os" + "path/filepath" + "testing" + "time" + + "github.com/stretchr/testify/require" +) + +func mustJSON(t *testing.T, v any) []byte { + t.Helper() + b, err := json.Marshal(v) + require.NoError(t, err) + return b +} + +func appendFile(t *testing.T, path string, b []byte) { + t.Helper() + f, err := os.OpenFile(path, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0o644) // nolint:gosec // test-only temp file + require.NoError(t, err) + _, err = f.Write(b) + require.NoError(t, err) + require.NoError(t, f.Close()) +} + +// A torn write waits in the buffer for its newline; it is never reported, and +// never repeated once completed. +func TestLogTailerReadsOnlyCompleteRecords(t *testing.T) { + path := filepath.Join(t.TempDir(), "log.jsonl") + + header := mustJSON(t, headerRecord{Type: RecordHeader, Header: Header{ + SchemaVersion: LogSchemaVersion, StartedAt: testNow, + }}) + poll1 := mustJSON(t, pollRecord{Type: RecordPoll, Poll: Poll{RuleUID: checkUID, GrafanaNow: testNow, Found: true}}) + poll2 := mustJSON(t, pollRecord{Type: RecordPoll, Poll: Poll{RuleUID: checkUID, GrafanaNow: testNow.Add(time.Second), Found: true}}) + + // Header and one poll complete; the second poll is torn mid-write. + split := len(poll2) / 2 + appendFile(t, path, []byte(string(header)+"\n")) + appendFile(t, path, []byte(string(poll1)+"\n")) + appendFile(t, path, poll2[:split]) + + tail, err := newLogTailer(path) + require.NoError(t, err) + defer tail.Close() + + polls, sentinel, err := tail.read() + require.NoError(t, err) + require.Len(t, polls, 1, "the torn record must not be reported") + require.Nil(t, sentinel) + + // A second read with nothing new must not repeat the first poll. + polls, _, err = tail.read() + require.NoError(t, err) + require.Empty(t, polls) + + // Finish the torn line; now it is reportable, exactly once. + appendFile(t, path, append(poll2[split:], '\n')) + polls, sentinel, err = tail.read() + require.NoError(t, err) + require.Len(t, polls, 1) + require.True(t, polls[0].GrafanaNow.Equal(testNow.Add(time.Second))) + require.Nil(t, sentinel) + + // The sentinel is reported when it lands. + stopped := mustJSON(t, stoppedRecord{Type: RecordStopped, At: testNow.Add(2 * time.Second)}) + appendFile(t, path, append(stopped, '\n')) + _, sentinel, err = tail.read() + require.NoError(t, err) + require.NotNil(t, sentinel) + require.True(t, sentinel.Equal(testNow.Add(2*time.Second))) +} + +// A complete but unparseable line is corruption: fail closed, never skip. +func TestLogTailerRejectsAnUnparseableCompleteLine(t *testing.T) { + path := filepath.Join(t.TempDir(), "log.jsonl") + appendFile(t, path, []byte("not json\n")) + + tail, err := newLogTailer(path) + require.NoError(t, err) + defer tail.Close() + + _, _, err = tail.read() + require.Error(t, err) + require.Contains(t, err.Error(), "unparseable") +} diff --git a/grafana-alertcheck/internal/gate/stop.go b/grafana-alertcheck/internal/gate/stop.go new file mode 100644 index 000000000..48a9b85a0 --- /dev/null +++ b/grafana-alertcheck/internal/gate/stop.go @@ -0,0 +1,235 @@ +package gate + +import ( + "context" + "errors" + "fmt" + "io" + "os" + "time" +) + +// recorderStopTimeout bounds the wait for the recorder's exit after SIGTERM. +// recorderStopPoll is how often the wait re-checks the lock. +const ( + recorderStopTimeout = 30 * time.Second + recorderStopPoll = 100 * time.Millisecond +) + +// StopConfig locates the recorder and selects the stop semantics. +type StopConfig struct { + // Log is the JSONL path whose flock proves a writer exists right now. + // PidFile defaults to .pid. + Log string + PidFile string + + Clock Clock + Notes io.Writer + + // Timeout bounds each wait (SIGTERM, then SIGKILL when Cleanup is set). + // Zero means recorderStopTimeout; a test overrides it to drive the SIGKILL + // path quickly. + Timeout time.Duration + + // Cleanup is the `stop` command's semantics rather than check's: SIGKILL a + // recorder that ignores SIGTERM, and remove the pidfile once no writer + // holds the log. check leaves it false — a writer that will not exit means + // the log cannot be trusted, which is a could-not-check, never a silent + // kill. + Cleanup bool +} + +func (c StopConfig) withDefaults() StopConfig { + if c.Clock == nil { + c.Clock = SystemClock{} + } + if c.Notes == nil { + c.Notes = io.Discard + } + if c.PidFile == "" && c.Log != "" { + c.PidFile = c.Log + ".pid" + } + if c.Timeout <= 0 { + c.Timeout = recorderStopTimeout + } + return c +} + +// StopRecorder signals the recorder holding cfg.Log and waits for it to go. +// It returns the log under a shared flock — the caller must keep it open across +// any read and close it when done — even when no writer was found. +// +// The recorder holds the log exclusively; readers and cleanup hold it shared. +// Shared is refused only by an exclusive holder, so contention always means the +// recorder: a pid can be reused, the lock cannot. +func StopRecorder(ctx context.Context, cfg StopConfig) (*os.File, error) { + cfg = cfg.withDefaults() + if cfg.Log == "" { + return nil, fmt.Errorf("stop recorder: no log path") + } + if cfg.PidFile == "" { + return nil, fmt.Errorf("stop recorder: no pidfile path") + } + + pid, pidErr := ReadPidFile(cfg.PidFile) + if pidErr == nil { + return cfg.stopNamed(ctx, pid) + } + if !cfg.Cleanup { + return nil, fmt.Errorf("cannot stop the recorder: %w; a pidfile is written only once a recorder reports that it is running, so an unreadable one means the recording never started", pidErr) + } + return cfg.stopUnnamed(pidErr) +} + +// stopNamed stops the recorder the pidfile names. +func (c StopConfig) stopNamed(ctx context.Context, pid int) (*os.File, error) { + log, err := os.Open(c.Log) + if err != nil { + return nil, fmt.Errorf("cannot check whether a recorder is still running: %w; pidfile %s names a recorder that may still hold the log", err, c.PidFile) + } + + // held is true when we TAKE the lock, i.e. when no writer holds it. + held, err := tryLockShared(log) + if err != nil { + log.Close() + return nil, err + } + if held { + // No writer, so no signal — the pid may belong to somebody else by now. + fmt.Fprintf(c.Notes, "note: no writer holds %s; the recorder has already finished\n", c.Log) + return c.finish(log) + } + return c.signalAndWait(ctx, pid, log) +} + +// stopUnnamed is cleanup with an unreadable pidfile: the log's lock still +// decides whether a writer exists, but a live one cannot be named, so it can +// only be reported, never signalled. +func (c StopConfig) stopUnnamed(pidErr error) (*os.File, error) { + log, err := os.Open(c.Log) + if err != nil { + // Nothing was recorded at all: no readable pidfile and no log. This is + // the only open failure cleanup may read as "nothing to stop". + if errors.Is(err, os.ErrNotExist) { + fmt.Fprintf(c.Notes, "note: no recorder to stop: %s is absent and its pidfile %s is unreadable\n", c.Log, c.PidFile) + return nil, nil + } + return nil, fmt.Errorf("cannot check whether a recorder is still running: %w; its unreadable pidfile %s: %w", err, c.PidFile, pidErr) + } + + held, err := tryLockShared(log) + if err != nil { + log.Close() + return nil, err + } + if held { + fmt.Fprintf(c.Notes, "note: no writer holds %s; nothing to stop\n", c.Log) + return c.finish(log) + } + log.Close() + return nil, fmt.Errorf("a writer holds %s but pidfile %s is unreadable: %w", c.Log, c.PidFile, pidErr) +} + +// finish removes the pidfile while the lock is still held — so a concurrent +// watch cannot take the lock and write a pidfile this then deletes — and +// returns the locked log. +func (c StopConfig) finish(log *os.File) (*os.File, error) { + if err := c.removePidFile(); err != nil { + log.Close() + return nil, err + } + return log, nil +} + +// errStillHeld is waitForRelease's timeout result. +var errStillHeld = errors.New("recorder still holds the log") + +// waitForRelease polls the lock until it is released, the timeout passes, or +// ctx is cancelled. +func (c StopConfig) waitForRelease(ctx context.Context, log *os.File, timeout time.Duration) error { + deadline := c.Clock.Now().Add(timeout) + for { + select { + case <-ctx.Done(): + return ctx.Err() + case <-c.Clock.After(recorderStopPoll): + } + held, err := tryLockShared(log) + if err != nil { + return err + } + if held { + return nil + } + if !c.Clock.Now().Before(deadline) { + return errStillHeld + } + } +} + +// signalAndWait sends SIGTERM to pid and waits for the log's lock to release. +// The lock, never the pid, is the proof the writer is gone. +func (c StopConfig) signalAndWait(ctx context.Context, pid int, log *os.File) (*os.File, error) { + fail := func(err error) (*os.File, error) { + log.Close() + return nil, err + } + + gone, err := signalRecorder(pid) + if err != nil { + return fail(err) + } + if gone { + // The pid died while the lock was held. A free lock now means the + // recorder finished in the race window; remaining contention means the + // pidfile does not name the process holding the log. + held, err := tryLockShared(log) + if err != nil { + return fail(err) + } + if held { + return c.finish(log) + } + return fail(fmt.Errorf("a writer holds %s but pidfile %s names pid %d, which does not exist: the pidfile does not name the process that holds the log", + c.Log, c.PidFile, pid)) + } + + err = c.waitForRelease(ctx, log, c.Timeout) + if err == nil { + return c.finish(log) + } + if !errors.Is(err, errStillHeld) { + return fail(err) + } + if !c.Cleanup { + return fail(fmt.Errorf("recorder pid %d still holds %s %s after SIGTERM; refusing to read a log a writer can still append to", + pid, c.Log, c.Timeout)) + } + + // The recorder ignored SIGTERM; SIGKILL cannot be caught, so the lock drops + // as soon as the process is gone. + if err := killRecorder(pid); err != nil { + return fail(err) + } + if err := c.waitForRelease(ctx, log, c.Timeout); err != nil { + if errors.Is(err, errStillHeld) { + err = fmt.Errorf("recorder pid %d still holds %s %s after SIGKILL", pid, c.Log, c.Timeout) + } + return fail(err) + } + return c.finish(log) +} + +// removePidFile is a no-op outside cleanup semantics: check must not remove the +// pidfile, because its presence is what a later run reads to learn a recording +// started. An already-absent pidfile is success; any other failure is reported, +// so cleanup cannot claim to have cleaned up when it did not. +func (c StopConfig) removePidFile() error { + if !c.Cleanup { + return nil + } + if err := os.Remove(c.PidFile); err != nil && !errors.Is(err, os.ErrNotExist) { + return fmt.Errorf("remove pidfile %s: %w", c.PidFile, err) + } + return nil +} diff --git a/grafana-alertcheck/internal/gate/stop_test.go b/grafana-alertcheck/internal/gate/stop_test.go new file mode 100644 index 000000000..087606e9a --- /dev/null +++ b/grafana-alertcheck/internal/gate/stop_test.go @@ -0,0 +1,162 @@ +package gate + +import ( + "context" + "fmt" + "os" + "os/exec" + "path/filepath" + "testing" + "time" + + "github.com/stretchr/testify/require" +) + +// An empty log is enough: StopRecorder opens it only to probe the flock. +func emptyLog(t *testing.T) string { + t.Helper() + path := filepath.Join(t.TempDir(), "log.jsonl") + require.NoError(t, os.WriteFile(path, nil, 0o644)) // nolint:gosec // test-only temp file + return path +} + +// Cleanup SIGKILLs a recorder that ignores SIGTERM. The lock holder is a real +// process that ignores SIGTERM (watch_daemon_test.go), so only a real SIGKILL +// releases the flock. +func TestStopRecorderCleanupKillsARecorderThatIgnoresSigterm(t *testing.T) { + logPath := emptyLog(t) + pid := startLockHolder(t, logPath) + writePid(t, logPath+".pid", fmt.Sprintf("%d\n", pid)) + + // Real clock, short timeout: the SIGTERM wait must elapse before SIGKILL. + held, err := StopRecorder(context.Background(), StopConfig{ + Log: logPath, + Clock: SystemClock{}, + Timeout: 200 * time.Millisecond, + Cleanup: true, + }) + require.NoError(t, err) + require.NotNil(t, held, "the log is returned held") + // A free lock is proof the SIGKILL landed. + free, err := tryLockExclusive(held) + require.NoError(t, err) + require.True(t, free, "the killed recorder still holds the lock") + require.NoError(t, held.Close()) + + require.NoFileExists(t, logPath+".pid", "cleanup removes the pidfile") +} + +// With no writer, cleanup sends no signal but still removes the pidfile. +func TestStopRecorderCleanupRemovesPidfileWhenNoWriterHoldsTheLog(t *testing.T) { + logPath := emptyLog(t) + writePid(t, logPath+".pid", "12345\n") // lock is free, so no signal is sent + + held, err := StopRecorder(context.Background(), StopConfig{ + Log: logPath, + Clock: newFakeClock(testNow), + Cleanup: true, + }) + require.NoError(t, err) + require.NotNil(t, held) + require.NoError(t, held.Close()) + require.NoFileExists(t, logPath+".pid") +} + +// A second stop while another operation still holds the returned lock must see +// a reader, not a writer; otherwise it signals the stale pid. +func TestStopRecorderDoesNotSignalPastAnotherOperationsLock(t *testing.T) { + logPath := emptyLog(t) + first, err := StopRecorder(context.Background(), StopConfig{Log: logPath, Clock: newFakeClock(testNow), Cleanup: true}) + require.NoError(t, err) + require.NotNil(t, first) + defer first.Close() + + bystander := exec.Command("sleep", "30") + require.NoError(t, bystander.Start()) + exited := make(chan struct{}) + go func() { _ = bystander.Wait(); close(exited) }() + t.Cleanup(func() { _ = bystander.Process.Kill() }) + writePid(t, logPath+".pid", fmt.Sprintf("%d\n", bystander.Process.Pid)) + + second, err := StopRecorder(context.Background(), StopConfig{Log: logPath, Clock: newFakeClock(testNow), Cleanup: true}) + require.NoError(t, err) + require.NotNil(t, second) + require.NoError(t, second.Close()) + + select { + case <-exited: + require.Fail(t, "the second stop signalled the stale pid") + case <-time.After(200 * time.Millisecond): + } +} + +// Cleanup is safe to run twice, or after check already stopped the recorder. +func TestStopRecorderCleanupIsIdempotent(t *testing.T) { + logPath := emptyLog(t) + held, err := StopRecorder(context.Background(), StopConfig{ + Log: logPath, + Clock: newFakeClock(testNow), + Cleanup: true, + }) + require.NoError(t, err) + require.NotNil(t, held, "the locked log is returned even when nothing was running") + require.NoError(t, held.Close()) +} + +// Cleanup must not read an unreadable pidfile as "nothing to stop" while a +// writer still holds the log: the flock is authoritative. +func TestStopRecorderCleanupRefusesToIgnoreALiveWriterWithNoPidfile(t *testing.T) { + logPath := emptyLog(t) + _ = startLockHolder(t, logPath) // deliberately no pidfile + + _, err := StopRecorder(context.Background(), StopConfig{ + Log: logPath, + Clock: newVirtualClock(testNow), + Cleanup: true, + }) + require.Error(t, err) + require.Contains(t, err.Error(), "unreadable") +} + +// Cleanup treats only a missing log as idempotent; a real open error (here +// ENOTDIR) must not be reported as "nothing to stop". +func TestStopRecorderCleanupPropagatesNonMissingOpenErrors(t *testing.T) { + dir := t.TempDir() + notADir := filepath.Join(dir, "file") + require.NoError(t, os.WriteFile(notADir, nil, 0o644)) // nolint:gosec // test-only temp file + logPath := filepath.Join(notADir, "log.jsonl") + + _, err := StopRecorder(context.Background(), StopConfig{ + Log: logPath, + Clock: newFakeClock(testNow), + Cleanup: true, + }) + require.Error(t, err) + require.Contains(t, err.Error(), "unreadable pidfile") +} + +// Without cleanup an absent pidfile is a hard error, as check requires. +func TestStopRecorderWithoutCleanupRefusesAMissingPidfile(t *testing.T) { + logPath := emptyLog(t) + _, err := StopRecorder(context.Background(), StopConfig{ + Log: logPath, + Clock: newFakeClock(testNow), + }) + require.Error(t, err) + require.Contains(t, err.Error(), "cannot stop the recorder") +} + +// Without cleanup a SIGTERM-ignoring holder is a hard error; the pidfile stays. +func TestStopRecorderWithoutCleanupLeavesThePidfileAndErrors(t *testing.T) { + logPath := emptyLog(t) + pid := startLockHolder(t, logPath) + writePid(t, logPath+".pid", fmt.Sprintf("%d\n", pid)) + + _, err := StopRecorder(context.Background(), StopConfig{ + Log: logPath, + Clock: newVirtualClock(testNow), + }) + require.Error(t, err) + require.Contains(t, err.Error(), "still holds") + require.FileExists(t, logPath+".pid") +} diff --git a/grafana-alertcheck/internal/gate/terminal.go b/grafana-alertcheck/internal/gate/terminal.go new file mode 100644 index 000000000..c8b4197bb --- /dev/null +++ b/grafana-alertcheck/internal/gate/terminal.go @@ -0,0 +1,120 @@ +package gate + +import ( + "fmt" + "time" +) + +// TerminationKind is the kind of terminal verdict. It reaches the JSON output. +type TerminationKind string + +const ( + // TerminationViolation is a post-`from` bad onset (new_failure or unstable). + TerminationViolation TerminationKind = "violation" + // TerminationNotVerified is an inability that has already happened. + TerminationNotVerified TerminationKind = "not_verified" +) + +// Termination is why a fail-fast run stopped before the window closed. At is +// the runner-domain detection time. +type Termination struct { + Kind TerminationKind `json:"kind"` + Alert string `json:"alert,omitempty"` + RuleUID string `json:"rule_uid,omitempty"` + Outcome Outcome `json:"outcome,omitempty"` + Reason UnobservableReason `json:"reason,omitempty"` + At time.Time `json:"at"` +} + +// terminalVerdict reports the first monotone terminal condition among the polls +// observed so far, treating at as the provisional end of the window. PURE. +// +// Only two conditions qualify, because only they can never become a pass: an +// inability that already happened, and a post-`from` bad onset (new_failure or +// unstable). `recovered` forgives an observed bad state and is reserved for +// bad-at-`from`, so a preexisting condition is deliberately not terminal. +// not_verified beats violation, as it does at the end of a full run. +func terminalVerdict(h Header, polls []Poll, defs []Definition, rt map[string]RuleTimings, + pol Policy, from, at time.Time) (Termination, bool) { + + badStates := badStateSet(pol.States) + pausedAtStart := h.pausedAtStart() + + var violation *Termination + for _, def := range defs { + if pausedAtStart[def.UID] { + continue + } + // The synthetic sentinel at `at` satisfies check 1, leaving only the + // checks decidable from the polls so far. The policy-specific nodata + // escalation is applied here too, or a configured terminal inability + // would never fail fast. + cov := proveCoverage(h, polls, &at, rt[def.UID], def, from, at, 0) + if pol.NodataIsUnobservable { + applyNodataPolicy(def, polls, &cov, rt[def.UID], from, at) + } + if cov.Unobservable { + return Termination{ + Kind: TerminationNotVerified, + Alert: def.Title, + RuleUID: def.UID, + Outcome: OutcomeNotVerified, + Reason: cov.Reason, + At: at, + }, true + } + + outcome, _, _ := classifyRule(def, polls, from, at, badStates, pol.Preexisting) + if outcome == OutcomeNewFailure || outcome == OutcomeUnstable { + if violation == nil { + v := Termination{ + Kind: TerminationViolation, + Alert: def.Title, + RuleUID: def.UID, + Outcome: outcome, + At: at, + } + violation = &v + } + } + } + if violation != nil { + return *violation, true + } + return Termination{}, false +} + +// earlyResult classifies the evidence observed when a terminal condition ended +// the run, and marks the Result as stopped early. It reuses decide unchanged by +// clamping the window to term.At and zeroing the grace (a synthetic sentinel +// at term.At satisfies check 1), then restores the requested window and real +// thresholds for reporting. +func earlyResult(h Header, polls []Poll, defs []Definition, rt map[string]RuleTimings, + gt GlobalTimings, pol Policy, term Termination) (Result, error) { + + earlyPol := pol + earlyPol.To = term.At + earlyGT := gt + earlyGT.transitionGrace = 0 + + sentinel := term.At + res, err := decide(h, polls, &sentinel, defs, rt, earlyGT, earlyPol) + + // An early Result can never be a pass; fail closed if classification lost + // the terminal verdict. + if err == nil && len(res.Violations) == 0 { + return res, fmt.Errorf("gate: fail-fast stopped on %s (%s) but classification found no violation; refusing to report a pass", + term.Kind, term.Alert) + } + + res.From = pol.From + res.To = pol.To + res.Global = GlobalThresholds{ + TransitionGrace: gt.transitionGrace, + GraceSource: graceSourceOrNone(gt.graceSource), + DrainTimeout: gt.drainTimeout, + } + t := term + res.TerminatedEarly = &t + return res, err +} diff --git a/grafana-alertcheck/internal/gate/terminal_test.go b/grafana-alertcheck/internal/gate/terminal_test.go new file mode 100644 index 000000000..a89b71437 --- /dev/null +++ b/grafana-alertcheck/internal/gate/terminal_test.go @@ -0,0 +1,205 @@ +package gate + +import ( + "encoding/json" + "testing" + "time" + + "github.com/stretchr/testify/require" +) + +func terminalHeader(from time.Time) Header { + return Header{ + SchemaVersion: LogSchemaVersion, + StartedAt: from.Add(-time.Minute), + Rules: []LoggedRule{{ + UID: checkUID, Title: checkTitle, IntervalSeconds: 60, + NoDataState: "OK", ExecErrState: "OK", + PollEverySeconds: checkPollEvery.Seconds(), + }}, + } +} + +// badFromPolls is healthy until onset, then carries the same firing instance. +func badFromPolls(uid string, from, at, onset time.Time) []Poll { + var out []Poll + for ts := from; !ts.After(at); ts = ts.Add(checkPollEvery) { + if ts.Before(onset) { + out = append(out, quietPoll(uid, ts)) + } else { + out = append(out, abnormalPoll(uid, ts, StateFiring, lbl("a"), onset)) + } + } + return out +} + +func TestTerminalVerdict(t *testing.T) { + from := testNow + at := from.Add(5 * time.Minute) + def := checkDef() + rt := map[string]RuleTimings{checkUID: newRuleTimings(checkPollEvery, 60)} + pol := Policy{From: from, To: at} + + tests := []struct { + name string + polls func() []Poll + want bool + wantKind TerminationKind + wantReason UnobservableReason + wantOut Outcome + }{ + { + name: "clean window is not terminal", + polls: func() []Poll { return denseHealthyPolls(checkUID, from, at, checkPollEvery) }, + want: false, + }, + { + name: "post-from onset is a terminal violation", + polls: func() []Poll { return badFromPolls(checkUID, from, at, from.Add(time.Minute)) }, + want: true, + wantKind: TerminationViolation, + wantOut: OutcomeNewFailure, + }, + { + // Bad before `from` can still become `recovered` (a pass). + name: "preexisting bad is not terminal", + polls: func() []Poll { + var out []Poll + for ts := from; !ts.After(at); ts = ts.Add(checkPollEvery) { + out = append(out, abnormalPoll(checkUID, ts, StateFiring, lbl("a"), from.Add(-10*time.Minute))) + } + return out + }, + want: false, + }, + { + name: "heartbeat gap is terminal not_verified", + polls: func() []Poll { + return append(denseHealthyPolls(checkUID, from, from.Add(time.Minute), checkPollEvery), + quietPoll(checkUID, at)) + }, + want: true, + wantKind: TerminationNotVerified, + wantReason: ReasonHeartbeatGap, + }, + { + name: "sustained health=error is terminal not_verified", + polls: func() []Poll { + var out []Poll + for ts := from; !ts.After(at); ts = ts.Add(checkPollEvery) { + p := quietPoll(checkUID, ts) + p.Health = "error" + out = append(out, p) + } + return out + }, + want: true, + wantKind: TerminationNotVerified, + wantReason: ReasonHealthError, + }, + { + name: "in-window pause is terminal not_verified", + polls: func() []Poll { + out := denseHealthyPolls(checkUID, from, at, checkPollEvery) + out = append(out, Poll{RuleUID: checkUID, GrafanaNow: from.Add(time.Minute), Found: true, Health: "ok", IsPaused: true}) + return out + }, + want: true, + wantKind: TerminationNotVerified, + wantReason: ReasonPausedInWindow, + }, + { + name: "absent rule is terminal not_verified", + polls: func() []Poll { + out := denseHealthyPolls(checkUID, from, at, checkPollEvery) + out = append(out, Poll{RuleUID: checkUID, GrafanaNow: from.Add(time.Minute)}) + return out + }, + want: true, + wantKind: TerminationNotVerified, + wantReason: ReasonRuleAbsent, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + term, ok := terminalVerdict(terminalHeader(from), tt.polls(), []Definition{def}, rt, pol, from, at) + require.Equal(t, tt.want, ok) + if !tt.want { + return + } + require.Equal(t, tt.wantKind, term.Kind) + require.Equal(t, tt.wantReason, term.Reason) + if tt.wantOut != "" { + require.Equal(t, tt.wantOut, term.Outcome) + } + require.True(t, term.At.Equal(at)) + }) + } +} + +// A sustained health=nodata run is terminal only when the policy escalates it, +// matching decide's end-of-run behavior. +func TestTerminalVerdictNodataIsTerminalOnlyWhenConfigured(t *testing.T) { + from := testNow + at := from.Add(5 * time.Minute) + def := checkDef() + rt := map[string]RuleTimings{checkUID: newRuleTimings(checkPollEvery, 60)} + + var polls []Poll + for ts := from; !ts.After(at); ts = ts.Add(checkPollEvery) { + p := quietPoll(checkUID, ts) + p.Health = "nodata" + polls = append(polls, p) + } + + _, ok := terminalVerdict(terminalHeader(from), polls, []Definition{def}, rt, Policy{From: from, To: at}, from, at) + require.False(t, ok, "health=nodata is only a note by default") + + term, ok := terminalVerdict(terminalHeader(from), polls, []Definition{def}, rt, + Policy{From: from, To: at, NodataIsUnobservable: true}, from, at) + require.True(t, ok) + require.Equal(t, TerminationNotVerified, term.Kind) + require.Equal(t, ReasonNodata, term.Reason) +} + +// Inability beats violation mid-window, as at the end of a full run. +func TestTerminalVerdictUnobservableBeatsViolation(t *testing.T) { + from := testNow + at := from.Add(5 * time.Minute) + + violating := checkDef() // checkUID/checkTitle, fires after from + gapped := Definition{UID: "rule-two", Title: "Rule Two", IntervalSeconds: 60} + rt := map[string]RuleTimings{ + checkUID: newRuleTimings(checkPollEvery, 60), + "rule-two": newRuleTimings(checkPollEvery, 60), + } + + polls := badFromPolls(checkUID, from, at, from.Add(time.Minute)) + // rule-two has a hole: nothing in the second half of the window. + polls = append(polls, denseHealthyPolls("rule-two", from, from.Add(time.Minute), checkPollEvery)...) + + term, ok := terminalVerdict(terminalHeader(from), polls, []Definition{violating, gapped}, rt, Policy{From: from, To: at}, from, at) + require.True(t, ok) + require.Equal(t, TerminationNotVerified, term.Kind) + require.Equal(t, "rule-two", term.RuleUID) + require.Equal(t, ReasonHeartbeatGap, term.Reason) +} + +// Termination kinds are published JSON too, and the assertions above compare +// against the constants. Pin the literals. +func TestTerminationKindJSONVocabulary(t *testing.T) { + want := map[TerminationKind]string{ + TerminationViolation: "violation", + TerminationNotVerified: "not_verified", + } + require.Len(t, want, 2, "every TerminationKind constant must be pinned here") + + for kind, literal := range want { + t.Run(string(kind), func(t *testing.T) { + raw, err := json.Marshal(kind) + require.NoError(t, err) + require.Equal(t, `"`+literal+`"`, string(raw)) + }) + } +}