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).
208 lines
7.0 KiB
Go
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...)
|
|
}
|