import asyncio import datetime import os import re import signal import subprocess import threading import time from pathlib import Path from typing import Optional from models import Recording, get_session, LogEntry def detect_stream_type(url: str) -> str: url_lower = url.lower().split("?")[0] if url_lower.startswith(("rtmp://", "rtmps://")): return "rtmp" if url_lower.endswith(".m3u8"): return "hls" if url_lower.endswith(".mp3"): return "mp3" if url_lower.endswith(".aac"): return "aac" if url_lower.endswith(".ogg"): return "ogg" if url_lower.endswith(".flac"): return "flac" return "http" def default_extension(stream_type: str) -> str: return {"mp3": "mp3", "aac": "aac", "ogg": "ogg", "flac": "flac"}.get(stream_type, "mp4") def safe_name(name: str) -> str: return re.sub(r"[^A-Za-z0-9_-]", "_", name).strip("_") def format_elapsed(seconds: int) -> str: d, rem = divmod(seconds, 86400) h, rem = divmod(rem, 3600) m, s = divmod(rem, 60) if d > 0: return f"{d}d {h}h {m}m" if h > 0: return f"{h}h {m}m {s}s" return f"{m}m {s}s" def format_size(size_bytes: int) -> str: if size_bytes >= 1073741824: return f"{size_bytes / 1073741824:.1f} GB" if size_bytes >= 1048576: return f"{size_bytes / 1048576:.1f} MB" if size_bytes >= 1024: return f"{size_bytes / 1024:.1f} KB" return f"{size_bytes} B" class RecordingManager: def __init__(self, config): self.config = config self._processes: dict[int, subprocess.Popen] = {} self._stop_events: dict[int, threading.Event] = {} self._threads: dict[int, threading.Thread] = {} self._metadata_monitor = None self._ntfy = None self._plik = None self._rec_options: dict[int, dict] = {} def _get_metadata_monitor(self): if self._metadata_monitor is None: from plugins.euer_radio import MetadataMonitor self._metadata_monitor = MetadataMonitor( log_fn=self._log, stop_fn=self.stop_recording, ) return self._metadata_monitor def _get_ntfy(self): if self._ntfy is None: from plugins.ntfy import NtfyNotifier self._ntfy = NtfyNotifier(self.config) return self._ntfy def _get_plik(self): if self._plik is None: from plugins.plik import PlikUploader self._plik = PlikUploader( self.config, log_fn=self._log, notify_fn=lambda event, title, msg, ntfy_en: self._get_ntfy().notify(event, title, msg, ntfy_en), ) return self._plik def _notify(self, event: str, title: str, message: str, rec_id: int = None, priority: str = "default"): ntfy_enabled = False if rec_id and rec_id in self._rec_options: ntfy_enabled = self._rec_options[rec_id].get("ntfy_enabled", False) desc = self._rec_options[rec_id].get("description", "") if desc: message = f"{message}\n{desc}" self._get_ntfy().notify(event, title, message, job_ntfy_enabled=ntfy_enabled, priority=priority) def _log(self, level: str, message: str, job_name: str = None): session = get_session() try: entry = LogEntry(level=level, message=message, job_name=job_name) session.add(entry) session.commit() finally: session.close() def start_recording( self, job_id: Optional[int], name: str, stream_url: str, stream_type: str = "auto", output_format: str = "", max_duration: Optional[int] = None, extra_ffmpeg_args: str = "", segment_duration: Optional[int] = None, is_scheduled: bool = False, description: str = "", metadata_monitor_enabled: bool = False, metadata_pattern: str = "", metadata_grace_period: int = 300, metadata_poll_interval: int = 30, ntfy_enabled: bool = False, plik_enabled: bool = False, delete_after_upload: bool = False, ) -> dict: if stream_type == "auto": stream_type = detect_stream_type(stream_url) ext = output_format if output_format and output_format.lower() != "auto" else default_extension(stream_type) sname = safe_name(name) timestamp = datetime.datetime.now().strftime("%Y%m%d_%H%M%S") download_path = self.config.get("recording", "download_path", default="/app/data/recordings") output_dir = os.path.join(download_path, sname) os.makedirs(output_dir, exist_ok=True) segment_mode = segment_duration is not None and segment_duration > 0 if segment_mode: output_file = os.path.join(output_dir, f"{sname}_{timestamp}_seg%03d.{ext}") else: output_file = os.path.join(output_dir, f"{sname}_{timestamp}.{ext}") scheduled_stop = None if max_duration and max_duration > 0: scheduled_stop = datetime.datetime.utcnow() + datetime.timedelta(seconds=max_duration) session = get_session() try: recording = Recording( job_id=job_id, job_name=name, stream_url=stream_url, stream_type=stream_type, output_dir=output_dir, output_file=output_file, status="starting", scheduled_stop=scheduled_stop, segment_mode=segment_mode, is_scheduled=is_scheduled, ) session.add(recording) session.commit() rec_id = recording.id finally: session.close() stop_event = threading.Event() self._stop_events[rec_id] = stop_event thread = threading.Thread( target=self._recording_loop, args=(rec_id, stream_url, stream_type, output_file, output_dir, extra_ffmpeg_args, segment_mode, segment_duration, stop_event, name), daemon=True, ) self._threads[rec_id] = thread thread.start() self._rec_options[rec_id] = { "ntfy_enabled": ntfy_enabled, "plik_enabled": plik_enabled, "delete_after_upload": delete_after_upload, "segment_mode": segment_mode, "description": description, } self._log("INFO", f"Aufnahme gestartet: {name}", name) self._notify("start", "Aufnahme gestartet", name, rec_id) if metadata_monitor_enabled and metadata_pattern: self._get_metadata_monitor().start_monitoring( rec_id=rec_id, stream_url=stream_url, pattern=metadata_pattern, grace_period=metadata_grace_period, poll_interval=metadata_poll_interval, job_name=name, ) return {"recording_id": rec_id, "output_dir": output_dir, "output_file": output_file} def _recording_loop( self, rec_id: int, stream_url: str, stream_type: str, output_file: str, output_dir: str, extra_ffmpeg_args: str, segment_mode: bool, segment_duration: Optional[int], stop_event: threading.Event, job_name: str, ): max_retries = self.config.get("recording", "max_retries", default=5) retry_delay = self.config.get("recording", "retry_delay", default=5) retry = 0 while not stop_event.is_set(): cmd = self._build_ffmpeg_cmd( stream_url, stream_type, output_file, extra_ffmpeg_args, segment_mode, segment_duration, ) try: log_path = os.path.join(output_dir, f".ffmpeg_{rec_id}.log") log_fh = open(log_path, "a") process = subprocess.Popen( cmd, stdout=log_fh, stderr=log_fh, stdin=subprocess.DEVNULL, ) except Exception as e: self._update_recording(rec_id, status="error", error_message=str(e)) self._log("ERROR", f"ffmpeg konnte nicht gestartet werden: {e}", job_name) self._notify("error", "Aufnahme fehlgeschlagen", f"{job_name}: {e}", rec_id, priority="high") return self._processes[rec_id] = process self._update_recording(rec_id, status="running", pid=process.pid) process.wait() retcode = process.returncode log_fh.close() self._processes.pop(rec_id, None) if stop_event.is_set() or retcode == 0: status = "completed" if retcode == 0 else "stopped" status_de = "abgeschlossen" if status == "completed" else "gestoppt" self._update_recording(rec_id, status=status) self._log("INFO", f"Aufnahme {status_de}: {job_name}", job_name) self._notify("stop", f"Aufnahme {status_de}", job_name, rec_id) self._handle_post_recording(rec_id, output_dir, output_file, job_name, segment_mode) return if retcode in (-2, 130, 143, 255): self._update_recording(rec_id, status="stopped") self._log("INFO", f"Aufnahme gestoppt: {job_name}", job_name) self._notify("stop", "Aufnahme gestoppt", job_name, rec_id) self._handle_post_recording(rec_id, output_dir, output_file, job_name, segment_mode) return retry += 1 if max_retries > 0 and retry >= max_retries: self._update_recording(rec_id, status="error", error_message=f"Max retries ({max_retries}) reached, exit code {retcode}") self._log("ERROR", f"Max. Versuche ({max_retries}) erreicht: {job_name}", job_name) self._notify("error", "Aufnahme fehlgeschlagen", f"{job_name}: Max. Versuche erreicht", rec_id, priority="high") return self._log("WARN", f"Verbindung verloren (Code: {retcode}), Reconnect in {retry_delay}s... [{job_name}]", job_name) self._notify("error", "Stream abgerissen", f"{job_name}: Reconnect in {retry_delay}s...", rec_id) if stop_event.wait(timeout=retry_delay): self._update_recording(rec_id, status="stopped") return sname = safe_name(job_name) ts = datetime.datetime.now().strftime("%Y%m%d_%H%M%S") ext = Path(output_file).suffix.lstrip(".") if segment_mode: output_file = os.path.join(output_dir, f"{sname}_{ts}_seg%03d.{ext}") else: output_file = os.path.join(output_dir, f"{sname}_{ts}.{ext}") self._update_recording(rec_id, status="stopped") def _build_ffmpeg_cmd(self, url, stream_type, output_file, extra_args, segment_mode, segment_duration): cmd = ["ffmpeg", "-nostdin", "-y", "-hide_banner", "-loglevel", "warning"] if stream_type == "hls": cmd += ["-i", url, "-c", "copy", "-bsf:a", "aac_adtstoasc"] else: cmd += ["-i", url, "-c", "copy"] if segment_mode and segment_duration: cmd += ["-f", "segment", "-segment_time", str(segment_duration), "-reset_timestamps", "1"] if extra_args: cmd += extra_args.split() cmd.append(output_file) return cmd def stop_recording(self, rec_id: int) -> bool: if self._metadata_monitor: self._metadata_monitor.stop_monitoring(rec_id) stop_event = self._stop_events.get(rec_id) if stop_event: stop_event.set() process = self._processes.get(rec_id) if process: try: process.send_signal(signal.SIGINT) try: process.wait(timeout=10) except subprocess.TimeoutExpired: process.kill() except (ProcessLookupError, OSError): pass self._update_recording(rec_id, status="stopped") self._log("INFO", f"Aufnahme gestoppt (ID: {rec_id})") self._cleanup(rec_id) return True def extend_recording(self, rec_id: int, extra_seconds: int) -> bool: session = get_session() try: rec = session.query(Recording).filter_by(id=rec_id).first() if not rec or rec.status != "running": return False if rec.scheduled_stop: rec.scheduled_stop = rec.scheduled_stop + datetime.timedelta(seconds=extra_seconds) else: rec.scheduled_stop = datetime.datetime.utcnow() + datetime.timedelta(seconds=extra_seconds) session.commit() self._log("INFO", f"Aufnahme verlängert um {format_elapsed(extra_seconds)}: {rec.job_name}", rec.job_name) return True finally: session.close() def remove_duration_limit(self, rec_id: int) -> bool: session = get_session() try: rec = session.query(Recording).filter_by(id=rec_id).first() if not rec or rec.status != "running": return False rec.scheduled_stop = None session.commit() self._log("INFO", f"Zeitlimit entfernt: {rec.job_name}", rec.job_name) return True finally: session.close() def get_status(self, rec_id: int) -> Optional[dict]: session = get_session() try: rec = session.query(Recording).filter_by(id=rec_id).first() if not rec: return None return self._recording_to_dict(rec) finally: session.close() def get_all_active(self) -> list[dict]: session = get_session() try: recs = session.query(Recording).filter(Recording.status.in_(["running", "starting"])).all() return [self._recording_to_dict(r) for r in recs] finally: session.close() def get_all_recordings(self, limit: int = 50) -> list[dict]: session = get_session() try: recs = session.query(Recording).order_by(Recording.id.desc()).limit(limit).all() return [self._recording_to_dict(r) for r in recs] finally: session.close() def check_scheduled_stops(self): session = get_session() try: now = datetime.datetime.utcnow() recs = session.query(Recording).filter( Recording.status == "running", Recording.scheduled_stop.isnot(None), Recording.scheduled_stop <= now, ).all() rec_ids = [r.id for r in recs] finally: session.close() for rec_id in rec_ids: self._log("INFO", f"Geplanter Stopp erreicht (ID: {rec_id})") self.stop_recording(rec_id) def cleanup_stale(self): session = get_session() try: recs = session.query(Recording).filter(Recording.status.in_(["running", "starting"])).all() for rec in recs: if rec.pid: try: os.kill(rec.pid, 0) except (ProcessLookupError, OSError): rec.status = "error" rec.error_message = "Process vanished" rec.stopped_at = datetime.datetime.utcnow() session.commit() finally: session.close() def _recording_to_dict(self, rec: Recording) -> dict: elapsed = 0 if rec.started_at: end = rec.stopped_at or datetime.datetime.utcnow() elapsed = int((end - rec.started_at).total_seconds()) total_size = 0 segment_count = 0 if rec.output_dir and os.path.isdir(rec.output_dir): for f in os.listdir(rec.output_dir): if f.startswith("."): continue fp = os.path.join(rec.output_dir, f) if os.path.isfile(fp): total_size += os.path.getsize(fp) segment_count += 1 elif rec.output_file and os.path.isfile(rec.output_file): total_size = os.path.getsize(rec.output_file) remaining = None if rec.scheduled_stop and rec.status == "running": diff = (rec.scheduled_stop - datetime.datetime.utcnow()).total_seconds() remaining = max(0, int(diff)) return { "id": rec.id, "job_id": rec.job_id, "job_name": rec.job_name, "stream_url": rec.stream_url, "stream_type": rec.stream_type, "output_dir": rec.output_dir, "output_file": rec.output_file, "pid": rec.pid, "status": rec.status, "started_at": rec.started_at.isoformat() if rec.started_at else None, "stopped_at": rec.stopped_at.isoformat() if rec.stopped_at else None, "scheduled_stop": rec.scheduled_stop.isoformat() if rec.scheduled_stop else None, "segment_mode": rec.segment_mode, "is_scheduled": rec.is_scheduled, "error_message": rec.error_message, "elapsed": elapsed, "elapsed_display": format_elapsed(elapsed), "total_size": total_size, "total_size_display": format_size(total_size), "segment_count": segment_count if rec.segment_mode else None, "remaining": remaining, "remaining_display": format_elapsed(remaining) if remaining is not None else None, } def _update_recording(self, rec_id: int, **kwargs): session = get_session() try: rec = session.query(Recording).filter_by(id=rec_id).first() if rec: for k, v in kwargs.items(): setattr(rec, k, v) if kwargs.get("status") in ("stopped", "completed", "error"): rec.stopped_at = datetime.datetime.utcnow() session.commit() finally: session.close() def _handle_post_recording(self, rec_id: int, output_dir: str, output_file: str, job_name: str, segment_mode: bool): opts = self._rec_options.get(rec_id, {}) if opts.get("plik_enabled"): self._get_plik().upload_recording( output_dir=output_dir, output_file=output_file, job_name=job_name, segment_mode=segment_mode, delete_after=opts.get("delete_after_upload", False), ntfy_enabled=opts.get("ntfy_enabled", False), ) def _cleanup(self, rec_id: int): self._processes.pop(rec_id, None) self._stop_events.pop(rec_id, None) self._threads.pop(rec_id, None) self._rec_options.pop(rec_id, None)