650 lines
27 KiB
Python
650 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")
|
|
|
|
@app.get("/live/{station}")
|
|
def live_stream_redirect(station: str, db: Session = Depends(db_session)):
|
|
is_m3u = False
|
|
|
|
# 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]
|
|
is_m3u = True
|
|
|
|
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"
|
|
if is_m3u:
|
|
m3u_content = f"#EXTM3U\n#EXTINF:-1,Radio 357\n{target_url}\n"
|
|
from fastapi.responses import PlainTextResponse
|
|
return PlainTextResponse(content=m3u_content, media_type="audio/x-mpegurl")
|
|
if "#" not in target_url:
|
|
target_url += "#.aac"
|
|
return RedirectResponse(url=target_url, status_code=302)
|
|
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"
|
|
if is_m3u:
|
|
m3u_content = f"#EXTM3U\n#EXTINF:-1,Radio Nowy Świat\n{target_url}\n"
|
|
from fastapi.responses import PlainTextResponse
|
|
return PlainTextResponse(content=m3u_content, media_type="audio/x-mpegurl")
|
|
if "#" not in target_url:
|
|
target_url += "#.mp3"
|
|
return RedirectResponse(url=target_url, status_code=302)
|
|
else:
|
|
raise HTTPException(status_code=404, detail="Brak streamu dla podanej stacji")
|