import os import signal import subprocess import threading import json import time from datetime import datetime from pathlib import Path from zoneinfo import ZoneInfo import requests from flask import Flask, render_template_string, send_from_directory, abort, Response, request, stream_with_context # --------------------------------------------------------------------------- # Configuration # --------------------------------------------------------------------------- ICECAST_STATUS_URL = os.environ.get( "ICECAST_STATUS_URL", "http://icecast:8000/status-json.xsl", ) ICECAST_STREAM_URL = os.environ.get( "ICECAST_STREAM_URL", "http://icecast:8000/stream.mp3", ) ICECAST_PUBLIC_URL = os.environ.get( "ICECAST_PUBLIC_URL", "https://audiostream.diereuthers.de/stream.mp3", ) STREAM_MOUNTPOINT = os.environ.get( "STREAM_MOUNTPOINT", "/stream.mp3", ) RECORDINGS_DIR = Path( os.environ.get("RECORDINGS_DIR", "/recordings") ) FINISHED_DIR = Path( os.environ.get("FINISHED_DIR", "/finished") ) WEB_PORT = int( os.environ.get("WEB_PORT", "8080") ) TIMEZONE = ZoneInfo( os.environ.get("TIMEZONE", "Europe/Berlin") ) START_DELAY = float( os.environ.get("START_DELAY", "3") ) STOP_DELAY = float( os.environ.get("STOP_DELAY", "10") ) POLL_INTERVAL = float( os.environ.get("POLL_INTERVAL", "2") ) NORMALIZE = os.environ.get("NORMALIZE", "true").lower() not in ("0", "false", "no") NORMALIZE_TARGET = os.environ.get("NORMALIZE_TARGET", "-16") NORMALIZE_TRUE_PEAK = os.environ.get("NORMALIZE_TRUE_PEAK", "-1.5") NORMALIZE_LRA = os.environ.get("NORMALIZE_LRA", "11") # --------------------------------------------------------------------------- # Application state # --------------------------------------------------------------------------- app = Flask(__name__) state_lock = threading.Lock() ffmpeg_process = None recording_filename = None recording_started = None normalizing_filename = None source_seen_since = None source_missing_since = None last_status_ok = False last_status_check = None last_status_error = None SSE_CLIENTS = set() SSE_LOCK = threading.Lock() # --------------------------------------------------------------------------- # Helpers # --------------------------------------------------------------------------- def now(): return datetime.now(TIMEZONE) def format_datetime(timestamp): if not timestamp: return "–" return datetime.fromtimestamp( timestamp, TIMEZONE ).strftime("%d.%m.%Y %H:%M:%S") def format_size(size): if size < 1024: return f"{size} B" if size < 1024 ** 2: return f"{size / 1024:.1f} KB" if size < 1024 ** 3: return f"{size / 1024 ** 2:.1f} MB" return f"{size / 1024 ** 3:.2f} GB" def format_duration(seconds): if seconds is None: return "–" seconds = int(seconds) hours = seconds // 3600 minutes = (seconds % 3600) // 60 seconds = seconds % 60 if hours: return f"{hours}:{minutes:02d}:{seconds:02d}" return f"{minutes}:{seconds:02d}" def read_metadata(path): """Read common ID3/FFmpeg metadata from an MP3.""" try: result = subprocess.run( ["ffprobe", "-v", "error", "-show_entries", "format_tags=date,artist,title,album", "-of", "json", str(path)], capture_output=True, text=True, timeout=10, ) if result.returncode != 0: return {} data = json.loads(result.stdout) tags = data.get("format", {}).get("tags", {}) or {} return {str(k).lower(): str(v) for k, v in tags.items() if v is not None} except Exception: return {} def metadata_datetime(path): tags = read_metadata(path) value = tags.get("date", "") if value: try: dt = datetime.fromisoformat(value.replace("Z", "+00:00")) if dt.tzinfo is None: dt = dt.replace(tzinfo=TIMEZONE) return dt.astimezone(TIMEZONE).strftime("%d.%m.%Y %H:%M") except ValueError: pass return format_datetime(path.stat().st_mtime) def metadata_sort_key(path): tags = read_metadata(path) value = tags.get("date", "") if value: try: dt = datetime.fromisoformat(value.replace("Z", "+00:00")) if dt.tzinfo is None: dt = dt.replace(tzinfo=TIMEZONE) return dt.timestamp() except ValueError: pass return path.stat().st_mtime def get_duration(filename, directory=RECORDINGS_DIR): path = directory / filename try: result = subprocess.run( [ "ffprobe", "-v", "error", "-show_entries", "format=duration", "-of", "default=noprint_wrappers=1:nokey=1", str(path), ], capture_output=True, text=True, timeout=5, ) if result.returncode != 0: return None return float(result.stdout.strip()) except Exception: return None def is_stream_active(data): """ Detect the configured Icecast mountpoint. Icecast normally returns a dict for one source and can return a list when multiple sources exist. """ try: source = data["icestats"]["source"] except (KeyError, TypeError): return False if not source: return False if isinstance(source, dict): sources = [source] elif isinstance(source, list): sources = source else: return False for item in sources: if not isinstance(item, dict): continue listenurl = item.get("listenurl", "") # Icecast returns the public listen URL, so compare the path. if listenurl: try: from urllib.parse import urlparse path = urlparse(listenurl).path if path == STREAM_MOUNTPOINT: return True except Exception: pass return False def query_icecast(): global last_status_ok global last_status_check global last_status_error try: response = requests.get( ICECAST_STATUS_URL, timeout=3, ) response.raise_for_status() data = response.json() with state_lock: last_status_ok = True last_status_check = time.time() last_status_error = None return is_stream_active(data) except Exception as exc: with state_lock: last_status_ok = False last_status_check = time.time() last_status_error = str(exc) print( f"[monitor] Icecast status check failed: {exc}", flush=True, ) # Important: # A temporary status failure must NOT be interpreted as # "stream stopped". return None def broadcast_event(event, filename=None): payload = json.dumps({ "event": event, "filename": filename, }, ensure_ascii=False) dead = [] with SSE_LOCK: for client_queue in list(SSE_CLIENTS): try: client_queue.put_nowait(payload) except Exception: dead.append(client_queue) for client_queue in dead: SSE_CLIENTS.discard(client_queue) def sse_stream(): import queue client_queue = queue.Queue() with SSE_LOCK: SSE_CLIENTS.add(client_queue) try: yield ": connected\n\n" while True: try: payload = client_queue.get(timeout=20) yield f"data: {payload}\n\n" except queue.Empty: yield ": keepalive\n\n" finally: with SSE_LOCK: SSE_CLIENTS.discard(client_queue) # --------------------------------------------------------------------------- # Recording # --------------------------------------------------------------------------- def start_recording(): global ffmpeg_process global recording_filename global recording_started timestamp = now().strftime("%Y-%m-%d_%H-%M-%S") filename = f"{timestamp}.mp3" output = RECORDINGS_DIR / filename RECORDINGS_DIR.mkdir( parents=True, exist_ok=True, ) print( f"[recorder] Starting recording: {filename}", flush=True, ) ffmpeg_process = subprocess.Popen( [ "ffmpeg", "-hide_banner", "-loglevel", "warning", "-i", ICECAST_STREAM_URL, "-c:a", "copy", "-y", str(output), ] ) recording_filename = filename recording_started = time.time() broadcast_event("recording_started", filename) def normalize_recording(filename): """Normalize a completed MP3 in place in RECORDINGS_DIR.""" source = RECORDINGS_DIR / filename if not source.exists(): return False if not NORMALIZE: print(f"[normalize] disabled: {filename}", flush=True) return True temp = RECORDINGS_DIR / (filename + ".normalizing.mp3") cmd = [ "ffmpeg", "-hide_banner", "-loglevel", "warning", "-y", "-i", str(source), "-af", ( f"loudnorm=I={NORMALIZE_TARGET}:" f"TP={NORMALIZE_TRUE_PEAK}:" f"LRA={NORMALIZE_LRA}" ), "-c:a", "libmp3lame", "-b:a", "128k", str(temp), ] print(f"[normalize] starting: {filename}", flush=True) try: result = subprocess.run( cmd, capture_output=True, text=True, timeout=1800, ) if result.returncode != 0: print( f"[normalize] ffmpeg failed ({result.returncode}): " f"{result.stderr[-2000:]}", flush=True, ) temp.unlink(missing_ok=True) return False temp.replace(source) print(f"[normalize] completed: {filename}", flush=True) return True except Exception as exc: print(f"[normalize] error: {exc}", flush=True) temp.unlink(missing_ok=True) return False def stop_recording(): global ffmpeg_process global recording_filename global recording_started global normalizing_filename if ffmpeg_process is None: return filename = recording_filename print( f"[recorder] Stopping recording: {filename}", flush=True, ) try: ffmpeg_process.send_signal(signal.SIGINT) ffmpeg_process.wait(timeout=15) except subprocess.TimeoutExpired: print( "[recorder] FFmpeg did not stop cleanly, terminating", flush=True, ) ffmpeg_process.terminate() try: ffmpeg_process.wait(timeout=5) except subprocess.TimeoutExpired: ffmpeg_process.kill() except Exception as exc: print( f"[recorder] Error stopping FFmpeg: {exc}", flush=True, ) ffmpeg_process = None recording_filename = None recording_started = None # The live recording is now complete. Normalize it in RECORDINGS_DIR. if filename: normalizing_filename = filename if NORMALIZE else None try: normalize_recording(filename) finally: normalizing_filename = None broadcast_event("recording_stopped", filename) # --------------------------------------------------------------------------- # Monitoring loop # --------------------------------------------------------------------------- def recorder_loop(): global ffmpeg_process global recording_filename global recording_started global source_seen_since global source_missing_since print("[monitor] Recorder monitor started", flush=True) while True: active = query_icecast() now_ts = time.time() with state_lock: process = ffmpeg_process current_filename = recording_filename if active is None: time.sleep(POLL_INTERVAL) continue # FFmpeg ending after the Icecast source disappeared is a normal # end of the recording. Finalize it so normalization is performed. if process is not None and process.poll() is not None: if not active: print( f"[recorder] FFmpeg ended because the Icecast stream stopped: {current_filename}", flush=True, ) stop_recording() source_seen_since = None source_missing_since = None time.sleep(POLL_INTERVAL) continue print( "[recorder] FFmpeg exited unexpectedly while Icecast is active", flush=True, ) ffmpeg_process = None recording_filename = None recording_started = None source_seen_since = now_ts time.sleep(POLL_INTERVAL) continue if active: source_missing_since = None if process is None: if source_seen_since is None: source_seen_since = now_ts elif now_ts - source_seen_since >= START_DELAY: start_recording() source_seen_since = None else: source_seen_since = None else: source_seen_since = None if process is not None: if source_missing_since is None: source_missing_since = now_ts elif now_ts - source_missing_since >= STOP_DELAY: stop_recording() source_missing_since = None else: source_missing_since = None time.sleep(POLL_INTERVAL) # --------------------------------------------------------------------------- # Web UI # --------------------------------------------------------------------------- HTML = """ Event Audio – Aufnahmen

EVENT AUDIO – Aufnahmen

Audiostream-Aufzeichnungen
{% if normalizing_filename %}
● AUFNAHME WIRD NORMALISIERT
{{ normalizing_filename }}
Lautstärke wird automatisch normalisiert …
{% elif recording %}
● AUFNAHME LÄUFT
🔴 Live-Stream
Direkt vom Icecast-Stream
▶ Aktuelle Aufnahme
Zeitversetzt · gestartet {{ recording.started }} · aktuell {{ recording.size }}
{% elif status_error %}
⚠ Icecast-Status momentan nicht erreichbar
{% else %}
● keine laufende Aufnahme
{% endif %}

Fertige Aufnahmen

{% if recordings %} {% for item in recordings %} {% endfor %}
Zeitpunkt Redner Beschreibung Veranstaltung Dauer Wiedergabe
{{ item.datetime }} {{ item.artist or "–" }} {{ item.title or "–" }} {{ item.album or "–" }} {{ item.duration }} Download
{% else %}
Noch keine fertigen Aufnahmen vorhanden.
{% endif %}
""" def render_index(): global ffmpeg_process global recording_filename global recording_started global normalizing_filename files = [] FINISHED_DIR.mkdir(parents=True, exist_ok=True) for path in FINISHED_DIR.glob("*.mp3"): try: stat = path.stat() tags = read_metadata(path) files.append({ "filename": path.name, "datetime": metadata_datetime(path), "artist": tags.get("artist", ""), "title": tags.get("title", ""), "album": tags.get("album", ""), "size": format_size(stat.st_size), "duration": format_duration(get_duration(path.name, FINISHED_DIR)), "sort_date": metadata_sort_key(path), }) except FileNotFoundError: continue files.sort(key=lambda item: item["sort_date"], reverse=True) with state_lock: process = ffmpeg_process current_filename = recording_filename current_started = recording_started status_error = not last_status_ok recording = None if process is not None and process.poll() is None: current_size = "–" if current_filename: current_path = RECORDINGS_DIR / current_filename try: current_size = format_size( current_path.stat().st_size ) except FileNotFoundError: pass recording = { "filename": current_filename, "started": format_datetime(current_started), "size": current_size, } return render_template_string( HTML, recordings=files, recording=recording, normalizing_filename=normalizing_filename, status_error=status_error, icecast_public_url=ICECAST_PUBLIC_URL, ) @app.route("/") def index(): return render_index() @app.route("/api/live") def api_live(): with state_lock: process = ffmpeg_process current_filename = recording_filename current_started = recording_started status_error = not last_status_ok recording = None if process is not None and process.poll() is None and current_filename: current_size = "–" try: current_size = format_size((RECORDINGS_DIR / current_filename).stat().st_size) except FileNotFoundError: pass recording = { "filename": current_filename, "started": format_datetime(current_started), "size": current_size, } return render_template_string( """{% if normalizing_filename %}
● AUFNAHME WIRD NORMALISIERT
{{ normalizing_filename }}
Lautstärke wird automatisch normalisiert …
{% elif recording %}
● AUFNAHME LÄUFT
🔴 Live-Stream
Direkt vom Icecast-Stream
▶ Aktuelle Aufnahme
Zeitversetzt · gestartet {{ recording.started }} · aktuell {{ recording.size }}
{% elif status_error %}
⚠ Icecast-Status momentan nicht erreichbar
{% else %}
● keine laufende Aufnahme
{% endif %}""", recording=recording, normalizing_filename=normalizing_filename, status_error=status_error, icecast_public_url=ICECAST_PUBLIC_URL, ) @app.route("/api/events") def api_events(): return Response( stream_with_context(sse_stream()), mimetype="text/event-stream", headers={ "Cache-Control": "no-cache, no-store, must-revalidate", "X-Accel-Buffering": "no", "Connection": "keep-alive", }, ) @app.route("/api/status") def api_status(): return render_index() @app.route("/finished/") def finished(filename): return send_from_directory(FINISHED_DIR, filename, as_attachment=False) @app.route("/recordings/") def recording(filename): path = RECORDINGS_DIR / filename if not path.is_file(): abort(404) # For the active recording, keep the HTTP connection open and wait for # additional MP3 data appended by FFmpeg. A normal static-file response # would expose the current Content-Length and the browser would stop # when it reaches the file size that existed when playback started. if request.args.get("live") == "1" and recording_filename == filename and ffmpeg_process is not None: def generate(): position = 0 while True: try: size = path.stat().st_size except FileNotFoundError: break if position < size: try: with path.open("rb") as f: f.seek(position) while True: chunk = f.read(64 * 1024) if not chunk: break position += len(chunk) yield chunk except (FileNotFoundError, OSError): break continue # The recording is still active: wait for FFmpeg to append # more data instead of closing the HTTP response. if recording_filename == filename and ffmpeg_process is not None: time.sleep(0.5) continue # Recording ended; all remaining bytes have been delivered. break return Response( generate(), mimetype="audio/mpeg", headers={ "Cache-Control": "no-cache, no-store, must-revalidate", "Accept-Ranges": "none", }, ) return send_from_directory( RECORDINGS_DIR, filename, as_attachment=False, ) # --------------------------------------------------------------------------- # Start monitor # --------------------------------------------------------------------------- monitor_thread = threading.Thread( target=recorder_loop, daemon=True, ) monitor_thread.start() if __name__ == "__main__": app.run( host="0.0.0.0", port=WEB_PORT, )