Say when the mesh is wrong: conditions, watchdogs, the bus's advisories, doctor (hq to-be 45 Phase 1)
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.
This commit is contained in:
@@ -0,0 +1,227 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/broker"
|
||||
)
|
||||
|
||||
// What the bus says about itself, in the mesh's words (novox/hq to-be 45 §3, S9).
|
||||
//
|
||||
// **The server already says it; nothing listened.** A durable consumer that hands a message over as
|
||||
// often as it may gives up on it and says so on `$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES`; one
|
||||
// deleted says so on `…CONSUMER.DELETED`. The controller's own connection is told when it falls behind
|
||||
// (a slow consumer: the client library dropped messages — issue 184, the controller deaf for 24
|
||||
// minutes) and when the bus refuses it a subject (a permissions violation — issues 183, 217, 265,
|
||||
// days each as a line in a client library's output). Each is recorded here, named in the mesh's words —
|
||||
// which consumer of whose, which subject — for the watchdog to say as a condition, as the refused
|
||||
// reply of issue 265 is said today.
|
||||
//
|
||||
// What the bus says about **other** principals' connections — a module's slow consumer, a module
|
||||
// refused a subject — the server publishes only to a system account, which the mesh's bus does not
|
||||
// have; that half is not heard yet (to-be 45 S9, recorded as deferred).
|
||||
|
||||
// The advisory kinds, as conditions name them.
|
||||
const (
|
||||
AdvisorySlowConsumer = "slow-consumer"
|
||||
AdvisoryMaxDeliveries = "max-deliveries"
|
||||
AdvisoryRefused = "refused"
|
||||
AdvisoryConsumerLost = "consumer-lost"
|
||||
)
|
||||
|
||||
// Advisory is one thing the bus said, kept as its newest word and how often it was said.
|
||||
type Advisory struct {
|
||||
Kind string
|
||||
// ID names what it is about, for the condition's key: `<stream>.<consumer>`, or `controller`.
|
||||
ID string
|
||||
// Stream and Consumer are the consumer it is about, when it is about one.
|
||||
Stream, Consumer string
|
||||
// Said is the newest saying, in the mesh's words.
|
||||
Said string
|
||||
First, Last time.Time
|
||||
Count int
|
||||
}
|
||||
|
||||
// AdvisoryLog keeps what the bus said lately.
|
||||
type AdvisoryLog struct {
|
||||
mu sync.Mutex
|
||||
seen map[string]*Advisory
|
||||
}
|
||||
|
||||
// Advisories is this process's log.
|
||||
var Advisories = &AdvisoryLog{seen: map[string]*Advisory{}}
|
||||
|
||||
// Heard records one advisory.
|
||||
func (l *AdvisoryLog) Heard(a Advisory, at time.Time) {
|
||||
l.mu.Lock()
|
||||
defer l.mu.Unlock()
|
||||
key := a.Kind + "/" + a.ID
|
||||
if had, ok := l.seen[key]; ok {
|
||||
had.Last, had.Said, had.Count = at, a.Said, had.Count+1
|
||||
return
|
||||
}
|
||||
a.First, a.Last, a.Count = at, at, 1
|
||||
l.seen[key] = &a
|
||||
}
|
||||
|
||||
// Since is every advisory said at or after a moment, and forgets the older ones.
|
||||
func (l *AdvisoryLog) Since(since time.Time) []Advisory {
|
||||
l.mu.Lock()
|
||||
defer l.mu.Unlock()
|
||||
var out []Advisory
|
||||
for key, a := range l.seen {
|
||||
if a.Last.Before(since) {
|
||||
delete(l.seen, key)
|
||||
continue
|
||||
}
|
||||
out = append(out, *a)
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// jsAdvisory is the part of a JetStream advisory the mesh reads.
|
||||
type jsAdvisory struct {
|
||||
Type string `json:"type"`
|
||||
Stream string `json:"stream"`
|
||||
Consumer string `json:"consumer"`
|
||||
StreamSeq uint64 `json:"stream_seq"`
|
||||
Deliveries uint64 `json:"deliveries"`
|
||||
Action string `json:"action"`
|
||||
}
|
||||
|
||||
// ReadAdvisory is one JetStream advisory as the mesh says it; false for one it does not watch.
|
||||
func ReadAdvisory(subject string, body []byte) (Advisory, bool) {
|
||||
var a jsAdvisory
|
||||
if err := json.Unmarshal(body, &a); err != nil || a.Stream == "" || a.Consumer == "" {
|
||||
return Advisory{}, false
|
||||
}
|
||||
who := ConsumerInWords(a.Stream, a.Consumer)
|
||||
id := a.Stream + "." + a.Consumer
|
||||
switch {
|
||||
case strings.HasPrefix(subject, "$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES."):
|
||||
return Advisory{Kind: AdvisoryMaxDeliveries, ID: id, Stream: a.Stream, Consumer: a.Consumer,
|
||||
Said: fmt.Sprintf("%s handed message %d over %d times and gave up on it: it will not be delivered "+
|
||||
"again, and what it asked for was not done", who, a.StreamSeq, a.Deliveries)}, true
|
||||
case strings.HasPrefix(subject, "$JS.EVENT.ADVISORY.CONSUMER.DELETED."):
|
||||
if !MeshNamed(a.Stream, a.Consumer) {
|
||||
// A reader's own consumer, gone when it finished — every watch of a bucket and every
|
||||
// read-back of a stream makes one, many a minute: not something the mesh defines.
|
||||
return Advisory{}, false
|
||||
}
|
||||
return Advisory{Kind: AdvisoryConsumerLost, ID: id, Stream: a.Stream, Consumer: a.Consumer,
|
||||
Said: fmt.Sprintf("%s was deleted from the bus", who)}, true
|
||||
}
|
||||
return Advisory{}, false
|
||||
}
|
||||
|
||||
// MeshNamed says a consumer is one of the durable consumers the mesh defines, by its name's shape:
|
||||
// the controller's own, a machine's declaration consumer, a module's (`<node>_<module>`), a seat's
|
||||
// worker. Every other consumer is a reader's own — an ordered consumer, a bucket's watcher — named at
|
||||
// random by the client library, deleted when the read is done, and not the mesh's to say anything of.
|
||||
func MeshNamed(stream, name string) bool {
|
||||
switch {
|
||||
case strings.HasPrefix(stream, "KV_") || stream == broker.AssignmentsStream:
|
||||
return false // the mesh defines no durable consumer on a bucket's stream, or the memberships'
|
||||
case stream == "NODES":
|
||||
return true
|
||||
case strings.HasPrefix(stream, "SEAT_"):
|
||||
return strings.HasSuffix(name, "_worker")
|
||||
case name == broker.ControllerName:
|
||||
return true
|
||||
case stream == broker.EventsStream:
|
||||
return strings.Contains(name, "_")
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// ConsumerInWords is a durable consumer as the mesh says it: whose, and for what.
|
||||
func ConsumerInWords(stream, name string) string {
|
||||
switch {
|
||||
case name == broker.ControllerName:
|
||||
return "the controller's consumer on " + stream
|
||||
case stream == "NODES":
|
||||
return "how " + name + " hears what it should be (its declaration consumer)"
|
||||
case strings.HasPrefix(stream, "SEAT_") && strings.HasSuffix(name, "_worker"):
|
||||
seat := strings.ToLower(strings.ReplaceAll(strings.TrimPrefix(stream, "SEAT_"), "_", "-"))
|
||||
return "the worker every holder of " + seat + " takes its asks from"
|
||||
case stream == broker.EventsStream:
|
||||
if node, module, ok := strings.Cut(name, "_"); ok {
|
||||
return "how " + module + " on " + node + " hears what it consumes"
|
||||
}
|
||||
}
|
||||
return "the consumer " + name + " on " + stream
|
||||
}
|
||||
|
||||
// HearAdvisories subscribes what the bus says about the mesh's account, and listens for what it
|
||||
// tells this connection, until the returned function is called. Read-only: nothing is published.
|
||||
func (s *Server) HearAdvisories(logf func(string, ...any)) (func(), error) {
|
||||
if s.js == nil {
|
||||
return nil, errors.New("this control plane is not on the bus, so it cannot hear what the bus says")
|
||||
}
|
||||
conn := s.js.Conn()
|
||||
var subs []*nats.Subscription
|
||||
stop := func() {
|
||||
for _, sub := range subs {
|
||||
_ = sub.Unsubscribe()
|
||||
}
|
||||
}
|
||||
for _, subject := range broker.BusAdvisories {
|
||||
sub, err := conn.Subscribe(subject, func(m *nats.Msg) {
|
||||
if a, ok := ReadAdvisory(m.Subject, m.Data); ok {
|
||||
Advisories.Heard(a, time.Now())
|
||||
logf("the bus says: %s", a.Said)
|
||||
}
|
||||
})
|
||||
if err != nil {
|
||||
stop()
|
||||
return nil, fmt.Errorf("listening to what the bus says (%s): %w", subject, err)
|
||||
}
|
||||
subs = append(subs, sub)
|
||||
}
|
||||
WatchConnection(conn, logf)
|
||||
return stop, nil
|
||||
}
|
||||
|
||||
// WatchConnection records what the bus tells this connection about itself: it fell behind, or a
|
||||
// subject was refused it. Chained before whatever handler the connection had, which still runs. A
|
||||
// refused answer to a call is the call log's to say (issue 265), and is not said twice.
|
||||
func WatchConnection(conn *nats.Conn, logf func(string, ...any)) {
|
||||
before := conn.ErrorHandler()
|
||||
conn.SetErrorHandler(func(c *nats.Conn, sub *nats.Subscription, err error) {
|
||||
if a, ok := connectionAdvisory(sub, err); ok {
|
||||
Advisories.Heard(a, time.Now())
|
||||
logf("the bus says: %s", a.Said)
|
||||
}
|
||||
if before != nil {
|
||||
before(c, sub, err)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
// connectionAdvisory is an error the bus handed the controller's connection, as an advisory.
|
||||
func connectionAdvisory(sub *nats.Subscription, err error) (Advisory, bool) {
|
||||
switch {
|
||||
case errors.Is(err, nats.ErrSlowConsumer):
|
||||
subject := "a subscription"
|
||||
if sub != nil {
|
||||
subject = sub.Subject
|
||||
}
|
||||
return Advisory{Kind: AdvisorySlowConsumer, ID: "controller",
|
||||
Said: fmt.Sprintf("the controller fell behind on %s and the client dropped messages it was sent", subject)}, true
|
||||
case errors.Is(err, nats.ErrPermissionViolation):
|
||||
if m := refusedPublish.FindStringSubmatch(err.Error()); m != nil && strings.HasPrefix(m[1], "_INBOX.") &&
|
||||
!strings.HasPrefix(m[1], "_INBOX.enrol.") {
|
||||
return Advisory{}, false // a refused answer to a call, which the call log says against its call
|
||||
}
|
||||
return Advisory{Kind: AdvisoryRefused, ID: "controller",
|
||||
Said: "the bus refused the controller: " + err.Error()}, true
|
||||
}
|
||||
return Advisory{}, false
|
||||
}
|
||||
@@ -0,0 +1,33 @@
|
||||
package link
|
||||
|
||||
import "testing"
|
||||
|
||||
// **A reader's own consumer gone is not a consumer lost**: every bucket watch and stream read-back
|
||||
// makes and deletes one, many a minute, and the bus says so each time (found running the controller
|
||||
// against a real bus: five a tick).
|
||||
func TestOnlyTheMeshsOwnConsumersAreSaidLost(t *testing.T) {
|
||||
deleted := func(stream, consumer string) bool {
|
||||
_, ok := ReadAdvisory("$JS.EVENT.ADVISORY.CONSUMER.DELETED."+stream+"."+consumer,
|
||||
[]byte(`{"type":"io.nats.jetstream.advisory.v1.consumer_action","stream":"`+stream+`","consumer":"`+consumer+`","action":"delete"}`))
|
||||
return ok
|
||||
}
|
||||
for _, c := range [][2]string{{"EVENTS", "controller"}, {"CONTROL", "controller"}, {"NODES", "anchor"},
|
||||
{"EVENTS", "anchor_shop"}, {"SEAT_NODE_BUILD_AGENT", "SEAT_NODE_BUILD_AGENT_worker"}} {
|
||||
if !deleted(c[0], c[1]) {
|
||||
t.Errorf("%s on %s deleted is not said", c[1], c[0])
|
||||
}
|
||||
}
|
||||
for _, c := range [][2]string{{"KV_mesh-controller_conditions", "381UWW5Y"}, {"EVENTS", "E7vVYoe6"},
|
||||
{"SEAT_NODE_BUILD_AGENT", "Rcnt6jla"}, {"ASSIGNMENTS", "x"}} {
|
||||
if deleted(c[0], c[1]) {
|
||||
t.Errorf("a reader's own consumer %s on %s is said lost", c[1], c[0])
|
||||
}
|
||||
}
|
||||
a, ok := ReadAdvisory("$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.EVENTS.anchor_shop",
|
||||
[]byte(`{"stream":"EVENTS","consumer":"anchor_shop","stream_seq":7,"deliveries":5}`))
|
||||
if !ok || a.Kind != AdvisoryMaxDeliveries || a.ID != "EVENTS.anchor_shop" ||
|
||||
a.Said != "how shop on anchor hears what it consumes handed message 7 over 5 times and gave up on it: "+
|
||||
"it will not be delivered again, and what it asked for was not done" {
|
||||
t.Fatalf("%+v", a)
|
||||
}
|
||||
}
|
||||
@@ -71,6 +71,10 @@ const (
|
||||
// next heartbeat, and a stream of them is the mesh's least valuable message competing for
|
||||
// retention with its most valuable (design 25 §3).
|
||||
AliveSubjects = "mesh.control.*.alive"
|
||||
|
||||
// ToolsAliveSubjects is every machine's node tools saying they are there (novox/hq to-be 45 §3,
|
||||
// S11): core NATS like the host's, for the same reason.
|
||||
ToolsAliveSubjects = "mesh.control.*.tools-alive"
|
||||
)
|
||||
|
||||
// ReportSubject is where one node says what it did. On the CONTROL stream, because it is the
|
||||
@@ -80,6 +84,9 @@ func ReportSubject(node string) string { return "mesh.control." + node + ".repor
|
||||
// AliveSubject is one node's heartbeat.
|
||||
func AliveSubject(node string) string { return "mesh.control." + node + ".alive" }
|
||||
|
||||
// ToolsAliveSubject is one machine's node tools' heartbeat.
|
||||
func ToolsAliveSubject(node string) string { return "mesh.control." + node + ".tools-alive" }
|
||||
|
||||
// EventSubject is where a module's event lands. Derived from the emitter, never taken from the
|
||||
// caller: a source that could differ from the subject is an envelope that can lie about its
|
||||
// origin, and on NATS the account's permissions make the subject the authority (design 29 §2).
|
||||
|
||||
@@ -271,6 +271,23 @@ func (l *CallLog) finish(c *Call, answer []byte, failed, answeredAlready bool) {
|
||||
// kept durably, every other the bus holds — a call a controller before this one served included.
|
||||
// Memory wins for a call in both, being the newer word on it. A bus that cannot be read is said in
|
||||
// the error beside what memory holds, never answered as no calls.
|
||||
// Running is every call this process is serving that has not finished, oldest first, without
|
||||
// answers: what the watchdog of a call's bound (novox/hq to-be 45 S7) reads. This process's own,
|
||||
// because a call another controller left running is said abandoned when this one starts.
|
||||
func (l *CallLog) Running() []Call {
|
||||
l.mu.Lock()
|
||||
defer l.mu.Unlock()
|
||||
var out []Call
|
||||
for _, c := range l.calls {
|
||||
if c.State == CallRunning {
|
||||
running := *c
|
||||
running.Answer = nil
|
||||
out = append(out, running)
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func (l *CallLog) Recent() ([]Call, error) {
|
||||
l.mu.Lock()
|
||||
out := make([]Call, 0, len(l.calls))
|
||||
|
||||
@@ -138,6 +138,20 @@ type Signed struct {
|
||||
// is current.
|
||||
type Alive struct {
|
||||
Node string `json:"node"`
|
||||
// IntervalSeconds is how often the node says it is there (novox/hq to-be 45 §3, S1): the
|
||||
// watchdog's bound is three of them. Zero from a host older than that, which is read as the
|
||||
// interval hosts have always used.
|
||||
IntervalSeconds int `json:"interval_seconds,omitempty"`
|
||||
}
|
||||
|
||||
// ToolsAlive is a machine's node tools saying they are there (novox/hq to-be 45 §3, S11): the
|
||||
// runtime every module's tools and every held seat's verbs are served by. Its own word, apart from the
|
||||
// host's, because a host heard and a runtime gone is a machine nobody can ask anything.
|
||||
type ToolsAlive struct {
|
||||
Node string `json:"node"`
|
||||
IntervalSeconds int `json:"interval_seconds,omitempty"`
|
||||
// Version is the runtime's build, as it says it.
|
||||
Version string `json:"version,omitempty"`
|
||||
}
|
||||
|
||||
// Report is what a node states after applying. It states; the owning context writes.
|
||||
|
||||
@@ -25,11 +25,13 @@ import (
|
||||
// on the wire: the wire is the transport's business, and a kind that travelled would be a third
|
||||
// name for the same thing.
|
||||
const (
|
||||
KindEnrolment = "enrolment"
|
||||
KindReport = "report"
|
||||
KindHeartbeat = "heartbeat"
|
||||
KindBuilt = "built"
|
||||
KindModuleMoved = "module-moved"
|
||||
KindEnrolment = "enrolment"
|
||||
KindReport = "report"
|
||||
KindHeartbeat = "heartbeat"
|
||||
// KindToolsHeartbeat is a machine's node tools saying they are there (novox/hq to-be 45 S11).
|
||||
KindToolsHeartbeat = "tools-heartbeat"
|
||||
KindBuilt = "built"
|
||||
KindModuleMoved = "module-moved"
|
||||
// KindSourceMoved is the forge announcing a merge: a source moved, and what it produces is
|
||||
// built without anybody telling the mesh (novox/hq 04-ISSUES/131).
|
||||
KindSourceMoved = "source-moved"
|
||||
|
||||
@@ -89,6 +89,10 @@ func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Con
|
||||
return nil // stopped while standing by
|
||||
}
|
||||
defer func() { _ = said.Unsubscribe() }()
|
||||
// This controller is the one acting now: its watchdogs and self-check may say what they see
|
||||
// (novox/hq to-be 45 §3). One standing by hears nothing, and would call every machine silent.
|
||||
holding.Store(true)
|
||||
defer holding.Store(false)
|
||||
|
||||
// Heartbeats, on core NATS and off any stream (design 25 §3). Their own subscription because
|
||||
// they are their own guarantee: a lost one is the next one.
|
||||
@@ -98,6 +102,12 @@ func (n *natsInbound) Receive(ctx context.Context, act func(context.Context, Con
|
||||
return fmt.Errorf("subscribing to heartbeats: %w", err)
|
||||
}
|
||||
defer func() { _ = alive.Unsubscribe() }()
|
||||
// And each machine's node tools, on the same channel: as cheap, and lost the same way.
|
||||
toolsAlive, err := conn.ChanSubscribe(ToolsAliveSubjects, beats)
|
||||
if err != nil {
|
||||
return fmt.Errorf("subscribing to the node tools' heartbeats: %w", err)
|
||||
}
|
||||
defer func() { _ = toolsAlive.Unsubscribe() }()
|
||||
|
||||
// The events the controller follows, when something is listening for them.
|
||||
var events chan *nats.Msg
|
||||
@@ -159,6 +169,8 @@ func (n *natsInbound) deliver(ctx context.Context, act func(context.Context, Con
|
||||
}
|
||||
m := &natsControl{kind: kind, msg: msg, on: n}
|
||||
if streamed {
|
||||
// The event loop took one: what S4 watches (novox/hq to-be 45 §3).
|
||||
Loop.Took(time.Now())
|
||||
// A message with no metadata is not from a stream, whatever it was delivered on, and the
|
||||
// window has nothing to hold it by. Said by leaving the sequence at zero.
|
||||
if meta, err := msg.Metadata(); err == nil {
|
||||
@@ -225,6 +237,8 @@ func kindOfSubject(subject string) (string, bool) {
|
||||
return KindReport, true
|
||||
case "alive":
|
||||
return KindHeartbeat, true
|
||||
case "tools-alive":
|
||||
return KindToolsHeartbeat, true
|
||||
}
|
||||
}
|
||||
switch subject {
|
||||
|
||||
@@ -170,6 +170,8 @@ func (s *Server) act(ctx context.Context, m Control) {
|
||||
s.reported(ctx, m)
|
||||
case KindHeartbeat:
|
||||
s.heartbeat(m)
|
||||
case KindToolsHeartbeat:
|
||||
s.toolsHeartbeat(m)
|
||||
case KindBuilt:
|
||||
s.wasBuilt(ctx, m)
|
||||
case KindModuleMoved:
|
||||
@@ -268,6 +270,8 @@ func (s *Server) heartbeat(m Control) {
|
||||
_ = m.Drop()
|
||||
return
|
||||
}
|
||||
// The interval it says, for the watchdog's bound (novox/hq to-be 45 S1).
|
||||
HostBeats.Heard(alive.Node, time.Now(), time.Duration(alive.IntervalSeconds)*time.Second, "")
|
||||
if s.listener != nil {
|
||||
if _, err := s.listener.Heard(context.Background(), Report{Node: alive.Node}); err != nil {
|
||||
s.log.Printf("could not record that %s is here: %v", alive.Node, err)
|
||||
@@ -276,6 +280,22 @@ func (s *Server) heartbeat(m Control) {
|
||||
_ = m.Took()
|
||||
}
|
||||
|
||||
// toolsHeartbeat records that a machine's node tools were heard from (novox/hq to-be 45 S11), and
|
||||
// nothing else: in this process's memory, for the watchdog — the next one is a minute away.
|
||||
func (s *Server) toolsHeartbeat(m Control) {
|
||||
var alive ToolsAlive
|
||||
if err := json.Unmarshal(m.Body(), &alive); err != nil || alive.Node == "" {
|
||||
_ = m.Drop()
|
||||
return
|
||||
}
|
||||
// The machine is the one in the subject the bus let the runtime publish on, never the body's.
|
||||
if node, ok := nodeOfToolsAlive(m.Subject()); ok {
|
||||
alive.Node = node
|
||||
}
|
||||
ToolsBeats.Heard(alive.Node, time.Now(), time.Duration(alive.IntervalSeconds)*time.Second, alive.Version)
|
||||
_ = m.Took()
|
||||
}
|
||||
|
||||
// reported records what a node says it did.
|
||||
//
|
||||
// A node states; nothing here writes anything the node claimed about itself beyond that it was
|
||||
|
||||
@@ -0,0 +1,227 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/broker"
|
||||
)
|
||||
|
||||
// What the serving controller hears that its watchdogs read (novox/hq to-be 45 §3).
|
||||
//
|
||||
// **In memory, on purpose.** A heartbeat is the least valuable message the mesh sends — the next one
|
||||
// is a minute away — and what a watchdog needs of it is the newest moment and the interval it said.
|
||||
// A machine's last word is kept in the store as well (`last_seen`), which is what S1 reads; these are
|
||||
// what the store does not keep: the interval, the node tools' word, and when the event loop last took
|
||||
// a message.
|
||||
|
||||
// Beat is one emitter's newest heartbeat.
|
||||
type Beat struct {
|
||||
At time.Time
|
||||
// Every is the interval it said; zero when it said none.
|
||||
Every time.Duration
|
||||
Version string
|
||||
}
|
||||
|
||||
// Beats is the newest heartbeat of each machine, for one kind of emitter.
|
||||
type Beats struct {
|
||||
mu sync.Mutex
|
||||
started time.Time
|
||||
heard map[string]Beat
|
||||
}
|
||||
|
||||
// NewBeats is an empty record, started now.
|
||||
func NewBeats() *Beats { return &Beats{started: time.Now(), heard: map[string]Beat{}} }
|
||||
|
||||
// HostBeats are the node-engines' heartbeats (S1); ToolsBeats the node tools' (S11).
|
||||
var (
|
||||
HostBeats = NewBeats()
|
||||
ToolsBeats = NewBeats()
|
||||
)
|
||||
|
||||
// Heard records one heartbeat.
|
||||
func (b *Beats) Heard(node string, at time.Time, every time.Duration, version string) {
|
||||
b.mu.Lock()
|
||||
defer b.mu.Unlock()
|
||||
b.heard[node] = Beat{At: at, Every: every, Version: version}
|
||||
}
|
||||
|
||||
// Of is one machine's newest heartbeat; false when none was heard since this process started.
|
||||
func (b *Beats) Of(node string) (Beat, bool) {
|
||||
b.mu.Lock()
|
||||
defer b.mu.Unlock()
|
||||
beat, ok := b.heard[node]
|
||||
return beat, ok
|
||||
}
|
||||
|
||||
// Started is when this record began: a machine not heard since is silent since then at the most.
|
||||
func (b *Beats) Started() time.Time {
|
||||
b.mu.Lock()
|
||||
defer b.mu.Unlock()
|
||||
return b.started
|
||||
}
|
||||
|
||||
// nodeOfToolsAlive is the machine a node tools' heartbeat names in its subject.
|
||||
func nodeOfToolsAlive(subject string) (string, bool) {
|
||||
rest, ok := strings.CutPrefix(subject, "mesh.control.")
|
||||
if !ok {
|
||||
return "", false
|
||||
}
|
||||
node, kind, ok := strings.Cut(rest, ".")
|
||||
if !ok || kind != "tools-alive" || node == "" {
|
||||
return "", false
|
||||
}
|
||||
return node, true
|
||||
}
|
||||
|
||||
// LoopActivity is when the controller's event loop last took a message from a stream (S4).
|
||||
type LoopActivity struct {
|
||||
mu sync.Mutex
|
||||
took time.Time
|
||||
}
|
||||
|
||||
// Loop is this process's event loop.
|
||||
var Loop = &LoopActivity{}
|
||||
|
||||
// Took records that the loop was handed a message.
|
||||
func (l *LoopActivity) Took(at time.Time) {
|
||||
l.mu.Lock()
|
||||
defer l.mu.Unlock()
|
||||
l.took = at
|
||||
}
|
||||
|
||||
// Last is when the loop last took one; zero when it has not since this process started.
|
||||
func (l *LoopActivity) Last() time.Time {
|
||||
l.mu.Lock()
|
||||
defer l.mu.Unlock()
|
||||
return l.took
|
||||
}
|
||||
|
||||
// PowerState is what a machine last said about its power (novox/hq ADR 0211): `sleeping` or
|
||||
// `shutting-down` until it says `booted` or `woke`.
|
||||
type PowerState struct {
|
||||
State string
|
||||
At time.Time
|
||||
}
|
||||
|
||||
// PowerModule is the module whose events say a machine's power (ADR 0211), and PowerEvents its
|
||||
// events that say whether the machine is about to be away or is back.
|
||||
const PowerModule = "power"
|
||||
|
||||
var PowerEvents = map[string]bool{"sleeping": true, "shutting-down": true, "booted": true, "woke": true}
|
||||
|
||||
// PowerStates is each machine's newest word about its power since a moment, read back from the
|
||||
// events stream on a consumer of its own that acknowledges nothing (as AnnouncedMerges reads). The
|
||||
// machine is the one the runtime stamped on the event (`x-node`); an event without one names none.
|
||||
func (s *Server) PowerStates(ctx context.Context, since time.Time) (map[string]PowerState, error) {
|
||||
if s.js == nil {
|
||||
return nil, errors.New("this control plane is not on the bus, so it cannot read what machines said of their power")
|
||||
}
|
||||
sub, err := s.js.Context().SubscribeSync(EventSubject(PowerModule, ">"), nats.OrderedConsumer(), nats.StartTime(since))
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("reading what machines said of their power: %w", err)
|
||||
}
|
||||
defer func() { _ = sub.Unsubscribe() }()
|
||||
out := map[string]PowerState{}
|
||||
for {
|
||||
wait, cancel := context.WithTimeout(ctx, readQuiet)
|
||||
msg, err := sub.NextMsgWithContext(wait)
|
||||
cancel()
|
||||
if err != nil {
|
||||
if ctx.Err() != nil {
|
||||
return nil, ctx.Err()
|
||||
}
|
||||
break
|
||||
}
|
||||
meta, err := msg.Metadata()
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
state := strings.TrimPrefix(msg.Subject, "mesh.mod."+PowerModule+".event.")
|
||||
node := msg.Header.Get("x-node")
|
||||
if PowerEvents[state] && node != "" {
|
||||
at := meta.Timestamp
|
||||
var body struct {
|
||||
At time.Time `json:"at"`
|
||||
}
|
||||
if json.Unmarshal(msg.Data, &body) == nil && !body.At.IsZero() {
|
||||
at = body.At
|
||||
}
|
||||
if before, ok := out[node]; !ok || !at.Before(before.At) {
|
||||
out[node] = PowerState{State: state, At: at}
|
||||
}
|
||||
}
|
||||
if meta.NumPending == 0 {
|
||||
break
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// Away says whether a power state is a machine that said it would be away.
|
||||
func (p PowerState) Away() bool { return p.State == "sleeping" || p.State == "shutting-down" }
|
||||
|
||||
// StaleRefusal is the node-engine's words for a declaration it refused because it is older than the
|
||||
// one it holds (mesh-host `refuseOlder`, novox/hq issue 107): the receiver's refusal S13 counts.
|
||||
const StaleRefusal = "is older than what the mesh last said to this node"
|
||||
|
||||
// IsStaleRefusal says a report's refusal is the node-engine refusing a declaration older than it holds.
|
||||
func IsStaleRefusal(refused string) bool { return strings.Contains(refused, StaleRefusal) }
|
||||
|
||||
// RefusalCount keeps when each machine refused a stale declaration, for S13 (novox/hq to-be 45 §3):
|
||||
// more than five from one writer in five minutes is a writer sending what it has moved past.
|
||||
type RefusalCount struct {
|
||||
mu sync.Mutex
|
||||
per map[string][]time.Time
|
||||
}
|
||||
|
||||
// StaleRefusals is this process's count.
|
||||
var StaleRefusals = &RefusalCount{per: map[string][]time.Time{}}
|
||||
|
||||
// Refused records one refusal by a machine.
|
||||
func (r *RefusalCount) Refused(node string, at time.Time) {
|
||||
r.mu.Lock()
|
||||
defer r.mu.Unlock()
|
||||
r.per[node] = append(r.per[node], at)
|
||||
}
|
||||
|
||||
// Within is how many refusals each machine made since a moment; older ones are forgotten.
|
||||
func (r *RefusalCount) Within(since time.Time) map[string]int {
|
||||
r.mu.Lock()
|
||||
defer r.mu.Unlock()
|
||||
out := map[string]int{}
|
||||
for node, times := range r.per {
|
||||
var kept []time.Time
|
||||
for _, t := range times {
|
||||
if !t.Before(since) {
|
||||
kept = append(kept, t)
|
||||
}
|
||||
}
|
||||
if len(kept) == 0 {
|
||||
delete(r.per, node)
|
||||
continue
|
||||
}
|
||||
r.per[node] = kept
|
||||
out[node] = len(kept)
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// holding is whether this process holds the controller's consumer of what nodes say: the controller
|
||||
// acting, not one standing by for another (issue 213). Until the lease (to-be 45 §6) it is how a
|
||||
// watchdog knows it is the one that hears.
|
||||
var holding atomic.Bool
|
||||
|
||||
// Holding says this process is the controller acting now.
|
||||
func Holding() bool { return holding.Load() }
|
||||
|
||||
// JetStream is the connection the server consumes on.
|
||||
func (s *Server) JetStream() *broker.JetStream { return s.js }
|
||||
Reference in New Issue
Block a user