Files

306 lines
14 KiB
Python

import requests
import json
import base64
import time
import os
from sqlalchemy.orm import Session
from ..database import Program, Episode, Config
from ..diagnostics import ScraperError, invalid_response
def decode_jwt_exp(token):
try:
parts = token.split('.')
if len(parts) != 3: return 0
payload_b64 = parts[1]
payload_b64 += "=" * ((4 - len(payload_b64) % 4) % 4)
return json.loads(base64.b64decode(payload_b64).decode('utf-8')).get("exp", 0)
except:
return 0
def get_audio_length(url):
try:
r = requests.head(url, timeout=5)
return int(r.headers.get('Content-Length', 0))
except:
return 0
class Radio357Scraper:
def __init__(self, db: Session, logger=print, progress_callback=None, stop_flag=None):
self.db = db
self.logger = logger
self.progress_callback = progress_callback
self.stop_flag = stop_flag
self.requests_made = 0
limit_conf = self.db.query(Config).filter_by(key="r357_limit").first()
self.backfill_limit = int(limit_conf.value) if limit_conf and limit_conf.value.isdigit() else 600
hard_limit_conf = self.db.query(Config).filter_by(key="r357_hard_limit").first()
self.hard_limit = int(hard_limit_conf.value) if hard_limit_conf and hard_limit_conf.value.isdigit() else 1000
self.start_time = time.time()
self.max_execution_time = 900 # 15 minut max na cały sync
if self.progress_callback: self.progress_callback(self.requests_made, self.hard_limit)
def check_timeout(self):
if self.stop_flag and self.stop_flag():
raise Exception("Zatrzymano na żądanie użytkownika.")
if time.time() - self.start_time > self.max_execution_time:
raise Exception(f"Przekroczono limit czasu wykonywania skryptu ({self.max_execution_time // 60} min). Zatrzymano awaryjnie.")
def _get(self, url, headers=None):
self.check_timeout()
if self.requests_made >= self.hard_limit:
self.logger(f"R357: Limit {self.hard_limit} zapytań osiągnięty. Przerywam.")
return None
self.requests_made += 1
if self.progress_callback: self.progress_callback(self.requests_made, self.hard_limit)
if not headers:
headers = {}
headers["User-Agent"] = "Mozilla/5.0 (Windows NT 10.0; Win64; x64)"
try:
r = requests.get(url, headers=headers, timeout=10)
if r.status_code == 200:
try:
data = r.json()
except ValueError:
invalid_response("radio357", "INVALID_JSON", "API zwróciło nieprawidłowy JSON", url, r)
embedded = data.get("_embedded") if isinstance(data, dict) else None
if not isinstance(embedded, dict) or not isinstance(embedded.get("podcasts"), list):
invalid_response("radio357", "API_SCHEMA_CHANGED", "Brak oczekiwanej listy podcastów w odpowiedzi API", url, r)
return data
else:
invalid_response("radio357", "HTTP_ERROR", "Błąd pobierania danych API", url, r)
except Exception as e:
if isinstance(e, ScraperError):
raise
raise ScraperError("radio357", "NETWORK_ERROR", f"Błąd zapytania: {e}", url) from e
return None
def get_token(self):
token_conf = self.db.query(Config).filter_by(key="r357_token").first()
if token_conf and decode_jwt_exp(token_conf.value) > time.time() + 3600:
return token_conf.value
email = self.db.query(Config).filter_by(key="r357_email").first()
password = self.db.query(Config).filter_by(key="r357_password").first()
if not email or not password:
self.logger("R357: Brak danych logowania w bazie.")
return None
self.logger("R357: Pobieranie nowego tokenu...")
url = "https://auth.r357.eu/api/auth/login"
payload = {"email": email.value, "password": password.value}
headers = {
"Content-Type": "application/json",
"Origin": "https://konto.radio357.pl",
"User-Agent": "Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36"
}
try:
r = requests.post(url, json=payload, headers=headers, timeout=10)
if r.status_code == 200:
new_token = r.json().get("accessToken")
if new_token:
if not token_conf:
token_conf = Config(key="r357_token", value=new_token)
self.db.add(token_conf)
else:
token_conf.value = new_token
self.db.commit()
return new_token
else:
self.logger(f"R357 Błąd logowania: {r.status_code} {r.text}")
except Exception as e:
self.logger(f"R357 Wyjątek logowania: {e}")
return None
def check_auth_status(self):
token_conf = self.db.query(Config).filter_by(key="r357_token").first()
if not token_conf:
return False, "Brak zapisanego tokenu."
exp = decode_jwt_exp(token_conf.value)
if exp < time.time():
return False, "Token wygasł."
return True, "Token aktywny."
def sync_episodes(self, specific_program_slug=None):
self.logger("R357: Rozpoczynam synchronizację...")
page = 0
new_found = 0
is_first_run = self.db.query(Episode).filter_by(station="radio357").count() == 0
# Przy pierwszym uruchomieniu indeks metadanych musi przejść przez całe API.
# W kolejnych uruchomieniach zostawiamy budżet na backfill URL-i audio.
catchup_limit = self.hard_limit if is_first_run else max(1, self.hard_limit - self.backfill_limit)
seen_slugs = set() # Guard przeciw duplikatom w tej sesji
while True:
if self.requests_made >= catchup_limit:
self.logger(f"R357: [{self.requests_made}/{self.hard_limit}] Limit Catchup osiągnięty. Zostawiam resztę dla Backfill.")
break
url = f"https://static.radio357.pl/api/content/v1/podcasts?page={page}"
self.logger(f"R357: [{self.requests_made}/{self.hard_limit}] Pobieram stronę {page}...")
data = self._get(url)
if not data:
raise Exception(f"Błąd sieci/API podczas pobierania strony {page} (Radio 357)")
episodes = data.get("_embedded", {}).get("podcasts", [])
if not episodes:
break
all_known = True
page_has_target = False
for ep in episodes:
ep_id = str(ep["id"])
# Extract program info
programs_list = ep.get("programs")
if not programs_list: continue
program = programs_list[0]
slug = program.get("slug", f"program_{program['id']}")
if specific_program_slug and slug != specific_program_slug:
continue # Skip if we are targeting a specific program
page_has_target = True
existing = self.db.query(Episode).filter_by(station="radio357", ep_id=ep_id).first()
if not existing:
all_known = False
# Ensure program exists (cache guard against duplicates)
if slug not in seen_slugs:
db_prog = self.db.query(Program).filter_by(station="radio357", slug=slug).first()
if not db_prog:
db_prog = Program(
station="radio357",
slug=slug,
name=program.get("name", f"Podcast {program['id']}"),
image=program.get("image", ""),
description=program.get("description", "")
)
self.db.add(db_prog)
self.db.flush() # Natychmiastowy flush żeby kolejne odcinki widziały program
seen_slugs.add(slug)
team = ep.get("team", [])
authors = ", ".join([t.get("name", "") for t in team]) or "Radio 357"
new_ep = Episode(
station="radio357",
program_slug=slug,
ep_id=ep_id,
title=ep.get("subTitle") or ep.get("title", ""),
duration_secs=ep.get("duration", 0),
image=ep.get("image", ""),
description=ep.get("description", ""),
pub_date=ep.get("published", ""),
authors=authors,
url=None,
is_broken=False
)
self.db.add(new_ep)
new_found += 1
self.db.commit()
if not is_first_run and all_known and not specific_program_slug:
self.logger(f"R357: Wszystkie odcinki na stronie {page} są znane. Zakończono catchup.")
break
if specific_program_slug and page_has_target and all_known:
self.logger(f"R357: Wszystkie odcinki programu {specific_program_slug} na stronie {page} są znane. Zakończono catchup programu.")
break
if len(episodes) < data.get("page_size", 250):
break
page += 1
self.logger(f"R357: Znaleziono {new_found} nowych odcinków w bazie.")
def backfill_urls(self, specific_program_slug=None):
backfill_count = 0
self.logger(f"R357: Rozpoczynam pobieranie linków audio (limit {self.backfill_limit})...")
query = self.db.query(Episode).filter_by(station="radio357", url=None, is_broken=False).order_by(Episode.pub_date.desc())
if specific_program_slug:
query = query.filter_by(program_slug=specific_program_slug)
missing_eps = query.all()
token = self.get_token()
if not token:
self.logger("R357: Brak tokena, przerywam backfill url")
return
fetched = 0
for ep in missing_eps:
self.check_timeout()
if self.requests_made >= self.hard_limit or backfill_count >= self.backfill_limit:
break
gateway_url = f"https://gateway.r357.eu/api/content/podcast/{ep.ep_id}/url"
self.requests_made += 1
if self.progress_callback: self.progress_callback(self.requests_made, self.hard_limit)
backfill_count += 1
self.logger(f"R357: [{self.requests_made}/{self.hard_limit}] Pobieram URL dla {ep.title[:30]}...")
try:
r = requests.get(gateway_url, headers={
"Authorization": f"Bearer {token}",
"Origin": "https://radio357.pl",
"User-Agent": "Mozilla/5.0"
}, timeout=10)
if r.status_code == 200:
audio_url = r.json().get("url")
if audio_url:
ep.url = audio_url
fetched += 1
elif r.status_code == 401 or r.status_code == 403:
self.logger(f"R357: [{self.requests_made}/{self.hard_limit}] Token odrzucony przez gateway. Wymagam reloginu.")
# Usuwamy token, zeby przy nastepnym runie pobral nowy
self.db.query(Config).filter_by(key="r357_token").delete()
self.db.commit()
raise Exception("R357 Token Expired")
else:
self.logger(f"R357: [{self.requests_made}/{self.hard_limit}] Błąd HTTP {r.status_code} podczas pobierania audio: {r.text[:100]}")
except Exception as e:
self.logger(f"R357: [{self.requests_made}/{self.hard_limit}] Błąd URL dla odcinka {ep.ep_id}: {e}")
if "Token Expired" in str(e):
raise
time.sleep(0.5)
self.db.commit()
self.logger(f"R357: [{self.requests_made}/{self.hard_limit}] Pomyślnie pobrano {fetched} nowych linków.")
def update_live_stream_url(self):
token = self.get_token()
if not token:
return
self.logger("R357: Aktualizacja adresu live stream...")
url = "https://stream.radio357.pl/?s=www"
try:
r = requests.get(url, headers={"User-Agent": "Mozilla/5.0"}, cookies={"token": token}, allow_redirects=False, timeout=5)
if r.status_code in (301, 302, 307) and "Location" in r.headers:
target_url = r.headers["Location"]
config_url = self.db.query(Config).filter_by(key="r357_live_url").first()
if not config_url:
self.db.add(Config(key="r357_live_url", value=target_url))
else:
config_url.value = target_url
self.db.commit()
self.logger(f"R357: Zapisano nowy strumień: {target_url}")
except Exception as e:
self.logger(f"R357 Błąd podczas pobierania streamu premium: {e}")
def run_full_sync(self, specific_program_slug=None):
self.sync_episodes(specific_program_slug)
self.backfill_urls(specific_program_slug)
if not specific_program_slug:
self.update_live_stream_url()
self.logger(f"R357: Koniec. Wysłano zapytania: {self.requests_made}/{self.hard_limit} (w tym backfill ograniczony do {self.backfill_limit}).")