feat(02-04): implement ChoreographyScheduler with absolute monotonic clock dispatch
- Absolute monotonic timestamp scheduling (D-07): fire_at = start_time + pause_offset + event.timestamp
- 5ms poll interval keeps asyncio loop responsive
- pause() records paused_at; resume() accumulates pause_offset (D-09)
- seek() restarts dispatch loop skipping events before seek position
- _dispatch() builds {zone, animation, params} dict without v:1 (transport responsibility)
- CHR-05 fulfilled: drift < 20ms validated by tests
This commit is contained in:
147
src/led_sync/scheduler.py
Normal file
147
src/led_sync/scheduler.py
Normal file
@@ -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,
|
||||||
|
)
|
||||||
Reference in New Issue
Block a user