Files
mesh-controller/internal/link/advisories.go
T
jochen 826dcb91b1 Keep what a consumer gives up on until a person delivers it again or drops it
Design 25 promised a dead-letter stream that did not exist: a message a
consumer gave up on stayed only in its source, which drops it after a week,
and its condition cleared when the advisories stopped (hq issue 330, ADR 0264).
2026-10-08 18:32:58 +02:00

295 lines
11 KiB
Go

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"
// AdvisoryNotKept is the token of a max-deliveries advisory whose message could not be kept in
// DEAD_LETTERS (novox/hq issue 330): said from the log, since the stream does not hold it.
AdvisoryNotKept = "not-kept"
)
// 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
// Token is the condition key's last part, when it is not Kind.
Token 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 + "/" + a.Token
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: what it asked for was not "+
"done, and the controller keeps it in %s until it is delivered again or dropped", who, a.StreamSeq,
a.Deliveries, broker.DeadLettersStream)}, 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 {
// A message given up on is said from DEAD_LETTERS, where it is kept, for as long as it is
// kept (novox/hq issue 330) — not for an hour after the server said it; one that could not
// be kept is recorded by whoever tried.
if a.Kind != AdvisoryMaxDeliveries {
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
}
// Refusals of a publish, by subject, for a caller waiting on an answer to it.
var refusalWaiters = struct {
sync.Mutex
hooked map[*nats.Conn]bool
by map[string][]chan error
}{hooked: map[*nats.Conn]bool{}, by: map[string][]chan error{}}
// refusalsOf is told when the bus refuses this connection a publish to subject, until stop is called.
// The connection's error handler is chained once, before whatever it had, which still runs.
func refusalsOf(conn *nats.Conn, subject string) (<-chan error, func()) {
ch := make(chan error, 1)
refusalWaiters.Lock()
defer refusalWaiters.Unlock()
if !refusalWaiters.hooked[conn] {
refusalWaiters.hooked[conn] = true
before := conn.ErrorHandler()
conn.SetErrorHandler(func(c *nats.Conn, sub *nats.Subscription, err error) {
if m := refusedPublish.FindStringSubmatch(errString(err)); m != nil {
refusalWaiters.Lock()
for _, w := range refusalWaiters.by[m[1]] {
select {
case w <- err:
default:
}
}
refusalWaiters.Unlock()
}
if before != nil {
before(c, sub, err)
}
})
}
refusalWaiters.by[subject] = append(refusalWaiters.by[subject], ch)
return ch, func() {
refusalWaiters.Lock()
defer refusalWaiters.Unlock()
waiting := refusalWaiters.by[subject]
for i, w := range waiting {
if w == ch {
refusalWaiters.by[subject] = append(waiting[:i], waiting[i+1:]...)
break
}
}
if len(refusalWaiters.by[subject]) == 0 {
delete(refusalWaiters.by, subject)
}
}
}
func errString(err error) string {
if err == nil {
return ""
}
return err.Error()
}