diff --git a/node-tools/cmd/node-tools/main.go b/node-tools/cmd/node-tools/main.go index a52242d..8a851bd 100644 --- a/node-tools/cmd/node-tools/main.go +++ b/node-tools/cmd/node-tools/main.go @@ -21,6 +21,7 @@ import ( "syscall" "time" + "github.com/novox/mesh-tools/node-tools/internal/alive" "github.com/novox/mesh-tools/node-tools/internal/bus" "github.com/novox/mesh-tools/node-tools/internal/console" "github.com/novox/mesh-tools/node-tools/internal/runtime" @@ -54,6 +55,12 @@ func main() { if err != nil { log.Fatalf("node-tools: %v", err) } + // And says it is there (novox/hq to-be 45 S11): the machine's runtime, on its machine's name. A + // runtime on another credential — a module's own, in a test — is not the machine's. + stopSaying := func() {} + if cred.Module == runtimeModule && cred.Node != "" { + stopSaying = alive.Keep(conn, cred.Node, alive.Every, log.Printf) + } listen := os.Getenv("MESH_CONSOLE_LISTEN") if listen == "" && cred.Module == runtimeModule { @@ -76,6 +83,7 @@ func main() { signals := make(chan os.Signal, 1) signal.Notify(signals, syscall.SIGTERM, syscall.SIGINT) <-signals + stopSaying() stop() if up != nil { _ = up.Close() diff --git a/node-tools/internal/alive/alive.go b/node-tools/internal/alive/alive.go new file mode 100644 index 0000000..f333457 --- /dev/null +++ b/node-tools/internal/alive/alive.go @@ -0,0 +1,80 @@ +// Package alive is the node tools saying they are there (novox/hq to-be 45 §3, S11). +// +// **A machine heard and its runtime gone is a machine nobody can ask anything.** The node-engine says +// every minute that the machine is there; nothing said whether the runtime serving every module's +// tools and every held seat's verbs was. Now it says so too, on its own subject, every interval, with +// the interval — the controller's watchdog raises `tools-silent` after three missed — and its build. +// Core NATS, kept by nobody, as the host's heartbeat: a lost one is the next one. +package alive + +import ( + "encoding/json" + "runtime/debug" + "sync" + "time" +) + +// Every is how often the runtime says it is there. +const Every = 60 * time.Second + +// Subject is where one machine's node tools say it. +func Subject(node string) string { return "mesh.control." + node + ".tools-alive" } + +// Beat is what is said: the machine, the interval, and the build. +type Beat struct { + Node string `json:"node"` + IntervalSeconds int `json:"interval_seconds"` + Version string `json:"version,omitempty"` +} + +// Sayer publishes one message, kept by nobody. +type Sayer interface { + Say(subject string, body []byte) error +} + +// Version is this build's commit, as the toolchain stamped it; empty when it did not. +func Version() string { + info, ok := debug.ReadBuildInfo() + if !ok { + return "" + } + for _, s := range info.Settings { + if s.Key == "vcs.revision" { + return s.Value + } + } + return "" +} + +// Keep says the runtime is there now and every interval until the returned function is called. A +// heartbeat that cannot be said is logged when that starts and when it stops — not every minute. +func Keep(bus Sayer, node string, every time.Duration, logf func(string, ...any)) func() { + body, _ := json.Marshal(Beat{Node: node, IntervalSeconds: int(every / time.Second), Version: Version()}) + done := make(chan struct{}) + var once sync.Once + go func() { + failing := "" + tick := time.NewTicker(every) + defer tick.Stop() + for { + why := "" + if err := bus.Say(Subject(node), body); err != nil { + why = err.Error() + } + if why != failing { + if why != "" { + logf("node-tools: could not say this runtime is there — the mesh will call it silent: %s", why) + } else { + logf("node-tools: says again that this runtime is there") + } + failing = why + } + select { + case <-done: + return + case <-tick.C: + } + } + }() + return func() { once.Do(func() { close(done) }) } +} diff --git a/node-tools/internal/alive/alive_test.go b/node-tools/internal/alive/alive_test.go new file mode 100644 index 0000000..f7cd680 --- /dev/null +++ b/node-tools/internal/alive/alive_test.go @@ -0,0 +1,74 @@ +package alive + +import ( + "encoding/json" + "errors" + "sync" + "testing" + "time" +) + +type said struct { + mu sync.Mutex + subs []string + body [][]byte + fail error +} + +func (s *said) Say(subject string, body []byte) error { + s.mu.Lock() + defer s.mu.Unlock() + if s.fail != nil { + return s.fail + } + s.subs, s.body = append(s.subs, subject), append(s.body, body) + return nil +} + +func (s *said) count() int { + s.mu.Lock() + defer s.mu.Unlock() + return len(s.subs) +} + +// **The runtime says it is there at once and every interval, as its machine, with the interval**: the +// controller's bound is three of them (to-be 45 S11). +func TestTheRuntimeSaysItIsThereEveryInterval(t *testing.T) { + s := &said{} + stop := Keep(s, "anchor", 20*time.Millisecond, t.Logf) + time.Sleep(70 * time.Millisecond) + stop() + stop() + if n := s.count(); n < 3 { + t.Fatalf("said %d times in three intervals", n) + } + var beat Beat + if err := json.Unmarshal(s.body[0], &beat); err != nil || beat.Node != "anchor" || beat.IntervalSeconds != 0 || + s.subs[0] != "mesh.control.anchor.tools-alive" { + t.Fatalf("%s on %s: %v", s.body[0], s.subs[0], err) + } + n := s.count() + time.Sleep(50 * time.Millisecond) + if s.count() != n { + t.Fatal("said after it was stopped") + } +} + +// **A heartbeat the bus will not take is said once, and again when it is taken**: not a line a minute. +func TestAHeartbeatThatCannotBeSaidIsLoggedOnce(t *testing.T) { + s := &said{fail: errors.New("permissions violation")} + var mu sync.Mutex + var lines []string + stop := Keep(s, "anchor", 10*time.Millisecond, func(f string, a ...any) { + mu.Lock() + lines = append(lines, f) + mu.Unlock() + }) + time.Sleep(60 * time.Millisecond) + stop() + mu.Lock() + defer mu.Unlock() + if len(lines) != 1 { + t.Fatalf("logged %d lines for one failure", len(lines)) + } +} diff --git a/node-tools/internal/bus/bus.go b/node-tools/internal/bus/bus.go index ba6cd8b..64f543d 100644 --- a/node-tools/internal/bus/bus.go +++ b/node-tools/internal/bus/bus.go @@ -628,6 +628,9 @@ func (c *Conn) PublishAs(module string, env Envelope) error { return err } +// Say publishes one message on core NATS, kept by nobody: a heartbeat, whose loss is the next one. +func (c *Conn) Say(subject string, body []byte) error { return c.nc.Publish(subject, body) } + // ServedOn is where a served module's tool is answered right now: the membership's subjects when // issued, the derived shape otherwise — what Handle subscribes, for what announces it (ADR 0197). func (c *Conn) ServedOn(module, tool string) []Served { return c.servedOn(module, tool) }