Files
radiosync/main.py
T

647 lines
27 KiB
Python

import os
import datetime
import threading
import hashlib
import re
import secrets
from fastapi import FastAPI, HTTPException, Depends, Request, Form, BackgroundTasks
from fastapi.responses import HTMLResponse, Response, RedirectResponse
from fastapi.staticfiles import StaticFiles
from fastapi.templating import Jinja2Templates
from sqlalchemy.orm import Session
from apscheduler.schedulers.background import BackgroundScheduler
from core.database import SessionLocal, init_db, Config, Program, Episode, ErrorEvent
from core.scrapers.radio357 import Radio357Scraper
from core.scrapers.rns import RNScraper
from core.scrapers.jazz import JazzScraper
from core.rss_generator import generate_master_opml, generate_podcast_rss
from core.notifier import send_notification, validate_notification_url
from starlette.middleware.sessions import SessionMiddleware
app = FastAPI(title="RadioSync Hub")
session_secret = os.environ.get("RADIOSYNC_SECRET_KEY")
if not session_secret:
raise RuntimeError("RADIOSYNC_SECRET_KEY must be set")
app.add_middleware(SessionMiddleware, secret_key=session_secret)
# Set up templates
base_dir = os.path.dirname(os.path.abspath(__file__))
templates = Jinja2Templates(directory=os.path.join(base_dir, "templates"))
app.mount("/static", StaticFiles(directory=os.path.join(base_dir, "static")), name="static")
init_db()
scheduler = BackgroundScheduler()
sync_locks = {station: threading.Lock() for station in ("radio357", "rns", "jazz")}
# State variables for logging
sync_status = {
"radio357": {"last_run": None, "status": "Oczekuje", "progress": "", "stop_requested": False},
"rns": {"last_run": None, "status": "Oczekuje", "progress": "", "stop_requested": False},
"jazz": {"last_run": None, "status": "Oczekuje", "progress": "", "stop_requested": False},
"logs": []
}
def add_log(msg: str):
ts = datetime.datetime.now().strftime("%Y-%m-%d %H:%M:%S")
sync_status["logs"].insert(0, f"[{ts}] {msg}")
if len(sync_status["logs"]) > 100:
sync_status["logs"].pop()
print(msg)
def db_session():
db = SessionLocal()
try:
yield db
finally:
db.close()
def notify_error(db: Session, station: str, message: str):
url = db.query(Config).filter_by(key="ntfy_url").first()
topic = db.query(Config).filter_by(key="ntfy_topic").first()
user = db.query(Config).filter_by(key="ntfy_user").first()
pw = db.query(Config).filter_by(key="ntfy_password").first()
if url and topic:
try:
validate_notification_url(url.value)
send_notification(url.value, topic.value, f"RadioSync Błąd: {station}", message,
username=user.value if user else None,
password=pw.value if pw else None)
except Exception as exc:
add_log(f"Nie udało się wysłać powiadomienia ntfy: {exc}")
def record_event(db: Session, station: str, message: str, source: str = "sync", severity: str = "error"):
now = datetime.datetime.now()
snapshot_match = re.search(r"Snapshot: (.+?)(?: URL:|$)", message)
snapshot_path = snapshot_match.group(1) if snapshot_match else None
stable_message = re.sub(r"\. Snapshot: .+?(?= URL:|$)", "", message)
fingerprint_message = re.sub(r" URL: https?://\S+", "", stable_message)
fingerprint_message = re.sub(r"\b\d+\b", "<count>", fingerprint_message)
code_match = re.search(r"\[([A-Z0-9_]+)\]", stable_message)
stable_source = code_match.group(1) if code_match else source
fingerprint = hashlib.sha256(f"{station}:{stable_source}:{fingerprint_message}".encode("utf-8")).hexdigest()
try:
event = db.query(ErrorEvent).filter_by(fingerprint=fingerprint).first()
if event:
event.occurrences += 1
event.last_seen = now
event.acknowledged = False
event.acknowledged_at = None
if snapshot_path:
event.snapshot_path = snapshot_path
else:
db.add(ErrorEvent(station=station, source=stable_source, message=stable_message,
severity=severity, fingerprint=fingerprint, first_seen=now,
last_seen=now, snapshot_path=snapshot_path))
db.commit()
except Exception as exc:
db.rollback()
print(f"Failed to persist error event: {exc}")
label = "BŁĄD" if severity == "error" else "WARNING"
add_log(f"[{label}][{station}] {stable_message}")
def record_error(db: Session, station: str, message: str, source: str = "sync"):
record_event(db, station, message, source, severity="error")
def record_warning(db: Session, station: str, message: str, source: str = "sync"):
record_event(db, station, message, source, severity="warning")
def is_user_stop(message: str):
return "Zatrzymano na żądanie użytkownika" in message
def csrf_token(request: Request):
token = request.session.get("csrf_token")
if not token:
token = secrets.token_urlsafe(32)
request.session["csrf_token"] = token
return token
def validate_csrf(request: Request, token: str):
if not token or token != request.session.get("csrf_token"):
raise HTTPException(status_code=403, detail="Nieprawidłowy token CSRF")
def validate_interval(value: str, field: str):
try:
interval = int(value)
except (TypeError, ValueError) as exc:
raise HTTPException(status_code=400, detail=f"{field}: podaj liczbę godzin") from exc
if not 1 <= interval <= 168:
raise HTTPException(status_code=400, detail=f"{field}: zakres to 1-168 godzin")
return interval
def validate_limit(value: str, field: str, maximum: int):
try:
limit = int(value)
except (TypeError, ValueError) as exc:
raise HTTPException(status_code=400, detail=f"{field}: podaj liczbę") from exc
if not 1 <= limit <= maximum:
raise HTTPException(status_code=400, detail=f"{field}: zakres to 1-{maximum}")
return limit
def configured_interval(db: Session, key: str, default: int):
row = db.query(Config).filter_by(key=key).first()
try:
value = int(row.value) if row else default
except (TypeError, ValueError):
return default
return value if 1 <= value <= 168 else default
def station_status_key(station: str):
return {
"Radio 357": "radio357",
"RNŚ": "rns",
"Radio Jazz FM": "jazz",
}.get(station)
def persist_last_run(station: str, dt: datetime.datetime):
"""Zapisuje datę ostatniego udanego syncu do bazy danych."""
db = SessionLocal()
try:
existing = db.query(Config).filter_by(key=f"{station}_last_run").first()
if existing:
existing.value = dt.isoformat()
else:
db.add(Config(key=f"{station}_last_run", value=dt.isoformat()))
db.commit()
finally:
db.close()
def persist_status(station: str, status: str):
db = SessionLocal()
try:
key = f"{station}_status"
existing = db.query(Config).filter_by(key=key).first()
if existing:
existing.value = status
else:
db.add(Config(key=key, value=status))
db.commit()
finally:
db.close()
def set_status(station: str, status: str):
sync_status[station]["status"] = status
persist_status(station, status)
def load_last_runs():
"""Wczytuje daty ostatnich syncronizacji z bazy danych przy starcie."""
db = SessionLocal()
try:
for station in ["radio357", "rns", "jazz"]:
row = db.query(Config).filter_by(key=f"{station}_last_run").first()
if row and row.value:
try:
sync_status[station]["last_run"] = datetime.datetime.fromisoformat(row.value)
except ValueError:
pass
status_row = db.query(Config).filter_by(key=f"{station}_status").first()
if status_row and status_row.value != "W trakcie":
sync_status[station]["status"] = status_row.value
finally:
db.close()
def update_progress(station: str, made: int, limit: int):
sync_status[station]["progress"] = f"{made} / {limit}"
# --- BACKGROUND JOBS ---
def job_sync_radio357(specific_program=None):
if not sync_locks["radio357"].acquire(blocking=False):
add_log("R357: Pomijam synchronizację, ponieważ już trwa.")
return
db = None
try:
set_status("radio357", "W trakcie")
sync_status["radio357"]["progress"] = "0 / ?"
sync_status["radio357"]["stop_requested"] = False
add_log(f"Rozpoczęto synchronizację Radia 357 (program: {specific_program or 'wszystkie'})")
db = SessionLocal()
stop_fn = lambda: sync_status["radio357"]["stop_requested"]
scraper = Radio357Scraper(db, logger=add_log,
progress_callback=lambda m, l: update_progress("radio357", m, l),
stop_flag=stop_fn)
scraper.run_full_sync(specific_program)
now = datetime.datetime.now()
sync_status["radio357"]["last_run"] = now
set_status("radio357", "OK")
persist_last_run("radio357", now)
except Exception as e:
msg = f"Błąd R357: {e}"
if db:
if is_user_stop(msg):
record_warning(db, "Radio 357", "Synchronizacja zatrzymana na żądanie użytkownika.", "user_stop")
else:
record_error(db, "Radio 357", msg)
set_status("radio357", "Przerwano" if is_user_stop(msg) else "Błąd")
if db and not is_user_stop(msg):
notify_error(db, "Radio 357", msg)
finally:
sync_status["radio357"]["progress"] = ""
sync_status["radio357"]["stop_requested"] = False
if db:
db.close()
sync_locks["radio357"].release()
def job_sync_rns(specific_program=None):
if not sync_locks["rns"].acquire(blocking=False):
add_log("RNŚ: Pomijam synchronizację, ponieważ już trwa.")
return
db = None
try:
set_status("rns", "W trakcie")
sync_status["rns"]["progress"] = "0 / ?"
sync_status["rns"]["stop_requested"] = False
add_log(f"Rozpoczęto synchronizację RNŚ (program: {specific_program or 'wszystkie'})")
db = SessionLocal()
stop_fn = lambda: sync_status["rns"]["stop_requested"]
scraper = RNScraper(db, logger=add_log,
progress_callback=lambda m, l: update_progress("rns", m, l),
warning_callback=lambda message: record_warning(db, "RNŚ", message, "missing_player"),
stop_flag=stop_fn)
scraper.run_full_sync(specific_program)
now = datetime.datetime.now()
sync_status["rns"]["last_run"] = now
set_status("rns", "OK")
persist_last_run("rns", now)
except Exception as e:
msg = f"Błąd RNŚ: {e}"
if db:
if is_user_stop(msg):
record_warning(db, "RNŚ", "Synchronizacja zatrzymana na żądanie użytkownika.", "user_stop")
else:
record_error(db, "RNŚ", msg)
set_status("rns", "Przerwano" if is_user_stop(msg) else "Błąd")
if db and not is_user_stop(msg):
notify_error(db, "RNŚ", msg)
finally:
sync_status["rns"]["progress"] = ""
sync_status["rns"]["stop_requested"] = False
if db:
db.close()
sync_locks["rns"].release()
def job_sync_jazz():
if not sync_locks["jazz"].acquire(blocking=False):
add_log("JAZZ: Pomijam synchronizację, ponieważ już trwa.")
return
db = None
try:
set_status("jazz", "W trakcie")
add_log("Rozpoczęto synchronizację Radio Jazz FM")
db = SessionLocal()
scraper = JazzScraper(db, logger=add_log)
scraper.run_full_sync()
sync_status["jazz"]["last_run"] = datetime.datetime.now()
now = datetime.datetime.now()
sync_status["jazz"]["last_run"] = now
set_status("jazz", "OK")
persist_last_run("jazz", now)
except Exception as e:
msg = f"Błąd JAZZ: {e}"
if db:
if is_user_stop(msg):
record_warning(db, "Radio Jazz FM", "Synchronizacja zatrzymana na żądanie użytkownika.", "user_stop")
else:
record_error(db, "Radio Jazz FM", msg)
set_status("jazz", "Przerwano" if is_user_stop(msg) else "Błąd")
if db and not is_user_stop(msg):
notify_error(db, "Radio Jazz FM", msg)
finally:
if db:
db.close()
sync_locks["jazz"].release()
@app.on_event("startup")
def startup_event():
db = SessionLocal()
r357_int = configured_interval(db, "r357_interval", 6)
rns_int = configured_interval(db, "rns_interval", 6)
jazz_int = configured_interval(db, "jazz_interval", 24)
db.close()
scheduler.add_job(job_sync_radio357, 'interval', hours=r357_int, id='sync_r357', max_instances=1, coalesce=True)
scheduler.add_job(job_sync_rns, 'interval', hours=rns_int, id='sync_rns', max_instances=1, coalesce=True)
scheduler.add_job(job_sync_jazz, 'interval', hours=jazz_int, id='sync_jazz', max_instances=1, coalesce=True)
scheduler.start()
load_last_runs()
add_log("Aplikacja i harmonogram uruchomione.")
@app.on_event("shutdown")
def shutdown_event():
scheduler.shutdown()
# --- ROUTES ---
@app.get("/", response_class=HTMLResponse)
def index(request: Request, db: Session = Depends(db_session)):
msg = request.session.pop("msg", None)
stats = {
"r357_progs": db.query(Program).filter_by(station="radio357").count(),
"r357_eps": db.query(Episode).filter_by(station="radio357").count(),
"r357_audio_eps": db.query(Episode).filter_by(station="radio357", is_broken=False).filter(Episode.url.isnot(None)).count(),
"rns_progs": db.query(Program).filter_by(station="rns").count(),
"rns_eps": db.query(Episode).filter_by(station="rns").count(),
"rns_audio_eps": db.query(Episode).filter_by(station="rns", is_broken=False).filter(Episode.url.isnot(None)).count(),
}
config = {c.key: c.value for c in db.query(Config).all()}
pending_errors = db.query(ErrorEvent).filter_by(acknowledged=False).order_by(ErrorEvent.last_seen.desc()).all()
acknowledged_errors = db.query(ErrorEvent).filter_by(acknowledged=True).order_by(ErrorEvent.last_seen.desc()).limit(20).all()
return templates.TemplateResponse(request, "index.html", {
"status": sync_status,
"stats": stats,
"config": config,
"msg": msg,
"pending_errors": pending_errors,
"acknowledged_errors": acknowledged_errors,
"csrf_token": csrf_token(request),
})
@app.post("/errors/{error_id}/acknowledge")
def acknowledge_error(request: Request, error_id: int, csrf: str = Form(...), db: Session = Depends(db_session)):
validate_csrf(request, csrf)
event = db.query(ErrorEvent).filter_by(id=error_id).first()
if event:
event.acknowledged = True
event.acknowledged_at = datetime.datetime.now()
db.commit()
add_log(f"Potwierdzono odczytanie błędu: {event.message}")
status_key = station_status_key(event.station)
pending_for_station = db.query(ErrorEvent).filter_by(
station=event.station, acknowledged=False
).count()
if status_key and pending_for_station == 0 and sync_status[status_key]["status"] == "Błąd":
set_status(status_key, "Oczekuje")
return RedirectResponse(url="/", status_code=303)
@app.post("/config")
def save_config(
request: Request,
csrf: str = Form(...),
r357_email: str = Form(""), r357_password: str = Form(""),
rns_email: str = Form(""), rns_password: str = Form(""),
ntfy_url: str = Form(""), ntfy_topic: str = Form(""),
ntfy_user: str = Form(""), ntfy_password: str = Form(""),
r357_limit: str = Form("600"), rns_limit: str = Form("50"),
r357_hard_limit: str = Form("1000"), rns_hard_limit: str = Form("500"),
r357_interval: str = Form("6"), rns_interval: str = Form("6"), jazz_interval: str = Form("24"),
db: Session = Depends(db_session)
):
validate_csrf(request, csrf)
r357_interval_value = validate_interval(r357_interval, "r357_interval")
rns_interval_value = validate_interval(rns_interval, "rns_interval")
jazz_interval_value = validate_interval(jazz_interval, "jazz_interval")
validate_limit(r357_limit, "r357_limit", 10000)
validate_limit(rns_limit, "rns_limit", 10000)
validate_limit(r357_hard_limit, "r357_hard_limit", 20000)
validate_limit(rns_hard_limit, "rns_hard_limit", 20000)
if ntfy_url:
try:
validate_notification_url(ntfy_url)
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc)) from exc
import urllib.parse
updates = {
"r357_email": r357_email, "r357_password": r357_password,
"rns_email": rns_email, "rns_password": rns_password,
"ntfy_url": ntfy_url, "ntfy_topic": ntfy_topic,
"ntfy_user": ntfy_user, "ntfy_password": ntfy_password,
"r357_limit": r357_limit, "rns_limit": rns_limit,
"r357_hard_limit": r357_hard_limit, "rns_hard_limit": rns_hard_limit,
"r357_interval": r357_interval, "rns_interval": rns_interval, "jazz_interval": jazz_interval
}
for k, v in updates.items():
if v:
c = db.query(Config).filter_by(key=k).first()
if not c:
db.add(Config(key=k, value=v))
else:
c.value = v
db.commit()
scheduler.reschedule_job('sync_r357', trigger='interval', hours=r357_interval_value)
scheduler.reschedule_job('sync_rns', trigger='interval', hours=rns_interval_value)
scheduler.reschedule_job('sync_jazz', trigger='interval', hours=jazz_interval_value)
if r357_email or r357_password:
db.query(Config).filter(Config.key == "r357_token").delete(synchronize_session=False)
if rns_email or rns_password:
db.query(Config).filter(Config.key == "rns_cookies").delete(synchronize_session=False)
db.commit()
m = "Zapisano nową konfigurację. Harmonogram i sesje zresetowane."
add_log(m)
request.session["msg"] = m
return RedirectResponse(url="/", status_code=303)
@app.post("/test-ntfy")
def test_ntfy(request: Request, csrf: str = Form(...), db: Session = Depends(db_session)):
validate_csrf(request, csrf)
import urllib.parse
notify_error(db, "TEST", "To jest wiadomość testowa z RadioSync Hub.")
m = "Wysłano testowe powiadomienie ntfy."
add_log(m)
request.session["msg"] = m
return RedirectResponse(url="/", status_code=303)
@app.post("/test-auth")
def test_auth(request: Request, station: str = Form(...), csrf: str = Form(...), db: Session = Depends(db_session)):
validate_csrf(request, csrf)
import urllib.parse
m = ""
if station == "radio357":
scraper = Radio357Scraper(db, logger=add_log)
ok, msg = scraper.check_auth_status()
m = f"R357: {'✅' if ok else '❌'} {msg}"
elif station == "rns":
scraper = RNScraper(db, logger=add_log)
ok, msg = scraper.check_auth_status()
m = f"RNŚ: {'✅' if ok else '❌'} {msg}"
add_log(m)
request.session["msg"] = m
return RedirectResponse(url="/", status_code=303)
@app.post("/force-relogin")
def force_relogin(request: Request, station: str = Form(...), csrf: str = Form(...), db: Session = Depends(db_session)):
validate_csrf(request, csrf)
import urllib.parse
add_log(f"Wymuszony relogin dla: {station}...")
m = ""
if station == "radio357":
db.query(Config).filter_by(key="r357_token").delete()
db.commit()
scraper = Radio357Scraper(db, logger=add_log)
if scraper.get_token():
m = "R357: ✅ Relogin poprawny!"
else:
m = "R357: ❌ Relogin nieudany (sprawdź dane)."
elif station == "rns":
db.query(Config).filter_by(key="rns_cookies").delete()
db.commit()
scraper = RNScraper(db, logger=add_log)
if scraper.perform_login():
m = "RNŚ: ✅ Relogin poprawny!"
else:
m = "RNŚ: ❌ Relogin nieudany (sprawdź dane)."
add_log(m)
request.session["msg"] = m
return RedirectResponse(url="/", status_code=303)
@app.post("/force-sync")
def force_sync(request: Request, station: str = Form(...), csrf: str = Form(...), bg_tasks: BackgroundTasks = BackgroundTasks()):
validate_csrf(request, csrf)
import urllib.parse
if station not in sync_locks:
raise HTTPException(status_code=400, detail="Unknown station")
if sync_locks[station].locked():
m = f"Synchronizacja dla {station} już trwa."
add_log(m)
request.session["msg"] = m
return RedirectResponse(url="/", status_code=303)
if station == "radio357":
bg_tasks.add_task(job_sync_radio357)
elif station == "rns":
bg_tasks.add_task(job_sync_rns)
elif station == "jazz":
bg_tasks.add_task(job_sync_jazz)
names = {"radio357": "Radia 357", "rns": "Radia Nowy Świat", "jazz": "Radio Jazz FM"}
m = f"Rozpoczęto w tle synchronizację: {names.get(station, station)}."
add_log(m)
request.session["msg"] = m
return RedirectResponse(url="/", status_code=303)
@app.post("/stop-sync")
def stop_sync(request: Request, station: str = Form(...), csrf: str = Form(...)):
validate_csrf(request, csrf)
if station in sync_status and sync_status[station]["status"] == "W trakcie":
sync_status[station]["stop_requested"] = True
m = f"Wysłano sygnał zatrzymania synchronizacji: {station}."
else:
m = f"Sync dla {station} nie jest aktualnie uruchomiony."
add_log(m)
request.session["msg"] = m
return RedirectResponse(url="/", status_code=303)
@app.post("/force-sync-program")
def force_sync_program(request: Request, station: str = Form(...), slug: str = Form(...), csrf: str = Form(...), bg_tasks: BackgroundTasks = BackgroundTasks()):
validate_csrf(request, csrf)
import urllib.parse
if station not in {"radio357", "rns"}:
raise HTTPException(status_code=400, detail="Station does not support program sync")
if sync_locks[station].locked():
m = f"Synchronizacja dla {station} już trwa."
add_log(m)
request.session["msg"] = m
return RedirectResponse(url=f"/programs/{station}", status_code=303)
if station == "radio357":
bg_tasks.add_task(job_sync_radio357, specific_program=slug)
else:
bg_tasks.add_task(job_sync_rns, specific_program=slug)
m = f"Zlecono wymuszoną synchronizację pojedynczego programu: {slug} ({station})"
add_log(m)
request.session["msg"] = m
return RedirectResponse(url=f"/programs/{station}", status_code=303)
@app.get("/programs/{station}", response_class=HTMLResponse)
def list_programs(request: Request, station: str, db: Session = Depends(db_session)):
msg = request.session.pop("msg", None)
import datetime
progs = db.query(Program).filter_by(station=station).order_by(Program.name).all()
for p in progs:
p.missing_urls_count = db.query(Episode).filter_by(station=station, program_slug=p.slug, url=None).count()
p.total_eps = db.query(Episode).filter_by(station=station, program_slug=p.slug).count()
if p.last_catchup:
p.last_catchup_str = datetime.datetime.fromtimestamp(p.last_catchup).strftime('%Y-%m-%d %H:%M')
else:
p.last_catchup_str = "Nigdy"
return templates.TemplateResponse(request, "programs.html", {"station": station, "programs": progs, "msg": msg, "csrf_token": csrf_token(request)})
@app.get("/episodes/{station}/{slug:path}", response_class=HTMLResponse)
def list_episodes(request: Request, station: str, slug: str, db: Session = Depends(db_session)):
msg = request.session.pop("msg", None)
prog = db.query(Program).filter_by(station=station, slug=slug).first()
if not prog:
raise HTTPException(status_code=404, detail="Program not found")
eps = db.query(Episode).filter_by(station=station, program_slug=slug).order_by(Episode.pub_date.desc()).all()
return templates.TemplateResponse(request, "episodes.html", {"station": station, "program": prog, "episodes": eps, "msg": msg, "csrf_token": csrf_token(request)})
@app.post("/episodes/delete-url")
def delete_episode_url(request: Request, ep_id: int = Form(...), csrf: str = Form(...), db: Session = Depends(db_session)):
validate_csrf(request, csrf)
import urllib.parse
ep = db.query(Episode).filter_by(id=ep_id).first()
if ep:
ep.url = None
ep.is_broken = True
db.commit()
m = f"Usunięto błędny link audio dla odcinka: {ep.title}"
add_log(m)
request.session["msg"] = m
return RedirectResponse(url=f"/episodes/{ep.station}/{ep.program_slug}", status_code=303)
return RedirectResponse(url="/", status_code=303)
# --- RSS ENDPOINTS ---
@app.get("/feeds/Podcasts.opml")
def get_master_opml(request: Request, db: Session = Depends(db_session)):
xml = generate_master_opml(db, str(request.base_url).rstrip('/'))
return Response(content=xml, media_type="application/xml")
@app.get("/feeds/{station}.opml")
def get_station_opml(request: Request, station: str, db: Session = Depends(db_session)):
if station not in {"radio357", "rns", "jazz"}:
raise HTTPException(status_code=404, detail="Station not found")
xml = generate_master_opml(db, str(request.base_url).rstrip('/'), station=station)
return Response(content=xml, media_type="application/xml")
@app.get("/feeds/{station}/{slug_with_ext:path}")
def podcast_rss(station: str, slug_with_ext: str, request: Request, db: Session = Depends(db_session)):
if not slug_with_ext.endswith(".xml"):
raise HTTPException(status_code=404, detail="Not Found")
slug = slug_with_ext[:-4]
xml = generate_podcast_rss(db, station, slug)
if not xml:
return Response(status_code=404)
return Response(content=xml, media_type="application/xml")
def iter_stream(target_url):
import requests
with requests.get(target_url, stream=True, timeout=10) as r:
for chunk in r.iter_content(chunk_size=8192):
if chunk:
yield chunk
@app.get("/live/{station}")
def live_stream_redirect(station: str, db: Session = Depends(db_session)):
# Usuwamy sztuczne rozszerzenia dodane dla Mopidy
if station.endswith(".aac"):
station = station[:-4]
elif station.endswith(".mp3"):
station = station[:-4]
elif station.endswith(".m3u"):
station = station[:-4]
from fastapi.responses import StreamingResponse
if station == "radio357":
config_url = db.query(Config).filter_by(key="r357_live_url").first()
target_url = config_url.value if config_url else "https://stream.rcs.revma.com/ye5kghkgcm0uv"
# Proxy strumienia! Zwracamy jako audio/mpeg żeby ominąć unwrapper
return StreamingResponse(iter_stream(target_url), media_type="audio/mpeg")
elif station == "rns":
config_url = db.query(Config).filter_by(key="rns_live_url").first()
target_url = config_url.value if config_url else "https://stream.rcs.revma.com/4md4m0a0fs8uv"
return StreamingResponse(iter_stream(target_url), media_type="audio/mpeg")
else:
raise HTTPException(status_code=404, detail="Brak streamu dla podanej stacji")