import base64 import contextlib import datetime import os import secrets import sys from fastapi import FastAPI, HTTPException, Request, Query from fastapi.middleware.cors import CORSMiddleware from fastapi.responses import HTMLResponse, JSONResponse from fastapi.staticfiles import StaticFiles from pydantic import BaseModel from typing import Optional from config import Config from models import init_db, get_session, Job, Station, Recording as RecordingModel, LogEntry from recorder import RecordingManager, format_elapsed from scheduler import StreamScheduler APP_VERSION = os.environ.get("APP_VERSION", "dev") DATA_DIR = os.environ.get("DATA_DIR", "/app/data") STATIC_DIR = os.environ.get("STATIC_DIR", "/app/frontend/static") os.makedirs(DATA_DIR, exist_ok=True) config = Config(os.path.join(DATA_DIR, "config.yml")) init_db(os.path.join(DATA_DIR, "stream-recorder.db")) recorder = RecordingManager(config) scheduler = StreamScheduler(recorder) @contextlib.asynccontextmanager async def lifespan(app): recorder.cleanup_stale() scheduler.start() yield scheduler.stop() app = FastAPI(title="stream-recorder", version=APP_VERSION, lifespan=lifespan) app.add_middleware( CORSMiddleware, allow_origins=["*"], allow_methods=["*"], allow_headers=["*"], ) @app.middleware("http") async def auth_middleware(request: Request, call_next): if not config.get("auth", "enabled", default=False): return await call_next(request) if request.url.path.startswith("/api/health"): return await call_next(request) auth_header = request.headers.get("Authorization", "") if auth_header.startswith("Basic "): try: decoded = base64.b64decode(auth_header[6:]).decode("utf-8") username, password = decoded.split(":", 1) expected_user = config.get("auth", "username", default="admin") expected_pass = config.get("auth", "password", default="stream-recorder") if secrets.compare_digest(username, expected_user) and secrets.compare_digest(password, expected_pass): return await call_next(request) except Exception: pass return JSONResponse( status_code=401, content={"detail": "Unauthorized"}, headers={"WWW-Authenticate": 'Basic realm="stream-recorder"'}, ) # --- Health --- @app.get("/api/health") async def health(): return {"status": "ok", "version": APP_VERSION} # --- Jobs CRUD --- class StationCreate(BaseModel): name: str stream_url: str stream_type: str = "auto" description: str = "" class StationUpdate(StationCreate): pass def _station_to_dict(station: Station) -> dict: return { "id": station.id, "name": station.name, "stream_url": station.stream_url, "stream_type": station.stream_type, "description": station.description, "created_at": station.created_at.isoformat() if station.created_at else None, "updated_at": station.updated_at.isoformat() if station.updated_at else None, } def _clean_station_data(data: StationCreate | StationUpdate) -> dict: values = data.model_dump() values["name"] = values["name"].strip() values["stream_url"] = values["stream_url"].strip() values["stream_type"] = (values.get("stream_type") or "auto").strip() or "auto" values["description"] = (values.get("description") or "").strip() if not values["name"] or not values["stream_url"]: raise HTTPException(400, "Name und URL sind Pflichtfelder") return values @app.get("/api/stations") async def list_stations(): session = get_session() try: stations = session.query(Station).order_by(Station.name.collate("NOCASE"), Station.id).all() return [_station_to_dict(s) for s in stations] finally: session.close() @app.post("/api/stations") async def create_station(data: StationCreate): values = _clean_station_data(data) session = get_session() try: station = Station(**values) session.add(station) session.commit() return {"id": station.id, "message": "Sender gespeichert"} finally: session.close() @app.get("/api/stations/{station_id}") async def get_station(station_id: int): session = get_session() try: station = session.query(Station).filter_by(id=station_id).first() if not station: raise HTTPException(404, "Sender nicht gefunden") return _station_to_dict(station) finally: session.close() @app.put("/api/stations/{station_id}") async def update_station(station_id: int, data: StationUpdate): values = _clean_station_data(data) session = get_session() try: station = session.query(Station).filter_by(id=station_id).first() if not station: raise HTTPException(404, "Sender nicht gefunden") for k, v in values.items(): setattr(station, k, v) station.updated_at = datetime.datetime.utcnow() session.commit() return {"message": "Sender aktualisiert"} finally: session.close() @app.delete("/api/stations/{station_id}") async def delete_station(station_id: int): session = get_session() try: station = session.query(Station).filter_by(id=station_id).first() if not station: raise HTTPException(404, "Sender nicht gefunden") session.delete(station) session.commit() return {"message": "Sender gelöscht"} finally: session.close() class JobCreate(BaseModel): 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 schedule_enabled: bool = False schedule_once: bool = False schedule_date: str = "" schedule_days: str = "*" schedule_start: str = "" schedule_stop: str = "" schedule_end_date: str = "" pre_buffer: int = 0 post_buffer: int = 0 description: str = "" ntfy_enabled: bool = False plik_enabled: bool = False delete_after_upload: bool = False station_id: Optional[int] = None class JobUpdate(JobCreate): pass @app.get("/api/jobs") async def list_jobs(): session = get_session() try: jobs = session.query(Job).order_by(Job.id).all() station_ids = {j.station_id for j in jobs if j.station_id} station_names = {} if station_ids: for s in session.query(Station).filter(Station.id.in_(station_ids)).all(): station_names[s.id] = s.name result = [] for j in jobs: is_recording = session.query( session.query(RecordingModel).filter_by(job_id=j.id, status="running").exists() ).scalar() result.append({ "id": j.id, "name": j.name, "stream_url": j.stream_url, "stream_type": j.stream_type, "output_format": j.output_format, "max_duration": j.max_duration, "extra_ffmpeg_args": j.extra_ffmpeg_args, "segment_duration": j.segment_duration, "schedule_enabled": j.schedule_enabled, "schedule_once": j.schedule_once, "schedule_date": j.schedule_date, "schedule_days": j.schedule_days, "schedule_start": j.schedule_start, "schedule_stop": j.schedule_stop, "schedule_end_date": j.schedule_end_date, "pre_buffer": j.pre_buffer or 0, "post_buffer": j.post_buffer or 0, "description": j.description, "ntfy_enabled": j.ntfy_enabled, "plik_enabled": j.plik_enabled, "delete_after_upload": j.delete_after_upload, "station_id": j.station_id, "station_name": station_names.get(j.station_id) if j.station_id else None, "is_recording": is_recording, "created_at": j.created_at.isoformat() if j.created_at else None, }) return result finally: session.close() @app.post("/api/jobs") async def create_job(data: JobCreate): session = get_session() try: job = Job(**data.model_dump()) session.add(job) session.commit() return {"id": job.id, "message": "Job erstellt"} finally: session.close() @app.get("/api/jobs/{job_id}") async def get_job(job_id: int): session = get_session() try: job = session.query(Job).filter_by(id=job_id).first() if not job: raise HTTPException(404, "Job nicht gefunden") return { "id": job.id, "name": job.name, "stream_url": job.stream_url, "stream_type": job.stream_type, "output_format": job.output_format, "max_duration": job.max_duration, "extra_ffmpeg_args": job.extra_ffmpeg_args, "segment_duration": job.segment_duration, "schedule_enabled": job.schedule_enabled, "schedule_once": job.schedule_once, "schedule_date": job.schedule_date, "schedule_days": job.schedule_days, "schedule_start": job.schedule_start, "schedule_stop": job.schedule_stop, "schedule_end_date": job.schedule_end_date, "pre_buffer": job.pre_buffer or 0, "post_buffer": job.post_buffer or 0, "description": job.description, "ntfy_enabled": job.ntfy_enabled, "plik_enabled": job.plik_enabled, "delete_after_upload": job.delete_after_upload, "station_id": job.station_id, } finally: session.close() @app.put("/api/jobs/{job_id}") async def update_job(job_id: int, data: JobUpdate): session = get_session() try: job = session.query(Job).filter_by(id=job_id).first() if not job: raise HTTPException(404, "Job nicht gefunden") for k, v in data.model_dump().items(): setattr(job, k, v) job.updated_at = datetime.datetime.utcnow() session.commit() return {"message": "Job aktualisiert"} finally: session.close() @app.delete("/api/jobs/{job_id}") async def delete_job(job_id: int): session = get_session() try: job = session.query(Job).filter_by(id=job_id).first() if not job: raise HTTPException(404, "Job nicht gefunden") session.delete(job) session.commit() return {"message": "Job gelöscht"} finally: session.close() @app.post("/api/jobs/{job_id}/record") async def start_job_recording(job_id: int): session = get_session() try: job = session.query(Job).filter_by(id=job_id).first() if not job: raise HTTPException(404, "Job nicht gefunden") duration = job.max_duration if not duration and job.schedule_start and job.schedule_stop: if job.schedule_end_date and job.schedule_date: start_dt = datetime.datetime.strptime(f"{job.schedule_date} {job.schedule_start}", "%Y-%m-%d %H:%M") end_dt = datetime.datetime.strptime(f"{job.schedule_end_date} {job.schedule_stop}", "%Y-%m-%d %H:%M") duration = max(0, int((end_dt - start_dt).total_seconds())) else: def _time_to_minutes(t: str) -> int: parts = t.split(":") return int(parts[0]) * 60 + int(parts[1]) start_min = _time_to_minutes(job.schedule_start) stop_min = _time_to_minutes(job.schedule_stop) if stop_min > start_min: duration = (stop_min - start_min) * 60 else: duration = (1440 - start_min + stop_min) * 60 job_data = { "name": job.name, "stream_url": job.stream_url, "stream_type": job.stream_type, "output_format": job.output_format, "max_duration": duration, "extra_ffmpeg_args": job.extra_ffmpeg_args, "segment_duration": job.segment_duration, "description": job.description, "ntfy_enabled": job.ntfy_enabled, "plik_enabled": job.plik_enabled, "delete_after_upload": job.delete_after_upload, } finally: session.close() result = recorder.start_recording(job_id=job_id, **job_data) return result # --- Quick Recording --- class QuickRecordRequest(BaseModel): url: str name: str = "" stream_type: str = "auto" max_duration: Optional[int] = None @app.post("/api/quick-record") async def quick_record(data: QuickRecordRequest): name = data.name or f"Schnellaufnahme_{datetime.datetime.now().strftime('%H%M%S')}" result = recorder.start_recording( job_id=None, name=name, stream_url=data.url, stream_type=data.stream_type, max_duration=data.max_duration, ) return result # --- Recordings --- @app.get("/api/recordings") async def list_recordings(active_only: bool = False, limit: int = 50): if active_only: return recorder.get_all_active() return recorder.get_all_recordings(limit=limit) @app.get("/api/recordings/{rec_id}") async def get_recording(rec_id: int): status = recorder.get_status(rec_id) if not status: raise HTTPException(404, "Aufnahme nicht gefunden") return status @app.post("/api/recordings/{rec_id}/stop") async def stop_recording(rec_id: int): status = recorder.get_status(rec_id) if status and status.get("job_id"): scheduler.suppress_current_occurrence(status["job_id"]) recorder.stop_recording(rec_id) return {"message": "Aufnahme gestoppt"} @app.post("/api/recordings/{rec_id}/extend") async def extend_recording(rec_id: int, minutes: int = Query(default=30)): ok = recorder.extend_recording(rec_id, minutes * 60) if not ok: raise HTTPException(400, "Aufnahme kann nicht verlängert werden") return {"message": f"Aufnahme um {minutes} Minuten verlängert"} @app.post("/api/recordings/{rec_id}/unlimited") async def remove_limit(rec_id: int): ok = recorder.remove_duration_limit(rec_id) if not ok: raise HTTPException(400, "Zeitlimit kann nicht entfernt werden") return {"message": "Zeitlimit entfernt"} # --- Logs --- @app.get("/api/logs") async def get_logs(limit: int = 100, level: str = ""): session = get_session() try: q = session.query(LogEntry).order_by(LogEntry.id.desc()) if level: q = q.filter(LogEntry.level == level.upper()) logs = q.limit(limit).all() return [ { "id": l.id, "timestamp": l.timestamp.isoformat() if l.timestamp else None, "level": l.level, "message": l.message, "job_name": l.job_name, } for l in logs ] finally: session.close() # --- Config --- @app.get("/api/config") async def get_config(): safe_config = {**config.data} if "auth" in safe_config: safe_auth = {**safe_config["auth"]} if safe_auth.get("password"): safe_auth["password"] = "***" safe_config["auth"] = safe_auth return safe_config class ConfigUpdate(BaseModel): data: dict @app.put("/api/config") async def update_config(payload: ConfigUpdate): from config import _deep_merge _deep_merge(config.data, payload.data) config.save() return {"message": "Konfiguration gespeichert"} # --- Static files (Frontend) --- if os.path.isdir(STATIC_DIR): app.mount("/", StaticFiles(directory=STATIC_DIR, html=True), name="static") if __name__ == "__main__": import uvicorn host = config.get("server", "host", default="0.0.0.0") port = config.get("server", "port", default=8484) uvicorn.run("app:app", host=host, port=port, reload=False)