from __future__ import annotations import asyncio from dataclasses import dataclass from time import monotonic from app.bpm.service import BpmService from app.dmx.engine import DmxEngine from app.dmx.frame import FrameLayer from app.models.schemas import EffectPayload @dataclass(slots=True) class ActiveEffect: slug: str started_at: float duration_ms: int sync_mode: str = "manual" class EffectService: def __init__(self, engine: DmxEngine, bpm: BpmService) -> None: self.engine = engine self.bpm = bpm self.effects: dict[str, EffectPayload] = {} self.active_effects: dict[str, ActiveEffect] = {} self._effect_tasks: dict[str, asyncio.Task[None]] = {} def save(self, payload: EffectPayload) -> EffectPayload: self.effects[payload.slug] = payload return payload def list(self) -> list[EffectPayload]: return list(self.effects.values()) def trigger(self, slug: str) -> EffectPayload: effect = self.effects[slug] self.stop(slug) if effect.effect_type == "beat-flash": self._effect_tasks[slug] = asyncio.create_task( self._run_beat_flash(effect), name=f"tuxdmx-effect-{slug}", ) self.active_effects[slug] = ActiveEffect( slug=slug, started_at=monotonic(), duration_ms=self._coerce_duration_ms(effect.parameters.get("duration_ms"), fallback=180), sync_mode="bpm", ) return effect values_by_universe, precedence_map_by_universe = self._resolve_channel_maps(effect) duration_ms = self._coerce_duration_ms(effect.parameters.get("duration_ms"), fallback=5000) self.engine.set_layer( FrameLayer( name=f"effect:{slug}", priority=effect.priority, values_by_universe=values_by_universe, precedence_map_by_universe=precedence_map_by_universe, ) ) self.active_effects[slug] = ActiveEffect( slug=slug, started_at=monotonic(), duration_ms=duration_ms, sync_mode="static", ) return effect def stop(self, slug: str) -> None: task = self._effect_tasks.pop(slug, None) if task is not None: task.cancel() self.engine.remove_layer(f"effect:{slug}") self.active_effects.pop(slug, None) async def shutdown(self) -> None: for slug in list(self._effect_tasks): self.stop(slug) tasks = list(self._effect_tasks.values()) self._effect_tasks.clear() for task in tasks: try: await task except asyncio.CancelledError: pass async def _run_beat_flash(self, effect: EffectPayload) -> None: slug = effect.slug values_by_universe, precedence_map_by_universe = self._resolve_channel_maps(effect) if not any(values_by_universe.values()): return pulse_ms = self._coerce_duration_ms(effect.parameters.get("duration_ms"), fallback=180) pulse_seconds = max(0.05, pulse_ms / 1000) last_seen_beat = self.bpm.beat_counter next_manual_pulse_at = monotonic() try: while True: bpm_value = max(40.0, min(220.0, float(self.bpm.current_bpm))) beat_interval = max(0.2, 60.0 / bpm_value) should_pulse = False if self.bpm.audio_connected: if self.bpm.beat_counter != last_seen_beat: last_seen_beat = self.bpm.beat_counter should_pulse = True else: await asyncio.sleep(0.02) continue else: now = monotonic() if now >= next_manual_pulse_at: next_manual_pulse_at = now + beat_interval should_pulse = True else: await asyncio.sleep(min(0.05, next_manual_pulse_at - now)) continue if not should_pulse: continue self.engine.set_layer( FrameLayer( name=f"effect:{slug}", priority=effect.priority, values_by_universe=values_by_universe, precedence_map_by_universe=precedence_map_by_universe, ) ) await asyncio.sleep(min(pulse_seconds, beat_interval * 0.8)) self.engine.remove_layer(f"effect:{slug}") except asyncio.CancelledError: self.engine.remove_layer(f"effect:{slug}") raise def _resolve_channel_maps(self, effect: EffectPayload) -> tuple[dict[int, dict[int, int]], dict[int, dict[int, str]]]: raw_channels = effect.parameters.get("channels", {}) raw_precedence = effect.parameters.get("precedence", {}) channels = raw_channels if isinstance(raw_channels, dict) else {} precedence = raw_precedence if isinstance(raw_precedence, dict) else {} values = { int(channel): int(value) for channel, value in channels.items() } precedence_map = { int(channel): str(value) for channel, value in precedence.items() } universe = int(effect.parameters.get("universe", 1) or 1) return {universe: values}, {universe: precedence_map} def _coerce_duration_ms(self, raw_value: object, fallback: int) -> int: if isinstance(raw_value, (int, float)): return max(50, int(raw_value)) if isinstance(raw_value, str): try: return max(50, int(float(raw_value))) except ValueError: return fallback return fallback