Files
mesh-controller/internal/conditions/store_test.go
T
jschoubben f83fcdc15f
mesh/merge-gate pass: builds build-agent, mesh-controller → ace, g14, novox, shanks; no bus step; every machine composes with the change as it did without …
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery superseded: a newer head of the same pull request
Give a consumer's max-deliveries key one watcher, counting what was given up (issue 440)
The dead-letter row and a max-deliveries advisory without a token both said
bus.<stream>.<consumer>.max-deliveries: each look added an observation and a
line of evidence through the row, and the two overwrote each other's words.
The row owns the key since issue 330, so the advisory watcher leaves it, and
the row now says when the newest held message was given up, so a letter held
for a day no longer reads as observed every 30 seconds. Times without Happened
is refused, since the raise ignored it and an update counted it.
2026-10-11 02:30:09 +02:00

513 lines
18 KiB
Go

package conditions
import (
"context"
"encoding/json"
"errors"
"strings"
"testing"
"time"
)
// clock is a time a test moves by hand.
type clock struct{ at time.Time }
func (c *clock) now() time.Time { return c.at }
func (c *clock) pass(d time.Duration) { c.at = c.at.Add(d) }
func newClock() *clock { return &clock{at: time.Date(2026, 10, 6, 12, 0, 0, 0, time.UTC)} }
func keeper(t *testing.T) (*Keeper, *InMemory, *Told, *clock) {
t.Helper()
store, told, c := NewInMemory(), &Told{}, newClock()
k := NewKeeper(t.Context(), Options{Store: store, History: store, Teller: told, Now: c.now,
Say: func(f string, a ...any) { t.Logf(f, a...) }})
t.Cleanup(func() { k.Close(context.Background()) })
return k, store, told, c
}
// settled waits until the teller has been told n events.
func settled(t *testing.T, told *Told, n int) []Event {
t.Helper()
deadline := time.Now().Add(5 * time.Second)
for {
said := told.Said()
if len(said) >= n {
return said
}
if time.Now().After(deadline) {
t.Fatalf("told %d event(s), want %d: %+v", len(said), n, said)
}
time.Sleep(5 * time.Millisecond)
}
}
func silent(node string) Observation {
return Observation{Scope: ScopeMachine, ID: node, Kind: "silent", Machine: node, Severity: Warning,
Summary: node + " has not been heard from", Source: "S1"}
}
// **A condition is raised once, observed many times, and said on the bus only when it changes**
// (to-be 45 §2): an observation that changes nothing is written and said nowhere, or the operator's
// channel would hear the same fault every thirty seconds.
func TestAConditionIsSaidWhenItChangesNotWhenItIsSeenAgain(t *testing.T) {
k, _, told, c := keeper(t)
ctx := t.Context()
for i := 0; i < 3; i++ {
if _, err := k.Observe(ctx, silent("ace")); err != nil {
t.Fatal(err)
}
c.pass(time.Minute)
}
urgent := silent("ace")
urgent.Severity = Urgent
got, err := k.Observe(ctx, urgent)
if err != nil {
t.Fatal(err)
}
if got.Key != "machine.ace.silent" || got.Observations != 4 || got.Count != 1 || got.Severity != Urgent {
t.Fatalf("held %+v", got)
}
if len(got.Evidence) != 4 || !got.Evidence[0].At.Equal(c.at) {
t.Fatalf("evidence is not newest first: %+v", got.Evidence)
}
said := settled(t, told, 2)
if said[0].Event != EventRaised || said[0].Change != ChangeRaised || said[1].Event != EventChanged ||
said[1].Change != ChangeSeverity || said[1].Was != string(Warning) {
t.Fatalf("said %+v", said)
}
time.Sleep(50 * time.Millisecond)
if n := len(told.Said()); n != 2 {
t.Fatalf("said %d events for one raising and one change", n)
}
for i, name := range told.Names {
if name != told.Events[i].Event {
t.Errorf("event %d published as %s and says it is %s", i, name, told.Events[i].Event)
}
}
}
// **Evidence is bounded**: a condition open for a week keeps its newest ten observations, not all.
func TestEvidenceKeepsTheNewestTen(t *testing.T) {
k, _, _, c := keeper(t)
var got Condition
for i := 0; i < 25; i++ {
var err error
if got, err = k.Observe(t.Context(), silent("ace")); err != nil {
t.Fatal(err)
}
c.pass(time.Minute)
}
if len(got.Evidence) != KeptEvidence || got.Observations != 25 {
t.Fatalf("kept %d evidence of %d observations", len(got.Evidence), got.Observations)
}
}
// **Cleared and raised again within ten minutes is the same condition again** (to-be 45 §2): its
// count goes up and it is said as reopened, not as news; a person's silence of it still holds.
func TestRaisedAgainSoonAfterClearingReopens(t *testing.T) {
k, store, told, c := keeper(t)
ctx := t.Context()
if _, err := k.Observe(ctx, silent("ace")); err != nil {
t.Fatal(err)
}
if _, err := k.Silence(ctx, "machine.ace.silent", time.Hour, "jochen", "the laptop is on the train"); err != nil {
t.Fatal(err)
}
if cleared, err := k.Clear(ctx, "machine.ace.silent", "heard again"); err != nil || !cleared {
t.Fatalf("cleared %v: %v", cleared, err)
}
c.pass(5 * time.Minute)
again, err := k.Observe(ctx, silent("ace"))
if err != nil {
t.Fatal(err)
}
if again.Count != 2 || again.Silenced == nil {
t.Fatalf("reopened as %+v", again)
}
said := settled(t, told, 4)
if said[3].Event != EventRaised || said[3].Change != ChangeReopened {
t.Fatalf("the reopening was said as %+v", said[3])
}
// And from a new keeper — the controller restarted between — reading what cleared from history.
if _, err := k.Clear(ctx, "machine.ace.silent", "heard again"); err != nil {
t.Fatal(err)
}
settled(t, told, 5)
k.Close(context.Background())
next := NewKeeper(ctx, Options{Store: store, History: store, Now: c.now})
defer next.Close(context.Background())
c.pass(time.Minute)
third, err := next.Observe(ctx, silent("ace"))
if err != nil {
t.Fatal(err)
}
if third.Count != 3 {
t.Fatalf("a controller restarted between cleared and raised said it as new: %+v", third)
}
// Past the window it is news.
if _, err := next.Clear(ctx, "machine.ace.silent", "heard again"); err != nil {
t.Fatal(err)
}
c.pass(ReopenWithin + time.Minute)
fourth, err := next.Observe(ctx, silent("ace"))
if err != nil {
t.Fatal(err)
}
if fourth.Count != 1 || fourth.Silenced != nil {
t.Fatalf("raised past the window as %+v", fourth)
}
}
// **A source's whole observation clears what it no longer observes, and only its own.** A watchdog
// that stops seeing a fault says it is resolved; it does not clear what another raised.
func TestReconcileClearsOnlyTheSourcesOwn(t *testing.T) {
k, _, told, _ := keeper(t)
ctx := t.Context()
other := Observation{Scope: ScopeProbe, ID: "D3", Kind: "probe-failed", Token: "failed", Severity: Warning,
Summary: "D3 did not answer", Source: "doctor"}
if _, err := k.Observe(ctx, other); err != nil {
t.Fatal(err)
}
if err := k.Reconcile(ctx, "S1", []Observation{silent("ace"), silent("g14")}); err != nil {
t.Fatal(err)
}
if err := k.Reconcile(ctx, "S1", []Observation{silent("g14")}); err != nil {
t.Fatal(err)
}
open, err := k.Open(ctx)
if err != nil {
t.Fatal(err)
}
var keys []string
for _, c := range open {
keys = append(keys, c.Key)
}
if strings.Join(keys, ",") != "machine.g14.silent,probe.D3.failed" {
t.Fatalf("open after the second observation: %v", keys)
}
said := settled(t, told, 4)
last := said[3]
if last.Event != EventCleared || last.Key != "machine.ace.silent" || last.Why == "" {
t.Fatalf("the clearing was said as %+v", last)
}
}
// **A store that cannot be read is never an empty one** (ADR 0227 rule 4): reconciling against it
// clears nothing and says why.
func TestAnUnreadableStoreClearsNothing(t *testing.T) {
k, store, _, _ := keeper(t)
ctx := t.Context()
if _, err := k.Observe(ctx, silent("ace")); err != nil {
t.Fatal(err)
}
store.SetFail(errors.New("the bus is away"))
if err := k.Reconcile(ctx, "S1", nil); err == nil {
t.Fatal("reconciled against a store it could not read")
}
if _, err := k.Open(ctx); err == nil {
t.Fatal("an unreadable store answered as read")
}
store.SetFail(nil)
store.mu.Lock()
store.values["machine.g14.silent"] = Entry{Value: []byte("{not a condition"), Revision: 99}
store.mu.Unlock()
if _, err := k.Open(ctx); err == nil || !strings.Contains(err.Error(), "machine.g14.silent") {
t.Fatalf("an unreadable condition was left out rather than said: %v", err)
}
if cleared, err := k.Clear(ctx, "machine.g14.silent", "x"); err == nil || cleared {
t.Fatal("an unreadable condition was cleared unread")
}
}
// **A silence is bounded, says why, and ends on its own** (to-be 45 §2): the condition stays open
// through it, and its messages start again when it ends.
func TestASilenceIsBoundedAndEnds(t *testing.T) {
k, _, told, c := keeper(t)
ctx := t.Context()
if _, err := k.Observe(ctx, silent("ace")); err != nil {
t.Fatal(err)
}
if _, err := k.Silence(ctx, "machine.ace.silent", 8*24*time.Hour, "jochen", "away"); err == nil {
t.Fatal("silenced for longer than a week")
}
if _, err := k.Silence(ctx, "machine.ace.silent", time.Hour, "jochen", " "); err == nil {
t.Fatal("silenced without a reason")
}
if _, err := k.Silence(ctx, "machine.nothing.silent", time.Hour, "jochen", "x"); err == nil {
t.Fatal("silenced a condition that is not open")
}
held, err := k.Silence(ctx, "machine.ace.silent", time.Hour, "jochen", "on the train")
if err != nil {
t.Fatal(err)
}
if !held.SilencedAt(c.at) || held.Silenced.By != "jochen" {
t.Fatalf("silenced as %+v", held.Silenced)
}
c.pass(30 * time.Minute)
if err := k.EndSilences(ctx); err != nil {
t.Fatal(err)
}
c.pass(31 * time.Minute)
if err := k.EndSilences(ctx); err != nil {
t.Fatal(err)
}
got, _, _ := k.Get(ctx, "machine.ace.silent")
if got.Silenced != nil {
t.Fatalf("a silence past its end still held: %+v", got.Silenced)
}
said := settled(t, told, 3)
if said[1].Change != ChangeSilenced || said[2].Change != ChangeUnsilenced || said[2].Event != EventChanged {
t.Fatalf("said %+v", said)
}
}
// **Two writers never lose each other's word.** The serving controller observes while a person's
// command silences: the write that lost the compare-and-set reads again and redoes itself.
func TestAWriteThatLostTheRaceRedoesItself(t *testing.T) {
k, store, _, _ := keeper(t)
ctx := t.Context()
if _, err := k.Observe(ctx, silent("ace")); err != nil {
t.Fatal(err)
}
racing := &racingStore{InMemory: store, before: func() {
// Another process silences between this keeper's read and its write.
other := NewKeeper(ctx, Options{Store: store})
defer other.Close(context.Background())
if _, err := other.Silence(ctx, "machine.ace.silent", time.Hour, "jochen", "known"); err != nil {
t.Error(err)
}
}}
k.store = racing
got, err := k.Observe(ctx, silent("ace"))
if err != nil {
t.Fatal(err)
}
if got.Silenced == nil || got.Observations != 2 {
t.Fatalf("the observation overwrote the silence: %+v", got)
}
}
// racingStore lets another writer in once, between a read and the write after it.
type racingStore struct {
*InMemory
before func()
done bool
}
func (r *racingStore) Update(ctx context.Context, key string, value []byte, revision uint64) error {
if !r.done {
r.done = true
r.before()
}
return r.InMemory.Update(ctx, key, value, revision)
}
// **An observation that could not be routed is refused**, naming what it lacks.
func TestAnObservationSaysWhatItIs(t *testing.T) {
k, _, _, _ := keeper(t)
for _, o := range []Observation{
{Scope: "elsewhere", ID: "x", Kind: "k", Severity: Warning, Summary: "s", Source: "S1"},
{Scope: ScopeMachine, ID: "x", Kind: "k", Severity: "loud", Summary: "s", Source: "S1"},
{Scope: ScopeMachine, ID: "x", Kind: "k", Severity: Warning, Source: "S1"},
{Scope: ScopeMachine, ID: "x", Kind: "k", Severity: Warning, Summary: "s"},
} {
if _, err := k.Observe(t.Context(), o); err == nil {
t.Errorf("observed %+v", o)
}
}
}
// **A key holds nothing the bus would refuse or read as a wildcard**, whatever the thing is called.
func TestAKeyIsSafeForTheBus(t *testing.T) {
if got := Key(ScopeProvider, "keycloak.novox.my app*", "failing"); got != "provider.keycloak.novox.my_app_.failing" {
t.Fatalf("key %q", got)
}
if got := Key(ScopeBus, "EVENTS.>", "consumer-lost"); got != "bus.EVENTS._.consumer-lost" {
t.Fatalf("key %q", got)
}
}
// **The event's shape is a contract** (to-be 45 §2): the operator-channel's holder is written against
// these field names. A rename here is a channel that reads nothing, so they are held still.
func TestTheEventShapeIsTheContract(t *testing.T) {
k, _, told, _ := keeper(t)
if _, err := k.Observe(t.Context(), silent("ace")); err != nil {
t.Fatal(err)
}
said := settled(t, told, 1)
body, err := json.Marshal(said[0])
if err != nil {
t.Fatal(err)
}
var shape map[string]any
if err := json.Unmarshal(body, &shape); err != nil {
t.Fatal(err)
}
// The condition at the top level, kebab-case, beside what happened to it.
for _, field := range []string{"event", "at", "change", "show", "key", "kind", "subject", "severity",
"summary", "evidence", "source", "raised", "last-observed", "observations", "count", "resolver",
"silenced", "epoch"} {
if _, ok := shape[field]; !ok {
t.Errorf("the event carries no %q: %s", field, body)
}
}
if shape["silenced"] != nil {
t.Errorf("an unsilenced condition says silenced %v, not null", shape["silenced"])
}
subject, _ := shape["subject"].(map[string]any)
if subject["scope"] != "machine" || subject["id"] != "ace" || subject["machine"] != "ace" {
t.Errorf("subject %v", shape["subject"])
}
if said[0].Show != "mesh-controller.conditions key=machine.ace.silent" {
t.Errorf("show is %q", said[0].Show)
}
}
// **A transition the bus will not take is offered again**, and said lost only after TellFor.
func TestATransitionIsOfferedAgainWhileTheBusIsAway(t *testing.T) {
store, told, c := NewInMemory(), &Told{}, newClock()
told.SetFail(errors.New("no responders"))
k := NewKeeper(t.Context(), Options{Store: store, History: store, Teller: told, Now: c.now})
defer k.Close(context.Background())
if _, err := k.Observe(t.Context(), silent("ace")); err != nil {
t.Fatal(err)
}
time.Sleep(300 * time.Millisecond)
told.SetFail(nil)
said := settled(t, told, 1)
if said[0].Key != "machine.ace.silent" || k.Unsaid() != 0 {
t.Fatalf("said %+v, unsaid %d", said, k.Unsaid())
}
}
// **A reopening keeps when the fault began** (novox/hq issue 348): Raised is the reopening, Began the
// first raising — also from a keeper that read what cleared from the history — and past the window a
// raising is a new fault that begins then.
func TestAReopeningKeepsWhenTheFaultBegan(t *testing.T) {
k, store, told, c := keeper(t)
ctx := t.Context()
first, err := k.Observe(ctx, silent("ace"))
if err != nil {
t.Fatal(err)
}
if !first.Began().Equal(first.Raised) || !first.First.IsZero() {
t.Fatalf("a first raising began at %s, raised %s, first %s", first.Began(), first.Raised, first.First)
}
c.pass(30 * time.Second)
if _, err := k.Clear(ctx, "machine.ace.silent", "heard again"); err != nil {
t.Fatal(err)
}
c.pass(time.Minute)
again, err := k.Observe(ctx, silent("ace"))
if err != nil {
t.Fatal(err)
}
if !again.Raised.After(first.Raised) || !again.Began().Equal(first.Raised) || len(again.Gaps) != 1 ||
!again.Gaps[0].Reopened.Equal(again.Raised) || !again.Gaps[0].Cleared.Before(again.Raised) {
t.Fatalf("reopened: raised %s, began %s, gaps %+v; first raised %s", again.Raised, again.Began(), again.Gaps, first.Raised)
}
if !again.OpenAt(first.Raised) || again.OpenAt(again.Gaps[0].Cleared) || !again.OpenAt(again.Raised) ||
again.OpenAt(first.Raised.Add(-time.Second)) {
t.Fatalf("open at the wrong moments: %+v", again)
}
// Through a restarted controller, reading the clearing from the history.
if _, err := k.Clear(ctx, "machine.ace.silent", "heard again"); err != nil {
t.Fatal(err)
}
settled(t, told, 4)
k.Close(context.Background())
next := NewKeeper(ctx, Options{Store: store, History: store, Now: c.now})
defer next.Close(context.Background())
c.pass(time.Minute)
third, err := next.Observe(ctx, silent("ace"))
if err != nil {
t.Fatal(err)
}
if !third.Began().Equal(first.Raised) || len(third.Gaps) != 2 {
t.Fatalf("after a restart the reopening began at %s, not %s, with gaps %+v", third.Began(), first.Raised, third.Gaps)
}
if _, err := next.Clear(ctx, "machine.ace.silent", "heard again"); err != nil {
t.Fatal(err)
}
c.pass(ReopenWithin + time.Minute)
fourth, err := next.Observe(ctx, silent("ace"))
if err != nil {
t.Fatal(err)
}
if !fourth.Began().Equal(fourth.Raised) || len(fourth.Gaps) != 0 {
t.Fatalf("past the window the fault began at %s, raised %s, gaps %+v", fourth.Began(), fourth.Raised, fourth.Gaps)
}
}
// Past KeptGaps the oldest gaps go, and the fault is said to have begun after them, never before.
func TestAFaultKeepsItsNewestGapsAndForgetsWhatWasBefore(t *testing.T) {
t0 := time.Date(2026, 10, 9, 10, 0, 0, 0, time.UTC)
c := Condition{Raised: t0}
first, gaps := t0, []Gap(nil)
for i := 1; i <= KeptGaps+2; i++ {
cleared, now := t0.Add(time.Duration(2*i)*time.Minute), t0.Add(time.Duration(2*i+1)*time.Minute)
c = Condition{Raised: now}
c.reopen(first, gaps, cleared, now)
first, gaps = c.Began(), c.Gaps
}
if len(c.Gaps) != KeptGaps || !c.Began().Equal(t0.Add(5*time.Minute)) || c.OpenAt(t0.Add(time.Minute)) {
t.Fatalf("after %d gaps: began %s, %d gaps", KeptGaps+2, c.Began(), len(c.Gaps))
}
// A first time not before the reopening is no earlier fault.
d := Condition{Raised: t0}
d.reopen(t0, nil, t0.Add(-time.Minute), t0)
if !d.First.IsZero() || len(d.Gaps) != 0 {
t.Fatalf("a fault reopened at its own first raising: %+v", d)
}
}
// Two reopenings in one keeper, each after ClearSaying: both gaps kept, the fault begun at its first raising,
// and open again at the second clearing's reopening.
func TestASecondReopeningKeepsBothGaps(t *testing.T) {
k, _, _, c := keeper(t)
ctx := t.Context()
first, err := k.Observe(ctx, silent("ace"))
if err != nil {
t.Fatal(err)
}
var cleared []time.Time
for i := 0; i < 2; i++ {
c.pass(30 * time.Second)
cleared = append(cleared, c.now().UTC())
if ok, err := k.ClearSaying(ctx, "machine.ace.silent", "heard again", "ace is heard again"); err != nil || !ok {
t.Fatalf("cleared %v: %v", ok, err)
}
c.pass(time.Minute)
if _, err := k.Observe(ctx, silent("ace")); err != nil {
t.Fatal(err)
}
}
got, err := k.Open(ctx)
if err != nil || len(got) != 1 {
t.Fatalf("%+v %v", got, err)
}
g := got[0]
if !g.Began().Equal(first.Raised) || len(g.Gaps) != 2 || !g.Gaps[0].Cleared.Equal(cleared[0]) ||
!g.Gaps[1].Cleared.Equal(cleared[1]) || !g.Gaps[1].Reopened.Equal(g.Raised) || g.Count != 3 {
t.Fatalf("after two reopenings: began %s, gaps %+v, count %d", g.Began(), g.Gaps, g.Count)
}
if g.OpenAt(cleared[0].Add(time.Second)) || !g.OpenAt(cleared[0].Add(-time.Second)) || g.OpenAt(cleared[1].Add(time.Second)) {
t.Fatalf("open at the wrong moments: %+v", g)
}
}
// **Times is how often a record says it happened, never alone** (novox/hq issue 440). The raise read
// Times only beside Happened while an update read it always, so a source setting Times alone would have
// jumped the count on its second look and not its first. The keeper refuses it, saying why.
func TestTimesWithoutWhenItHappenedIsRefused(t *testing.T) {
k, store, _, _ := keeper(t)
o := Observation{Scope: ScopeMachine, ID: "ace", Kind: "silent", Machine: "ace", Severity: Warning,
Summary: "ace has not been heard from", Source: "S1", Times: 5}
_, err := k.Observe(t.Context(), o)
if err == nil || !strings.Contains(err.Error(), "when") {
t.Fatalf("Times without Happened was taken: %v", err)
}
if _, found, _ := ReadOne(t.Context(), store, o.Key()); found {
t.Fatal("a refused observation raised its condition")
}
}