Files
mesh-controller/internal/link/standing_test.go
T
jochen 68009b16fe Hear what providers retire, and let a person approve, reject and delete (hq ADR 0230)
A provider now waits for a person before retiring more than three consumers
or half of what it holds, and deletes only when asked. The controller is that
person's way in: it keeps waiting and rejected sets as conditions, answers them
with retire approve|reject, lists and deletes retired consumers through the
provider's own tools on its machine, records each act in the hand-act log, and
probes for anything retired longer than thirty days (D11).
2026-10-06 14:41:15 +02:00

208 lines
7.0 KiB
Go

package link
import (
"context"
"sync"
"testing"
"time"
"github.com/novox/mesh-controller/internal/broker"
)
// A provider's standing (novox/hq ADR 0224): who said it is read from the subject the bus let it
// publish on, a failing one is kept and a recovery cleared, and a recovery is never lost to a store
// that is away — said once, it would leave status naming a consumer that is fine.
type keptStandings struct {
kept []Standing
err error
}
func (k *keptStandings) Stood(_ context.Context, s Standing) (bool, error) {
if k.err != nil {
return false, k.err
}
k.kept = append(k.kept, s)
return !s.Failing, nil
}
func standingSays(t *testing.T, in *fakeInbound, to *settled, subject string, body map[string]any) Control {
t.Helper()
m := in.sends(t, to, KindProvisioner, body).(*fakeControl)
m.subject = subject
return m
}
func TestTheControllerFollowsEveryProvidersStandingAndNothingElse(t *testing.T) {
for subject, want := range map[string]string{
"mesh.mod.keycloak.event.provisioner.failing": "keycloak",
"mesh.mod.postgres.event.provisioner.recovered": "postgres",
"mesh.mod.minio.event.provisioner.retirement": "minio",
"mesh.mod.*.event.provisioner.failing": "*",
} {
got, ok := ProvisionerEmitter(subject)
if !ok || got != want {
t.Errorf("%s: %q %v", subject, got, ok)
}
if kind, _ := kindOfSubject(subject); kind != KindProvisioner {
t.Errorf("%s decodes to %q", subject, kind)
}
}
for _, subject := range []string{
"mesh.mod.keycloak.event.provisioner.other",
"mesh.mod.keycloak.event.client.created",
"mesh.mod.a.b.event.provisioner.failing",
"mesh.seat.keycloak.event.provisioner.failing",
} {
if _, ok := ProvisionerEmitter(subject); ok {
t.Errorf("%s read as a provider's standing", subject)
}
}
var follows int
for _, s := range broker.ControllerFollows {
if kind, _ := kindOfSubject(s); kind == KindProvisioner {
follows++
}
}
if follows != 3 {
t.Fatalf("the controller follows %d provider subjects, want failing, recovered and retirement (ADR 0230)", follows)
}
if !IsRetirement("mesh.mod.postgres.event.provisioner.retirement") || IsRetirement("mesh.mod.postgres.event.provisioner.failing") {
t.Fatal("a retirement word is not told from a standing")
}
}
func TestAFailingStandingIsKeptNamingTheEmitterFromTheSubject(t *testing.T) {
s, in := serving()
kept := &keptStandings{}
if err := s.Watches(kept); err != nil {
t.Fatal(err)
}
to := &settled{}
since := time.Date(2026, 10, 5, 0, 49, 0, 0, time.UTC)
s.act(t.Context(), standingSays(t, in, to, "mesh.mod.keycloak.event.provisioner.failing", map[string]any{
"provider": "oidc-client", "provider-node": "anchor", "consumer": "mesh_home_grafana",
"node": "home-server", "class": "credentials-rejected", "error": "401 invalid_grant",
"since": since, "attempts": 31000,
// A body naming another module is not believed: the subject is the bus's word.
"module": "postgres",
}))
if !to.acked || len(kept.kept) != 1 {
t.Fatalf("settled %+v, kept %+v", to, kept.kept)
}
got := kept.kept[0]
if got.Module != "keycloak" || !got.Failing || got.Consumer != "mesh_home_grafana" || got.Node != "home-server" ||
got.ProviderNode != "anchor" || got.Class != "credentials-rejected" || got.Attempts != 31000 || !got.Since.Equal(since) {
t.Fatalf("%+v", got)
}
to = &settled{}
s.act(t.Context(), standingSays(t, in, to, "mesh.mod.keycloak.event.provisioner.recovered", map[string]any{
"provider": "oidc-client", "provider-node": "anchor", "consumer": "mesh_home_grafana",
}))
if !to.acked || len(kept.kept) != 2 || kept.kept[1].Failing {
t.Fatalf("settled %+v, kept %+v", to, kept.kept)
}
}
func TestARecoveryIsHeldWhileTheStoreIsAway(t *testing.T) {
s, in := serving()
kept := &keptStandings{err: restarting}
if err := s.Watches(kept); err != nil {
t.Fatal(err)
}
to := &settled{}
s.act(t.Context(), standingSays(t, in, to, "mesh.mod.keycloak.event.provisioner.recovered",
map[string]any{"consumer": "mesh_home_grafana"}))
if !to.unsettled() || len(in.held) != 1 {
t.Fatalf("a recovery was settled while the store was away: %+v", to)
}
kept.err = nil
in.retries(t.Context(), s)
if !to.acked || len(kept.kept) != 1 {
t.Fatalf("the held recovery was not kept when the store came back: %+v %+v", to, kept.kept)
}
}
func TestAStandingThatNamesNoConsumerIsTakenAndForgotten(t *testing.T) {
s, in := serving()
kept := &keptStandings{}
if err := s.Watches(kept); err != nil {
t.Fatal(err)
}
to := &settled{}
s.act(t.Context(), standingSays(t, in, to, "mesh.mod.keycloak.event.provisioner.failing", map[string]any{}))
if !to.acked || len(kept.kept) != 0 {
t.Fatalf("%+v %+v", to, kept.kept)
}
}
func TestAStandingWithNothingKeepingItIsTaken(t *testing.T) {
s, in := serving()
to := &settled{}
s.act(t.Context(), standingSays(t, in, to, "mesh.mod.keycloak.event.provisioner.failing",
map[string]any{"consumer": "x"}))
if !to.acked {
t.Fatal("a standing nothing keeps was left for the bus to hand over again")
}
}
// Over a real bus: a provider's standing published under its own module's namespace reaches the
// controller through the events consumer's filter — the one wildcard filter on it — names the emitter
// from the subject, and is acknowledged.
func TestNatsAProvidersStandingReachesTheController(t *testing.T) {
js := aBus(t)
kept := &lockedStandings{}
s := &Server{inbound: Nats(js), bus: OverNATS{Conn: js.Conn(), JS: js.Context()}, log: quiet()}
if err := s.Follows(&toldAbout{}); err != nil {
t.Fatal(err)
}
if err := s.Watches(kept); err != nil {
t.Fatal(err)
}
ctx, stop := context.WithCancel(context.Background())
defer stop()
go func() { _ = s.Serve(ctx) }()
eventually(t, "the controller's event consumer being made", func() bool {
_, err := js.Context().ConsumerInfo("EVENTS", broker.ControllerName)
return err == nil
})
for _, event := range []string{broker.ProvisionerFailing, broker.ProvisionerRecovered} {
if _, err := js.Context().Publish("mesh.mod.keycloak.event."+event,
[]byte(`{"consumer":"mesh_home_grafana","provider-node":"anchor","class":"credentials-rejected"}`)); err != nil {
t.Fatal(err)
}
}
// Somebody else's event under the same prefix is not the controller's to hear.
if _, err := js.Context().Publish("mesh.mod.keycloak.event.client.created", []byte(`{}`)); err != nil {
t.Fatal(err)
}
eventually(t, "both standings being kept, in order, naming the emitter", func() bool {
got := kept.all()
return len(got) == 2 && got[0].Module == "keycloak" && got[0].Failing && !got[1].Failing
})
eventually(t, "both being acknowledged and nothing else delivered", func() bool {
info, err := js.Context().ConsumerInfo("EVENTS", broker.ControllerName)
return err == nil && info.NumAckPending == 0 && info.Delivered.Consumer == 2
})
}
type lockedStandings struct {
mu sync.Mutex
kept []Standing
}
func (l *lockedStandings) Stood(_ context.Context, s Standing) (bool, error) {
l.mu.Lock()
defer l.mu.Unlock()
l.kept = append(l.kept, s)
return !s.Failing, nil
}
func (l *lockedStandings) all() []Standing {
l.mu.Lock()
defer l.mu.Unlock()
return append([]Standing(nil), l.kept...)
}