Merge pull request 'Say the node tools are there, every minute (hq to-be 45 Phase 1, S11)' (#16) from feat/a-core-that-cannot-fail-silently-phase-1 into main
mesh/delivery held for a person: merged without a passing check: only a person decides that it goes on
mesh/delivery held for a person: merged without a passing check: only a person decides that it goes on
This commit was merged in pull request #16.
This commit is contained in:
@@ -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()
|
||||
|
||||
@@ -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) }) }
|
||||
}
|
||||
@@ -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))
|
||||
}
|
||||
}
|
||||
@@ -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) }
|
||||
|
||||
Reference in New Issue
Block a user