from fastapi import FastAPI, HTTPException, Request, Response from fastapi.middleware.cors import CORSMiddleware from fastapi.responses import StreamingResponse from pydantic import BaseModel from typing import List, Optional, Dict from pathlib import Path import subprocess import asyncio import json from fastapi.security import OAuth2PasswordBearer from config import settings from config_validation import validate_config, ConfigValidationError from cache import init_cache, set as cache_set from auth import ( create_access_token, create_refresh_token, decode_token, is_blacklisted, add_to_blacklist, ) from ratelimit import ( check_rate_limit, get_rate_limit_remaining, validate_api_key, api_keys, create_api_key as ratelimit_create_api_key, delete_api_key as ratelimit_delete_api_key, ) from prescan import PreScan from nfo_generator import NFOGenerator from image_downloader import ImageDownloader app = FastAPI( title="Rippy API", description="API für das automatische Ripping-System", version="1.0.0" ) # OAuth2 Scheme oauth2_scheme = OAuth2PasswordBearer(tokenUrl="token") # SSE-Connections sse_connections: List = [] @app.on_event("startup") async def startup_event(): """Initialisiere Cache beim Start und validiere Konfiguration.""" init_cache() try: validate_config() except ConfigValidationError as e: print(f"⚠️ Konfigurations-Warnung: {e}") # Middleware für Rate-Limiting @app.middleware("http") async def rate_limit_middleware(request: Request, call_next): """Rate-Limiting Middleware.""" client_ip = request.client.host api_key = request.headers.get("X-API-Key") # Prüfe API Key if api_key: key_info = validate_api_key(api_key) if not key_info: raise HTTPException(status_code=401, detail="Ungültiger API Key") # Rate Limit prüfen if not check_rate_limit(client_ip): return Response( content=json.dumps({"error": "Rate limit exceeded"}), status_code=429, media_type="application/json" ) response = await call_next(request) # Füge Rate-Limit Header hinzu remaining = get_rate_limit_remaining(client_ip) response.headers["X-RateLimit-Remaining"] = str(remaining) return response # CORS hinzufügen app.add_middleware( CORSMiddleware, allow_origins=["*"], allow_credentials=True, allow_methods=["*"], allow_headers=["*"], ) class Job(BaseModel): id: str type: str status: str device: str startTime: str endTime: Optional[str] = None progress: int = 0 class Device(BaseModel): id: str name: str type: str path: str status: str @app.get("/health") async def health_check(): return {"status": "ok", "service": "api"} @app.get("/") async def root(): return { "name": "Rippy", "version": "1.0.0", "description": "Automatisches Ripping-System für CD, DVD und Blu-ray" } @app.get("/jobs", response_model=List[Job]) async def get_jobs(): """Holt alle Jobs.""" return [] @app.get("/devices", response_model=List[Device]) async def get_devices(): """Holt alle Geräte.""" devices = [] try: # 1. Zuerst Symlinks in /dev/disc/ prüfen result = subprocess.run( ["ls", "-la", "/dev/disc/"], capture_output=True, text=True, timeout=5 ) for line in result.stdout.strip().split('\n')[1:]: if line and 'total' not in line: parts = line.split() if len(parts) >= 9: name = parts[-1] devices.append(Device( id=name, name=f"Laufwerk {name}", type="dvd", path=f"/dev/disc/{name}", status="ready" )) # 2. Direkte Geräte /dev/sr0, /dev/cdrom prüfen import os direct_devices = ['/dev/sr0', '/dev/cdrom', '/dev/dvd'] for device_path in direct_devices: if os.path.exists(device_path): # Prüfen ob es ein block device ist import stat mode = os.stat(device_path).st_mode if stat.S_ISBLK(mode): # Geräte-Info aus udevadm holen result = subprocess.run( ["udevadm", "info", "-q", "property", "-n", device_path], capture_output=True, text=True, timeout=5 ) info = {} for line in result.stdout.strip().split('\n'): if '=' in line: key, value = line.split('=', 1) info[key] = value device_type = "dvd" if info.get("ID_CDROM_BD") == "1": device_type = "bluray" elif info.get("ID_CDROM_CD") == "1": device_type = "cd" model = info.get("ID_MODEL", "Laufwerk") serial = info.get("ID_SERIAL", "unknown") devices.append(Device( id=f"{model}_{serial}".replace(" ", "_"), name=model, type=device_type, path=device_path, status="ready", serial=serial, model=model )) except Exception as e: print(f"Fehler beim Lesen der Geräte: {e}") pass return devices # SSE-Stream für Echtzeit-Updates @app.get("/stream/jobs") async def job_stream(): """SSE-Stream für Job-Updates.""" async def event_generator(): while True: if sse_connections: # Job-Status aktualisieren jobs = await get_jobs() yield f"data: {json.dumps([j.dict() for j in jobs])}\n\n" await asyncio.sleep(1) return StreamingResponse(event_generator(), media_type="text/event-stream") # Metadaten-Lookup Endpoints class MetadataLookupRequest(BaseModel): title: str year: Optional[int] = None disc_type: str = "dvd" @app.post("/metadata/lookup") async def lookup_metadata(request: MetadataLookupRequest): """Suche Metadaten für Disc.""" prescan = PreScan() # Dummy device für Pre-Scan device = "/dev/dvd" if request.disc_type in ["dvd", "bluray"] else "/dev/cdrom" result = prescan.scan(device) return { "title": result.title, "year": result.year, "confidence": result.confidence, "metadata": result.metadata, "tracks": result.tracks } @app.post("/metadata/confirm") async def confirm_metadata(title: str, year: Optional[int] = None, metadata: Dict = None): """Bestätige Metadaten.""" from cache.keys import generate_confirmed_key # In Cache speichern cache_key = generate_confirmed_key(title, year) cache_set(cache_key, {"title": title, "year": year, "metadata": metadata or {}}) return {"status": "confirmed", "key": cache_key} # Pre-Scan Endpoint class PreScanRequest(BaseModel): device_path: str @app.post("/prescan") async def run_prescan(request: PreScanRequest): """Führe Pre-Scan durch.""" try: prescan = PreScan() result = prescan.scan(request.device_path) return result.to_dict() except Exception as e: raise HTTPException(status_code=500, detail=str(e)) # Jellyfin-Formatierung Endpoints class JellyfinFormatRequest(BaseModel): title: str year: Optional[int] metadata: Dict disc_type: str output_dir: str @app.post("/jellyfin/format") async def jellyfin_format(request: JellyfinFormatRequest): """Formatiere für Jellyfin (NFO + Images).""" try: nfo_gen = NFOGenerator() img_downloader = ImageDownloader() # Ordnerstruktur erstellen output_path = Path(request.output_dir) if request.disc_type in ["dvd", "bluray"]: # Film-Formatierung title = request.metadata.get("title", request.title) year = request.year or request.metadata.get("year") # movie.nfo movie_nfo = nfo_gen.generate_movie_nfo( title=title, year=year or 2000, overview=request.metadata.get("overview", ""), rating=request.metadata.get("rating", 0), runtime=request.metadata.get("runtime", 0), genres=request.metadata.get("genres", []), director=request.metadata.get("director", ""), actors=request.metadata.get("actors", []) ) nfo_path = output_path / "movie.nfo" nfo_gen.save_nfo(movie_nfo, nfo_path) # Poster und Fanart img_downloader.download_poster(title, output_path, 500) img_downloader.download_fanart(title, output_path, 1920) return { "status": "formatted", "nfo_path": str(nfo_path), "poster_path": str(output_path / "poster.jpg"), "fanart_path": str(output_path / "fanart.jpg") } else: # Audio-Formatierung artist = request.metadata.get("artist", "Unknown Artist") album = request.metadata.get("title", request.title) # Review-Fix 22.07.: `year` war hier undefiniert (existierte nur im Film-Zweig) year = request.year or request.metadata.get("year") # album.nfo album_nfo = nfo_gen.generate_album_nfo( title=album, artist=artist, year=year or 2000, genres=request.metadata.get("genres", []) ) nfo_path = output_path / "album.nfo" nfo_gen.save_nfo(album_nfo, nfo_path) # Album-Cover img_downloader.download_music_images(artist, album, output_path) return { "status": "formatted", "nfo_path": str(nfo_path), "album_cover_path": str(output_path / "album.jpg") } except Exception as e: raise HTTPException(status_code=500, detail=str(e)) # Auth Endpoints class LoginRequest(BaseModel): username: str password: str @app.post("/token") async def login(request: LoginRequest): """Login und Token generieren.""" # Einfache Auth für MVP (in Produktion mit Datenbank); Zugangsdaten aus .env if request.username == settings.admin_username and request.password == settings.admin_password: access_token = create_access_token( data={"sub": request.username, "scopes": ["admin"]} ) refresh_token = create_refresh_token( data={"sub": request.username} ) return { "access_token": access_token, "refresh_token": refresh_token, "token_type": "bearer" } raise HTTPException(status_code=401, detail="Ungültige Anmeldedaten") @app.post("/token/refresh") async def refresh_token(refresh_token: str): """Refresh Access Token.""" payload = decode_token(refresh_token) if not payload or payload.get("type") != "refresh": raise HTTPException(status_code=401, detail="Ungültiges Refresh Token") access_token = create_access_token( data={"sub": payload.get("sub"), "scopes": payload.get("scopes", [])} ) return {"access_token": access_token, "token_type": "bearer"} @app.post("/token/invalidate") async def invalidate_token(token: str): """Invalidate Token (Logout).""" if is_blacklisted(token): raise HTTPException(status_code=400, detail="Token bereits invalidiert") # Review-Fix 22.07.: vorher wurde hier NICHTS geblacklistet (Placebo-Logout) add_to_blacklist(token) if not is_blacklisted(token): raise HTTPException(status_code=400, detail="Ungültiger Token") return {"status": "invalidated"} # API Key Endpoints class APIKeyCreateRequest(BaseModel): name: str # Review-Fix 22.07.: diese Endpoints nutzten `secrets` und `api_keys`, die in # diesem Modul NIE existierten (Crash bei jedem Aufruf) — der echte Key-Store # lebt in ratelimit.py und wird jetzt benutzt. @app.post("/api-keys") async def create_api_key(request: APIKeyCreateRequest): """Erstelle API Key.""" # In Produktion mit Auth prüfen return ratelimit_create_api_key(request.name) @app.get("/api-keys") async def list_api_keys(): """Liste API Keys.""" return list(api_keys.values()) @app.delete("/api-keys/{key}") async def delete_api_key(key: str): """Lösche API Key.""" # In Produktion mit Auth prüfen if ratelimit_delete_api_key(key): return {"status": "deleted"} raise HTTPException(status_code=404, detail="API Key nicht gefunden")