Add MPRIS2 volume routing and PA event watcher; fix Spotify song transition reset
- Add MPRISController class: wraps playerctl with 3s player-list cache, fuzzy app-name matching, and thread-safe get/set_volume via MPRIS2 D-Bus - Add PulseController event watcher thread: separate pulsectl connection listens for sink-input events and queues them so reapply fires within ~0.1s for PA-only apps instead of waiting up to 1s - Route set_app_volume and get_app_volume_norm through MPRIS first, falling back to PA stream for non-MPRIS apps (e.g. Brave/Chromium) - Update reapply_app_volumes to prefer MPRIS path, preventing turnupd from fighting Spotify on the PA layer during song transitions - Wire MPRISController and start_watching() into main(); drain_events() triggers immediate reapply on any PA sink-input eventmain
parent
f1fa3080b5
commit
eba8e66879
|
|
@ -14,9 +14,11 @@ even when no knob is being moved.
|
|||
|
||||
import logging
|
||||
import os
|
||||
import queue
|
||||
import signal
|
||||
import subprocess
|
||||
import sys
|
||||
import threading
|
||||
import time
|
||||
|
||||
import pulsectl
|
||||
|
|
@ -131,17 +133,156 @@ def all_led_colors(
|
|||
]
|
||||
|
||||
|
||||
# ── MPRIS2 controller (playerctl-backed) ──────────────────────────────────────
|
||||
|
||||
class MPRISController:
|
||||
"""Uses *playerctl* to read/write per-app volume via the MPRIS2 D-Bus interface.
|
||||
|
||||
Caches the player list for ``_CACHE_TTL`` seconds to avoid spawning a new
|
||||
subprocess on every single call.
|
||||
"""
|
||||
|
||||
_CACHE_TTL: float = 3.0 # seconds between ``playerctl --list-all`` calls
|
||||
|
||||
def __init__(self) -> None:
|
||||
self._players: list[str] = []
|
||||
self._players_ts: float = 0.0
|
||||
self._lock = threading.Lock()
|
||||
|
||||
# ── internal helpers ──────────────────────────────────────────────────────
|
||||
|
||||
def _run(self, *args: str, timeout: float = 2.0) -> tuple[bool, str]:
|
||||
"""Run ``playerctl <args>`` and return ``(success, stdout.strip())``."""
|
||||
try:
|
||||
result = subprocess.run(
|
||||
["playerctl", *args],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
timeout=timeout,
|
||||
)
|
||||
return result.returncode == 0, result.stdout.strip()
|
||||
except (FileNotFoundError, subprocess.TimeoutExpired) as exc:
|
||||
log.debug("playerctl call failed: %s", exc)
|
||||
return False, ""
|
||||
|
||||
def _refresh_players(self, *, force: bool = False) -> None:
|
||||
"""Refresh the cached player list if it has expired (or *force* is set)."""
|
||||
now = time.monotonic()
|
||||
if not force and (now - self._players_ts) < self._CACHE_TTL:
|
||||
return
|
||||
ok, out = self._run("--list-all")
|
||||
with self._lock:
|
||||
self._players = [p.strip() for p in out.splitlines() if p.strip()] if ok else []
|
||||
self._players_ts = now
|
||||
|
||||
# ── public API ────────────────────────────────────────────────────────────
|
||||
|
||||
def find_player(self, app_name: str) -> str | None:
|
||||
"""Return the first cached player whose name contains *app_name* (case-insensitive)."""
|
||||
self._refresh_players()
|
||||
needle = app_name.lower()
|
||||
with self._lock:
|
||||
for player in self._players:
|
||||
if needle in player.lower():
|
||||
return player
|
||||
return None
|
||||
|
||||
def get_volume(self, app_name: str) -> float | None:
|
||||
"""Return the MPRIS volume (0.0–1.0) for *app_name*, or ``None`` if unavailable."""
|
||||
player = self.find_player(app_name)
|
||||
if player is None:
|
||||
return None
|
||||
ok, out = self._run("--player", player, "volume")
|
||||
if not ok or not out:
|
||||
return None
|
||||
try:
|
||||
return max(0.0, min(1.0, float(out)))
|
||||
except ValueError:
|
||||
return None
|
||||
|
||||
def set_volume(self, app_name: str, volume: float) -> bool:
|
||||
"""Set the MPRIS volume for *app_name*. Returns ``True`` on success."""
|
||||
player = self.find_player(app_name)
|
||||
if player is None:
|
||||
return False
|
||||
volume = max(0.0, min(1.0, volume))
|
||||
ok, _ = self._run("--player", player, "volume", f"{volume:.4f}")
|
||||
if ok:
|
||||
log.debug("MPRIS: %r volume → %.4f", player, volume)
|
||||
return ok
|
||||
|
||||
|
||||
# ── PulseAudio / PipeWire controller ──────────────────────────────────────────
|
||||
|
||||
class PulseController:
|
||||
"""Thin wrapper around :class:`pulsectl.Pulse` for volume and mute control."""
|
||||
"""Thin wrapper around :class:`pulsectl.Pulse` for volume and mute control.
|
||||
|
||||
def __init__(self) -> None:
|
||||
When *mpris* is supplied, ``set_app_volume`` and ``get_app_volume_norm``
|
||||
will prefer the MPRIS2 path for any app that has a live playerctl player,
|
||||
falling back to the PulseAudio stream only when MPRIS is unavailable.
|
||||
|
||||
A background watcher thread (started by :meth:`start_watching`) listens
|
||||
for PulseAudio sink-input events and pushes indices onto ``_event_q`` so
|
||||
the main loop can trigger an immediate reapply for PA-only apps instead of
|
||||
waiting for the 1-second timer.
|
||||
"""
|
||||
|
||||
def __init__(self, mpris: MPRISController | None = None) -> None:
|
||||
self._pulse = pulsectl.Pulse("turnupd")
|
||||
self._mpris = mpris
|
||||
self._event_q: queue.Queue[int] = queue.Queue()
|
||||
self._watcher_thread: threading.Thread | None = None
|
||||
self._stop_event = threading.Event()
|
||||
|
||||
def close(self) -> None:
|
||||
self._stop_event.set()
|
||||
self._pulse.close()
|
||||
|
||||
# ── PA event watcher ──────────────────────────────────────────────────────
|
||||
|
||||
def start_watching(self) -> None:
|
||||
"""Spawn the background PA event listener thread (idempotent)."""
|
||||
if self._watcher_thread and self._watcher_thread.is_alive():
|
||||
return
|
||||
self._stop_event.clear()
|
||||
t = threading.Thread(target=self._event_loop, daemon=True, name="pa-watcher")
|
||||
t.start()
|
||||
self._watcher_thread = t
|
||||
log.debug("PA watcher thread started")
|
||||
|
||||
def _event_loop(self) -> None:
|
||||
"""Background thread: open a *separate* Pulse connection and listen for events."""
|
||||
try:
|
||||
with pulsectl.Pulse("turnupd-watcher") as watch_pulse:
|
||||
def _cb(ev: pulsectl.PulseEventInfo) -> None: # type: ignore[name-defined]
|
||||
if ev.facility == "sink_input":
|
||||
self._event_q.put(int(ev.index))
|
||||
raise pulsectl.PulseLoopStop
|
||||
|
||||
watch_pulse.event_mask_set("sink_input")
|
||||
watch_pulse.event_callback_set(_cb)
|
||||
while not self._stop_event.is_set():
|
||||
try:
|
||||
watch_pulse.event_listen(timeout=1.0)
|
||||
except pulsectl.PulseLoopStop:
|
||||
pass
|
||||
except Exception as exc:
|
||||
log.debug("PA event loop error: %s", exc)
|
||||
time.sleep(0.5)
|
||||
except Exception as exc:
|
||||
log.warning("PA watcher thread exiting: %s", exc)
|
||||
|
||||
def drain_events(self) -> bool:
|
||||
"""Drain all pending PA events. Returns ``True`` if any events were present."""
|
||||
had_events = False
|
||||
while True:
|
||||
try:
|
||||
self._event_q.get_nowait()
|
||||
had_events = True
|
||||
except queue.Empty:
|
||||
break
|
||||
return had_events
|
||||
|
||||
def set_sink_volume(self, sink_name: str, volume: float) -> None:
|
||||
volume = max(0.0, min(VOLUME_MAX, volume))
|
||||
try:
|
||||
|
|
@ -192,14 +333,24 @@ class PulseController:
|
|||
|
||||
def set_app_volume(self, app_name: str, volume: float) -> None:
|
||||
volume = max(0.0, min(VOLUME_MAX, volume))
|
||||
|
||||
# Prefer the MPRIS2 path — it writes to the app's internal slider so the
|
||||
# volume survives song transitions (e.g. Spotify resetting on new tracks).
|
||||
if self._mpris and self._mpris.set_volume(app_name, volume):
|
||||
log.debug("MPRIS set_volume: %r = %.4f", app_name, volume)
|
||||
return
|
||||
|
||||
# Fall back to PulseAudio stream volume.
|
||||
needle = app_name.lower()
|
||||
found = False
|
||||
try:
|
||||
for inp in self._pulse.sink_input_list():
|
||||
name = inp.proplist.get("application.name", "")
|
||||
binary = inp.proplist.get("application.process.binary", "")
|
||||
if needle in name.lower() or needle in binary.lower():
|
||||
self._pulse.volume_set_all_chans(inp, volume)
|
||||
return
|
||||
found = True
|
||||
if not found:
|
||||
log.debug("App %r not found in sink inputs", app_name)
|
||||
except Exception as exc:
|
||||
log.warning("set_app_volume(%r) failed: %s", app_name, exc)
|
||||
|
|
@ -230,6 +381,13 @@ class PulseController:
|
|||
|
||||
def get_app_volume_norm(self, app_name: str) -> float | None:
|
||||
"""Return the current app volume normalised to 0.0–1.0, or None if not found."""
|
||||
# Prefer MPRIS — more accurate for apps like Spotify.
|
||||
if self._mpris:
|
||||
vol = self._mpris.get_volume(app_name)
|
||||
if vol is not None:
|
||||
return vol
|
||||
|
||||
# Fall back to PulseAudio stream.
|
||||
needle = app_name.lower()
|
||||
try:
|
||||
for inp in self._pulse.sink_input_list():
|
||||
|
|
@ -278,6 +436,83 @@ def init_knob_norms(config: dict, pulse: PulseController) -> list[float]:
|
|||
return norms
|
||||
|
||||
|
||||
def build_app_volume_map(config: dict, knob_norms: list[float]) -> dict[str, float]:
|
||||
"""Return ``{app_name_lower: volume}`` for every app/group knob in *config*.
|
||||
|
||||
Used by :func:`reapply_app_volumes` to know what volume each configured
|
||||
application should currently be at, based on the last knob positions.
|
||||
"""
|
||||
app_volumes: dict[str, float] = {}
|
||||
for knob_id_str, knob_cfg in config.get("knobs", {}).items():
|
||||
try:
|
||||
knob_id = int(knob_id_str)
|
||||
except ValueError:
|
||||
continue
|
||||
if knob_id >= NUM_KNOBS:
|
||||
continue
|
||||
action = knob_cfg.get("action", "")
|
||||
vol = round(knob_norms[knob_id] * VOLUME_MAX, 4)
|
||||
if action == "app_volume":
|
||||
t = knob_cfg.get("target", "")
|
||||
if t:
|
||||
app_volumes[t.lower()] = vol
|
||||
elif action == "group_volume":
|
||||
for t in knob_cfg.get("targets", []):
|
||||
if t:
|
||||
app_volumes[t.lower()] = vol
|
||||
return app_volumes
|
||||
|
||||
|
||||
def reapply_app_volumes(config: dict, pulse: PulseController, knob_norms: list[float]) -> None:
|
||||
"""Re-apply stored knob volumes to every matching active sink input.
|
||||
|
||||
Called on a 1-second timer and whenever a PA sink-input event fires so
|
||||
that new streams (e.g. Spotify starting a new song) are brought back to
|
||||
the last knob position rather than being left at the 100 % default that
|
||||
``module-stream-restore`` restores them to.
|
||||
|
||||
For MPRIS-capable apps the volume is written via playerctl (which updates
|
||||
the app's own internal slider). For PA-only apps (e.g. Brave) the volume
|
||||
is corrected on the PulseAudio stream level.
|
||||
"""
|
||||
app_volumes = build_app_volume_map(config, knob_norms)
|
||||
if not app_volumes:
|
||||
return
|
||||
|
||||
# Split targets into MPRIS-handled vs PA-only.
|
||||
mpris = pulse._mpris
|
||||
pa_only: dict[str, float] = {}
|
||||
|
||||
for app_name, vol in app_volumes.items():
|
||||
if mpris and mpris.set_volume(app_name, vol):
|
||||
log.debug("reapply MPRIS: %r → %.4f", app_name, vol)
|
||||
else:
|
||||
pa_only[app_name] = vol
|
||||
|
||||
if not pa_only:
|
||||
return
|
||||
|
||||
# PA stream correction for non-MPRIS apps.
|
||||
try:
|
||||
for inp in pulse._pulse.sink_input_list():
|
||||
name = inp.proplist.get("application.name", "").lower()
|
||||
binary = inp.proplist.get("application.process.binary", "").lower()
|
||||
for needle, vol in pa_only.items():
|
||||
if needle in name or needle in binary:
|
||||
current = inp.volume.value_flat
|
||||
if abs(current - vol) > 0.01:
|
||||
pulse._pulse.volume_set_all_chans(inp, vol)
|
||||
log.debug(
|
||||
"reapply PA: %r volume %.2f → %.2f",
|
||||
inp.proplist.get("application.name", needle),
|
||||
current,
|
||||
vol,
|
||||
)
|
||||
break # one needle is enough per stream
|
||||
except Exception as exc:
|
||||
log.debug("reapply_app_volumes (PA) failed: %s", exc)
|
||||
|
||||
|
||||
# ── Event handlers ─────────────────────────────────────────────────────────────
|
||||
|
||||
def handle_knob(
|
||||
|
|
@ -356,7 +591,9 @@ def main() -> None:
|
|||
|
||||
log.info("Turn Up daemon starting — %s @ %d baud", port, baud)
|
||||
|
||||
pulse = PulseController()
|
||||
mpris = MPRISController()
|
||||
pulse = PulseController(mpris)
|
||||
pulse.start_watching()
|
||||
knob_norms = init_knob_norms(config, pulse)
|
||||
buf = bytearray()
|
||||
|
||||
|
|
@ -366,6 +603,9 @@ def main() -> None:
|
|||
except OSError:
|
||||
config_mtime = None
|
||||
last_config_check = time.monotonic()
|
||||
# Re-apply app volumes periodically so new streams (e.g. Spotify new song)
|
||||
# are brought to the last knob position rather than resetting to 100 %.
|
||||
last_reapply = time.monotonic()
|
||||
|
||||
def _shutdown(sig: int, _frame: object) -> None:
|
||||
log.info("Received signal %d — shutting down", sig)
|
||||
|
|
@ -413,6 +653,14 @@ def main() -> None:
|
|||
except OSError:
|
||||
pass
|
||||
|
||||
# Re-apply configured app volumes every 1 s to catch new
|
||||
# streams (e.g. Spotify starting a new song resets to 100 %).
|
||||
# Also trigger immediately on any PA sink-input event so
|
||||
# PA-only apps (e.g. Brave) are corrected within ~0.1 s.
|
||||
if pulse.drain_events() or now - last_reapply >= 1.0:
|
||||
last_reapply = now
|
||||
reapply_app_volumes(config, pulse, knob_norms)
|
||||
|
||||
except serial.SerialException as exc:
|
||||
log.warning("Serial error: %s — retrying in 3 s", exc)
|
||||
time.sleep(3)
|
||||
|
|
|
|||
Loading…
Reference in New Issue