306 lines
14 KiB
Python
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}).")
|