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