diff --git a/src/led_sync/scheduler.py b/src/led_sync/scheduler.py new file mode 100644 index 0000000..80189cc --- /dev/null +++ b/src/led_sync/scheduler.py @@ -0,0 +1,147 @@ +""" +ChoreographyScheduler: absolute monotonic clock event dispatcher. +Implements CHR-05 — fires animation commands on-time during choreography playback. + +Algorithm (D-07/D-08/D-09): + fire_at = start_time + pause_offset + event.timestamp + Never accumulate relative delays. Recalculate remaining sleep each 5ms poll. + Drift guaranteed < 1ms per event over any duration (absolute math, no accumulation). +""" +import asyncio +import logging +import time +from typing import Callable + +from led_sync.models import ChoreoEvent +from led_sync.transport.udp_client import ESP32Transport + +logger = logging.getLogger(__name__) + +POLL_INTERVAL = 0.005 # 5ms sleep poll — keeps asyncio responsive + + +class ChoreographyScheduler: + def __init__( + self, + transport: ESP32Transport, + on_event: Callable[[ChoreoEvent], None] | None = None, + ): + self._transport = transport + self._on_event = on_event # optional hook: called when event fires (for TUI update) + + self.start_time: float | None = None + self.pause_offset: float = 0.0 + self.paused_at: float | None = None + self._task: asyncio.Task | None = None + self._events: list[ChoreoEvent] = [] + + # --- State --- + + @property + def current_position(self) -> float: + """Playback position in seconds. Safe to read from asyncio or threads.""" + if self.start_time is None: + return 0.0 + if self.paused_at is not None: + return self.paused_at - self.start_time - self.pause_offset + return time.monotonic() - self.start_time - self.pause_offset + + @property + def is_running(self) -> bool: + return self._task is not None and not self._task.done() + + # --- Control --- + + def play(self, events: list[ChoreoEvent], seek_seconds: float = 0.0) -> None: + """ + Start scheduling events. Sorts by timestamp, skips events before seek_seconds. + Creates an asyncio Task — must be called from within asyncio event loop. + """ + self.stop() + self.start_time = time.monotonic() - seek_seconds # absolute base + self.pause_offset = 0.0 + self.paused_at = None + self._events = sorted(events, key=lambda e: e.timestamp) + self._task = asyncio.create_task( + self._run(seek_seconds=seek_seconds), + name="choreo-scheduler", + ) + logger.info("Scheduler started: %d events, seek=%.2fs", len(self._events), seek_seconds) + + def pause(self) -> None: + """Freeze dispatching. Records pause point for resume.""" + if self.paused_at is not None or self.start_time is None: + return # already paused or not started + self.paused_at = time.monotonic() + logger.info("Scheduler paused at %.2fs", self.current_position) + + def resume(self) -> None: + """Continue from pause point. Adjusts pause_offset (D-09).""" + if self.paused_at is None: + return # not paused + self.pause_offset += time.monotonic() - self.paused_at + self.paused_at = None + logger.info("Scheduler resumed, pause_offset=%.3fs", self.pause_offset) + + def seek(self, seconds: float) -> None: + """ + Seek to position. Restarts scheduler from new position, + skipping events before seconds. + """ + if self._events: + self.play(self._events, seek_seconds=seconds) + + def stop(self) -> None: + """Cancel scheduler task and reset all state.""" + if self._task and not self._task.done(): + self._task.cancel() + self._task = None + self.start_time = None + self.pause_offset = 0.0 + self.paused_at = None + logger.info("Scheduler stopped") + + # --- Internal --- + + async def _run(self, seek_seconds: float = 0.0) -> None: + """ + Core dispatch loop. Absolute monotonic scheduling (D-07). + 5ms poll interval ensures asyncio loop stays responsive. + """ + pending = [e for e in self._events if e.timestamp >= seek_seconds] + logger.debug("Dispatch loop: %d events pending after seek=%.2fs", len(pending), seek_seconds) + + for event in pending: + # Wait until fire time, polling every 5ms + while True: + if self.paused_at is not None: + # Paused: just poll and wait + await asyncio.sleep(POLL_INTERVAL) + continue + # Absolute fire time — recalculated fresh each iteration (D-07) + fire_at = self.start_time + self.pause_offset + event.timestamp + remaining = fire_at - time.monotonic() + if remaining <= 0: + break + await asyncio.sleep(min(remaining, POLL_INTERVAL)) + + self._dispatch(event) + + logger.info("Scheduler: all events dispatched") + + def _dispatch(self, event: ChoreoEvent) -> None: + """Build protocol command and send via transport (CHR-05).""" + cmd: dict = { + "zone": event.zone, + "animation": event.animation, + "params": event.params, + } + self._transport.send_command(cmd) + if self._on_event: + self._on_event(event) + logger.debug( + "Dispatched: zone=%s animation=%s at position=%.3fs", + event.zone, + event.animation, + self.current_position, + )