diff --git a/modules/dbus/README.md b/modules/dbus/README.md new file mode 100644 index 0000000..692081d --- /dev/null +++ b/modules/dbus/README.md @@ -0,0 +1,111 @@ +# dbus + +A machine's message bus (novox/hq ADR 0215). This module holds `node-message-bus` on every machine, +servers included, because every machine runs a D-Bus system bus. It owns the bus implementation's +packages and declares its system service running. It publishes what matters about the bus on the +mesh's bus, and serves tools to look at the system bus and at the operator's session bus. + +## What it owns + +| | | +|---|---| +| package `dbus-broker` | the bus every machine runs | +| package `dbus-broker-units` | `dbus.service` as an alias of `dbus-broker.service`, for the system and for each user | +| package `dbus` | the bus's configuration (`system.conf`, `session.conf`), `dbus.socket`, which starts the bus at boot, and libdbus | +| service `dbus-broker.service` | declared running, with no boot state and no restart or reload trigger (below) | + +All four machines were found with dbus-broker 37 behind `dbus.service` and dbus 1.16.2, from the +distribution. One server also has an old `dbus-units` package, an empty package that depends on +`dbus-broker-units`. It is not declared: nothing needs it, and removing it is a choice for its +operator. + +## The bus is never restarted live + +On 2026-10-04 a full upgrade on a workstation restarted the system bus while the upgrade was still +running. From then on every login hung, sshd answered nothing and the machine's host stopped +reporting, until someone rebooted it at its keyboard. Every program that speaks on the bus (the +service manager, logind, the network manager, the keyring, the power module's sleep lock) holds a +connection to it. A restart takes all of them away at once, and not all of them come back. + +So the module never restarts or reloads the bus, for any change. Its service has no `restart-on` and +no `reload-on`, and the host never restarts a service that has neither. A new bus takes effect at +the next boot. `dbus_check` says when the running bus is older than an installed package, which is +the sign that a reboot is due. The distribution's own upgrade can still restart the bus. This rule +covers what the mesh does, not what pacman does. + +**Why `state: running` and no `boot`.** `dbus.service` is an alias (`systemctl is-enabled` says +`alias`), so it cannot be declared enabled: the host would run `systemctl enable` and read back +something other than `enabled`. The real unit, `dbus-broker.service`, reads `disabled` on three +machines and `enabled` on one. It needs no enabling. At boot `dbus.socket`, which the `dbus` package +links into `sockets.target`, pulls in `dbus.service`, and `dbus-broker-units` ships that name as a +link to `dbus-broker.service`. On the three machines where it reads `disabled`, enabling it would +only write a second alias link into `/etc`, and a module's apply would change something on a running +machine for no gain. The service is declared so that the mesh knows the bus is the module's and says +so when it is not running. A host finding it stopped would start it, which can only help a machine +whose bus is down. ADR 0215 §1 says "running and enabled". On these machines the package's own +socket link is what makes it start at boot. + +**A policy change** a module needs in the future uses the bus's own reload (`dbus-broker.service` is +`Type=notify-reload`), which keeps every connection. No module needs one today (below). + +## Events + +The watcher runs beside the tools in the same process, on its own connection to the system bus as the +operator's account. It reads only the bus driver's answers, the driver's `NameOwnerChanged` signal and +the bus unit's journal. + +| event | when | carries | +|---|---|---| +| `bus.stalled` | the bus does not answer the driver's `GetId` within 3 s, or cannot be reached | the reason | +| `bus.recovered` | a stalled bus answers again | since when it was stalled, and for how long | +| `bus.restarted` | after reconnecting, the bus's id or the driver's pid has changed within the same boot (after a boot both change, and `power` says `booted`) | the old and new id and pid, the unit | +| `service.appeared` | a well-known name appears and is still there 10 s later | the name, the owner's pid, process and unit, whether it is activatable | +| `service.left` | a well-known name is gone and still gone 10 s later | the name and what held it | +| `policy.denied` | the bus logged a policy denial, at most once per 5 minutes | how many since the last one, and up to 5 distinct examples: the refused message's type, sender, destination, path, interface and member | + +Unique names (`:1.42`) come and go with every client and are never published. A service restarted +within the 10 s, or one that came and went, publishes nothing. Activatable services that exit when +idle, such as hostnamed, do appear and leave, and their events say `activatable: true`. + +**Why traffic never leaves the machine.** The bus carries secrets from the keyring, notification text +and the clipboard. The watcher never becomes a monitor, so it never sees another peer's messages. Its +events are built from names, pids, units and the header fields the broker logs with a denial. The log +line itself is not passed on. The tests hold every event's body to those fields. + +Events the mesh's bus does not take wait in order and go out when it answers again, as `power`'s do. +At most 1000 events wait. When there are more, the oldest are dropped and counted. + +The watcher's ping is `GetId`, not `org.freedesktop.DBus.Peer.Ping`. The system policy refuses the +peer ping to an account that is not root and logs a denial each time, which the watcher would then +publish. + +## Tools + +| tool | | +|---|---| +| `dbus_names` | `bus`: system or session. Every well-known name with its owner's pid, process, user and unit, the activatable names that are not running, and the number of connections | +| `dbus_introspect` | `bus`, `service`, `path`. The object's interfaces with their methods, properties (type and access, never values) and signals, and its children. Uses `--auto-start=no`, so looking never starts a service | +| `dbus_monitor` | `bus`, `seconds` (at most 15), optional `match` rule and `names`. The headers of what passed (type, sender, destination, path, interface, member, error name), at most 500. Never a body: each line is decoded into the header alone. On the system bus only root may monitor, so it runs through `sudo -n` | +| `dbus_check` | the packages are installed. `dbus.service` is `dbus-broker.service` and active. The running bus is not older than the installed `dbus-broker` or `dbus` (otherwise a reboot is due). Every activatable service file's `SystemdService` exists. Policy denials in the last hour. The watcher is connected, not stalled, and has nothing stuck | +| `dbus_health` | the round trip of a ping on the watcher's connection, the connection count (from `Debug.Stats` through `sudo -n`, else the unique names), and the bus's unit, pid, start and uptime | + +Every command is bounded at 20 s and its output read up to 1 MiB. Lists stop at 500 entries. +`Debug.Stats` is asked only as root: asked as the account it is refused and logged as a denial. + +The session bus is the runtime's `DBUS_SESSION_BUS_ADDRESS`, else `$XDG_RUNTIME_DIR/bus`, else +`/run/user//bus`. On a machine where the account has no session, the tools say so. + +The check notes, without failing, systemd's own services whose `dbus-org.*` alias is missing because +they are not enabled, such as resolved, networkd and homed on the workstations. Activating them fails +by design. + +## Assignment + +On every machine, as the first holder of `node-message-bus`. The module needs nothing set: no +settings, no secrets, no ports. + +## What comes later + +ADR 0215 §4: the seat receives nothing yet. Packages ship their own D-Bus policy and service files, +and no module writes one of its own. When one does, that file will be a contribution to this seat, +and this module will list the kind and reload the bus's policy for it, without a restart. diff --git a/modules/dbus/cmd/dbus-tools/bus.go b/modules/dbus/cmd/dbus-tools/bus.go new file mode 100644 index 0000000..afa8c61 --- /dev/null +++ b/modules/dbus/cmd/dbus-tools/bus.go @@ -0,0 +1,103 @@ +package main + +import ( + "context" + + "github.com/godbus/dbus/v5" +) + +const ( + busName = "org.freedesktop.DBus" + busPath = "/org/freedesktop/DBus" +) + +// systemBus is the machine's system bus as the watcher uses it, over its own private connection: the +// bus driver's answers and its NameOwnerChanged signal, never another peer's messages. +type systemBus struct { + conn *dbus.Conn + out chan NameChange +} + +// DialSystemBus connects to the system bus as this account and listens for names changing owner. +func DialSystemBus() (Bus, error) { + conn, err := dbus.SystemBusPrivate() + if err != nil { + return nil, err + } + if err := conn.Auth(nil); err != nil { + conn.Close() + return nil, err + } + if err := conn.Hello(); err != nil { + conn.Close() + return nil, err + } + if err := conn.AddMatchSignal(dbus.WithMatchSender(busName), dbus.WithMatchInterface(busName), + dbus.WithMatchMember("NameOwnerChanged")); err != nil { + conn.Close() + return nil, err + } + raw := make(chan *dbus.Signal, 256) + conn.Signal(raw) + b := &systemBus{conn: conn, out: make(chan NameChange, 256)} + go func() { + // The signal channel closes when the connection is lost; so does this one, which is how the + // watcher learns the bus went away. + defer close(b.out) + for { + select { + case s, open := <-raw: + if !open { + return + } + if s.Name != busName+".NameOwnerChanged" || len(s.Body) != 3 { + continue + } + name, _ := s.Body[0].(string) + old, _ := s.Body[1].(string) + nw, _ := s.Body[2].(string) + b.out <- NameChange{Name: name, Old: old, New: nw} + case <-conn.Context().Done(): + return + } + } + }() + return b, nil +} + +func (b *systemBus) driver() dbus.BusObject { return b.conn.Object(busName, busPath) } + +// Ping asks the bus driver its id. Not org.freedesktop.DBus.Peer.Ping: dbus-broker's system policy +// refuses that to an account that is not root, and logs the refusal as a denial. +func (b *systemBus) Ping(ctx context.Context) error { + var id string + return b.driver().CallWithContext(ctx, busName+".GetId", 0).Store(&id) +} + +func (b *systemBus) ID(ctx context.Context) (string, error) { + var id string + err := b.driver().CallWithContext(ctx, busName+".GetId", 0).Store(&id) + return id, err +} + +func (b *systemBus) PID(ctx context.Context, name string) (uint32, error) { + var pid uint32 + err := b.driver().CallWithContext(ctx, busName+".GetConnectionUnixProcessID", 0, name).Store(&pid) + return pid, err +} + +func (b *systemBus) Names(ctx context.Context) ([]string, error) { + var names []string + err := b.driver().CallWithContext(ctx, busName+".ListNames", 0).Store(&names) + return names, err +} + +func (b *systemBus) Activatable(ctx context.Context) ([]string, error) { + var names []string + err := b.driver().CallWithContext(ctx, busName+".ListActivatableNames", 0).Store(&names) + return names, err +} + +func (b *systemBus) Changes() <-chan NameChange { return b.out } + +func (b *systemBus) Close() { b.conn.Close() } diff --git a/modules/dbus/cmd/dbus-tools/check.go b/modules/dbus/cmd/dbus-tools/check.go new file mode 100644 index 0000000..3fa3006 --- /dev/null +++ b/modules/dbus/cmd/dbus-tools/check.go @@ -0,0 +1,290 @@ +package main + +import ( + "context" + "fmt" + "sort" + "strconv" + "strings" + "time" +) + +// Packages are the bus implementation's packages, as the manifest declares them: dbus-broker (the +// bus), its units (the dbus.service alias), and the reference package, which ships the bus's +// configuration, the socket that starts it at boot, and libdbus. +var Packages = []string{"dbus", "dbus-broker", "dbus-broker-units"} + +// RunningPackages are the packages whose files the running bus loaded when it started: a newer one +// installed since takes effect only at the next boot (novox/hq ADR 0215 §2). +var RunningPackages = []string{"dbus-broker", "dbus"} + +// SystemUnit is the bus's real unit; dbus.service is its alias. +const SystemUnit = "dbus-broker.service" + +// ServiceDirs are where activatable system services are described. +var ServiceDirs = []string{"/usr/share/dbus-1/system-services", "/usr/local/share/dbus-1/system-services", + "/usr/lib/dbus-1/system-services"} + +// Package is one installed package as the package manager's local database records it. +type Package struct { + Name string + Version string + Installed time.Time +} + +// ParseDesc reads one package's desc file in the local database. +func ParseDesc(desc string) Package { + var p Package + lines := strings.Split(desc, "\n") + for i := 0; i+1 < len(lines); i++ { + v := strings.TrimSpace(lines[i+1]) + switch strings.TrimSpace(lines[i]) { + case "%NAME%": + p.Name = v + case "%VERSION%": + p.Version = v + case "%INSTALLDATE%": + if n, err := strconv.ParseInt(v, 10, 64); err == nil { + p.Installed = time.Unix(n, 0) + } + } + } + return p +} + +// InstalledPackage reads a package from the local database, without running the package manager. +func (m *Machine) InstalledPackage(name string) (Package, bool) { + for _, d := range m.glob("/var/lib/pacman/local/" + name + "-*/desc") { + if p := ParseDesc(m.read(d) + "\n"); p.Name == name { + return p, true + } + } + return Package{}, false +} + +// ParseShow reads `systemctl show` blocks: one map per unit, in the order asked. +func ParseShow(out string) []map[string]string { + var blocks []map[string]string + cur := map[string]string{} + for _, line := range strings.Split(out, "\n") { + line = strings.TrimSpace(line) + if line == "" { + if len(cur) > 0 { + blocks = append(blocks, cur) + cur = map[string]string{} + } + continue + } + if k, v, ok := strings.Cut(line, "="); ok { + cur[k] = v + } + } + if len(cur) > 0 { + blocks = append(blocks, cur) + } + return blocks +} + +// unixStamp reads systemd's "@" timestamp. +func unixStamp(s string) (time.Time, bool) { + n, err := strconv.ParseInt(strings.TrimPrefix(s, "@"), 10, 64) + if err != nil || n == 0 { + return time.Time{}, false + } + return time.Unix(n, 0), true +} + +// BusUnit is the system bus's unit as the service manager has it. +type BusUnit struct { + Unit string `json:"unit"` + Active string `json:"active"` + PID uint32 `json:"pid,omitempty"` + Since time.Time `json:"-"` + Started string `json:"started,omitempty"` +} + +// SystemBusUnit asks the service manager about the system bus through its alias, so the answer is +// the implementation's unit whichever it is. +func (m *Machine) SystemBusUnit(ctx context.Context) (BusUnit, error) { + out, err := m.Run(ctx, "systemctl", "show", "dbus.service", "-p", "Id,ActiveState,MainPID,ActiveEnterTimestamp", + "--timestamp=unix") + if err != nil { + return BusUnit{}, err + } + b := ParseShow(out) + if len(b) == 0 { + return BusUnit{}, fmt.Errorf("systemctl show answered nothing for dbus.service") + } + u := BusUnit{Unit: b[0]["Id"], Active: b[0]["ActiveState"], PID: parsePID(b[0]["MainPID"])} + if t, ok := unixStamp(b[0]["ActiveEnterTimestamp"]); ok { + u.Since, u.Started = t, t.UTC().Format(time.RFC3339) + } + return u, nil +} + +// ServiceFile is one activatable system service's description. +type ServiceFile struct { + File string + Name string + Unit string +} + +// ParseServiceFile reads the Name and SystemdService of a D-Bus service file. +func ParseServiceFile(file, content string) ServiceFile { + s := ServiceFile{File: file} + for _, line := range strings.Split(content, "\n") { + k, v, ok := strings.Cut(strings.TrimSpace(line), "=") + if !ok { + continue + } + switch strings.TrimSpace(k) { + case "Name": + s.Name = strings.TrimSpace(v) + case "SystemdService": + s.Unit = strings.TrimSpace(v) + } + } + return s +} + +// Check is one thing the module expects of the machine. +type Check struct { + Name string `json:"name"` + OK bool `json:"ok"` + Detail string `json:"detail"` +} + +// RebootDue says, per package the running bus loaded, whether a newer one was installed after the bus +// started: the bus is never restarted live, so that package waits for a boot. +func RebootDue(started time.Time, pkgs []Package) (bool, string) { + var newer []string + for _, p := range pkgs { + if !p.Installed.IsZero() && p.Installed.After(started) { + newer = append(newer, fmt.Sprintf("%s %s installed %s", p.Name, p.Version, p.Installed.UTC().Format(time.RFC3339))) + } + } + if len(newer) == 0 { + return false, "the running bus started " + started.UTC().Format(time.RFC3339) + ", after every package it loaded was installed" + } + return true, "the running bus started " + started.UTC().Format(time.RFC3339) + " and is older than " + + strings.Join(newer, ", ") + ": a reboot is due (the bus is never restarted live, novox/hq ADR 0215)" +} + +// Check says what this module expects and whether the machine meets it. +func (m *Machine) Check(ctx context.Context, w *Watcher) map[string]any { + var checks []Check + var notes []string + add := func(name string, ok bool, format string, args ...any) { + checks = append(checks, Check{Name: name, OK: ok, Detail: fmt.Sprintf(format, args...)}) + } + var running []Package + for _, name := range Packages { + p, ok := m.InstalledPackage(name) + add("package "+name, ok, "installed: %v %s", ok, p.Version) + for _, r := range RunningPackages { + if ok && r == name { + running = append(running, p) + } + } + } + unit, err := m.SystemBusUnit(ctx) + if err != nil { + add("system bus", false, "%v", err) + } else { + add("system bus", unit.Unit == SystemUnit && unit.Active == "active", + "dbus.service is %s, %s, pid %d, since %s (want %s active)", orWord(unit.Unit, "unknown"), + orWord(unit.Active, "unknown"), unit.PID, orWord(unit.Started, "unknown"), SystemUnit) + if !unit.Since.IsZero() { + due, detail := RebootDue(unit.Since, running) + add("running bus is the installed one", !due, "%s", detail) + } + } + + var files []ServiceFile + for _, dir := range ServiceDirs { + for _, f := range m.glob(dir + "/*.service") { + if s := ParseServiceFile(f, m.read(f)); s.Unit != "" { + files = append(files, s) + } + } + } + sort.Slice(files, func(i, j int) bool { return files[i].Name < files[j].Name }) + if len(files) > 0 { + args := []string{"show", "-p", "Id,LoadState"} + for _, f := range files { + args = append(args, f.Unit) + } + out, err := m.Run(ctx, "systemctl", args...) + blocks := ParseShow(out) + switch { + case err != nil: + add("activatable services", false, "systemctl show: %v", err) + case len(blocks) != len(files): + add("activatable services", false, "systemctl show answered %d units for %d service files", len(blocks), len(files)) + default: + var broken, disabled []string + for i, f := range files { + if blocks[i]["LoadState"] != "not-found" { + continue + } + // systemd's own bus services are reached through a dbus-org.* alias that exists only + // while the service is enabled: a disabled one is a choice, not a fault. + if strings.HasPrefix(f.Unit, "dbus-org.") { + disabled = append(disabled, f.Name+" → "+f.Unit) + } else { + broken = append(broken, f.Name+" → "+f.Unit+" ("+f.File+")") + } + } + add("activatable services", len(broken) == 0, "%d service files; whose unit does not exist: %s", + len(files), orWord(strings.Join(broken, ", "), "none")) + if len(disabled) > 0 { + notes = append(notes, "activation fails for "+strings.Join(disabled, ", ")+ + ": the service is not enabled, so its dbus-org alias does not exist") + } + } + } + + denials, _, err := m.Denials(ctx, "", m.Now().Add(-time.Hour)) + if err != nil { + add("policy denials in the last hour", false, "reading the journal: %v", err) + } else { + seen := map[string]bool{} + var ex []string + for _, d := range denials { + k := strings.TrimSpace(d.Type + " " + d.Interface + "." + d.Member + " to " + d.Destination) + if !seen[k] && len(ex) < DenialExamples { + seen[k] = true + ex = append(ex, k) + } + } + add("policy denials in the last hour", len(denials) == 0, "%d: %s", len(denials), orWord(strings.Join(ex, "; "), "none")) + } + + if a, err := m.SessionAddress(); err == nil { + notes = append(notes, "this account's session bus: "+a) + } else { + notes = append(notes, err.Error()) + } + + if w != nil { + s := w.Snapshot() + add("watcher", s.Connected && !s.Stalled, "connected %v, stalled %v%s, last ping %.1f ms, %d well-known names", + s.Connected, s.Stalled, orNote(s.StallReason), s.LastPingMS, s.Services) + add("events reach the mesh's bus", s.Pending == 0 && s.Problem == "", "%d event(s) waiting, %d dropped%s", + s.Pending, s.Dropped, orNote(s.Problem)) + } + failing := 0 + for _, c := range checks { + if !c.OK { + failing++ + } + } + return map[string]any{"checks": checks, "failing": failing, "notes": notes} +} + +func orNote(s string) string { + if s == "" { + return "" + } + return " (" + s + ")" +} diff --git a/modules/dbus/cmd/dbus-tools/health.go b/modules/dbus/cmd/dbus-tools/health.go new file mode 100644 index 0000000..60fb81f --- /dev/null +++ b/modules/dbus/cmd/dbus-tools/health.go @@ -0,0 +1,91 @@ +package main + +import ( + "context" + "encoding/json" + "strings" + "time" +) + +// CountConnections reads the bus driver's Debug.Stats answer (`busctl call … GetStats --json=short`): +// dbus-broker lists one accounting entry per peer, the reference daemon says ActiveConnections. +func CountConnections(out string) (int, bool) { + var reply struct { + Data []map[string]struct { + Data json.RawMessage `json:"data"` + } `json:"data"` + } + if json.Unmarshal([]byte(strings.TrimSpace(out)), &reply) != nil || len(reply.Data) != 1 { + return 0, false + } + if v, ok := reply.Data[0]["org.bus1.DBus.Debug.Stats.PeerAccounting"]; ok { + var peers []json.RawMessage + if json.Unmarshal(v.Data, &peers) == nil { + return len(peers), true + } + } + if v, ok := reply.Data[0]["ActiveConnections"]; ok { + var n int + if json.Unmarshal(v.Data, &n) == nil { + return n, true + } + } + return 0, false +} + +// CountUniqueNames counts the connections among ListNames' answer (`busctl call … ListNames`). +func CountUniqueNames(out string) (int, bool) { + var reply struct { + Data [][]string `json:"data"` + } + if json.Unmarshal([]byte(strings.TrimSpace(out)), &reply) != nil || len(reply.Data) != 1 { + return 0, false + } + n := 0 + for _, name := range reply.Data[0] { + if strings.HasPrefix(name, ":") { + n++ + } + } + return n, true +} + +// Health is the system bus's health now: a ping through the watcher's connection, how many +// connections the bus has, and how long it has run. +func (m *Machine) Health(ctx context.Context, w *Watcher) map[string]any { + h := map[string]any{"bus": "system"} + if w != nil { + took, err := w.PingNow() + if err != nil { + h["ping"] = "no answer: " + err.Error() + } else { + h["ping_ms"] = float64(took.Microseconds()) / 1000 + } + h["watcher"] = w.Snapshot() + } + // Debug.Stats is root's on the system bus; asked as the account it is refused and logged as a + // denial, which the watcher would then publish. So it is asked through sudo -n or not at all. + call := []string{"busctl", "--system", "--json=short", "--no-pager", "call", busName, busPath} + if out, err := m.privileged(ctx, call[0], append(call[1:], busName+".Debug.Stats", "GetStats")...); err == nil { + if n, ok := CountConnections(out); ok { + h["connections"], h["connections_from"] = n, "the bus's Debug.Stats" + } + } + if _, ok := h["connections"]; !ok { + if out, err := m.Run(ctx, call[0], append(call[1:], busName, "ListNames")...); err == nil { + if n, ok := CountUniqueNames(out); ok { + h["connections"], h["connections_from"] = n, "the unique names on the bus (Debug.Stats needs sudo -n)" + } + } + } + if u, err := m.SystemBusUnit(ctx); err != nil { + h["unit"] = err.Error() + } else { + h["unit"], h["active"], h["pid"] = u.Unit, u.Active, u.PID + if !u.Since.IsZero() { + h["started"] = u.Started + h["uptime"] = m.Now().Sub(u.Since).Round(time.Second).String() + } + } + return h +} diff --git a/modules/dbus/cmd/dbus-tools/introspect.go b/modules/dbus/cmd/dbus-tools/introspect.go new file mode 100644 index 0000000..8b8de4e --- /dev/null +++ b/modules/dbus/cmd/dbus-tools/introspect.go @@ -0,0 +1,133 @@ +package main + +import ( + "context" + "encoding/json" + "encoding/xml" + "fmt" + "regexp" + "strings" + + "github.com/godbus/dbus/v5/introspect" +) + +// A bus name and an object path as the specification allows them; anything else is refused before +// it reaches busctl, so no argument is ever taken for an option. +var ( + busNameRE = regexp.MustCompile(`^(:[A-Za-z0-9_-]+(\.[A-Za-z0-9_-]+)+|[A-Za-z_-][A-Za-z0-9_-]*(\.[A-Za-z_-][A-Za-z0-9_-]*)+)$`) + objectPathRE = regexp.MustCompile(`^/([A-Za-z0-9_]+(/[A-Za-z0-9_]+)*)?$`) +) + +// Member is one method or signal, with its arguments as "name type". +type Member struct { + Name string `json:"name"` + In []string `json:"in,omitempty"` + Out []string `json:"out,omitempty"` + Args []string `json:"args,omitempty"` +} + +// Property is one property's name, type and access, never its value: a value can be anything a +// service holds, and reading it is a call of its own. +type Property struct { + Name string `json:"name"` + Type string `json:"type"` + Access string `json:"access"` +} + +// Interface is one interface of an object. +type Interface struct { + Name string `json:"name"` + Methods []Member `json:"methods,omitempty"` + Properties []Property `json:"properties,omitempty"` + Signals []Member `json:"signals,omitempty"` +} + +// Object is dbus_introspect's answer. +type Object struct { + Bus string `json:"bus"` + Service string `json:"service"` + Path string `json:"path"` + Interfaces []Interface `json:"interfaces"` + Children []string `json:"children,omitempty"` +} + +// ParseIntrospection reads `busctl call … Introspect --json=short` ({"type":"s","data":[""]}). +func ParseIntrospection(out string) (introspect.Node, error) { + var reply struct { + Data []string `json:"data"` + } + var node introspect.Node + if err := json.Unmarshal([]byte(strings.TrimSpace(out)), &reply); err != nil || len(reply.Data) != 1 { + return node, fmt.Errorf("the introspection answer is not a single string") + } + if err := xml.Unmarshal([]byte(reply.Data[0]), &node); err != nil { + return node, fmt.Errorf("the introspection document does not parse: %w", err) + } + return node, nil +} + +// Shape turns an introspection document into the tool's answer. +func Shape(bus, service, path string, node introspect.Node) Object { + o := Object{Bus: bus, Service: service, Path: path, Interfaces: []Interface{}} + for _, i := range node.Interfaces { + iface := Interface{Name: i.Name} + for _, m := range i.Methods { + mem := Member{Name: m.Name} + for _, a := range m.Args { + s := strings.TrimSpace(a.Name + " " + a.Type) + if a.Direction == "out" { + mem.Out = append(mem.Out, s) + } else { + mem.In = append(mem.In, s) + } + } + iface.Methods = append(iface.Methods, mem) + } + for _, p := range i.Properties { + iface.Properties = append(iface.Properties, Property{Name: p.Name, Type: p.Type, Access: p.Access}) + } + for _, s := range i.Signals { + mem := Member{Name: s.Name} + for _, a := range s.Args { + mem.Args = append(mem.Args, strings.TrimSpace(a.Name+" "+a.Type)) + } + iface.Signals = append(iface.Signals, mem) + } + o.Interfaces = append(o.Interfaces, iface) + } + for _, c := range node.Children { + if len(o.Children) >= AnswerCap { + break + } + o.Children = append(o.Children, c.Name) + } + return o +} + +// Introspect asks one object what it offers, as this account, without starting a service that is not +// running (--auto-start=no): looking must not change what runs. +func (m *Machine) Introspect(ctx context.Context, bus, service, path string) (Object, error) { + if !busNameRE.MatchString(service) { + return Object{}, fmt.Errorf("%q is not a bus name", service) + } + if path == "" { + path = "/" + } + if !objectPathRE.MatchString(path) { + return Object{}, fmt.Errorf("%q is not an object path", path) + } + args, err := m.busArgs(bus) + if err != nil { + return Object{}, err + } + out, err := m.Run(ctx, "busctl", append(args, "--json=short", "--no-pager", "--auto-start=no", "call", + service, path, "org.freedesktop.DBus.Introspectable", "Introspect")...) + if err != nil { + return Object{}, err + } + node, err := ParseIntrospection(out) + if err != nil { + return Object{}, err + } + return Shape(orWord(bus, "system"), service, path, node), nil +} diff --git a/modules/dbus/cmd/dbus-tools/journal.go b/modules/dbus/cmd/dbus-tools/journal.go new file mode 100644 index 0000000..dbb3e6a --- /dev/null +++ b/modules/dbus/cmd/dbus-tools/journal.go @@ -0,0 +1,114 @@ +package main + +import ( + "context" + "encoding/json" + "strconv" + "strings" + "time" +) + +// BusUnits are the system bus's units as the journal knows them: dbus-broker's own name, and the +// alias every implementation answers to. +var BusUnits = []string{"dbus-broker.service", "dbus.service"} + +// JournalLines is the most journal entries one read takes. +const JournalLines = 2000 + +// Denial is one policy denial as the bus logged it: the header of the message it refused, never its +// body, which the bus does not log either. +type Denial struct { + At string `json:"at"` + Action string `json:"action,omitempty"` + Type string `json:"type,omitempty"` + Sender string `json:"sender,omitempty"` + Destination string `json:"destination,omitempty"` + Path string `json:"path,omitempty"` + Interface string `json:"interface,omitempty"` + Member string `json:"member,omitempty"` + Policy string `json:"policy,omitempty"` +} + +// key is what makes two denials the same example: who was refused what, ignoring the sender's unique +// name, which differs at every connection. +func (d Denial) key() string { + return d.Action + "|" + d.Type + "|" + d.Destination + "|" + d.Interface + "|" + d.Member +} + +// journalEntry is the part of a journal entry the module reads. MESSAGE is read only to recognise a +// denial; it is never answered or published. +type journalEntry struct { + Cursor string `json:"__CURSOR"` + Realtime string `json:"__REALTIME_TIMESTAMP"` + Message any `json:"MESSAGE"` + Action string `json:"DBUS_BROKER_TRANSMIT_ACTION"` + Type string `json:"DBUS_BROKER_MESSAGE_TYPE"` + Sender string `json:"DBUS_BROKER_SENDER_UNIQUE_NAME"` + Destination string `json:"DBUS_BROKER_MESSAGE_DESTINATION"` + Path string `json:"DBUS_BROKER_MESSAGE_PATH"` + Interface string `json:"DBUS_BROKER_MESSAGE_INTERFACE"` + Member string `json:"DBUS_BROKER_MESSAGE_MEMBER"` + Policy string `json:"DBUS_BROKER_POLICY_TYPE"` +} + +// IsDenial is whether a bus's log line is a policy denial: dbus-broker's "A security policy denied", +// or the reference daemon's "Rejected send message". +func IsDenial(message string) bool { + return strings.Contains(message, "security policy denied") || strings.Contains(message, "Rejected send message") || + strings.Contains(message, "Rejected receive message") +} + +// ParseDenials reads journalctl's JSON lines and answers the denials among them and the last cursor. +func ParseDenials(out string) ([]Denial, string) { + var got []Denial + cursor := "" + for _, line := range strings.Split(out, "\n") { + line = strings.TrimSpace(line) + if line == "" || line[0] != '{' { + continue + } + var e journalEntry + if json.Unmarshal([]byte(line), &e) != nil { + continue + } + if e.Cursor != "" { + cursor = e.Cursor + } + msg, _ := e.Message.(string) // a binary MESSAGE comes as an array of bytes and is no denial + if !IsDenial(msg) { + continue + } + at := "" + if us, err := strconv.ParseInt(e.Realtime, 10, 64); err == nil { + at = time.UnixMicro(us).UTC().Format(time.RFC3339) + } + got = append(got, Denial{At: at, Action: e.Action, Type: e.Type, Sender: e.Sender, + Destination: e.Destination, Path: e.Path, Interface: e.Interface, Member: e.Member, Policy: e.Policy}) + } + return got, cursor +} + +func journalArgs() []string { + args := []string{"--no-pager", "-o", "json", "-n", strconv.Itoa(JournalLines)} + for _, u := range BusUnits { + args = append(args, "-u", u) + } + return args +} + +// Denials is the watcher's Journal on this machine: journalctl as the operator's account, which +// reads the system journal through its group. +func (m *Machine) Denials(ctx context.Context, after string, since time.Time) ([]Denial, string, error) { + args := journalArgs() + if after != "" { + args = append(args, "--after-cursor", after) + } else { + args = append(args, "--since", "@"+strconv.FormatInt(since.Unix(), 10)) + } + out, err := m.Run(ctx, "journalctl", args...) + if err != nil { + return nil, "", err + } + d, cursor := ParseDenials(out) + return d, cursor, nil +} diff --git a/modules/dbus/cmd/dbus-tools/live_test.go b/modules/dbus/cmd/dbus-tools/live_test.go new file mode 100644 index 0000000..a075158 --- /dev/null +++ b/modules/dbus/cmd/dbus-tools/live_test.go @@ -0,0 +1,67 @@ +package main + +import ( + "context" + "encoding/json" + "os" + "path/filepath" + "testing" + "time" +) + +// Run on a real machine with MESH_LIVE=1: every tool reads this machine's buses, and the watcher's +// connection answers. Nothing is changed. +func TestLiveTools(t *testing.T) { + if os.Getenv("MESH_LIVE") == "" { + t.Skip("set MESH_LIVE=1 on a machine with a system bus") + } + ctx := context.Background() + m := Here() + b, err := DialSystemBus() + if err != nil { + t.Fatal(err) + } + defer b.Close() + if err := b.Ping(ctx); err != nil { + t.Fatal(err) + } + id, _ := b.ID(ctx) + pid, _ := b.PID(ctx, busName) + t.Logf("bus %s, driver pid %d in %s", id, pid, m.UnitOf(pid)) + show := func(name string, v any, err error) { + raw, _ := json.Marshal(v) + if len(raw) > 1500 { + raw = append(raw[:1500], "…"...) + } + t.Logf("%s: %v\n%s", name, err, raw) + } + n, err := m.Names(ctx, "system") + show("names system", n, err) + n, err = m.Names(ctx, "session") + show("names session", n, err) + o, err := m.Introspect(ctx, "system", "org.freedesktop.login1", "/org/freedesktop/login1") + show("introspect", o, err) + w, err := m.Monitor(ctx, "system", 2, "", nil) + show("monitor", w, err) + show("health", m.Health(ctx, nil), nil) + show("check", m.Check(ctx, nil), nil) +} + +// Run with MESH_LIVE=1: the watcher connects, pings and baselines the names, and says nothing. +func TestLiveWatcher(t *testing.T) { + if os.Getenv("MESH_LIVE") == "" { + t.Skip("set MESH_LIVE=1 on a machine with a system bus") + } + m := Here() + var said []string + w := NewWatcher(m, func(e string, b any) error { said = append(said, e); return nil }, DialSystemBus, m.Denials) + w.state = filepath.Join(t.TempDir(), "bus") + ctx, cancel := context.WithTimeout(context.Background(), PingEvery+2*time.Second) + defer cancel() + w.Run(ctx) + s := w.Snapshot() + t.Logf("%+v said %v", s, said) + if s.BusID == "" || s.LastPingAt == "" || s.Stalled || s.Services == 0 { + t.Fatalf("%+v", s) + } +} diff --git a/modules/dbus/cmd/dbus-tools/machine.go b/modules/dbus/cmd/dbus-tools/machine.go new file mode 100644 index 0000000..eab2cec --- /dev/null +++ b/modules/dbus/cmd/dbus-tools/machine.go @@ -0,0 +1,247 @@ +package main + +import ( + "bufio" + "bytes" + "context" + "errors" + "fmt" + "os" + "os/exec" + "path/filepath" + "strconv" + "strings" + "syscall" + "time" +) + +// CommandTimeout bounds every command a tool runs: a bus that hangs must cost a tool call twenty +// seconds, never the runtime's thirty. +const CommandTimeout = 20 * time.Second + +// ReadCap is the most of one command's output the module reads. An introspection document of the +// service manager is about 100 KiB; nothing the module asks is near a mebibyte. +const ReadCap = 1024 * 1024 + +// AnswerCap is the most entries a tool answers in one list (names, messages, denials). +const AnswerCap = 500 + +// Runner runs one command and answers its standard output. Injected, so every tool is tested against +// recorded answers rather than this machine's bus. +type Runner func(ctx context.Context, name string, args ...string) (string, error) + +// Streamer runs one command for at most the context's time and hands each line of its output to +// line, which says whether it wants more. The end of the time is the normal end, not an error. +// Injected, so dbus_monitor is tested without a bus. +type Streamer func(ctx context.Context, line func(string) bool, name string, args ...string) error + +// ExecRunner runs a command, bounded by CommandTimeout. A failure carries what it said on stderr. +func ExecRunner(ctx context.Context, name string, args ...string) (string, error) { + ctx, cancel := context.WithTimeout(ctx, CommandTimeout) + defer cancel() + cmd := exec.CommandContext(ctx, name, args...) + cmd.Env = append(os.Environ(), "LC_ALL=C", "SYSTEMD_PAGER=", "SYSTEMD_COLORS=0") + var stdout, stderr bytes.Buffer + cmd.Stdout, cmd.Stderr = &stdout, &stderr + err := cmd.Run() + out := capped(stdout.String(), ReadCap) + if ctx.Err() == context.DeadlineExceeded { + return out, fmt.Errorf("%s did not answer within %s", name, CommandTimeout) + } + if err != nil { + said := strings.TrimSpace(stderr.String()) + if said == "" { + said = strings.TrimSpace(stdout.String()) + } + return out, fmt.Errorf("%s %s: %w: %s", name, strings.Join(args, " "), err, capped(said, 2048)) + } + return out, nil +} + +// ExecStreamer runs a command until the context ends, line by line. A line longer than ReadCap is +// skipped, never held: a monitored message's line carries its body, which the module never keeps. +func ExecStreamer(ctx context.Context, line func(string) bool, name string, args ...string) error { + cmd := exec.CommandContext(ctx, name, args...) + cmd.Env = append(os.Environ(), "LC_ALL=C", "SYSTEMD_PAGER=", "SYSTEMD_COLORS=0") + // A terminate, which sudo passes on to what it runs; a kill would leave a root monitor behind. + cmd.Cancel = func() error { return cmd.Process.Signal(syscall.SIGTERM) } + cmd.WaitDelay = 2 * time.Second + var stderr bytes.Buffer + cmd.Stderr = &stderr + out, err := cmd.StdoutPipe() + if err != nil { + return err + } + if err := cmd.Start(); err != nil { + return err + } + r := bufio.NewReaderSize(out, 64*1024) + var cur []byte + tooLong := false + for { + chunk, isPrefix, err := r.ReadLine() + if err != nil { + break + } + if !tooLong { + cur = append(cur, chunk...) + if len(cur) > ReadCap { + cur, tooLong = cur[:0], true + } + } + if isPrefix { + continue + } + if !tooLong && !line(string(cur)) { + break + } + cur, tooLong = cur[:0], false + } + _ = cmd.Process.Signal(syscall.SIGTERM) + werr := cmd.Wait() + if ctx.Err() != nil { + return nil // the time ran out: the normal end of a bounded watch + } + if werr != nil && stderr.Len() > 0 { + return fmt.Errorf("%s: %w: %s", name, werr, capped(strings.TrimSpace(stderr.String()), 2048)) + } + return nil +} + +func capped(s string, n int) string { + if len(s) <= n { + return s + } + return s[:n] + "\n… (cut)" +} + +// Machine is what the module reads and acts on: a filesystem root (the real one, or a test's tree), +// a way to run commands, a way to stream one, and the environment the runtime gave it. +type Machine struct { + Root string + Run Runner + Stream Streamer + Env func(string) string + UID int + Now func() time.Time +} + +// Here is the machine this process runs on. +func Here() *Machine { + return &Machine{Root: "/", Run: ExecRunner, Stream: ExecStreamer, Env: os.Getenv, UID: os.Getuid(), Now: time.Now} +} + +func (m *Machine) path(p string) string { return filepath.Join(m.Root, p) } + +func (m *Machine) read(p string) string { + b, err := os.ReadFile(m.path(p)) + if err != nil { + return "" + } + return strings.TrimSpace(string(b)) +} + +func (m *Machine) exists(p string) bool { + _, err := os.Stat(m.path(p)) + return err == nil +} + +func (m *Machine) glob(pattern string) []string { + got, _ := filepath.Glob(m.path(pattern)) + out := make([]string, 0, len(got)) + for _, g := range got { + rel, err := filepath.Rel(m.Root, g) + if err != nil { + continue + } + out = append(out, "/"+rel) + } + return out +} + +// privileged runs a command as root without asking for a password (sudo -n), as the other modules' +// tools do: the runtime runs as the operator's account, and watching the system bus is root's. +func (m *Machine) privileged(ctx context.Context, name string, args ...string) (string, error) { + return m.Run(ctx, "sudo", append([]string{"-n", name}, args...)...) +} + +// Buses are the two buses a tool can be pointed at. +var Buses = []string{"system", "session"} + +// ErrNoSession is the answer on a machine where this account has no session bus: the servers. +var ErrNoSession = errors.New("no session bus for this account on this machine") + +// SessionAddress is the account's session bus: the runtime's DBUS_SESSION_BUS_ADDRESS, else the +// user manager's socket under XDG_RUNTIME_DIR or /run/user/. Only a socket that exists counts. +func (m *Machine) SessionAddress() (string, error) { + if a := m.Env("DBUS_SESSION_BUS_ADDRESS"); strings.HasPrefix(a, "unix:path=") { + p := strings.TrimPrefix(a, "unix:path=") + if i := strings.IndexByte(p, ','); i >= 0 { + p = p[:i] + } + if m.exists(p) { + return a, nil + } + } else if a != "" { + return a, nil // an abstract or other address: taken as given + } + dir := m.Env("XDG_RUNTIME_DIR") + if dir == "" { + dir = "/run/user/" + strconv.Itoa(m.UID) + } + if p := dir + "/bus"; m.exists(p) { + return "unix:path=" + p, nil + } + return "", fmt.Errorf("%w (no socket at %s/bus)", ErrNoSession, dir) +} + +// busArgs are busctl's words for one bus. +func (m *Machine) busArgs(bus string) ([]string, error) { + switch bus { + case "", "system": + return []string{"--system"}, nil + case "session": + a, err := m.SessionAddress() + if err != nil { + return nil, err + } + return []string{"--address=" + a}, nil + } + return nil, fmt.Errorf("bus is system or session, not %q", bus) +} + +// UnitOf is the systemd unit a process runs in, from its cgroup: readable for any process. +func (m *Machine) UnitOf(pid uint32) string { + if pid == 0 { + return "" + } + return UnitFromCgroup(m.read(fmt.Sprintf("/proc/%d/cgroup", pid))) +} + +// ProcessName is a process's command name. +func (m *Machine) ProcessName(pid uint32) string { + if pid == 0 { + return "" + } + return m.read(fmt.Sprintf("/proc/%d/comm", pid)) +} + +// UnitFromCgroup is the innermost service, socket or scope in a cgroup v2 path. +func UnitFromCgroup(cgroup string) string { + for _, line := range strings.Split(cgroup, "\n") { + if !strings.HasPrefix(line, "0::") { + continue + } + parts := strings.Split(strings.TrimPrefix(line, "0::"), "/") + for i := len(parts) - 1; i >= 0; i-- { + p := parts[i] + if strings.HasSuffix(p, ".service") || strings.HasSuffix(p, ".scope") { + return p + } + } + } + return "" +} + +// BootID is this boot's id: a bus that changed across a boot did not restart, the machine did. +func (m *Machine) BootID() string { return m.read("/proc/sys/kernel/random/boot_id") } diff --git a/modules/dbus/cmd/dbus-tools/main.go b/modules/dbus/cmd/dbus-tools/main.go new file mode 100644 index 0000000..0c09329 --- /dev/null +++ b/modules/dbus/cmd/dbus-tools/main.go @@ -0,0 +1,26 @@ +// The dbus module's Go bundle (novox/hq ADR 0215): the holder of node-message-bus. One process the +// node's runtime launches as the operator's account, serving the module's tools over MCP on stdio and +// running its watcher beside them, which publishes what matters about the system bus, never its +// traffic. +package main + +import ( + "context" + "fmt" + "os" + + stdio "git.novox.be/novox/mesh-sdk/go" +) + +func main() { + m := Here() + w := NewWatcher(m, func(eventType string, body any) error { return stdio.Emit(eventType, body) }, + DialSystemBus, m.Denials) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + go w.Run(ctx) + if err := stdio.Serve("", Tools(m, w)); err != nil { + fmt.Fprintln(os.Stderr, err) + os.Exit(1) + } +} diff --git a/modules/dbus/cmd/dbus-tools/manifest_test.go b/modules/dbus/cmd/dbus-tools/manifest_test.go new file mode 100644 index 0000000..d6f4134 --- /dev/null +++ b/modules/dbus/cmd/dbus-tools/manifest_test.go @@ -0,0 +1,159 @@ +package main + +// The dbus module's shape (novox/hq ADR 0215): it holds node-message-bus, owns the bus's packages, +// declares the bus running and never restarts or reloads it, and publishes only its curated events. + +import ( + "encoding/json" + "os" + "reflect" + "regexp" + "sort" + "strings" + "testing" +) + +type manifestShape struct { + Module string `json:"module"` + Capabilities []string `json:"capabilities"` + Claims []map[string]any `json:"claims"` + Emits []string `json:"emits"` + Consumes []string `json:"consumes"` + Tools []string `json:"tools"` + Resources []map[string]any `json:"resources"` + Build struct { + Artifacts []map[string]any `json:"artifacts"` + } `json:"build"` +} + +func manifest(t *testing.T) (manifestShape, string) { + t.Helper() + raw, err := os.ReadFile("../../module.json") + if err != nil { + t.Fatal(err) + } + var m manifestShape + if err := json.Unmarshal(raw, &m); err != nil { + t.Fatal(err) + } + return m, string(raw) +} + +func TestItHoldsTheMessageBusSeatWithOnlyWhatItNeeds(t *testing.T) { + m, _ := manifest(t) + if m.Module != "dbus" || len(m.Claims) != 1 || m.Claims[0]["name"] != "node-message-bus" || m.Claims[0]["scope"] != "node" { + t.Fatalf("claims: %+v", m.Claims) + } + if _, serves := m.Claims[0]["serves"]; serves { + t.Error("the seat receives nothing yet (ADR 0215 §4) and serves no verbs") + } + if !reflect.DeepEqual(m.Capabilities, []string{"package-manager", "service-manager"}) { + t.Fatalf("capabilities: %v", m.Capabilities) + } +} + +func TestItOwnsTheBusImplementationsPackages(t *testing.T) { + m, _ := manifest(t) + var pkgs []string + for _, r := range m.Resources { + if r["type"] == "package" { + pkgs = append(pkgs, r["package"].(string)) + if r["absent"] == true { + t.Errorf("%v is declared absent", r["package"]) + } + } + } + sort.Strings(pkgs) + if !reflect.DeepEqual(pkgs, Packages) { + t.Fatalf("packages %v, the check reads %v", pkgs, Packages) + } +} + +// TestTheBusIsNeverRestartedOrReloaded is ADR 0215 §2: the bus is declared running, with no trigger +// that would restart or reload it, and no boot state: dbus.service is an alias and dbus-broker.service +// is started at boot through dbus.socket, so enabling would write an alias link the machines do not +// have today. +func TestTheBusIsNeverRestartedOrReloaded(t *testing.T) { + m, raw := manifest(t) + if strings.Contains(raw, "restart-on") || strings.Contains(raw, "reload-on") { + t.Fatal("the manifest carries a restart or reload trigger") + } + var services []map[string]any + for _, r := range m.Resources { + if r["type"] == "service" { + services = append(services, r) + } + } + if len(services) != 1 { + t.Fatalf("services: %v", services) + } + s := services[0] + if s["unit"] != SystemUnit || s["state"] != "running" { + t.Fatalf("%v", s) + } + if _, ok := s["boot"]; ok { + t.Fatal("a boot state on the bus would have the host enable it") + } + for k := range s { + if k != "id" && k != "type" && k != "unit" && k != "state" { + t.Errorf("the bus's service carries %q", k) + } + } +} + +func TestItEmitsTheCuratedEventsAndConsumesNothing(t *testing.T) { + m, _ := manifest(t) + if !reflect.DeepEqual(m.Emits, Emits) { + t.Fatalf("emits %v, the watcher publishes %v", m.Emits, Emits) + } + if len(m.Consumes) != 0 { + t.Fatalf("consumes %v", m.Consumes) + } +} + +func TestToolsAreTheManifests(t *testing.T) { + m, _ := manifest(t) + served := map[string]bool{} + for _, tool := range Tools(testMachine(t, nil), nil) { + if served[tool.Name] { + t.Errorf("%s is served twice", tool.Name) + } + served[tool.Name] = true + if !strings.HasPrefix(tool.Name, "dbus_") { + t.Errorf("%s is not named dbus_", tool.Name) + } + } + for _, want := range m.Tools { + if !served[want] { + t.Errorf("the manifest lists %s and the bundle does not serve it", want) + } + delete(served, want) + } + for extra := range served { + t.Errorf("the bundle serves %s, which the manifest does not list", extra) + } +} + +func TestTheBundleIsTheBuildersShape(t *testing.T) { + m, _ := manifest(t) + if len(m.Build.Artifacts) != 1 { + t.Fatalf("%v", m.Build.Artifacts) + } + a := m.Build.Artifacts[0] + for k, v := range map[string]string{"kind": "bundle", "language": "go", "system": "arch", "from": "cmd/dbus-tools", "binary": "dbus-tools"} { + if a[k] != v { + t.Errorf("%s is %v, want %s", k, a[k], v) + } + } +} + +// TestNothingInstallationSpecific holds the manifest to the catalogue's rule: no node names, domains +// or home paths. +func TestNothingInstallationSpecific(t *testing.T) { + _, raw := manifest(t) + for _, re := range []*regexp.Regexp{regexp.MustCompile(`/home/`), regexp.MustCompile(`\.(be|internal|com)\b`)} { + if loc := re.FindString(raw); loc != "" { + t.Errorf("the manifest carries %q", loc) + } + } +} diff --git a/modules/dbus/cmd/dbus-tools/monitor.go b/modules/dbus/cmd/dbus-tools/monitor.go new file mode 100644 index 0000000..6454393 --- /dev/null +++ b/modules/dbus/cmd/dbus-tools/monitor.go @@ -0,0 +1,115 @@ +package main + +import ( + "context" + "encoding/json" + "fmt" + "strconv" + "strings" + "time" +) + +// MaxMonitorSeconds bounds a watch of the bus: long enough to catch an exchange, short enough that +// nobody leaves a monitor running. +const MaxMonitorSeconds = 15 + +// MonitorLimit is the most messages busctl reads in one watch before it stops by itself. +const MonitorLimit = 20000 + +// Header is what dbus_monitor answers of one message: who sent what to whom. The body is never read +// into it: the struct has no field for it, so a payload is dropped as each line is decoded. +type Header struct { + Time string `json:"time,omitempty"` + Type string `json:"type"` + Sender string `json:"sender,omitempty"` + Destination string `json:"destination,omitempty"` + Path string `json:"path,omitempty"` + Interface string `json:"interface,omitempty"` + Member string `json:"member,omitempty"` + ErrorName string `json:"error_name,omitempty"` +} + +type monitorLine struct { + Type string `json:"type"` + Sender string `json:"sender"` + Destination string `json:"destination"` + Path string `json:"path"` + Interface string `json:"interface"` + Member string `json:"member"` + ErrorName string `json:"error_name"` + Realtime int64 `json:"timestamp-realtime"` +} + +// ParseHeader reads one line of `busctl monitor --json=short` into its header alone. +func ParseHeader(line string) (Header, bool) { + var l monitorLine + if json.Unmarshal([]byte(line), &l) != nil || l.Type == "" { + return Header{}, false + } + h := Header{Type: l.Type, Sender: l.Sender, Destination: l.Destination, Path: l.Path, + Interface: l.Interface, Member: l.Member, ErrorName: l.ErrorName} + if l.Realtime > 0 { + h.Time = time.UnixMicro(l.Realtime).UTC().Format("15:04:05.000000") + } + return h, true +} + +// Watch is dbus_monitor's answer. +type Watch struct { + Bus string `json:"bus"` + Seconds int `json:"seconds"` + Match string `json:"match,omitempty"` + Names []string `json:"names,omitempty"` + Total int `json:"total"` + Messages []Header `json:"messages"` + Cut bool `json:"cut,omitempty"` + Note string `json:"note"` +} + +// Monitor watches one bus for a few seconds and answers the headers of what passed. The system bus is +// watched as root (sudo -n): only root may become a monitor there. The session bus is the account's +// own. The bodies, which carry secrets, notification text and the clipboard, are never kept. +func (m *Machine) Monitor(ctx context.Context, bus string, seconds int, match string, names []string) (Watch, error) { + if seconds < 1 || seconds > MaxMonitorSeconds { + return Watch{}, fmt.Errorf("seconds is 1 to %d", MaxMonitorSeconds) + } + if len(match) > 512 || strings.ContainsAny(match, "\n\x00") { + return Watch{}, fmt.Errorf("the match rule is not one line of at most 512 characters") + } + for _, n := range names { + if !busNameRE.MatchString(n) { + return Watch{}, fmt.Errorf("%q is not a bus name", n) + } + } + args, err := m.busArgs(bus) + if err != nil { + return Watch{}, err + } + cmd := append([]string{"timeout", strconv.Itoa(seconds) + "s", "busctl"}, args...) + cmd = append(cmd, "--json=short", "--no-pager", "--limit-messages="+strconv.Itoa(MonitorLimit), "monitor") + if match != "" { + cmd = append(cmd, "--match="+match) + } + cmd = append(cmd, names...) + if orWord(bus, "system") == "system" { + cmd = append([]string{"sudo", "-n"}, cmd...) + } + w := Watch{Bus: orWord(bus, "system"), Seconds: seconds, Match: match, Names: names, Messages: []Header{}, + Note: "headers only; message bodies are never read into the answer"} + ctx, cancel := context.WithTimeout(ctx, time.Duration(seconds)*time.Second+5*time.Second) + defer cancel() + err = m.Stream(ctx, func(line string) bool { + h, ok := ParseHeader(line) + if !ok { + return true + } + w.Total++ + if len(w.Messages) < AnswerCap { + w.Messages = append(w.Messages, h) + } else { + w.Cut = true + } + return true + }, cmd[0], cmd[1:]...) + return w, err +} diff --git a/modules/dbus/cmd/dbus-tools/names.go b/modules/dbus/cmd/dbus-tools/names.go new file mode 100644 index 0000000..5f0fdcc --- /dev/null +++ b/modules/dbus/cmd/dbus-tools/names.go @@ -0,0 +1,103 @@ +package main + +import ( + "context" + "encoding/json" + "fmt" + "sort" + "strings" +) + +// busctlName is one line of `busctl list --json`. +type busctlName struct { + Name string `json:"name"` + PID *uint32 `json:"pid"` + Process *string `json:"process"` + User *string `json:"user"` + Connection string `json:"connection"` + Unit *string `json:"unit"` +} + +// Name is one well-known name on a bus, with what holds it. +type Name struct { + Name string `json:"name"` + PID uint32 `json:"pid,omitempty"` + Process string `json:"process,omitempty"` + User string `json:"user,omitempty"` + Unit string `json:"unit,omitempty"` + Connection string `json:"connection,omitempty"` +} + +// Names is dbus_names's answer. +type Names struct { + Bus string `json:"bus"` + Running []Name `json:"running"` + ActivatableNotRunning []string `json:"activatable_not_running"` + Connections int `json:"connections"` + Cut bool `json:"cut,omitempty"` +} + +// ParseNames reads `busctl list --json=short`: the well-known names running, the activatable names +// that are not, and how many connections (unique names) the bus has. +func ParseNames(bus, out string) (Names, error) { + var raw []busctlName + if err := json.Unmarshal([]byte(strings.TrimSpace(out)), &raw); err != nil { + return Names{}, fmt.Errorf("busctl list answered what is not its JSON: %w", err) + } + n := Names{Bus: bus, Running: []Name{}, ActivatableNotRunning: []string{}} + for _, r := range raw { + switch { + case r.Name == busName: // the bus itself, which busctl attributes to the service manager + case strings.HasPrefix(r.Name, ":"): + n.Connections++ + case r.Connection == "(activatable)" || r.PID == nil: + n.ActivatableNotRunning = append(n.ActivatableNotRunning, r.Name) + default: + n.Running = append(n.Running, Name{Name: r.Name, PID: deref(r.PID), Process: derefS(r.Process), + User: derefS(r.User), Unit: derefS(r.Unit), Connection: r.Connection}) + } + } + sort.Slice(n.Running, func(i, j int) bool { return n.Running[i].Name < n.Running[j].Name }) + sort.Strings(n.ActivatableNotRunning) + if len(n.Running) > AnswerCap { + n.Running, n.Cut = n.Running[:AnswerCap], true + } + if len(n.ActivatableNotRunning) > AnswerCap { + n.ActivatableNotRunning, n.Cut = n.ActivatableNotRunning[:AnswerCap], true + } + return n, nil +} + +// Names lists one bus's names as this account sees them. +func (m *Machine) Names(ctx context.Context, bus string) (Names, error) { + args, err := m.busArgs(bus) + if err != nil { + return Names{}, err + } + out, err := m.Run(ctx, "busctl", append(args, "--json=short", "--no-pager", "list")...) + if err != nil { + return Names{}, err + } + return ParseNames(orWord(bus, "system"), out) +} + +func deref(p *uint32) uint32 { + if p == nil { + return 0 + } + return *p +} + +func derefS(p *string) string { + if p == nil { + return "" + } + return *p +} + +func orWord(s, word string) string { + if strings.TrimSpace(s) == "" { + return word + } + return s +} diff --git a/modules/dbus/cmd/dbus-tools/parse_test.go b/modules/dbus/cmd/dbus-tools/parse_test.go new file mode 100644 index 0000000..071ad10 --- /dev/null +++ b/modules/dbus/cmd/dbus-tools/parse_test.go @@ -0,0 +1,304 @@ +package main + +import ( + "context" + "encoding/json" + "errors" + "reflect" + "strings" + "testing" + "time" +) + +// Recorded answers, trimmed, from a workstation of 2026-10-05. +const busctlList = `[{"name":":1.0","pid":456,"process":"systemd-timesyn","user":"systemd-timesync","connection":":1.0","unit":"systemd-timesyncd.service","session":null,"description":null},` + + `{"name":"fi.w1.wpa_supplicant1","pid":1442,"process":"wpa_supplicant","user":"root","connection":":1.23","unit":"wpa_supplicant.service","session":null,"description":null},` + + `{"name":"org.blueman.Mechanism","pid":null,"process":null,"user":null,"connection":"(activatable)","unit":null,"session":null,"description":null},` + + `{"name":":1.11","pid":975,"process":"polkitd","user":"polkitd","connection":":1.11","unit":"polkit.service","session":null,"description":null},` + + `{"name":"org.freedesktop.DBus","pid":807,"process":"dbus-broker-lau","user":"root","connection":"org.freedesktop.DBus","unit":"dbus-broker.service","session":null,"description":null}]` + +func TestNamesAreSplitIntoRunningActivatableAndConnections(t *testing.T) { + n, err := ParseNames("system", busctlList) + if err != nil { + t.Fatal(err) + } + if n.Connections != 2 || len(n.Running) != 1 || n.Running[0].Name != "fi.w1.wpa_supplicant1" || + n.Running[0].Unit != "wpa_supplicant.service" || n.Running[0].PID != 1442 { + t.Fatalf("%+v", n) + } + if !reflect.DeepEqual(n.ActivatableNotRunning, []string{"org.blueman.Mechanism"}) { + t.Fatalf("%v", n.ActivatableNotRunning) + } +} + +const introspection = `{"type":"s","data":["\n\n \n \n \n \n \n \n \n \n \n \n \n \n\n"]}` + +func TestAnIntrospectionIsShapedWithoutValues(t *testing.T) { + node, err := ParseIntrospection(introspection) + if err != nil { + t.Fatal(err) + } + o := Shape("system", "org.freedesktop.login1", "/org/freedesktop/login1", node) + if len(o.Interfaces) != 1 || !reflect.DeepEqual(o.Children, []string{"session", "seat"}) { + t.Fatalf("%+v", o) + } + i := o.Interfaces[0] + if !reflect.DeepEqual(i.Methods[0], Member{Name: "Inhibit", In: []string{"what s"}, Out: []string{"pipe_fd h"}}) || + !reflect.DeepEqual(i.Properties[0], Property{Name: "IdleHint", Type: "b", Access: "read"}) || + !reflect.DeepEqual(i.Signals[0], Member{Name: "PrepareForSleep", Args: []string{"start b"}}) { + t.Fatalf("%+v", i) + } +} + +func TestIntrospectRefusesWhatIsNotANameOrAPathAndNeverStartsAService(t *testing.T) { + var asked []string + m := testMachine(t, nil) + m.Run = func(_ context.Context, name string, args ...string) (string, error) { + asked = append([]string{name}, args...) + return introspection, nil + } + if _, err := m.Introspect(context.Background(), "system", "--address=x", "/"); err == nil { + t.Fatal("an option was taken for a bus name") + } + if _, err := m.Introspect(context.Background(), "system", "org.freedesktop.login1", "relative"); err == nil { + t.Fatal("a relative path was taken") + } + if _, err := m.Introspect(context.Background(), "system", "org.freedesktop.login1", ""); err != nil { + t.Fatal(err) + } + if !strings.Contains(strings.Join(asked, " "), "--auto-start=no") || asked[len(asked)-3] != "/" { + t.Fatalf("%v", asked) + } +} + +const monitorLines = `{"type":"method_call","endian":"l","flags":0,"version":1,"cookie":186058,"timestamp-realtime":1791193058632202,"sender":":1.797","destination":"org.freedesktop.UPower","path":"/org/freedesktop/UPower","interface":"org.freedesktop.UPower","member":"GetDisplayDevice","payload":{"type":"","data":[]}} +{"type":"signal","endian":"l","flags":1,"version":1,"cookie":7,"timestamp-realtime":1791193058632725,"sender":":1.40","path":"/org/freedesktop/Notifications","interface":"org.freedesktop.Notifications","member":"Notify","payload":{"type":"susssasa{sv}i","data":["app",0,"","the secret notification text","the clipboard's password",[],{},5000]}} +not json at all +{"type":"error","endian":"l","flags":1,"version":1,"cookie":9,"sender":":1.9","destination":":1.797","error_name":"org.freedesktop.DBus.Error.AccessDenied","payload":{"type":"s","data":["the reason, with a token"]}}` + +func TestAMonitoredMessageIsItsHeaderOnly(t *testing.T) { + var asked []string + m := testMachine(t, nil) + m.Stream = func(ctx context.Context, line func(string) bool, name string, args ...string) error { + asked = append([]string{name}, args...) + if _, ok := ctx.Deadline(); !ok { + t.Error("the watch has no deadline") + } + for _, l := range strings.Split(monitorLines, "\n") { + if !line(l) { + break + } + } + return nil + } + w, err := m.Monitor(context.Background(), "system", 3, "type='signal'", []string{"org.freedesktop.Notifications"}) + if err != nil { + t.Fatal(err) + } + raw, _ := json.Marshal(w) + for _, secret := range []string{"secret notification", "password", "token", "payload", "data"} { + if strings.Contains(string(raw), secret) { + t.Errorf("the answer carries %q: %s", secret, raw) + } + } + if w.Total != 3 || w.Messages[1].Member != "Notify" || w.Messages[2].ErrorName != "org.freedesktop.DBus.Error.AccessDenied" { + t.Fatalf("%+v", w) + } + cmd := strings.Join(asked, " ") + if !strings.HasPrefix(cmd, "sudo -n timeout 3s busctl --system") || !strings.Contains(cmd, "--match=type='signal'") || + !strings.HasSuffix(cmd, "org.freedesktop.Notifications") { + t.Fatalf("%s", cmd) + } +} + +func TestAMonitorIsBoundedAndTheSessionBusIsTheAccounts(t *testing.T) { + m := testMachine(t, map[string]string{"/run/user/1000/bus": ""}) + var asked string + m.Stream = func(_ context.Context, _ func(string) bool, name string, args ...string) error { + asked = name + " " + strings.Join(args, " ") + return nil + } + for _, s := range []int{0, MaxMonitorSeconds + 1} { + if _, err := m.Monitor(context.Background(), "system", s, "", nil); err == nil { + t.Errorf("%d seconds were accepted", s) + } + } + if _, err := m.Monitor(context.Background(), "session", 2, "", []string{"-x"}); err == nil { + t.Error("an option was taken for a name") + } + if _, err := m.Monitor(context.Background(), "session", 2, "", nil); err != nil { + t.Fatal(err) + } + if strings.HasPrefix(asked, "sudo") || !strings.Contains(asked, "--address=unix:path=/run/user/1000/bus") { + t.Fatalf("%s", asked) + } +} + +func TestASessionBusIsFoundOrSaidMissing(t *testing.T) { + m := testMachine(t, nil) + if _, err := m.SessionAddress(); !errors.Is(err, ErrNoSession) { + t.Fatalf("a server without a session said %v", err) + } + m = testMachine(t, map[string]string{"/run/user/1000/bus": ""}) + if a, err := m.SessionAddress(); err != nil || a != "unix:path=/run/user/1000/bus" { + t.Fatalf("%q %v", a, err) + } + m.Env = func(k string) string { + if k == "DBUS_SESSION_BUS_ADDRESS" { + return "unix:path=/run/user/1000/bus,guid=abc" + } + return "" + } + if a, _ := m.SessionAddress(); a != "unix:path=/run/user/1000/bus,guid=abc" { + t.Fatalf("%q", a) + } +} + +const journalLines = `{"__CURSOR":"s=1;i=1","__REALTIME_TIMESTAMP":"1791192960341045","MESSAGE":"Ready","_PID":"807"} +{"__CURSOR":"s=1;i=2","__REALTIME_TIMESTAMP":"1791192960341045","MESSAGE":"A security policy denied :1.1199 to send method call /org/freedesktop/DBus:org.freedesktop.DBus.Debug.Stats.GetStats to org.freedesktop.DBus.","DBUS_BROKER_TRANSMIT_ACTION":"send","DBUS_BROKER_MESSAGE_TYPE":"method_call","DBUS_BROKER_SENDER_UNIQUE_NAME":":1.1199","DBUS_BROKER_MESSAGE_DESTINATION":"org.freedesktop.DBus","DBUS_BROKER_MESSAGE_PATH":"/org/freedesktop/DBus","DBUS_BROKER_MESSAGE_INTERFACE":"org.freedesktop.DBus.Debug.Stats","DBUS_BROKER_MESSAGE_MEMBER":"GetStats","DBUS_BROKER_POLICY_TYPE":"internal"} +{"__CURSOR":"s=1;i=3","MESSAGE":[65,66]}` + +func TestDenialsAreReadFromTheJournalByTheirFields(t *testing.T) { + d, cursor := ParseDenials(journalLines) + if cursor != "s=1;i=3" || len(d) != 1 { + t.Fatalf("%q %+v", cursor, d) + } + want := Denial{At: "2026-10-05T09:36:00Z", Action: "send", Type: "method_call", Sender: ":1.1199", + Destination: "org.freedesktop.DBus", Path: "/org/freedesktop/DBus", Interface: "org.freedesktop.DBus.Debug.Stats", + Member: "GetStats", Policy: "internal"} + if d[0] != want { + t.Fatalf("%+v", d[0]) + } +} + +func TestTheJournalIsReadFromTheCursorOn(t *testing.T) { + m := testMachine(t, nil) + var asked []string + m.Run = func(_ context.Context, name string, args ...string) (string, error) { + asked = append([]string{name}, args...) + return journalLines, nil + } + _, next, _ := m.Denials(context.Background(), "", time.Unix(100, 0)) + if !strings.Contains(strings.Join(asked, " "), "--since @100") { + t.Fatalf("%v", asked) + } + m.Denials(context.Background(), next, time.Time{}) + if !strings.Contains(strings.Join(asked, " "), "--after-cursor s=1;i=3") || !strings.Contains(strings.Join(asked, " "), "-u dbus-broker.service") { + t.Fatalf("%v", asked) + } +} + +func TestAUnitIsReadFromACgroup(t *testing.T) { + for in, want := range map[string]string{ + "0::/system.slice/bluetooth.service\n": "bluetooth.service", + "0::/user.slice/user-1000.slice/user@1000.service/app.slice/dunst.service": "dunst.service", + "0::/user.slice/user-1000.slice/session-2.scope": "session-2.scope", + "0::/init.scope": "init.scope", + "": "", + } { + if got := UnitFromCgroup(in); got != want { + t.Errorf("%q: %q, want %q", in, got, want) + } + } +} + +func TestShowBlocksAreReadInOrder(t *testing.T) { + b := ParseShow("Id=dbus-org.freedesktop.resolve1.service\nLoadState=not-found\n\nId=systemd-hostnamed.service\nLoadState=loaded\n") + if len(b) != 2 || b[0]["LoadState"] != "not-found" || b[1]["Id"] != "systemd-hostnamed.service" { + t.Fatalf("%v", b) + } +} + +func TestARunningBusOlderThanItsPackageIsSaid(t *testing.T) { + p := ParseDesc("%NAME%\ndbus-broker\n\n%VERSION%\n37-3\n\n%INSTALLDATE%\n1772097880\n") + if p.Name != "dbus-broker" || p.Version != "37-3" || p.Installed.Unix() != 1772097880 { + t.Fatalf("%+v", p) + } + if due, _ := RebootDue(time.Unix(1791123478, 0), []Package{p}); due { + t.Fatal("a bus started after its package was said to be older") + } + due, detail := RebootDue(time.Unix(1772000000, 0), []Package{p}) + if !due || !strings.Contains(detail, "reboot is due") || !strings.Contains(detail, "dbus-broker 37-3") { + t.Fatalf("%v %s", due, detail) + } +} + +func TestConnectionsAreCountedFromStatsOrNames(t *testing.T) { + stats := `{"type":"a{sv}","data":[{"org.bus1.DBus.Debug.Stats.PeerAccounting":{"type":"a(sa{sv}a{su})","data":[[":1.0",{},{}],[":1.1",{},{}],[":1.7",{},{}]]}}]}` + if n, ok := CountConnections(stats); !ok || n != 3 { + t.Fatalf("%d %v", n, ok) + } + if n, ok := CountConnections(`{"type":"a{sv}","data":[{"ActiveConnections":{"type":"u","data":12}}]}`); !ok || n != 12 { + t.Fatalf("%d %v", n, ok) + } + if n, ok := CountUniqueNames(`{"type":"as","data":[["org.freedesktop.DBus",":1.0","org.bluez",":1.5"]]}`); !ok || n != 2 { + t.Fatalf("%d %v", n, ok) + } +} + +func TestStatsAreNeverAskedAsTheAccount(t *testing.T) { + m := testMachine(t, nil) + var calls []string + m.Run = func(_ context.Context, name string, args ...string) (string, error) { + line := name + " " + strings.Join(args, " ") + calls = append(calls, line) + if name == "sudo" { + return "", errors.New("sudo: a password is required") + } + if strings.Contains(line, "ListNames") { + return `{"type":"as","data":[[":1.0",":1.1"]]}`, nil + } + return "Id=dbus-broker.service\nActiveState=active\nMainPID=807\nActiveEnterTimestamp=@1791123478\n", nil + } + h := m.Health(context.Background(), nil) + for _, c := range calls { + if strings.Contains(c, "GetStats") && !strings.HasPrefix(c, "sudo -n ") { + t.Fatalf("Debug.Stats asked as the account, which the bus logs as a denial: %s", c) + } + } + if h["connections"] != 2 || h["unit"] != "dbus-broker.service" { + t.Fatalf("%v", h) + } +} + +func TestTheCheckReadsPackagesUnitsServiceFilesAndDenials(t *testing.T) { + m := testMachine(t, map[string]string{ + "/var/lib/pacman/local/dbus-1.16.2-1/desc": "%NAME%\ndbus\n\n%VERSION%\n1.16.2-1\n\n%INSTALLDATE%\n1741340000\n", + "/var/lib/pacman/local/dbus-broker-37-3/desc": "%NAME%\ndbus-broker\n\n%VERSION%\n37-3\n\n%INSTALLDATE%\n1800000000\n", + "/var/lib/pacman/local/dbus-broker-units-37-3/desc": "%NAME%\ndbus-broker-units\n\n%VERSION%\n37-3\n\n%INSTALLDATE%\n1772097880\n", + "/usr/share/dbus-1/system-services/org.bluez.service": "[D-BUS Service]\nName=org.bluez\nExec=/bin/false\nUser=root\nSystemdService=dbus-org.bluez.service\n", + "/usr/share/dbus-1/system-services/org.example.service": "[D-BUS Service]\nName=org.example\nSystemdService=example.service\n", + "/usr/share/dbus-1/system-services/org.freedesktop.systemd1.service": "[D-BUS Service]\nName=org.freedesktop.systemd1\nExec=/bin/false\n", + }) + m.Run = func(_ context.Context, name string, args ...string) (string, error) { + line := strings.Join(args, " ") + switch { + case name == "journalctl": + return journalLines, nil + case strings.Contains(line, "Id,LoadState"): + return "Id=dbus-org.bluez.service\nLoadState=not-found\n\nId=example.service\nLoadState=not-found\n", nil + default: + return "Id=dbus-broker.service\nActiveState=active\nMainPID=807\nActiveEnterTimestamp=@1791123478\n", nil + } + } + got := m.Check(context.Background(), nil) + byName := map[string]Check{} + for _, c := range got["checks"].([]Check) { + byName[c.Name] = c + } + if !byName["package dbus-broker-units"].OK || !byName["system bus"].OK { + t.Fatalf("%+v", got) + } + if c := byName["running bus is the installed one"]; c.OK || !strings.Contains(c.Detail, "dbus-broker 37-3") { + t.Fatalf("an upgraded broker was not said: %+v", c) + } + if c := byName["activatable services"]; c.OK || !strings.Contains(c.Detail, "org.example → example.service") || strings.Contains(c.Detail, "bluez") { + t.Fatalf("%+v", c) + } + if c := byName["policy denials in the last hour"]; c.OK || !strings.Contains(c.Detail, "GetStats") { + t.Fatalf("%+v", c) + } + if !strings.Contains(strings.Join(got["notes"].([]string), " "), "org.bluez → dbus-org.bluez.service") { + t.Fatalf("%v", got["notes"]) + } +} diff --git a/modules/dbus/cmd/dbus-tools/tools.go b/modules/dbus/cmd/dbus-tools/tools.go new file mode 100644 index 0000000..e2f6576 --- /dev/null +++ b/modules/dbus/cmd/dbus-tools/tools.go @@ -0,0 +1,99 @@ +package main + +import ( + "context" + "fmt" + + stdio "git.novox.be/novox/mesh-sdk/go" +) + +var busArg = map[string]any{"type": "string", "enum": Buses, + "description": "system (the default) or session (this account's login session bus, where one exists)"} + +func busOf(args map[string]any) string { + b, _ := args["bus"].(string) + return b +} + +// Tools are the module's tools over MCP (novox/hq ADR 0215): the names on either bus, what an object +// offers, a bounded watch of message headers, the module's check, and the bus's health. +func Tools(m *Machine, w *Watcher) []stdio.Tool { + return []stdio.Tool{ + { + Name: "dbus_names", + Description: "Every well-known name on the system or session bus: its owner's pid, process, user and " + + "systemd unit, the activatable names that are not running, and how many connections the bus has. " + + "Read as the operator's account.", + Input: map[string]any{"type": "object", "properties": map[string]any{"bus": busArg}}, + Run: func(args map[string]any) (any, error) { + return m.Names(context.Background(), busOf(args)) + }, + }, + { + Name: "dbus_introspect", + Description: "What one object on a bus offers: its interfaces with their methods (arguments in and out), " + + "properties (type and access, never their values) and signals, and its child objects. A service " + + "that is not running is not started for it.", + Input: map[string]any{ + "type": "object", + "properties": map[string]any{ + "bus": busArg, + "service": map[string]any{"type": "string", "description": "the bus name, e.g. org.freedesktop.login1"}, + "path": map[string]any{"type": "string", "description": "the object path (default /)"}, + }, + "required": []string{"service"}, + }, + Run: func(args map[string]any) (any, error) { + service, _ := args["service"].(string) + path, _ := args["path"].(string) + return m.Introspect(context.Background(), busOf(args), service, path) + }, + }, + { + Name: "dbus_monitor", + Description: fmt.Sprintf("Watch a bus for a few seconds (at most %d) and answer the headers of the "+ + "messages that passed: type, sender, destination, path, interface, member. Never a message's body, "+ + "which carries secrets, notification text and the clipboard. Optional match rule and names to "+ + "narrow it. The system bus is watched as root through sudo -n.", MaxMonitorSeconds), + Input: map[string]any{ + "type": "object", + "properties": map[string]any{ + "bus": busArg, + "seconds": map[string]any{"type": "integer", "description": fmt.Sprintf("how long (default 5, at most %d)", MaxMonitorSeconds)}, + "match": map[string]any{"type": "string", "description": "a D-Bus match rule, e.g. type='signal',interface='org.freedesktop.login1.Manager'"}, + "names": map[string]any{"type": "array", "items": map[string]any{"type": "string"}, "description": "only messages to or from these bus names"}, + }, + }, + Run: func(args map[string]any) (any, error) { + s := 5 + if v, ok := args["seconds"].(float64); ok { + s = int(v) + } + match, _ := args["match"].(string) + var names []string + if list, ok := args["names"].([]any); ok { + for _, n := range list { + if str, ok := n.(string); ok { + names = append(names, str) + } + } + } + return m.Monitor(context.Background(), busOf(args), s, match, names) + }, + }, + { + Name: "dbus_check", + Description: "Whether this machine's message bus is as the mesh expects: the bus's packages installed, " + + "dbus-broker running behind dbus.service, whether the running bus is older than its installed " + + "package (then a reboot is due: the bus is never restarted live), activatable services whose " + + "unit does not exist, policy denials in the last hour, and the watcher's state.", + Run: func(map[string]any) (any, error) { return m.Check(context.Background(), w), nil }, + }, + { + Name: "dbus_health", + Description: "The system bus's health now: a ping's round trip, how many connections it has, its unit, " + + "pid and uptime, and what the watcher last saw.", + Run: func(map[string]any) (any, error) { return m.Health(context.Background(), w), nil }, + }, + } +} diff --git a/modules/dbus/cmd/dbus-tools/watcher.go b/modules/dbus/cmd/dbus-tools/watcher.go new file mode 100644 index 0000000..ea4bd39 --- /dev/null +++ b/modules/dbus/cmd/dbus-tools/watcher.go @@ -0,0 +1,619 @@ +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:]) +} diff --git a/modules/dbus/cmd/dbus-tools/watcher_test.go b/modules/dbus/cmd/dbus-tools/watcher_test.go new file mode 100644 index 0000000..ad641ee --- /dev/null +++ b/modules/dbus/cmd/dbus-tools/watcher_test.go @@ -0,0 +1,446 @@ +package main + +import ( + "context" + "encoding/json" + "errors" + "os" + "path/filepath" + "reflect" + "strings" + "sync" + "testing" + "time" +) + +// fakeBus is a system bus in memory: names with their owners' pids, a ping that can hang, and the +// NameOwnerChanged stream. +type fakeBus struct { + mu sync.Mutex + id string + pid uint32 + names map[string]uint32 + hang bool + changes chan NameChange +} + +func newFakeBus(id string, pid uint32, names map[string]uint32) *fakeBus { + return &fakeBus{id: id, pid: pid, names: names, changes: make(chan NameChange, 64)} +} + +func (f *fakeBus) Ping(ctx context.Context) error { + f.mu.Lock() + hang := f.hang + f.mu.Unlock() + if hang { + <-ctx.Done() + return ctx.Err() + } + return nil +} +func (f *fakeBus) ID(context.Context) (string, error) { return f.id, nil } +func (f *fakeBus) PID(_ context.Context, name string) (uint32, error) { + if name == busName { + return f.pid, nil + } + f.mu.Lock() + defer f.mu.Unlock() + if p, ok := f.names[name]; ok { + return p, nil + } + return 0, errors.New("no such name") +} +func (f *fakeBus) Names(context.Context) ([]string, error) { + f.mu.Lock() + defer f.mu.Unlock() + out := []string{busName, ":1.0", ":1.1"} + for n := range f.names { + out = append(out, n) + } + return out, nil +} +func (f *fakeBus) Activatable(context.Context) ([]string, error) { + return []string{"org.freedesktop.hostname1"}, nil +} +func (f *fakeBus) Changes() <-chan NameChange { return f.changes } +func (f *fakeBus) Close() {} + +// meshBus records what was published, and can refuse. +type meshBus struct { + mu sync.Mutex + down bool + types []string + bodies []map[string]any +} + +func (b *meshBus) emit(t string, body any) error { + b.mu.Lock() + defer b.mu.Unlock() + if b.down { + return errors.New("no bus") + } + b.types = append(b.types, t) + m, _ := body.(map[string]any) + b.bodies = append(b.bodies, m) + return nil +} + +func (b *meshBus) seen() []string { + b.mu.Lock() + defer b.mu.Unlock() + return append([]string(nil), b.types...) +} + +// clock is a time the test moves. +type clock struct{ t time.Time } + +func (c *clock) now() time.Time { return c.t } +func (c *clock) advance(d time.Duration) { c.t = c.t.Add(d) } + +func testMachine(t *testing.T, files map[string]string) *Machine { + t.Helper() + root := t.TempDir() + for p, c := range files { + full := filepath.Join(root, p) + os.MkdirAll(filepath.Dir(full), 0o755) + os.WriteFile(full, []byte(c), 0o644) + } + return &Machine{Root: root, Env: func(string) string { return "" }, UID: 1000, Now: time.Now} +} + +func testWatcher(t *testing.T, b *meshBus, c *clock) *Watcher { + m := testMachine(t, map[string]string{ + "/proc/sys/kernel/random/boot_id": "boot-1\n", + "/proc/700/comm": "systemd-logind\n", + "/proc/700/cgroup": "0::/system.slice/systemd-logind.service\n", + "/proc/900/comm": "bluetoothd\n", + "/proc/900/cgroup": "0::/system.slice/bluetooth.service\n", + }) + w := NewWatcher(m, b.emit, nil, nil) + w.now = c.now + w.state = filepath.Join(t.TempDir(), "bus") + return w +} + +func start() *clock { return &clock{t: time.Date(2026, 10, 5, 12, 0, 0, 0, time.UTC)} } + +func TestTheFirstConnectionSaysNothingAboutNames(t *testing.T) { + mb, c := &meshBus{}, start() + w := testWatcher(t, mb, c) + fb := newFakeBus("id-1", 500, map[string]uint32{"org.freedesktop.login1": 700}) + w.identity("id-1", 500) + w.connectedTo(fb) + c.advance(Debounce) + w.settle(fb) + w.flush() + if got := mb.seen(); len(got) != 0 { + t.Fatalf("the baseline was announced: %v", got) + } + if w.Snapshot().Services != 1 { + t.Fatalf("%+v", w.Snapshot()) + } +} + +func TestAServiceAppearingIsSaidOnceItStaysWithItsUnit(t *testing.T) { + mb, c := &meshBus{}, start() + w := testWatcher(t, mb, c) + fb := newFakeBus("id-1", 500, map[string]uint32{}) + w.connectedTo(fb) + fb.names["org.bluez"] = 900 + w.observe(NameChange{Name: "org.bluez", New: ":1.9"}) + c.advance(Debounce / 2) + w.settle(fb) + w.flush() + if got := mb.seen(); len(got) != 0 { + t.Fatalf("said before the debounce: %v", got) + } + c.advance(Debounce) + w.settle(fb) + w.flush() + if got := mb.seen(); !reflect.DeepEqual(got, []string{ServiceAppeared}) { + t.Fatalf("%v", got) + } + body := mb.bodies[0] + if body["name"] != "org.bluez" || body["unit"] != "bluetooth.service" || body["process"] != "bluetoothd" || body["pid"] != uint32(900) { + t.Fatalf("%v", body) + } +} + +func TestAFlapIsDebouncedAway(t *testing.T) { + mb, c := &meshBus{}, start() + w := testWatcher(t, mb, c) + fb := newFakeBus("id-1", 500, map[string]uint32{"org.freedesktop.login1": 700}) + w.connectedTo(fb) + c.advance(Debounce) + w.settle(fb) + // logind restarted: left and back within the debounce, and a newcomer that came and went. + w.observe(NameChange{Name: "org.freedesktop.login1", Old: ":1.5"}) + c.advance(time.Second) + w.observe(NameChange{Name: "org.freedesktop.login1", New: ":1.80"}) + w.observe(NameChange{Name: "org.example.Brief", New: ":1.81"}) + w.observe(NameChange{Name: "org.example.Brief", Old: ":1.81"}) + c.advance(Debounce) + w.settle(fb) + w.flush() + if got := mb.seen(); len(got) != 0 { + t.Fatalf("a flap was said: %v", got) + } +} + +func TestAServiceLeavingIsSaidWithTheUnitItHad(t *testing.T) { + mb, c := &meshBus{}, start() + w := testWatcher(t, mb, c) + fb := newFakeBus("id-1", 500, map[string]uint32{"org.freedesktop.login1": 700}) + w.connectedTo(fb) + c.advance(Debounce) + w.settle(fb) + delete(fb.names, "org.freedesktop.login1") + w.observe(NameChange{Name: "org.freedesktop.login1", Old: ":1.5"}) + c.advance(Debounce) + w.settle(fb) + w.flush() + if got := mb.seen(); !reflect.DeepEqual(got, []string{ServiceLeft}) { + t.Fatalf("%v", got) + } + if mb.bodies[0]["unit"] != "systemd-logind.service" { + t.Fatalf("%v", mb.bodies[0]) + } +} + +func TestUniqueNamesAndTheDriverAreNeverSaid(t *testing.T) { + mb, c := &meshBus{}, start() + w := testWatcher(t, mb, c) + fb := newFakeBus("id-1", 500, map[string]uint32{}) + w.connectedTo(fb) + w.observe(NameChange{Name: ":1.42", New: ":1.42"}) + w.observe(NameChange{Name: ":1.42", Old: ":1.42"}) + w.observe(NameChange{Name: busName, New: busName}) + c.advance(Debounce) + w.settle(fb) + w.flush() + if got := mb.seen(); len(got) != 0 { + t.Fatalf("%v", got) + } + for _, n := range []string{":1.1", busName, ""} { + if IsWellKnown(n) { + t.Errorf("%q counted as a service", n) + } + } +} + +func TestAPingThatHangsIsAStallAndAnAnswerARecovery(t *testing.T) { + mb, c := &meshBus{}, start() + w := testWatcher(t, mb, c) + fb := newFakeBus("id-1", 500, nil) + fb.hang = true + begun := time.Now() + w.ping(fb) + w.ping(fb) + if time.Since(begun) > 2*StallAfter+time.Second { + t.Fatal("a ping waited longer than its bound") + } + c.advance(42 * time.Second) + fb.hang = false + w.ping(fb) + w.ping(fb) + w.flush() + if got := mb.seen(); !reflect.DeepEqual(got, []string{BusStalled, BusRecovered}) { + t.Fatalf("%v", got) + } + if mb.bodies[1]["stalled_seconds"] != 42 || w.Snapshot().Stalled { + t.Fatalf("%v %+v", mb.bodies[1], w.Snapshot()) + } +} + +func TestARestartWithinABootIsSaidAndABootIsNot(t *testing.T) { + mb, c := &meshBus{}, start() + w := testWatcher(t, mb, c) + w.identity("id-1", 500) + w.identity("id-1", 500) // a reconnect to the same bus + w.identity("id-2", 501) // the bus came back as another + w.flush() + if got := mb.seen(); !reflect.DeepEqual(got, []string{BusRestarted}) { + t.Fatalf("%v", got) + } + if mb.bodies[0]["previous_bus_id"] != "id-1" || mb.bodies[0]["pid"] != uint32(501) { + t.Fatalf("%v", mb.bodies[0]) + } + + // The runtime restarts: the bus it remembers is the one still running, so nothing is said. + again := NewWatcher(w.m, mb.emit, nil, nil) + again.state, again.now = w.state, c.now + again.identity("id-2", 501) + // After a boot both change, and that is not the bus's restart. + os.WriteFile(filepath.Join(w.m.Root, "/proc/sys/kernel/random/boot_id"), []byte("boot-2\n"), 0o644) + third := NewWatcher(w.m, mb.emit, nil, nil) + third.state, third.now = w.state, c.now + third.identity("id-3", 400) + again.flush() + third.flush() + if got := mb.seen(); len(got) != 1 { + t.Fatalf("a runtime restart or a boot was taken for the bus's restart: %v", got) + } +} + +func TestAServiceThatDidNotComeBackAfterAReconnectHasLeft(t *testing.T) { + mb, c := &meshBus{}, start() + w := testWatcher(t, mb, c) + fb := newFakeBus("id-1", 500, map[string]uint32{"org.freedesktop.login1": 700, "org.bluez": 900}) + w.connectedTo(fb) + c.advance(Debounce) + w.settle(fb) + back := newFakeBus("id-2", 501, map[string]uint32{"org.freedesktop.login1": 700}) + w.connectedTo(back) + c.advance(Debounce) + w.settle(back) + w.flush() + if got := mb.seen(); !reflect.DeepEqual(got, []string{ServiceLeft}) || mb.bodies[0]["name"] != "org.bluez" { + t.Fatalf("%v %v", got, mb.bodies) + } +} + +func TestEventsWaitInOrderWhileTheMeshBusIsGone(t *testing.T) { + mb, c := &meshBus{down: true}, start() + w := testWatcher(t, mb, c) + w.markStalled("test") + c.advance(time.Minute) + w.answered() + w.flush() + if s := w.Snapshot(); s.Pending != 2 || s.Problem == "" { + t.Fatalf("what the mesh's bus did not take is not kept and said: %+v", s) + } + mb.mu.Lock() + mb.down = false + mb.mu.Unlock() + w.flush() + if got := mb.seen(); !reflect.DeepEqual(got, []string{BusStalled, BusRecovered}) { + t.Fatalf("%v", got) + } + if s := w.Snapshot(); s.Pending != 0 || s.Problem != "" { + t.Fatalf("%+v", s) + } +} + +func TestAFullQueueLetsTheOldestGo(t *testing.T) { + mb, c := &meshBus{down: true}, start() + w := testWatcher(t, mb, c) + for i := 0; i < MaxQueue+5; i++ { + w.enqueue(ServiceAppeared, map[string]any{"i": i}) + } + s := w.Snapshot() + if s.Pending != MaxQueue || s.Dropped != 5 || w.queue[0].Body["i"] != 5 { + t.Fatalf("%+v first %v", s, w.queue[0].Body) + } +} + +func TestDenialsAreSaidAtMostOncePerWindowWithACount(t *testing.T) { + mb, c := &meshBus{}, start() + w := testWatcher(t, mb, c) + batches := [][]Denial{ + {{Type: "method_call", Interface: "org.example.A", Member: "Do", Destination: "org.example"}}, + {{Type: "method_call", Interface: "org.example.A", Member: "Do", Destination: "org.example"}, + {Type: "method_call", Interface: "org.example.B", Member: "Other", Destination: "org.example"}}, + {{Type: "method_call", Interface: "org.example.A", Member: "Do", Destination: "org.example"}}, + } + n := 0 + w.journal = func(ctx context.Context, after string, since time.Time) ([]Denial, string, error) { + if n > 0 && after != "c"+string(rune('0'+n-1)) { + t.Errorf("read %d did not continue from the cursor: %q", n, after) + } + d := batches[n] + n++ + return d, "c" + string(rune('0'+n-1)), nil + } + w.readDenials(context.Background(), c.t) + c.advance(DenialsEvery) + w.readDenials(context.Background(), c.t) + c.advance(DenialEventEvery) + w.readDenials(context.Background(), c.t) + w.flush() + if got := mb.seen(); !reflect.DeepEqual(got, []string{PolicyDenied, PolicyDenied}) { + t.Fatalf("%v", got) + } + if mb.bodies[0]["count"] != 1 || mb.bodies[1]["count"] != 3 { + t.Fatalf("%v", mb.bodies) + } + if ex := mb.bodies[1]["examples"].([]Denial); len(ex) != 2 { + t.Fatalf("the same denial is one example: %v", ex) + } +} + +// TestNoEventCarriesTraffic holds every event body to names, pids, units, times, counts and the +// header fields of a denial: nothing in the watcher can carry a message's body. +func TestNoEventCarriesTraffic(t *testing.T) { + mb, c := &meshBus{}, start() + w := testWatcher(t, mb, c) + fb := newFakeBus("id-1", 500, map[string]uint32{}) + w.connectedTo(fb) + fb.names["org.bluez"] = 900 + w.observe(NameChange{Name: "org.bluez", New: ":1.9"}) + c.advance(Debounce) + w.settle(fb) + w.markStalled("x") + w.answered() + w.identity("a", 1) + w.identity("b", 2) + w.journal = func(context.Context, string, time.Time) ([]Denial, string, error) { + d, cur := ParseDenials(`{"__CURSOR":"c","MESSAGE":"A security policy denied :1.9 to send method call /p:i.m to d.","DBUS_BROKER_MESSAGE_MEMBER":"m","SECRET_BODY":"hunter2"}`) + return d, cur, nil + } + w.readDenials(context.Background(), c.t) + w.flush() + allowed := map[string]bool{"at": true, "name": true, "pid": true, "process": true, "unit": true, "activatable": true, + "reason": true, "stalled_since": true, "stalled_seconds": true, "previous_bus_id": true, "bus_id": true, + "previous_pid": true, "count": true, "since": true, "examples": true} + if len(mb.types) != 5 { + t.Fatalf("%v", mb.types) + } + for i, b := range mb.bodies { + for k := range b { + if !allowed[k] { + t.Errorf("%s carries %q", mb.types[i], k) + } + } + raw, _ := json.Marshal(b) + if strings.Contains(string(raw), "hunter2") || strings.Contains(string(raw), "security policy") { + t.Errorf("%s carries what the bus logged verbatim: %s", mb.types[i], raw) + } + } +} + +func TestRunWithoutABusSaysItStalledAndReconnects(t *testing.T) { + mb := &meshBus{} + m := testMachine(t, map[string]string{"/proc/sys/kernel/random/boot_id": "b\n"}) + fb := newFakeBus("id-1", 500, map[string]uint32{}) + var dials int + w := NewWatcher(m, mb.emit, func() (Bus, error) { + dials++ + if dials == 1 { + return nil, errors.New("no socket") + } + return fb, nil + }, nil) + w.state = filepath.Join(t.TempDir(), "bus") + w.retry = 20 * time.Millisecond + ctx, cancel := context.WithTimeout(context.Background(), 6*time.Second) + defer cancel() + done := make(chan struct{}) + go func() { w.Run(ctx); close(done) }() + deadline := time.Now().Add(6 * time.Second) + for time.Now().Before(deadline) && !w.Snapshot().Connected { + time.Sleep(50 * time.Millisecond) + } + if !w.Snapshot().Connected { + t.Fatal("the watcher did not reconnect") + } + close(fb.changes) // the bus goes away + for time.Now().Before(deadline) && w.Snapshot().Connected { + time.Sleep(10 * time.Millisecond) + } + cancel() + <-done + got := mb.seen() + if len(got) < 2 || got[0] != BusStalled || got[1] != BusRecovered { + t.Fatalf("%v", got) + } +} diff --git a/modules/dbus/go.mod b/modules/dbus/go.mod new file mode 100644 index 0000000..8f659da --- /dev/null +++ b/modules/dbus/go.mod @@ -0,0 +1,8 @@ +module dbus + +go 1.22 + +require ( + git.novox.be/novox/mesh-sdk/go v0.1.7 + github.com/godbus/dbus/v5 v5.1.0 +) diff --git a/modules/dbus/go.sum b/modules/dbus/go.sum new file mode 100644 index 0000000..8c37f7d --- /dev/null +++ b/modules/dbus/go.sum @@ -0,0 +1,4 @@ +git.novox.be/novox/mesh-sdk/go v0.1.7 h1:C0sTQmtTiyYH7bnqZb7PusXnqA37gKuT7Nqjn9gG47w= +git.novox.be/novox/mesh-sdk/go v0.1.7/go.mod h1:GFuZUElBZ9A++mxgIKo97aXXo+kV0uJ/UkbhQPPIbrY= +github.com/godbus/dbus/v5 v5.1.0 h1:4KLkAxT3aOY8Li4FRJe/KvhoNFFxo0m6fNuFUO8QJUk= +github.com/godbus/dbus/v5 v5.1.0/go.mod h1:xhWf0FNVPg57R7Z0UbKHbJfkEywrmjJnf7w5xrFpKfA= diff --git a/modules/dbus/module.json b/modules/dbus/module.json new file mode 100644 index 0000000..c0e28b3 --- /dev/null +++ b/modules/dbus/module.json @@ -0,0 +1,67 @@ +{ + "module": "dbus", + "version": "1", + "capabilities": [ + "package-manager", + "service-manager" + ], + "claims": [ + { + "name": "node-message-bus", + "scope": "node" + } + ], + "emits": [ + "bus.stalled", + "bus.recovered", + "bus.restarted", + "service.appeared", + "service.left", + "policy.denied" + ], + "tools": [ + "dbus_names", + "dbus_introspect", + "dbus_monitor", + "dbus_check", + "dbus_health" + ], + "resources": [ + { + "id": "dbus", + "type": "package", + "package": "dbus" + }, + { + "id": "broker", + "type": "package", + "package": "dbus-broker" + }, + { + "id": "broker-units", + "type": "package", + "package": "dbus-broker-units" + }, + { + "id": "system-bus", + "type": "service", + "unit": "dbus-broker.service", + "state": "running" + } + ], + "build": { + "artifacts": [ + { + "name": "tools-go", + "kind": "bundle", + "language": "go", + "system": "arch", + "from": "cmd/dbus-tools", + "binary": "dbus-tools", + "loads": [ + "dbus-tools" + ] + } + ] + } +}