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:]) }