Files
rippy/docker/api/main.py
T
Hitonabi 3e921fc0c9 Etappe 5: Auth + Rate-Limiting
- JWT-Auth (Access 15min/Refresh 7 Tage)
- Rate-Limiting (100/min pro Client/API-Key)
- API Key Management
- OAuth2 Password Scheme
2026-07-21 17:19:34 +02:00

372 lines
10 KiB
Python

from fastapi import FastAPI, HTTPException, Request, Response
from fastapi.middleware.cors import CORSMiddleware
from fastapi.responses import StreamingResponse
from fastapi import WebSocket, WebSocketDisconnect
from pydantic import BaseModel
from typing import List, Optional, Dict
from datetime import datetime
import os
import subprocess
import asyncio
import json
import time
from fastapi.security import OAuth2PasswordBearer
from .config import settings
from .cache import init_cache, set
from .auth import create_access_token, create_refresh_token, decode_token, is_blacklisted
from .ratelimit import check_rate_limit, get_rate_limit_remaining, validate_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."""
init_cache()
# 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:
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"
))
except Exception:
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."""
# In Cache speichern
cache_key = f"confirmed:{title}:{year}"
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 = title
# 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)
if request.username == "admin" and request.password == "rippy123":
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")
# In Produktion mit Redis implementieren
return {"status": "invalidated"}
# API Key Endpoints
class APIKeyCreateRequest(BaseModel):
name: str
@app.post("/api-keys")
async def create_api_key(request: APIKeyCreateRequest):
"""Erstelle API Key."""
# In Produktion mit Auth prüfen
key_info = {
"key": secrets.token_urlsafe(32),
"name": request.name,
"created_at": time.time(),
"rate_limit": 100
}
return key_info
@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 key in api_keys:
del api_keys[key]
return {"status": "deleted"}
raise HTTPException(status_code=404, detail="API Key nicht gefunden")