Every one of the 48 core failures of research 031 was found by a person looking; the mesh's answers carried the fact for whoever asked and told nobody. - The condition store (to-be 45 §2): mesh-controller_conditions, one key per open condition, written by compare-and-set so a person's silence and the watchdogs never lose each other's word; every transition kept ninety days in mesh-controller_condition-history and said as the seat's events condition-raised / condition-changed / condition-cleared (the condition at the top level, with event, at, change, why, show), offered again while the bus is away. Raised and cleared by observation only; a clearing reopened within ten minutes is the same condition with its count up, its silence kept. Verbs: conditions, conditions show, conditions silence (a hand act, at most a week), conditions history. - ADR 0224's provider standing is the first kind, provider-failing, held by the provider's events; the provider_standing table is no longer read or written (left in place: dropping it is the operator's word). - status leads with the open conditions, urgent first, and says all well only with none open; conditions it cannot read are said and not well. - The signals table compiled in, one watchdog loop over it every 30s: S1 heartbeat (3 intervals, asleep machines excepted, control node urgent after 30 min), S2 report after a send, S3 plan tier, S4 event loop deaf, S5 merge not acted, S6 ask lost, S7 call hung, S8 provider silent, S9 advisories, S10 self-check silent, S11 node tools silent, S13 stale refusals; S12, S14, S15 deferred with their reasons. A row that cannot see raises probe-failed and clears nothing. A test generated from the table suppresses each signal inside and past its bound. - The bus's advisories (maximum deliveries, a mesh consumer deleted) and the controller's own slow consumer and refused subjects, said in the mesh's words. - doctor: the probe registry D1-D10 (D5 deferred) and DW, every five minutes, each in thirty seconds; a probe that cannot run is never a pass. D1 validates with mesh-host's own validator. Every run ends with the doctor-heartbeat event mesh-watcher listens for. - The controller is granted its new buckets, events, the two advisories and $SRV.INFO; the node tools their tools-alive heartbeat. The streams and consumers the controller asserts and the ones D6/D7 expect are one derivation.
348 lines
12 KiB
Go
348 lines
12 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"net"
|
|
"os"
|
|
"slices"
|
|
"strings"
|
|
"sync/atomic"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
"github.com/nats-io/nats.go/jetstream"
|
|
"golang.org/x/net/dns/dnsmessage"
|
|
|
|
"github.com/novox/mesh-controller/internal/broker"
|
|
"github.com/novox/mesh-controller/internal/conditions"
|
|
"github.com/novox/mesh-controller/internal/link"
|
|
)
|
|
|
|
// The self-check (novox/hq to-be 45 §4): a probe that fails raises its condition, one that cannot run
|
|
// raises probe-failed for itself and is never a pass, and every run ends with its heartbeat.
|
|
|
|
// withProbes runs the test with a registry of its own.
|
|
func withProbes(t *testing.T, probes ...probe) {
|
|
t.Helper()
|
|
before := probeRegistry
|
|
probeRegistry = probes
|
|
t.Cleanup(func() { probeRegistry = before })
|
|
}
|
|
|
|
func TestARunKeepsWhatEachProbeFoundAndSaysItsHeartbeat(t *testing.T) {
|
|
failing := conditions.Observation{Scope: conditions.ScopeMachine, ID: "anchor", Token: "refused",
|
|
Severity: conditions.Urgent, Summary: "anchor's node-engine would refuse its declaration"}
|
|
var broken atomic.Bool
|
|
broken.Store(true)
|
|
withProbes(t,
|
|
probe{ID: "P1", Asserts: "passes", Kind: "never", Phase: 1,
|
|
run: func(context.Context, *doctor) ([]conditions.Observation, error) { return nil, nil }},
|
|
probe{ID: "P2", Asserts: "finds a fault", Kind: "declaration-refused", Phase: 1,
|
|
run: func(context.Context, *doctor) ([]conditions.Observation, error) {
|
|
if broken.Load() {
|
|
return []conditions.Observation{failing}, nil
|
|
}
|
|
return nil, nil
|
|
}},
|
|
probe{ID: "P3", Asserts: "cannot run", Kind: "x", Phase: 1,
|
|
run: func(context.Context, *doctor) ([]conditions.Observation, error) {
|
|
if broken.Load() {
|
|
return nil, errors.New("the store is away")
|
|
}
|
|
return nil, nil
|
|
}},
|
|
probe{ID: "P4", Asserts: "hangs", Kind: "x", Phase: 1,
|
|
run: func(ctx context.Context, _ *doctor) ([]conditions.Observation, error) {
|
|
if broken.Load() {
|
|
<-ctx.Done()
|
|
time.Sleep(50 * time.Millisecond)
|
|
}
|
|
return nil, nil
|
|
}},
|
|
probe{ID: "P5", Asserts: "later", Kind: "x", Phase: 2, Deferred: "not yet"},
|
|
)
|
|
before := probeWithin
|
|
probeWithin = 200 * time.Millisecond
|
|
t.Cleanup(func() { probeWithin = before })
|
|
|
|
store := conditions.NewInMemory()
|
|
told := &conditions.Told{}
|
|
k := conditions.NewKeeper(t.Context(), conditions.Options{Store: store, History: store})
|
|
defer k.Close(context.Background())
|
|
d := &doctor{keeper: k, teller: told, host: "anchor"}
|
|
run := d.runOnce(t.Context(), "a test")
|
|
if run.Counts != (doctorCounts{Passed: 1, Failed: 1, FailedToRun: 2, Deferred: 1}) {
|
|
t.Fatalf("counted %+v", run.Counts)
|
|
}
|
|
open, err := k.Open(t.Context())
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
var keys []string
|
|
for _, c := range open {
|
|
keys = append(keys, c.Key+"="+c.Kind)
|
|
}
|
|
for _, want := range []string{"machine.anchor.refused=declaration-refused", "probe.P3.failed=probe-failed",
|
|
"probe.P4.failed=probe-failed"} {
|
|
if !slices.Contains(keys, want) {
|
|
t.Errorf("%s is not open: %v", want, keys)
|
|
}
|
|
}
|
|
// The heartbeat, in the shape mesh-watcher reads (the contract with the operator's channel).
|
|
if len(told.Names) != 1 || told.Names[0] != conditions.HeartbeatEvent {
|
|
t.Fatalf("said %v", told.Names)
|
|
}
|
|
if d.lastRunEnded().IsZero() || d.lastRun().Run != run.Run {
|
|
t.Fatal("the run is not the last verdict")
|
|
}
|
|
body, _ := json.Marshal(run)
|
|
var shape map[string]any
|
|
_ = json.Unmarshal(body, &shape)
|
|
for _, field := range []string{"run", "at", "interval-seconds", "counts", "probes", "controller"} {
|
|
if _, ok := shape[field]; !ok {
|
|
t.Errorf("the heartbeat carries no %q: %s", field, body)
|
|
}
|
|
}
|
|
|
|
// Mended: the next run clears every one of them.
|
|
broken.Store(false)
|
|
d.runOnce(t.Context(), "a test")
|
|
if open, _ := k.Open(t.Context()); len(open) != 0 {
|
|
t.Fatalf("a passing run left open %+v", open)
|
|
}
|
|
}
|
|
|
|
// **The registry says what each probe asserts**, and a probe not built says why and when.
|
|
func TestTheRegistryIsTheDesignsLiveForm(t *testing.T) {
|
|
seen := map[string]bool{}
|
|
for _, p := range probeRegistry {
|
|
if seen[p.ID] {
|
|
t.Errorf("%s twice", p.ID)
|
|
}
|
|
seen[p.ID] = true
|
|
if p.Asserts == "" || p.From == "" || p.Kind == "" {
|
|
t.Errorf("%s does not say what it asserts, where from, or what it raises", p.ID)
|
|
}
|
|
if (p.run == nil) != (p.Deferred != "") || (p.Deferred != "" && p.Phase <= 1) {
|
|
t.Errorf("%s is run and deferred, or neither, or deferred out of Phase 1: %+v", p.ID, p)
|
|
}
|
|
}
|
|
for _, id := range []string{"D1", "D2", "D3", "D4", "D5", "D6", "D7", "D8", "D9", "D10"} {
|
|
if !seen[id] {
|
|
t.Errorf("to-be 45 §4 has %s and the registry does not", id)
|
|
}
|
|
}
|
|
}
|
|
|
|
// **The doctor and the watchdogs watch each other**: watchdogs that stopped are DW; a self-check that
|
|
// stopped is S10 (signals_test.go).
|
|
func TestWatchdogsThatStoppedAreSaid(t *testing.T) {
|
|
w := &watchdogs{started: time.Now().Add(-time.Hour)}
|
|
d := &doctor{watchdogs: w, host: "anchor"}
|
|
got, err := probeWatchdogs(t.Context(), d)
|
|
if err != nil || len(got) != 1 || got[0].Severity != conditions.Urgent {
|
|
t.Fatalf("%+v %v", got, err)
|
|
}
|
|
w.ticked = time.Now()
|
|
if got, _ := probeWatchdogs(t.Context(), d); len(got) != 0 {
|
|
t.Fatalf("%+v", got)
|
|
}
|
|
}
|
|
|
|
// **D1 composes every machine of a healthy mesh and the host's own validator takes each.**
|
|
func TestEveryMachineOfAHealthyMeshComposesAndValidates(t *testing.T) {
|
|
open := aMesh(t)
|
|
got, err := probeDeclarations(t.Context(), &doctor{open: open})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if len(got) != 0 {
|
|
t.Fatalf("a healthy mesh failed D1: %+v", got)
|
|
}
|
|
}
|
|
|
|
// **D1 names a machine nothing can be sent to**, and the network that cannot be computed for it.
|
|
func TestAMachineWhoseDeclarationDoesNotComposeIsSaid(t *testing.T) {
|
|
open := aMesh(t)
|
|
ctx := t.Context()
|
|
one, two := rivals()
|
|
register(t, open, one)
|
|
register(t, open, two)
|
|
for _, m := range []string{"rival-one", "rival-two"} {
|
|
if _, err := assign(ctx, open, "laptop", m); err != nil && m == "rival-one" {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
got, err := probeDeclarations(ctx, &doctor{open: open})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if len(got) != 1 || got[0].Key() != "machine.laptop.uncomposable" || !strings.Contains(got[0].Summary, "the-seat") {
|
|
t.Fatalf("%+v", got)
|
|
}
|
|
}
|
|
|
|
// **D2: a resolver answering NXDOMAIN for IPv6 is wrong** — musl takes it as no such name (issue 262).
|
|
func TestAResolverAnsweringNoSuchNameForIPv6IsWrong(t *testing.T) {
|
|
answerAs := func(rcode dnsmessage.RCode) string {
|
|
conn, err := net.ListenPacket("udp", "127.0.0.1:0")
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
t.Cleanup(func() { _ = conn.Close() })
|
|
go func() {
|
|
buf := make([]byte, 1500)
|
|
for {
|
|
n, from, err := conn.ReadFrom(buf)
|
|
if err != nil {
|
|
return
|
|
}
|
|
var q dnsmessage.Message
|
|
if q.Unpack(buf[:n]) != nil {
|
|
continue
|
|
}
|
|
reply := dnsmessage.Message{Header: dnsmessage.Header{ID: q.ID, Response: true}, Questions: q.Questions}
|
|
if q.Questions[0].Type == dnsmessage.TypeA {
|
|
reply.Answers = []dnsmessage.Resource{{Header: dnsmessage.ResourceHeader{Name: q.Questions[0].Name,
|
|
Type: dnsmessage.TypeA, Class: dnsmessage.ClassINET}, Body: &dnsmessage.AResource{A: [4]byte{10, 77, 0, 1}}}}
|
|
} else {
|
|
reply.RCode = rcode
|
|
}
|
|
packed, _ := reply.Pack()
|
|
_, _ = conn.WriteTo(packed, from)
|
|
}
|
|
}()
|
|
_, port, _ := net.SplitHostPort(conn.LocalAddr().String())
|
|
return port
|
|
}
|
|
before := resolverPort
|
|
t.Cleanup(func() { resolverPort = before })
|
|
|
|
resolverPort = answerAs(dnsmessage.RCodeSuccess)
|
|
v4, rcode, err := askResolver(t.Context(), "127.0.0.1", "anchor.internal", dnsmessage.TypeA)
|
|
if err != nil || rcode != dnsmessage.RCodeSuccess || !slices.Equal(v4, []string{"10.77.0.1"}) {
|
|
t.Fatalf("%v %v %v", v4, rcode, err)
|
|
}
|
|
v6, rcode, err := askResolver(t.Context(), "127.0.0.1", "anchor.internal", dnsmessage.TypeAAAA)
|
|
if err != nil || rcode != dnsmessage.RCodeSuccess || len(v6) != 0 {
|
|
t.Fatalf("NODATA read as %v %v %v", v6, rcode, err)
|
|
}
|
|
resolverPort = answerAs(dnsmessage.RCodeNameError)
|
|
if _, rcode, _ := askResolver(t.Context(), "127.0.0.1", "anchor.internal", dnsmessage.TypeAAAA); rcode != dnsmessage.RCodeNameError {
|
|
t.Fatalf("NXDOMAIN read as %v", rcode)
|
|
}
|
|
}
|
|
|
|
// **D6, D7: what the controller defines is what it finds**, and a consumer deleted or a stream
|
|
// redefined is said — against a real bus, raised by the same derivation the controller starts with.
|
|
func TestNatsTheBusIsWhatTheControllerDefines(t *testing.T) {
|
|
url := os.Getenv("MESH_TEST_NATS")
|
|
if url == "" {
|
|
t.Skip("MESH_TEST_NATS unset")
|
|
}
|
|
open := aMesh(t)
|
|
js, err := broker.Dial(url)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
t.Cleanup(js.Close)
|
|
for _, s := range []string{"CONTROL", "NODES", "ASSIGNMENTS", "EVENTS"} {
|
|
_ = js.Context().DeleteStream(s)
|
|
}
|
|
if _, err := assertBusObjects(t.Context(), open.inventory, js); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := js.EnsureControllerBuckets(); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
d := &doctor{open: open, js: js}
|
|
for _, p := range []func(context.Context, *doctor) ([]conditions.Observation, error){probeConsumers, probeStreams} {
|
|
got, err := p(t.Context(), d)
|
|
if err != nil || len(got) != 0 {
|
|
t.Fatalf("a bus just raised fails: %+v %v", got, err)
|
|
}
|
|
}
|
|
if err := js.Context().DeleteConsumer("NODES", "laptop"); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
info, err := js.Context().StreamInfo("EVENTS")
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
cfg := info.Config
|
|
cfg.MaxMsgsPerSubject = 3
|
|
if _, err := js.Context().UpdateStream(&cfg); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
consumers, err := probeConsumers(t.Context(), d)
|
|
if err != nil || len(consumers) != 1 || consumers[0].Key() != "bus.NODES.laptop.missing" {
|
|
t.Fatalf("the deleted consumer: %+v %v", consumers, err)
|
|
}
|
|
streams, err := probeStreams(t.Context(), d)
|
|
if err != nil || len(streams) != 1 || !strings.Contains(streams[0].Summary, "per subject") {
|
|
t.Fatalf("the redefined stream: %+v %v", streams, err)
|
|
}
|
|
}
|
|
|
|
// **S9 hears the bus**: a consumer that gives up on a message, and one deleted, as the server says.
|
|
func TestNatsTheBusSaysAConsumerGaveUpAndOneWasDeleted(t *testing.T) {
|
|
url := os.Getenv("MESH_TEST_NATS")
|
|
if url == "" {
|
|
t.Skip("MESH_TEST_NATS unset")
|
|
}
|
|
conn, err := nats.Connect(url)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer conn.Close()
|
|
heard := make(chan *nats.Msg, 16)
|
|
for _, subject := range broker.BusAdvisories {
|
|
if _, err := conn.ChanSubscribe(subject, heard); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
api, _ := jetstream.New(conn)
|
|
_ = api.DeleteStream(t.Context(), "SEAT_ADVISED")
|
|
stream, err := api.CreateStream(t.Context(), jetstream.StreamConfig{Name: "SEAT_ADVISED", Subjects: []string{"advised.>"}})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer func() { _ = api.DeleteStream(context.Background(), "SEAT_ADVISED") }()
|
|
consumer, err := stream.CreateConsumer(t.Context(), jetstream.ConsumerConfig{Durable: "SEAT_ADVISED_worker",
|
|
AckPolicy: jetstream.AckExplicitPolicy, MaxDeliver: 1, AckWait: 100 * time.Millisecond})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := api.Publish(t.Context(), "advised.x", []byte("x")); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := consumer.Fetch(1, jetstream.FetchMaxWait(time.Second)); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
// A seat's worker, so the deletion is of a consumer the mesh names (link.MeshNamed).
|
|
// Not acknowledged: after its one delivery the consumer gives up on it — on the next fetch.
|
|
time.Sleep(300 * time.Millisecond)
|
|
_, _ = consumer.Fetch(1, jetstream.FetchMaxWait(300*time.Millisecond))
|
|
if err := stream.DeleteConsumer(t.Context(), "SEAT_ADVISED_worker"); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
kinds := map[string]string{}
|
|
deadline := time.After(5 * time.Second)
|
|
for len(kinds) < 2 {
|
|
select {
|
|
case m := <-heard:
|
|
if a, ok := link.ReadAdvisory(m.Subject, m.Data); ok {
|
|
kinds[a.Kind] = a.Said
|
|
}
|
|
case <-deadline:
|
|
t.Fatalf("the bus said only %v", kinds)
|
|
}
|
|
}
|
|
if !strings.Contains(kinds["max-deliveries"], "gave up") || !strings.Contains(kinds["consumer-lost"], "was deleted") {
|
|
t.Fatalf("%v", kinds)
|
|
}
|
|
}
|