Files
mesh-catalog/modules/dbus/cmd/dbus-tools/watcher.go
T
jochen 7b5e1d3362 dbus: hold node-message-bus, and never restart the bus live (hq ADR 0215)
A live restart of the system bus during an upgrade hung every login on a
workstation until a reboot. The module owns the bus's packages, declares
the bus running with no restart or reload trigger, publishes only curated
events (health, services, denials; never traffic) and serves tools to look
at both buses.
2026-10-05 11:50:52 +02:00

620 lines
16 KiB
Go

package main
import (
"context"
"os"
"path/filepath"
"sort"
"strings"
"sync"
"time"
)
// The watcher publishes what matters about the machine's system bus as this module's events
// (novox/hq ADR 0215 §3): its health, a well-known service appearing or leaving, and a policy
// denial. It reads only the bus driver's own answers and signals and the bus unit's journal. It never
// becomes a monitor and never sees another peer's messages, so traffic cannot leave the machine
// through it: the event bodies are built from names, pids, units and the denial's header fields.
// Event types, as the module's manifest declares them in `emits`.
const (
BusStalled = "bus.stalled"
BusRecovered = "bus.recovered"
BusRestarted = "bus.restarted"
ServiceAppeared = "service.appeared"
ServiceLeft = "service.left"
PolicyDenied = "policy.denied"
)
// Emits is every event the watcher publishes, in the manifest's order.
var Emits = []string{BusStalled, BusRecovered, BusRestarted, ServiceAppeared, ServiceLeft, PolicyDenied}
const (
// PingEvery is how often the bus is pinged.
PingEvery = 10 * time.Second
// StallAfter is how long a ping may take before the bus counts as stalled.
StallAfter = 3 * time.Second
// Debounce is how long a name must stay as it is before its change is said: a service restarted
// within it, or one that flaps, says nothing.
Debounce = 10 * time.Second
// DenialsEvery is how often the journal is read for policy denials.
DenialsEvery = 30 * time.Second
// DenialEventEvery is the least time between two policy.denied events; denials in between are
// counted into the next one.
DenialEventEvery = 5 * time.Minute
// MaxQueue is the most events kept while the mesh's bus does not take them; the oldest go first.
MaxQueue = 1000
// DenialExamples is the most distinct denials one policy.denied names.
DenialExamples = 5
)
// NameChange is the bus driver's NameOwnerChanged: a name, its old owner and its new one.
type NameChange struct{ Name, Old, New string }
// Bus is the system bus as the watcher uses it, behind an interface so it is tested without one.
type Bus interface {
Ping(ctx context.Context) error
ID(ctx context.Context) (string, error)
PID(ctx context.Context, name string) (uint32, error)
Names(ctx context.Context) ([]string, error)
Activatable(ctx context.Context) ([]string, error)
// Changes closes when the connection is lost.
Changes() <-chan NameChange
Close()
}
// Emitter publishes one event and returns once the mesh's bus has it.
type Emitter func(eventType string, body any) error
// Journal answers the policy denials in the bus unit's journal after a cursor (or, with none, since
// a time), and the cursor to continue from.
type Journal func(ctx context.Context, after string, since time.Time) ([]Denial, string, error)
type queued struct {
Type string
Body map[string]any
}
type owner struct {
PID uint32
Process string
Unit string
}
// Watcher is the long-running half of the module.
type Watcher struct {
m *Machine
emit Emitter
dial func() (Bus, error)
journal Journal
now func() time.Time
state string // where the last seen bus identity is kept, across restarts of the runtime
retry time.Duration
mu sync.Mutex
bus Bus
queue []queued
dropped int
issue string
connected bool
stalled bool
stallSince time.Time
stallReason string
lastPing time.Duration
lastPingAt time.Time
busID string
brokerPID uint32
baselined bool
current map[string]bool
published map[string]owner
dirty map[string]time.Time
activatable map[string]bool
cursor string
denials int
denialSince time.Time
examples []Denial
lastDenial time.Time
kick chan struct{}
}
// NewWatcher is a watcher for this machine, emitting through emit, reaching the bus through dial and
// the journal through journal.
func NewWatcher(m *Machine, emit Emitter, dial func() (Bus, error), journal Journal) *Watcher {
home, _ := os.UserHomeDir()
now := time.Now
if m != nil && m.Now != nil {
now = m.Now
}
return &Watcher{m: m, emit: emit, dial: dial, journal: journal, now: now,
state: filepath.Join(home, ".local", "state", "mesh-dbus", "bus"), retry: 5 * time.Second,
current: map[string]bool{}, published: map[string]owner{}, dirty: map[string]time.Time{},
activatable: map[string]bool{}, kick: make(chan struct{}, 1)}
}
// Snapshot is what dbus_check and dbus_health show of the watcher.
type Snapshot struct {
Connected bool `json:"connected"`
Stalled bool `json:"stalled"`
StalledSince string `json:"stalled_since,omitempty"`
StallReason string `json:"stall_reason,omitempty"`
LastPingMS float64 `json:"last_ping_ms"`
LastPingAt string `json:"last_ping_at,omitempty"`
BusID string `json:"bus_id,omitempty"`
BrokerPID uint32 `json:"broker_pid,omitempty"`
Services int `json:"well_known_names"`
Pending int `json:"pending_events"`
Dropped int `json:"dropped_events,omitempty"`
Problem string `json:"problem,omitempty"`
}
func (w *Watcher) Snapshot() Snapshot {
w.mu.Lock()
defer w.mu.Unlock()
s := Snapshot{Connected: w.connected, Stalled: w.stalled, StallReason: w.stallReason,
LastPingMS: float64(w.lastPing.Microseconds()) / 1000, BusID: w.busID, BrokerPID: w.brokerPID,
Services: len(w.published), Pending: len(w.queue), Dropped: w.dropped, Problem: w.issue}
if w.stalled {
s.StalledSince = w.stallSince.UTC().Format(time.RFC3339)
}
if !w.lastPingAt.IsZero() {
s.LastPingAt = w.lastPingAt.UTC().Format(time.RFC3339)
}
return s
}
func (w *Watcher) problem(s string) {
w.mu.Lock()
w.issue = s
w.mu.Unlock()
}
// enqueue adds an event in order, stamped with when it happened. A full queue lets the oldest go and
// counts it: a mesh bus gone for a day must not grow the process without bound.
func (w *Watcher) enqueue(eventType string, body map[string]any) {
if body == nil {
body = map[string]any{}
}
body["at"] = w.now().UTC().Format(time.RFC3339)
w.mu.Lock()
w.queue = append(w.queue, queued{eventType, body})
if len(w.queue) > MaxQueue {
w.dropped += len(w.queue) - MaxQueue
w.queue = w.queue[len(w.queue)-MaxQueue:]
}
w.mu.Unlock()
select {
case w.kick <- struct{}{}:
default:
}
}
// flush publishes what waits, in order, and stops at the first the mesh's bus does not take.
func (w *Watcher) flush() {
for {
w.mu.Lock()
if len(w.queue) == 0 {
w.mu.Unlock()
return
}
next := w.queue[0]
w.mu.Unlock()
if err := w.emit(next.Type, next.Body); err != nil {
w.problem("the mesh's bus did not take " + next.Type + ": " + err.Error())
return
}
w.mu.Lock()
if len(w.queue) > 0 {
w.queue = w.queue[1:]
}
if len(w.queue) == 0 && strings.HasPrefix(w.issue, "the mesh's bus") {
w.issue = ""
}
w.mu.Unlock()
}
}
// flusher publishes on its own, so an emit waiting on the runtime never delays a ping.
func (w *Watcher) flusher(ctx context.Context) {
tick := time.NewTicker(5 * time.Second)
defer tick.Stop()
for {
select {
case <-ctx.Done():
return
case <-w.kick:
case <-tick.C:
}
w.flush()
}
}
// markStalled says bus.stalled once, until the bus answers again.
func (w *Watcher) markStalled(reason string) {
w.mu.Lock()
if w.stalled {
w.stallReason = reason
w.mu.Unlock()
return
}
w.stalled, w.stallSince, w.stallReason = true, w.now(), reason
w.mu.Unlock()
w.enqueue(BusStalled, map[string]any{"reason": reason})
}
// answered says bus.recovered when a stalled bus answers again.
func (w *Watcher) answered() {
w.mu.Lock()
if !w.stalled {
w.mu.Unlock()
return
}
since := w.stallSince
w.stalled, w.stallReason = false, ""
w.mu.Unlock()
w.enqueue(BusRecovered, map[string]any{"stalled_since": since.UTC().Format(time.RFC3339),
"stalled_seconds": int(w.now().Sub(since).Seconds())})
}
// ping asks the bus driver to answer within StallAfter.
func (w *Watcher) ping(b Bus) {
ctx, cancel := context.WithTimeout(context.Background(), StallAfter)
defer cancel()
start := time.Now()
err := b.Ping(ctx)
took := time.Since(start)
if err != nil {
if ctx.Err() != nil {
w.markStalled("the bus did not answer a ping within " + StallAfter.String())
} else {
w.markStalled("the bus answered a ping with an error: " + err.Error())
}
return
}
w.mu.Lock()
w.lastPing, w.lastPingAt = took, w.now()
w.mu.Unlock()
w.answered()
}
// PingNow pings the bus on the watcher's connection, for dbus_health.
func (w *Watcher) PingNow() (time.Duration, error) {
w.mu.Lock()
b := w.bus
w.mu.Unlock()
if b == nil {
return 0, errNotConnected
}
ctx, cancel := context.WithTimeout(context.Background(), StallAfter)
defer cancel()
start := time.Now()
err := b.Ping(ctx)
return time.Since(start), err
}
type watcherError string
func (e watcherError) Error() string { return string(e) }
const errNotConnected = watcherError("the watcher is not connected to the system bus")
// identity notes the bus's id and the bus driver's pid, and says bus.restarted when either changed
// within one boot: after a boot both change, and that is the machine's news, not the bus's.
func (w *Watcher) identity(id string, pid uint32) {
boot := w.m.BootID()
w.mu.Lock()
prevID, prevPID := w.busID, w.brokerPID
w.busID, w.brokerPID = id, pid
w.mu.Unlock()
prevBoot := boot
if prevID == "" {
if b, err := os.ReadFile(w.state); err == nil {
f := strings.Fields(string(b))
if len(f) == 3 {
prevBoot, prevID = f[0], f[1]
prevPID = parsePID(f[2])
}
}
}
if prevID != "" && prevBoot == boot && (prevID != id || prevPID != pid) {
w.enqueue(BusRestarted, map[string]any{"previous_bus_id": prevID, "bus_id": id,
"previous_pid": prevPID, "pid": pid, "unit": w.m.UnitOf(pid)})
}
if err := os.MkdirAll(filepath.Dir(w.state), 0o755); err == nil {
_ = os.WriteFile(w.state, []byte(boot+" "+id+" "+itoa(pid)+"\n"), 0o644)
}
}
// IsWellKnown is whether a name is a service's name rather than a connection's: unique names (":1.42")
// come and go with every client and are never said, nor is the bus driver's own.
func IsWellKnown(name string) bool {
return name != "" && !strings.HasPrefix(name, ":") && name != busName
}
// connected baselines the names after a (re)connect. The first time it says nothing; after a lost
// connection the difference with what was said is debounced like any other change, so a service that
// did not come back with a restarted bus is said to have left.
func (w *Watcher) connectedTo(b Bus) {
ctx, cancel := context.WithTimeout(context.Background(), StallAfter)
defer cancel()
names, err := b.Names(ctx)
if err != nil {
w.problem("listing the bus's names: " + err.Error())
return
}
act, _ := b.Activatable(ctx)
now := w.now()
w.mu.Lock()
defer w.mu.Unlock()
w.activatable = map[string]bool{}
for _, n := range act {
w.activatable[n] = true
}
cur := map[string]bool{}
for _, n := range names {
if IsWellKnown(n) {
cur[n] = true
}
}
if !w.baselined {
w.baselined = true
w.current = cur
for n := range cur {
w.published[n] = owner{}
w.dirty[n] = time.Time{} // resolved silently at the next settle
}
return
}
for n := range cur {
if _, said := w.published[n]; !said {
w.dirty[n] = now
}
}
for n := range w.published {
if !cur[n] {
w.dirty[n] = now
}
}
w.current = cur
}
// observe takes one NameOwnerChanged. Only well-known names count.
func (w *Watcher) observe(c NameChange) {
if !IsWellKnown(c.Name) {
return
}
w.mu.Lock()
defer w.mu.Unlock()
if c.New != "" {
w.current[c.Name] = true
} else {
delete(w.current, c.Name)
}
w.dirty[c.Name] = w.now()
}
// settle says what changed and stayed changed for Debounce. A name's owner is resolved when it is
// said, so the event names the process and the unit that holds it.
func (w *Watcher) settle(b Bus) {
now := w.now()
w.mu.Lock()
var due []string
for n, at := range w.dirty {
if now.Sub(at) >= Debounce {
due = append(due, n)
}
}
sort.Strings(due)
w.mu.Unlock()
for _, n := range due {
w.mu.Lock()
present := w.current[n]
was, said := w.published[n]
silent := w.dirty[n].IsZero()
delete(w.dirty, n)
activatable := w.activatable[n]
w.mu.Unlock()
switch {
case present && (!said || silent):
o := w.resolve(b, n)
w.mu.Lock()
w.published[n] = o
w.mu.Unlock()
if !silent {
w.enqueue(ServiceAppeared, o.body(n, activatable))
}
case !present && said:
w.mu.Lock()
delete(w.published, n)
w.mu.Unlock()
w.enqueue(ServiceLeft, was.body(n, activatable))
}
}
}
func (w *Watcher) resolve(b Bus, name string) owner {
if b == nil {
return owner{}
}
ctx, cancel := context.WithTimeout(context.Background(), StallAfter)
defer cancel()
pid, err := b.PID(ctx, name)
if err != nil {
return owner{}
}
return owner{PID: pid, Process: w.m.ProcessName(pid), Unit: w.m.UnitOf(pid)}
}
func (o owner) body(name string, activatable bool) map[string]any {
body := map[string]any{"name": name, "activatable": activatable}
if o.PID != 0 {
body["pid"] = o.PID
}
if o.Process != "" {
body["process"] = o.Process
}
if o.Unit != "" {
body["unit"] = o.Unit
}
return body
}
// readDenials takes the denials logged since the last read, and says policy.denied at most once per
// DenialEventEvery, with the count and a few distinct examples: a client denied in a loop must not
// flood the mesh's bus.
func (w *Watcher) readDenials(ctx context.Context, start time.Time) {
if w.journal == nil {
return
}
w.mu.Lock()
cursor := w.cursor
w.mu.Unlock()
got, next, err := w.journal(ctx, cursor, start)
if err != nil {
w.problem("reading the bus's journal: " + err.Error())
return
}
now := w.now()
w.mu.Lock()
if next != "" {
w.cursor = next
}
if strings.HasPrefix(w.issue, "reading the bus's journal") {
w.issue = ""
}
for _, d := range got {
if w.denials == 0 {
w.denialSince = now
}
w.denials++
if len(w.examples) < DenialExamples && !containsDenial(w.examples, d) {
w.examples = append(w.examples, d)
}
}
due := w.denials > 0 && (w.lastDenial.IsZero() || now.Sub(w.lastDenial) >= DenialEventEvery)
var body map[string]any
if due {
body = map[string]any{"count": w.denials, "since": w.denialSince.UTC().Format(time.RFC3339),
"examples": w.examples}
w.denials, w.examples, w.lastDenial = 0, nil, now
}
w.mu.Unlock()
if due {
w.enqueue(PolicyDenied, body)
}
}
// Run watches until ctx ends. Without the system bus it says the bus stalled, and tries again every
// few seconds; a lost connection is followed at once by a new one.
func (w *Watcher) Run(ctx context.Context) {
go w.flusher(ctx)
start := w.now()
tick := time.NewTicker(PingEvery)
defer tick.Stop()
denials := time.NewTicker(DenialsEvery)
defer denials.Stop()
retry := time.NewTimer(0)
defer retry.Stop()
var changes <-chan NameChange
for {
select {
case <-ctx.Done():
w.mu.Lock()
b := w.bus
w.bus, w.connected = nil, false
w.mu.Unlock()
if b != nil {
b.Close()
}
w.flush()
return
case <-retry.C:
b, err := w.dial()
if err != nil {
w.markStalled("the system bus is not reachable: " + err.Error())
retry.Reset(w.retry)
continue
}
idCtx, cancel := context.WithTimeout(ctx, StallAfter)
id, idErr := b.ID(idCtx)
pid, _ := b.PID(idCtx, busName)
cancel()
if idErr != nil {
b.Close()
w.markStalled("the system bus did not say its id: " + idErr.Error())
retry.Reset(w.retry)
continue
}
w.mu.Lock()
w.bus, w.connected = b, true
w.mu.Unlock()
changes = b.Changes()
w.identity(id, pid)
w.connectedTo(b)
w.answered()
w.settle(b)
case c, open := <-changes:
if !open {
w.mu.Lock()
b := w.bus
w.bus, w.connected = nil, false
w.mu.Unlock()
if b != nil {
b.Close()
}
changes = nil
retry.Reset(time.Second)
continue
}
w.observe(c)
case <-tick.C:
w.mu.Lock()
b := w.bus
w.mu.Unlock()
if b != nil {
w.ping(b)
w.settle(b)
}
case <-denials.C:
w.readDenials(ctx, start)
}
}
}
func containsDenial(list []Denial, d Denial) bool {
for _, x := range list {
if x.key() == d.key() {
return true
}
}
return false
}
func parsePID(s string) uint32 {
var n uint32
for _, c := range s {
if c < '0' || c > '9' {
return 0
}
n = n*10 + uint32(c-'0')
}
return n
}
func itoa(n uint32) string {
if n == 0 {
return "0"
}
var b [10]byte
i := len(b)
for n > 0 {
i--
b[i] = byte('0' + n%10)
n /= 10
}
return string(b[i:])
}