From dc9c4282e74257350cf9a6b1719878a431259f2e Mon Sep 17 00:00:00 2001 From: Zedwil Date: Fri, 2 Oct 2026 10:12:31 -0500 Subject: [PATCH] fix: read sink state once per tick instead of three pactl calls per dial --- dev.openwave.sdPlugin/owdeck/graph.py | 47 +++++++----- dev.openwave.sdPlugin/plugin.py | 29 +++++++- tests/test_plugin.py | 102 +++++++++++++++++++++++++- 3 files changed, 154 insertions(+), 24 deletions(-) diff --git a/dev.openwave.sdPlugin/owdeck/graph.py b/dev.openwave.sdPlugin/owdeck/graph.py index d26e577..d486f70 100644 --- a/dev.openwave.sdPlugin/owdeck/graph.py +++ b/dev.openwave.sdPlugin/owdeck/graph.py @@ -54,26 +54,37 @@ def node_id(node_name): return best -def sink_exists(name): - out = _run(["pactl", "list", "short", "sinks"]) - if not out: - return False - return any(line.split("\t")[1] == name - for line in out.splitlines() if "\t" in line) - - -def get_volume(name): - """Sink volume as 0..1, or None if it is not there.""" - out = _run(["pactl", "get-sink-volume", name]) +def sinks(): + """Every sink's (volume 0..1 or None, muted), by name, from one call. + + One pactl for the whole graph rather than three per sink: the deck reads + every dial several times a second, and a subprocess per field per dial is + a process storm on pipewire-pulse. JSON rather than the text listing, + whose field labels are translated. Volume is the first channel's, which + is what `pactl get-sink-volume` leads with. {} when pactl cannot be read, + which reads as every sink being gone. + """ + out = _run(["pactl", "-f", "json", "list", "sinks"]) if not out: - return None - for token in out.split(): - if token.endswith("%"): + return {} + try: + listing = json.loads(out) + except json.JSONDecodeError: + return {} + found = {} + for sink in listing if isinstance(listing, list) else []: + if not isinstance(sink, dict) or not sink.get("name"): + continue + volume = None + channels = sink.get("volume") + if isinstance(channels, dict) and channels: + first = next(iter(channels.values())) try: - return int(token.rstrip("%")) / 100.0 - except ValueError: - return None - return None + volume = int(str(first["value_percent"]).rstrip("%")) / 100.0 + except (KeyError, TypeError, ValueError): + volume = None + found[sink["name"]] = (volume, bool(sink.get("mute"))) + return found def set_volume(name, value): diff --git a/dev.openwave.sdPlugin/plugin.py b/dev.openwave.sdPlugin/plugin.py index 0fa913f..edca64b 100644 --- a/dev.openwave.sdPlugin/plugin.py +++ b/dev.openwave.sdPlugin/plugin.py @@ -172,6 +172,8 @@ def __init__(self, port, uuid, register_event): self._last_refresh = 0.0 self._snapshot = None self._snapshot_at = 0.0 + self._sinks = {} + self._sinks_at = 0.0 self._theme = render.DEFAULT_THEME # ------------------------------------------------------------ outbound @@ -195,6 +197,21 @@ def _openwave(self, force=False): self._snapshot = ipc.snapshot() return self._snapshot + def _sink_states(self): + """Every sink's (volume, muted), refetched at most once a tick. + + Shared by every mix key the same way the snapshot is: the meters + redraw each dial at ~8 Hz, and asking pactl per dial per redraw ran + dozens of processes a second. Our own writes drop it (_sinks_at = 0) + so the next read sees them; anything else changing a sink shows + within SNAPSHOT_SECONDS. + """ + now = time.monotonic() + if self._sinks_at == 0.0 or now - self._sinks_at >= SNAPSHOT_SECONDS: + self._sinks_at = now + self._sinks = graph.sinks() + return self._sinks + def _read(self, settings): """What a volume key should show: name, percent, mute, glyph. @@ -210,14 +227,15 @@ def _read(self, settings): name = owstate.mix_name(ident) glyph = _MIX_GLYPHS.get( (owstate.mixes().get(ident) or {}).get("icon_name"), "speaker") - if not sink or not graph.sink_exists(sink): + found = self._sink_states().get(sink) if sink else None + if found is None: return {"name": name, "percent": 0, "muted": False, "glyph": glyph, "ok": False, "context": ""} - volume = graph.get_volume(sink) + volume, muted = found return { "name": name, "percent": 0 if volume is None else round(volume * 100), - "muted": bool(graph.get_mute(sink)), + "muted": muted, "glyph": glyph, "ok": volume is not None, "context": "", @@ -502,6 +520,7 @@ def _set_level(self, settings, value): sink = owstate.mix_sink(ident) if sink: graph.set_volume(sink, value) + self._sinks_at = 0.0 elif kind == "src": ipc.set_source_level(ident, value) self._openwave(force=True) @@ -541,6 +560,7 @@ def _toggle_mute(self, context): self._send("showAlert", context) return graph.toggle_mute(sink) + self._sinks_at = 0.0 elif kind == "cell": if ipc.toggle_cell_mute(*split_cell(ident)) is False: self._send("showAlert", context) @@ -788,6 +808,9 @@ def _on_push(self, states): self._snapshot_at = time.monotonic() except (TypeError, ValueError): pass + # A push can follow a mix master moved in OpenWave's window, which + # lives in the sink rather than the snapshot. + self._sinks_at = 0.0 self._render_all() def _tick_levels(self, now): diff --git a/tests/test_plugin.py b/tests/test_plugin.py index f590670..3ae6985 100644 --- a/tests/test_plugin.py +++ b/tests/test_plugin.py @@ -20,6 +20,8 @@ import plugin as P # noqa: E402 from owdeck import graph, ipc, owstate, render # noqa: E402 +REAL_SINKS = graph.sinks + SNAPSHOT = { "groups": ["Mic"], @@ -66,6 +68,7 @@ class PluginCase(unittest.TestCase): def setUp(self): self.calls = [] + self.sink_reads = 0 self.sinks = {"openwave_personal_mix": [1.0, False], "openwave_chat_mix": [0.6, True]} self.snapshot = json.loads(json.dumps(SNAPSHOT)) @@ -78,9 +81,7 @@ def setUp(self): lambda i: (MIXES.get(i) or {}).get("sink")) self._patch(owstate, "mix_name", lambda i, d=None: (MIXES.get(i) or {}).get("name", i)) - self._patch(graph, "sink_exists", lambda n: n in self.sinks) - self._patch(graph, "get_volume", - lambda n: self.sinks[n][0] if n in self.sinks else None) + self._patch(graph, "sinks", self._list_sinks) self._patch(graph, "get_mute", lambda n: self.sinks[n][1] if n in self.sinks else None) self._patch(graph, "set_volume", self._set_volume) @@ -102,6 +103,8 @@ def setUp(self): self.plugin._last_refresh = 0.0 self.plugin._snapshot = None self.plugin._snapshot_at = 0.0 + self.plugin._sinks = {} + self.plugin._sinks_at = 0.0 self.plugin._scene_pressed = {} self.plugin._levels = {} self.plugin._levels_at = 0.0 @@ -117,6 +120,10 @@ def _patch(self, module, name, replacement): setattr(module, name, replacement) # -- stub behaviours -------------------------------------------------- + def _list_sinks(self): + self.sink_reads += 1 + return {n: (v, m) for n, (v, m) in self.sinks.items()} + def _set_volume(self, name, value): self.calls.append(("sink-volume", name, round(value, 3))) self.sinks[name][0] = max(0.0, min(1.0, value)) @@ -567,6 +574,7 @@ def test_an_unchanged_key_is_not_redrawn(self): def test_a_changed_level_is_redrawn(self): self.place("c", P.VOLUME, {"target": "mix:personal"}) self.sinks["openwave_personal_mix"][0] = 0.2 + self.plugin._sinks_at = 0.0 # the next tick, when it is reread self.plugin._render("c") self.assertEqual(len(self.events("setImage")), 1) @@ -574,6 +582,7 @@ def test_a_level_change_under_a_mute_is_not_redrawn(self): """A muted key shows MUTED, not a number, so it has not changed.""" self.place("c", P.VOLUME, {"target": "mix:chat"}) self.sinks["openwave_chat_mix"][0] = 0.2 + self.plugin._sinks_at = 0.0 # the next tick, when it is reread self.plugin._render("c") self.assertEqual(self.events("setImage"), []) @@ -1084,6 +1093,93 @@ def test_garbage_push_is_ignored(self): self.assertIn("stale", self.plugin._snapshot) +class TestSinkPolling(PluginCase): + """The meters redraw every dial at ~8 Hz; reading the sinks must not.""" + + def setUp(self): + super().setUp() + self.clock = 100.0 + self.spawned = [] + self._patch(P.time, "monotonic", lambda: self.clock) + self._patch(ipc, "levels", lambda: {}) + # The real reader, over a fake pactl, so what is counted is + # processes and not calls into a stub. + self._patch(graph, "sinks", REAL_SINKS) + self._patch(graph.subprocess, "run", self._fake_run) + + def _fake_run(self, argv, **_kwargs): + self.spawned.append(argv) + listing = [ + {"name": name, "mute": muted, + "volume": {"front-left": {"value_percent": f"{round(v * 100)}%"}, + "front-right": {"value_percent": "0%"}}} + for name, (v, muted) in self.sinks.items() + ] + return graph.subprocess.CompletedProcess( + argv, 0, stdout=json.dumps(listing), stderr="") + + def _dials(self): + for i, mix in enumerate(("personal", "chat", "personal", "chat")): + self.place(f"d{i}", P.VOLUME, {"target": f"mix:{mix}"}, + controller="Encoder") + + def _run_for(self, seconds): + """Drive the meter tick and the refresh timer as run() does.""" + end = self.clock + seconds + while self.clock < end: + self.clock += 0.125 + self.plugin._tick_levels(self.clock) + if self.clock - self.plugin._last_refresh >= P.REFRESH_SECONDS: + self.plugin._last_refresh = self.clock + self.plugin._render_all() + + def test_many_dials_cost_one_pactl_per_refresh(self): + self._dials() + self.spawned.clear() + self._run_for(3.0) + # 24 meter ticks and 3 redraws over four dials: one read per + # SNAPSHOT_SECONDS window, not three per dial per tick (~300). + self.assertLessEqual(len(self.spawned), 4) + self.assertTrue(all(argv[0] == "pactl" for argv in self.spawned)) + + def test_the_dial_still_shows_the_sink(self): + self._dials() + state = self.plugin._read({"target": "mix:chat"}) + self.assertEqual((state["percent"], state["muted"], state["ok"]), + (60, True, True)) + + def test_an_outside_change_shows_within_a_second(self): + self._dials() + self.sinks["openwave_chat_mix"] = [0.25, False] + self._run_for(1.0) + state = self.plugin._read({"target": "mix:chat"}) + self.assertEqual((state["percent"], state["muted"]), (25, False)) + + def test_a_vanished_sink_reads_unavailable(self): + self._dials() + self.sinks.pop("openwave_chat_mix") + self._run_for(1.0) + self.assertFalse(self.plugin._read({"target": "mix:chat"})["ok"]) + + def test_pactl_failing_reads_unavailable(self): + self._patch(graph.subprocess, "run", + lambda argv, **_: graph.subprocess.CompletedProcess( + argv, 1, stdout="", stderr="")) + self.assertFalse(self.plugin._read({"target": "mix:chat"})["ok"]) + + def test_turning_fast_steps_from_its_own_write(self): + """Back-to-back detents inside one cache window must each step from + the last, not from the level the cache held before the turn.""" + self.place("d", P.VOLUME, {"target": "mix:chat"}, + controller="Encoder") + for _ in range(3): + self.plugin._handle({"event": "dialRotate", "context": "d", + "payload": {"ticks": 1}}) + self.assertEqual( + [c[2] for c in self.calls if c[0] == "sink-volume"], + [0.62, 0.64, 0.66]) + + @unittest.skipUnless(ipc._HAVE_GI, "needs PyGObject") class TestChangedSignal(unittest.TestCase): """The relay must unpack org.gtk.Actions.Changed as the bus sends it."""