feat: #1 telemetry record/playback with auto-race detection
- TelemetryRecorder: manual + auto mode (IsRaceOn edge detection) - TelemetryPlayer: timed packet replay into _handle_packet pipeline - .tsr format: [uint32 elapsed_ms][uint32 len][bytes] - Settings UI: start/stop, auto toggle, recording list, play button - API: /record/*, /playback/*, /recordings
This commit is contained in:
@@ -422,3 +422,50 @@ async def api_update_forward_targets(data: dict[str, Any]):
|
||||
save_config(cfg)
|
||||
telemetry_listener.forwarder.reload_targets()
|
||||
return {"ok": True, "active": telemetry_listener.forwarder.active}
|
||||
|
||||
|
||||
# ---- Record & Playback ----
|
||||
from server.recorder import TelemetryRecorder, TelemetryPlayer
|
||||
|
||||
_recorder = None
|
||||
_player = None
|
||||
|
||||
def _get_recorder():
|
||||
global _recorder
|
||||
if _recorder is None: _recorder = TelemetryRecorder(telemetry_listener)
|
||||
return _recorder
|
||||
|
||||
def _get_player():
|
||||
global _player
|
||||
if _player is None: _player = TelemetryPlayer(telemetry_listener)
|
||||
return _player
|
||||
|
||||
@router.post("/record/start")
|
||||
async def api_start_record():
|
||||
_get_recorder().start()
|
||||
return {"ok": True, "recording": True}
|
||||
|
||||
@router.post("/record/stop")
|
||||
async def api_stop_record():
|
||||
_get_recorder().stop()
|
||||
return {"ok": True, "recording": False}
|
||||
|
||||
@router.get("/record/status")
|
||||
async def api_record_status():
|
||||
return {"recording": _get_recorder().is_recording, "playing": _get_player().is_playing}
|
||||
|
||||
@router.get("/recordings")
|
||||
async def api_list_recordings():
|
||||
return TelemetryPlayer.list_recordings()
|
||||
|
||||
@router.post("/playback/start")
|
||||
async def api_start_playback(data: dict):
|
||||
filename = data.get("filename", "")
|
||||
if not filename: raise HTTPException(400, "filename required")
|
||||
_get_player().start(filename)
|
||||
return {"ok": True, "playing": True}
|
||||
|
||||
@router.post("/playback/stop")
|
||||
async def api_stop_playback():
|
||||
_get_player().stop()
|
||||
return {"ok": True, "playing": False}
|
||||
|
||||
@@ -0,0 +1,134 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import os
|
||||
import struct
|
||||
import time
|
||||
from pathlib import Path
|
||||
|
||||
from utils.logger import get_logger
|
||||
|
||||
logger = get_logger(__name__)
|
||||
|
||||
RECORDINGS_DIR = Path(__file__).resolve().parent.parent / "data" / "recordings"
|
||||
RECORDINGS_DIR.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
|
||||
class TelemetryRecorder:
|
||||
def __init__(self, listener):
|
||||
self._listener = listener
|
||||
self._file = None
|
||||
self._recording = False
|
||||
self._auto_mode = False
|
||||
self._start_time = 0.0
|
||||
self._last_race_on = False
|
||||
|
||||
@property
|
||||
def is_recording(self) -> bool:
|
||||
return self._recording
|
||||
|
||||
@property
|
||||
def current_file(self) -> str | None:
|
||||
return str(self._file.name) if self._file else None
|
||||
|
||||
def start(self):
|
||||
if self._recording:
|
||||
return
|
||||
ts = time.strftime("%Y%m%d_%H%M%S")
|
||||
self._file = open(RECORDINGS_DIR / f"rec_{ts}.tsr", "wb")
|
||||
self._recording = True
|
||||
self._start_time = time.time()
|
||||
logger.info("Recording started: %s", self._file.name)
|
||||
|
||||
def stop(self):
|
||||
if not self._recording:
|
||||
return
|
||||
self._recording = False
|
||||
if self._file:
|
||||
self._file.close()
|
||||
self._file = None
|
||||
logger.info("Recording stopped")
|
||||
|
||||
def set_auto(self, enabled: bool):
|
||||
self._auto_mode = enabled
|
||||
self._last_race_on = False
|
||||
if not enabled:
|
||||
self.stop()
|
||||
|
||||
def write_packet(self, data: bytes):
|
||||
if self._auto_mode:
|
||||
try:
|
||||
is_race = struct.unpack_from("<i", data, 0)[0] == 1
|
||||
except Exception:
|
||||
is_race = True
|
||||
|
||||
if is_race and not self._last_race_on:
|
||||
self.start()
|
||||
elif not is_race and self._last_race_on and self._recording:
|
||||
self.stop()
|
||||
self._last_race_on = is_race
|
||||
|
||||
if not self._recording or not self._file:
|
||||
return
|
||||
elapsed_ms = int((time.time() - self._start_time) * 1000)
|
||||
self._file.write(struct.pack("<II", elapsed_ms, len(data)))
|
||||
self._file.write(data)
|
||||
|
||||
|
||||
class TelemetryPlayer:
|
||||
def __init__(self, listener):
|
||||
self._listener = listener
|
||||
self._task: asyncio.Task | None = None
|
||||
self._playing = False
|
||||
|
||||
@property
|
||||
def is_playing(self) -> bool:
|
||||
return self._playing
|
||||
|
||||
def start(self, filename: str):
|
||||
if self._playing:
|
||||
return
|
||||
path = RECORDINGS_DIR / filename
|
||||
if not path.exists():
|
||||
logger.error("Recording not found: %s", filename)
|
||||
return
|
||||
loop = asyncio.get_event_loop()
|
||||
self._task = loop.create_task(self._play_loop(path))
|
||||
self._playing = True
|
||||
logger.info("Playback started: %s", filename)
|
||||
|
||||
def stop(self):
|
||||
self._playing = False
|
||||
if self._task and not self._task.done():
|
||||
self._task.cancel()
|
||||
self._task = None
|
||||
logger.info("Playback stopped")
|
||||
|
||||
async def _play_loop(self, path: Path):
|
||||
try:
|
||||
data = path.read_bytes()
|
||||
offset = 0
|
||||
play_start = time.time()
|
||||
while offset < len(data) and self._playing:
|
||||
elapsed_ms, length = struct.unpack_from("<II", data, offset)
|
||||
offset += 8
|
||||
packet = data[offset:offset + length]
|
||||
offset += length
|
||||
|
||||
target_time = play_start + elapsed_ms / 1000.0
|
||||
delay = target_time - time.time()
|
||||
if delay > 0:
|
||||
await asyncio.sleep(min(delay, 0.5))
|
||||
|
||||
self._listener._handle_packet(packet, ("127.0.0.1", 0))
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
except Exception as e:
|
||||
logger.error("Playback error: %s", e)
|
||||
finally:
|
||||
self._playing = False
|
||||
|
||||
@staticmethod
|
||||
def list_recordings() -> list[dict]:
|
||||
return [{"name": f.name, "size": f.stat().st_size, "time": f.stat().st_mtime}
|
||||
for f in sorted(RECORDINGS_DIR.glob("*.tsr"), key=lambda x: x.stat().st_mtime, reverse=True)]
|
||||
@@ -151,6 +151,10 @@ class TelemetryListener:
|
||||
|
||||
self.forwarder.forward(data)
|
||||
|
||||
from server.api import _recorder
|
||||
if _recorder and _recorder.is_recording:
|
||||
_recorder.write_packet(data)
|
||||
|
||||
td = None
|
||||
if self._parser:
|
||||
td = self._parser.parse(data, addr)
|
||||
|
||||
Reference in New Issue
Block a user