@@ -0,0 +1,110 @@
|
||||
from sqlalchemy import create_engine, Column, String, Integer, Boolean, Float, Text, DateTime, TypeDecorator
|
||||
from sqlalchemy.orm import declarative_base, sessionmaker
|
||||
from .security import encrypt_config_value, decrypt_config_value
|
||||
|
||||
DATABASE_URL = "sqlite:///./data/radiosync.db"
|
||||
|
||||
engine = create_engine(DATABASE_URL, connect_args={"check_same_thread": False})
|
||||
SessionLocal = sessionmaker(autocommit=False, autoflush=False, bind=engine)
|
||||
|
||||
Base = declarative_base()
|
||||
|
||||
|
||||
class EncryptedString(TypeDecorator):
|
||||
impl = String
|
||||
cache_ok = True
|
||||
|
||||
def process_bind_param(self, value, dialect):
|
||||
return encrypt_config_value(value)
|
||||
|
||||
def process_result_value(self, value, dialect):
|
||||
return decrypt_config_value(value)
|
||||
|
||||
class Config(Base):
|
||||
__tablename__ = "config"
|
||||
key = Column(String, primary_key=True, index=True)
|
||||
value = Column(EncryptedString)
|
||||
|
||||
class Program(Base):
|
||||
__tablename__ = "programs"
|
||||
id = Column(Integer, primary_key=True, index=True)
|
||||
station = Column(String, index=True) # "radio357", "rns", "jazz"
|
||||
slug = Column(String, index=True)
|
||||
name = Column(String)
|
||||
description = Column(Text, nullable=True)
|
||||
image = Column(String, nullable=True)
|
||||
|
||||
# Sync tracking
|
||||
backfill_page = Column(Integer, default=1)
|
||||
backfill_complete = Column(Boolean, default=False)
|
||||
total_pages = Column(Integer, default=0)
|
||||
last_catchup = Column(Float, default=0.0)
|
||||
next_catchup_after = Column(Float, default=0.0) # timestamp – nie sprawdzaj przed tym czasem
|
||||
|
||||
class Episode(Base):
|
||||
__tablename__ = "episodes"
|
||||
id = Column(Integer, primary_key=True, index=True)
|
||||
station = Column(String, index=True)
|
||||
program_slug = Column(String, index=True)
|
||||
ep_id = Column(String, index=True) # ID from the radio API
|
||||
|
||||
title = Column(String)
|
||||
authors = Column(String, nullable=True)
|
||||
url = Column(String, nullable=True)
|
||||
image = Column(String, nullable=True)
|
||||
pub_date = Column(String, nullable=True) # ISO format or YYYY-MM-DD
|
||||
duration_secs = Column(Integer, default=0)
|
||||
description = Column(Text, nullable=True)
|
||||
|
||||
is_broken = Column(Boolean, default=False)
|
||||
|
||||
|
||||
class ErrorEvent(Base):
|
||||
__tablename__ = "error_events"
|
||||
id = Column(Integer, primary_key=True, index=True)
|
||||
station = Column(String, index=True, nullable=False)
|
||||
source = Column(String, nullable=False)
|
||||
message = Column(Text, nullable=False)
|
||||
severity = Column(String, default="error", nullable=False)
|
||||
fingerprint = Column(String, index=True, nullable=False)
|
||||
occurrences = Column(Integer, default=1, nullable=False)
|
||||
first_seen = Column(DateTime, nullable=False)
|
||||
last_seen = Column(DateTime, nullable=False)
|
||||
acknowledged = Column(Boolean, default=False, nullable=False)
|
||||
acknowledged_at = Column(DateTime, nullable=True)
|
||||
snapshot_path = Column(String, nullable=True)
|
||||
|
||||
def init_db():
|
||||
Base.metadata.create_all(bind=engine)
|
||||
import sqlite3, os
|
||||
db_path = DATABASE_URL.replace("sqlite:///", "")
|
||||
if os.path.exists(db_path):
|
||||
with sqlite3.connect(db_path) as conn:
|
||||
cursor = conn.cursor()
|
||||
existing_columns = {
|
||||
row[1] for row in cursor.execute("PRAGMA table_info(programs)")
|
||||
}
|
||||
for col, typedef in [
|
||||
("next_catchup_after", "REAL DEFAULT 0.0"),
|
||||
("total_pages", "INTEGER DEFAULT 0"),
|
||||
]:
|
||||
if col not in existing_columns:
|
||||
cursor.execute(f"ALTER TABLE programs ADD COLUMN {col} {typedef}")
|
||||
error_columns = {row[1] for row in cursor.execute("PRAGMA table_info(error_events)")}
|
||||
if "snapshot_path" not in error_columns:
|
||||
cursor.execute("ALTER TABLE error_events ADD COLUMN snapshot_path TEXT")
|
||||
|
||||
# Encrypt legacy plaintext values using the raw SQLite connection;
|
||||
# assigning the same decrypted value through the ORM is not dirty.
|
||||
rows = cursor.execute("SELECT key, value FROM config").fetchall()
|
||||
for key, value in rows:
|
||||
if value is not None and not value.startswith("enc:v1:"):
|
||||
cursor.execute("UPDATE config SET value = ? WHERE key = ?",
|
||||
(encrypt_config_value(value), key))
|
||||
|
||||
def get_db():
|
||||
db = SessionLocal()
|
||||
try:
|
||||
yield db
|
||||
finally:
|
||||
db.close()
|
||||
@@ -0,0 +1,48 @@
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
import hashlib
|
||||
from datetime import datetime
|
||||
|
||||
|
||||
class ScraperError(Exception):
|
||||
def __init__(self, station, code, message, url=None, status_code=None):
|
||||
self.station = station
|
||||
self.code = code
|
||||
self.url = url
|
||||
self.status_code = status_code
|
||||
details = f"[{code}] {message}"
|
||||
if status_code is not None:
|
||||
details += f" (HTTP {status_code})"
|
||||
if url:
|
||||
details += f" URL: {url}"
|
||||
super().__init__(details)
|
||||
|
||||
|
||||
def save_response_snapshot(station, code, url, content, status_code=None):
|
||||
directory = os.path.join("data", "diagnostics")
|
||||
try:
|
||||
os.makedirs(directory, exist_ok=True)
|
||||
digest = hashlib.sha256(f"{station}:{code}:{url}".encode("utf-8")).hexdigest()[:16]
|
||||
safe_station = re.sub(r"[^a-zA-Z0-9_-]", "_", station)
|
||||
base_path = os.path.join(directory, f"{safe_station}-{code}-{digest}")
|
||||
|
||||
with open(f"{base_path}.json", "w", encoding="utf-8") as metadata_file:
|
||||
json.dump({"station": station, "code": code, "url": url, "status_code": status_code,
|
||||
"updated_at": datetime.now().isoformat()}, metadata_file, ensure_ascii=False, indent=2)
|
||||
with open(f"{base_path}.html", "w", encoding="utf-8") as response_file:
|
||||
response_file.write((content or "")[:2_000_000])
|
||||
return base_path
|
||||
except OSError as exc:
|
||||
return f"snapshot unavailable: {exc}"
|
||||
|
||||
|
||||
def invalid_response(station, code, message, url, response=None, content=None):
|
||||
snapshot = save_response_snapshot(
|
||||
station,
|
||||
url,
|
||||
response.text if response is not None else content,
|
||||
response.status_code if response is not None else None,
|
||||
)
|
||||
raise ScraperError(station, code, f"{message}. Snapshot: {snapshot}", url,
|
||||
response.status_code if response is not None else None)
|
||||
@@ -0,0 +1,45 @@
|
||||
import ipaddress
|
||||
import socket
|
||||
from urllib.parse import urlparse
|
||||
import requests
|
||||
|
||||
|
||||
def validate_notification_url(url: str):
|
||||
parsed = urlparse(url or "")
|
||||
if parsed.scheme not in {"http", "https"} or not parsed.hostname or parsed.username or parsed.password:
|
||||
raise ValueError("Adres ntfy musi być adresem HTTP(S) bez danych logowania")
|
||||
try:
|
||||
addresses = {info[4][0] for info in socket.getaddrinfo(parsed.hostname, parsed.port, type=socket.SOCK_STREAM)}
|
||||
if any(ipaddress.ip_address(address).is_private or ipaddress.ip_address(address).is_loopback or
|
||||
ipaddress.ip_address(address).is_link_local for address in addresses):
|
||||
raise ValueError("Adres ntfy wskazuje na prywatną lub lokalną sieć")
|
||||
except socket.gaierror as exc:
|
||||
raise ValueError("Nie można rozpoznać hosta ntfy") from exc
|
||||
|
||||
def send_notification(url: str, topic: str, title: str, message: str, username: str = None, password: str = None):
|
||||
if not url or not topic:
|
||||
return False
|
||||
validate_notification_url(url)
|
||||
|
||||
full_url = f"{url.rstrip('/')}/{topic}"
|
||||
headers = {
|
||||
"Title": title.encode('utf-8')
|
||||
}
|
||||
|
||||
auth = None
|
||||
if username and password:
|
||||
auth = (username, password)
|
||||
|
||||
try:
|
||||
response = requests.post(
|
||||
full_url,
|
||||
data=message.encode('utf-8'),
|
||||
headers=headers,
|
||||
auth=auth,
|
||||
timeout=5
|
||||
)
|
||||
response.raise_for_status()
|
||||
return True
|
||||
except Exception as e:
|
||||
print(f"Failed to send ntfy notification: {e}")
|
||||
return False
|
||||
@@ -0,0 +1,111 @@
|
||||
from sqlalchemy.orm import Session
|
||||
from xml.etree import ElementTree as ET
|
||||
from xml.sax.saxutils import escape
|
||||
from email.utils import format_datetime
|
||||
from datetime import datetime
|
||||
from .database import Program, Episode
|
||||
|
||||
def get_base_url(request):
|
||||
# Retrieve base URL for constructing absolute paths (from FastAPI Request)
|
||||
return str(request.base_url).rstrip('/')
|
||||
|
||||
def generate_master_opml(db: Session, base_url: str, station: str = None):
|
||||
opml_outlines = []
|
||||
|
||||
# Radio 357
|
||||
r357 = db.query(Program).filter_by(station="radio357").order_by(Program.name).all() if not station or station == "radio357" else []
|
||||
if r357:
|
||||
opml_outlines.append(' <outline text="Radio 357">')
|
||||
for prog in r357:
|
||||
title = escape(prog.name, {'"': """, "'": "'"})
|
||||
xmlUrl = f"{base_url}/feeds/radio357/{prog.slug}.xml"
|
||||
opml_outlines.append(f' <outline text="{title}" type="rss" xmlUrl="{xmlUrl}" />')
|
||||
opml_outlines.append(' </outline>')
|
||||
|
||||
# RNŚ
|
||||
rns = db.query(Program).filter_by(station="rns").order_by(Program.name).all() if not station or station == "rns" else []
|
||||
if rns:
|
||||
opml_outlines.append(' <outline text="Radio Nowy Świat">')
|
||||
for prog in rns:
|
||||
title = escape(prog.name, {'"': """, "'": "'"})
|
||||
xmlUrl = f"{base_url}/feeds/rns/{prog.slug}.xml"
|
||||
opml_outlines.append(f' <outline text="{title}" type="rss" xmlUrl="{xmlUrl}" />')
|
||||
opml_outlines.append(' </outline>')
|
||||
|
||||
# Radio Jazz FM (Direct to source feeds!)
|
||||
jazz = db.query(Program).filter_by(station="jazz").order_by(Program.name).all() if not station or station == "jazz" else []
|
||||
if jazz:
|
||||
opml_outlines.append(' <outline text="Radio Jazz FM">')
|
||||
for prog in jazz:
|
||||
title = escape(prog.name, {'"': """, "'": "'"})
|
||||
# Link directly to their server!
|
||||
xmlUrl = f"https://podkasty.radiojazz.fm/@{prog.slug}/feed.xml"
|
||||
opml_outlines.append(f' <outline text="{title}" type="rss" xmlUrl="{xmlUrl}" />')
|
||||
opml_outlines.append(' </outline>')
|
||||
|
||||
xml_content = '<?xml version="1.0" encoding="UTF-8"?>\n'
|
||||
xml_content += '<opml version="1.0">\n <body>\n'
|
||||
xml_content += '\n'.join(opml_outlines) + '\n'
|
||||
xml_content += ' </body>\n</opml>'
|
||||
return xml_content
|
||||
|
||||
def generate_podcast_rss(db: Session, station: str, slug: str):
|
||||
prog = db.query(Program).filter_by(station=station, slug=slug).first()
|
||||
if not prog:
|
||||
return None
|
||||
|
||||
eps = db.query(Episode).filter_by(station=station, program_slug=slug, is_broken=False).filter(Episode.url.isnot(None)).order_by(Episode.pub_date.desc()).all()
|
||||
|
||||
rss = ET.Element("rss", {"version": "2.0", "xmlns:itunes": "http://www.itunes.com/dtds/podcast-1.0.dtd"})
|
||||
channel = ET.SubElement(rss, "channel")
|
||||
|
||||
ET.SubElement(channel, "title").text = prog.name
|
||||
ET.SubElement(channel, "description").text = prog.description or ""
|
||||
|
||||
if prog.image:
|
||||
ET.SubElement(channel, "itunes:image", {"href": prog.image})
|
||||
image_el = ET.SubElement(channel, "image")
|
||||
ET.SubElement(image_el, "url").text = prog.image
|
||||
ET.SubElement(image_el, "title").text = prog.name
|
||||
|
||||
for ep in eps:
|
||||
item = ET.SubElement(channel, "item")
|
||||
|
||||
date_prefix = f"({ep.pub_date[:10]}) " if ep.pub_date else ""
|
||||
ET.SubElement(item, "title").text = f"{date_prefix}{ep.title}"
|
||||
ET.SubElement(item, "itunes:author").text = ep.authors or ""
|
||||
|
||||
if ep.description:
|
||||
ET.SubElement(item, "description").text = ep.description
|
||||
ET.SubElement(item, "itunes:summary").text = ep.description
|
||||
|
||||
if ep.pub_date:
|
||||
try:
|
||||
# radio357 format: 2024-03-22T08:00:00+00:00 or rns: 2024-03-22
|
||||
if 'T' in ep.pub_date:
|
||||
dt = datetime.fromisoformat(ep.pub_date)
|
||||
else:
|
||||
dt = datetime.strptime(ep.pub_date, "%Y-%m-%d")
|
||||
ET.SubElement(item, "pubDate").text = format_datetime(dt)
|
||||
except:
|
||||
ET.SubElement(item, "pubDate").text = ep.pub_date
|
||||
|
||||
if ep.duration_secs:
|
||||
ET.SubElement(item, "itunes:duration").text = str(ep.duration_secs)
|
||||
|
||||
ET.SubElement(item, "guid", {"isPermaLink": "false"}).text = f"{station}-{ep.ep_id}"
|
||||
ET.SubElement(item, "enclosure", {"url": ep.url, "type": "audio/mpeg", "length": "0"})
|
||||
|
||||
if ep.image or prog.image:
|
||||
ET.SubElement(item, "itunes:image", {"href": ep.image or prog.image})
|
||||
|
||||
tree = ET.ElementTree(rss)
|
||||
try:
|
||||
ET.indent(tree, space=" ", level=0)
|
||||
except AttributeError:
|
||||
pass # older python versions
|
||||
|
||||
from io import BytesIO
|
||||
f = BytesIO()
|
||||
tree.write(f, encoding="utf-8", xml_declaration=True)
|
||||
return f.getvalue()
|
||||
@@ -0,0 +1,57 @@
|
||||
import requests
|
||||
from bs4 import BeautifulSoup
|
||||
from sqlalchemy.orm import Session
|
||||
from ..database import Program
|
||||
|
||||
BASE_URL = "https://podkasty.radiojazz.fm"
|
||||
|
||||
class JazzScraper:
|
||||
def __init__(self, db: Session, logger=print):
|
||||
self.db = db
|
||||
self.logger = logger
|
||||
|
||||
def sync_programs(self):
|
||||
self.logger("JAZZ: Pobieranie audycji...")
|
||||
try:
|
||||
r = requests.get(BASE_URL, headers={"User-Agent": "Mozilla/5.0"}, timeout=15)
|
||||
r.raise_for_status()
|
||||
html = r.text
|
||||
except Exception as e:
|
||||
self.logger(f"JAZZ: Błąd pobierania bazy: {e}")
|
||||
return
|
||||
|
||||
soup = BeautifulSoup(html, 'html.parser')
|
||||
links = soup.find_all('a', href=True)
|
||||
|
||||
found = 0
|
||||
for link in links:
|
||||
href = link['href']
|
||||
if '/@' in href:
|
||||
slug = href.split('/@')[-1].split('/')[0]
|
||||
if not slug: continue
|
||||
|
||||
text = link.get_text(strip=True)
|
||||
if not text or text.startswith('@') or "Recent activity" in text:
|
||||
continue
|
||||
|
||||
clean_title = text
|
||||
if clean_title.endswith(f"@{slug}"):
|
||||
clean_title = clean_title[:-len(f"@{slug}")].strip()
|
||||
if clean_title.startswith("Ż "):
|
||||
clean_title = clean_title[2:].strip()
|
||||
elif clean_title.startswith("Ż"):
|
||||
clean_title = clean_title[1:].strip()
|
||||
|
||||
prog = self.db.query(Program).filter_by(station="jazz", slug=slug).first()
|
||||
if not prog:
|
||||
self.db.add(Program(
|
||||
station="jazz", slug=slug, name=clean_title,
|
||||
description="Radio Jazz FM", image=""
|
||||
))
|
||||
found += 1
|
||||
|
||||
self.db.commit()
|
||||
self.logger(f"JAZZ: Znaleziono {found} nowych audycji (łącznie zaktualizowano).")
|
||||
|
||||
def run_full_sync(self):
|
||||
self.sync_programs()
|
||||
@@ -0,0 +1,282 @@
|
||||
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
|
||||
|
||||
self.logger(f"R357: Pobieram URL dla {ep.title[:30]}...")
|
||||
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
|
||||
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("R357: 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 Błąd HTTP {r.status_code} podczas pobierania audio: {r.text[:100]}")
|
||||
except Exception as e:
|
||||
self.logger(f"R357: 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: Pomyślnie pobrano {fetched} nowych linków.")
|
||||
|
||||
def run_full_sync(self, specific_program_slug=None):
|
||||
self.sync_episodes(specific_program_slug)
|
||||
self.backfill_urls(specific_program_slug)
|
||||
self.logger(f"R357: Koniec. Wysłano zapytania: {self.requests_made}/{self.hard_limit} (w tym backfill ograniczony do {self.backfill_limit}).")
|
||||
@@ -0,0 +1,508 @@
|
||||
import requests
|
||||
import time
|
||||
import os
|
||||
import re
|
||||
import urllib.parse
|
||||
from bs4 import BeautifulSoup
|
||||
from sqlalchemy.orm import Session
|
||||
from ..database import Program, Episode, Config
|
||||
from ..diagnostics import ScraperError, invalid_response
|
||||
|
||||
MAX_REQUESTS_LIMIT = 50
|
||||
BASE_URL = "https://nowyswiat.online"
|
||||
|
||||
|
||||
def get_cookie_for_url(cookie_jar, name, url):
|
||||
"""Wybiera najbardziej szczegółowe cookie, gdy jar zawiera duplikaty nazwy."""
|
||||
parsed_url = urllib.parse.urlparse(url)
|
||||
host = parsed_url.hostname.lower()
|
||||
path = parsed_url.path or "/"
|
||||
candidates = []
|
||||
for cookie in cookie_jar:
|
||||
if cookie.name != name or (cookie.secure and parsed_url.scheme != "https"):
|
||||
continue
|
||||
domain = cookie.domain.lstrip(".").lower()
|
||||
if domain and not (host == domain or host.endswith(f".{domain}")):
|
||||
continue
|
||||
cookie_path = cookie.path or "/"
|
||||
if not path.startswith(cookie_path.rstrip("/") or "/"):
|
||||
continue
|
||||
candidates.append(cookie)
|
||||
|
||||
if not candidates:
|
||||
return None
|
||||
selected = max(candidates, key=lambda cookie: (len(cookie.path or "/"), len(cookie.domain or "")))
|
||||
for cookie in candidates:
|
||||
if cookie is not selected:
|
||||
cookie_jar.clear(cookie.domain, cookie.path, cookie.name)
|
||||
return selected.value
|
||||
|
||||
def parse_polish_date(date_str):
|
||||
months = {
|
||||
"stycznia": "01", "lutego": "02", "marca": "03", "kwietnia": "04",
|
||||
"maja": "05", "czerwca": "06", "lipca": "07", "sierpnia": "08",
|
||||
"września": "09", "października": "10", "listopada": "11", "grudnia": "12",
|
||||
"styczeń": "01", "luty": "02", "marzec": "03", "kwiecień": "04",
|
||||
"maj": "05", "czerwiec": "06", "lipiec": "07", "sierpień": "08",
|
||||
"wrzesień": "09", "październik": "10", "listopad": "11", "grudzień": "12"
|
||||
}
|
||||
parts = date_str.lower().split()
|
||||
if len(parts) == 3:
|
||||
day = parts[0].zfill(2)
|
||||
month = months.get(parts[1], "01")
|
||||
year = parts[2]
|
||||
return f"{year}-{month}-{day}"
|
||||
return date_str
|
||||
|
||||
class RNScraper:
|
||||
def __init__(self, db: Session, logger=print, progress_callback=None, warning_callback=None, stop_flag=None):
|
||||
self.db = db
|
||||
self.logger = logger
|
||||
self.progress_callback = progress_callback
|
||||
self.warning_callback = warning_callback
|
||||
self.stop_flag = stop_flag
|
||||
self.requests_made = 0
|
||||
self.start_time = time.time()
|
||||
self.max_execution_time = 900 # 15 minut max na cały sync
|
||||
limit_conf = self.db.query(Config).filter_by(key="rns_limit").first()
|
||||
self.backfill_limit = int(limit_conf.value) if limit_conf and limit_conf.value.isdigit() else 50
|
||||
hard_limit_conf = self.db.query(Config).filter_by(key="rns_hard_limit").first()
|
||||
self.hard_limit = int(hard_limit_conf.value) if hard_limit_conf and hard_limit_conf.value.isdigit() else 500
|
||||
self.session = requests.Session()
|
||||
self.session.headers.update({"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64)"})
|
||||
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_html(self, url):
|
||||
self.check_timeout()
|
||||
if self.requests_made >= self.hard_limit:
|
||||
return None
|
||||
for attempt in range(1, 4):
|
||||
if self.requests_made >= self.hard_limit:
|
||||
return None
|
||||
self.requests_made += 1
|
||||
if self.progress_callback: self.progress_callback(self.requests_made, self.hard_limit)
|
||||
try:
|
||||
r = self.session.get(url, timeout=15)
|
||||
if r.status_code == 200:
|
||||
return r.text
|
||||
invalid_response("rns", "HTTP_ERROR", "Błąd pobierania strony", url, r)
|
||||
except requests.exceptions.Timeout as e:
|
||||
if attempt == 3:
|
||||
raise ScraperError("rns", "NETWORK_TIMEOUT", "Trzy próby pobrania strony zakończyły się timeoutem", url) from e
|
||||
self.logger(f"RNŚ: Timeout dla {url}, próba {attempt}/3. Ponawiam za 30 sekund.")
|
||||
self.check_timeout()
|
||||
time.sleep(30)
|
||||
except Exception as e:
|
||||
if isinstance(e, ScraperError):
|
||||
raise
|
||||
raise ScraperError("rns", "NETWORK_ERROR", f"Błąd zapytania: {e}", url) from e
|
||||
return None
|
||||
|
||||
def perform_login(self):
|
||||
email = self.db.query(Config).filter_by(key="rns_email").first()
|
||||
password = self.db.query(Config).filter_by(key="rns_password").first()
|
||||
if not email or not password:
|
||||
self.logger("RNŚ: Brak danych logowania w bazie.")
|
||||
return False
|
||||
|
||||
self.logger("RNŚ: Inicjalizacja logowania (pobieranie CSRF)...")
|
||||
try:
|
||||
r1 = self.session.get("https://nowyswiat.online/konto/zaloguj", timeout=15)
|
||||
except Exception as e:
|
||||
self.logger(f"RNŚ: Sieć zablokowała pobieranie CSRF: {e}")
|
||||
return False
|
||||
|
||||
csrf_token = get_cookie_for_url(
|
||||
self.session.cookies,
|
||||
"csrf_cookie_neocms",
|
||||
"https://nowyswiat.online/konto/zaloguj",
|
||||
)
|
||||
if not csrf_token:
|
||||
self.logger("RNŚ: Nie udało się pobrać tokenu CSRF.")
|
||||
return False
|
||||
|
||||
self.logger("RNŚ: Wysyłanie formularza...")
|
||||
payload = {"csrf_neocms": csrf_token, "login": email.value, "password": password.value, "ufd_data": "{}"}
|
||||
try:
|
||||
r2 = self.session.post("https://nowyswiat.online/konto/zaloguj", data=payload, headers={
|
||||
"X-Requested-With": "XMLHttpRequest"
|
||||
}, timeout=15)
|
||||
except Exception as e:
|
||||
self.logger(f"RNŚ: Błąd sieci przy wysyłaniu formularza: {e}")
|
||||
return False
|
||||
|
||||
# Check if login succeeded by looking for a session cookie or a success response
|
||||
if '"status":"OK"' in r2.text or "logowanie udane" in r2.text.lower():
|
||||
# Save cookies to DB (serialize)
|
||||
cookies_dict = requests.utils.dict_from_cookiejar(self.session.cookies)
|
||||
import json
|
||||
cookie_conf = self.db.query(Config).filter_by(key="rns_cookies").first()
|
||||
if not cookie_conf:
|
||||
self.db.add(Config(key="rns_cookies", value=json.dumps(cookies_dict)))
|
||||
else:
|
||||
cookie_conf.value = json.dumps(cookies_dict)
|
||||
self.db.commit()
|
||||
return True
|
||||
else:
|
||||
self.logger(f"RNŚ: Błędne dane logowania (lub zmiana mechanizmu). Odpowiedź: {r2.text[:100]}")
|
||||
return False
|
||||
|
||||
def ensure_auth(self):
|
||||
cookie_conf = self.db.query(Config).filter_by(key="rns_cookies").first()
|
||||
if cookie_conf:
|
||||
import json
|
||||
try:
|
||||
cookies_dict = json.loads(cookie_conf.value)
|
||||
self.session.cookies = requests.utils.cookiejar_from_dict(cookies_dict)
|
||||
html = self._get_html(f"{BASE_URL}/")
|
||||
if html and "wyloguj" in html.lower():
|
||||
return True
|
||||
self.logger("RNŚ: Zapisana sesja wygasła, wykonuję ponowne logowanie.")
|
||||
except:
|
||||
pass
|
||||
return self.perform_login()
|
||||
|
||||
def check_auth_status(self):
|
||||
cookie_conf = self.db.query(Config).filter_by(key="rns_cookies").first()
|
||||
if not cookie_conf:
|
||||
return False, "Brak zapisanych ciasteczek."
|
||||
import json
|
||||
try:
|
||||
self.session.cookies = requests.utils.cookiejar_from_dict(json.loads(cookie_conf.value))
|
||||
html = self._get_html("https://nowyswiat.online/")
|
||||
if html and "wyloguj" in html.lower():
|
||||
return True, "Ciasteczka aktywne, sesja poprawna."
|
||||
return False, "Ciasteczka nieaktywne lub wygasły (brak dostępu do profilu)."
|
||||
except Exception as e:
|
||||
return False, f"Błąd: {e}"
|
||||
|
||||
def update_programs(self):
|
||||
self.logger("RNŚ: Aktualizacja listy programów...")
|
||||
html = self._get_html("https://nowyswiat.online/podcasty")
|
||||
if not html: return
|
||||
soup = BeautifulSoup(html, "html.parser")
|
||||
links = soup.find_all("a", href=True)
|
||||
program_links = [link for link in links if "rbroadcast=" in link["href"]]
|
||||
if not program_links:
|
||||
invalid_response("rns", "HTML_SCHEMA_CHANGED", "Nie znaleziono programów w stronie podcastów", "https://nowyswiat.online/podcasty", content=html)
|
||||
for link in program_links:
|
||||
href = link["href"]
|
||||
if "rbroadcast=" in href:
|
||||
parsed = urllib.parse.urlparse(href)
|
||||
slug = urllib.parse.parse_qs(parsed.query).get("rbroadcast", [None])[0]
|
||||
if not slug: continue
|
||||
|
||||
title_el = link.find("h2", class_="rns-search-dropdown-title")
|
||||
title = title_el.text.strip() if title_el else slug
|
||||
|
||||
img_el = link.find("img")
|
||||
img_url = img_el["src"] if img_el and img_el.has_attr("src") else ""
|
||||
if img_url and not img_url.startswith("http"): img_url = f"{BASE_URL}/{img_url.lstrip('/')}"
|
||||
|
||||
prog = self.db.query(Program).filter_by(station="rns", slug=slug).first()
|
||||
if not prog:
|
||||
self.db.add(Program(
|
||||
station="rns", slug=slug, name=title, image=img_url,
|
||||
description=f"Radio Nowy Świat: {title}", backfill_page=2, backfill_complete=False
|
||||
))
|
||||
self.db.commit()
|
||||
|
||||
def _verify_audio_teaser(self, ep: Episode):
|
||||
"""Wykrywa 1-minutowy teaser (rozmiar mniejszy niż ~2MB, mimo że audycja trwa > 5 min)."""
|
||||
self.check_timeout()
|
||||
if not ep.url: return False
|
||||
|
||||
if self.requests_made >= self.hard_limit: return False
|
||||
self.requests_made += 1
|
||||
if self.progress_callback: self.progress_callback(self.requests_made, self.hard_limit)
|
||||
|
||||
try:
|
||||
r = self.session.head(ep.url, timeout=5)
|
||||
cl = int(r.headers.get("Content-Length", 0))
|
||||
if cl > 0 and cl < 2_500_000 and ep.duration_secs > 300:
|
||||
self.logger(f"RNŚ: Wykryto uszkodzony link (Teaser 1-min) dla {ep.title}. Usuwam URL.")
|
||||
ep.url = None
|
||||
ep.is_broken = True
|
||||
return True
|
||||
except Exception as e:
|
||||
pass
|
||||
return False
|
||||
|
||||
def fetch_program_page(self, program_slug, page):
|
||||
url = f"https://nowyswiat.online/podcasty?rbroadcast={program_slug}&page={page}"
|
||||
html = self._get_html(url)
|
||||
if not html:
|
||||
raise Exception(f"Błąd sieci podczas pobierania strony {page} audycji {program_slug}")
|
||||
|
||||
soup = BeautifulSoup(html, "html.parser")
|
||||
cards = soup.find_all("a", class_="rns-grid-podcast-card")
|
||||
known_episode_count = self.db.query(Episode).filter_by(
|
||||
station="rns", program_slug=program_slug
|
||||
).count()
|
||||
page_has_podcast_content = bool(soup.find(string=re.compile(r"podcast|podkast", re.IGNORECASE)))
|
||||
if not cards and page > 1:
|
||||
return 0, 0, page - 1
|
||||
if not cards and (known_episode_count or not page_has_podcast_content):
|
||||
invalid_response("rns", "HTML_SCHEMA_CHANGED", "Nie znaleziono kart odcinków dla istniejącego programu", url, content=html)
|
||||
|
||||
new_found = 0
|
||||
missing_player_count = 0
|
||||
valid_cards = 0
|
||||
|
||||
for card in cards:
|
||||
href = card.get("href", "")
|
||||
if not href or "/podcasty/" not in href: continue
|
||||
valid_cards += 1
|
||||
ep_id = href.split("/podcasty/")[-1].split("?")[0]
|
||||
|
||||
player_box = card.find("div", class_="rns-play-btn-box")
|
||||
audio_url = player_box.get("data-neo-player-src") if player_box else None
|
||||
|
||||
if not player_box:
|
||||
missing_player_count += 1
|
||||
|
||||
if audio_url and not audio_url.startswith("http"):
|
||||
audio_url = f"{BASE_URL}/{audio_url.lstrip('/')}"
|
||||
|
||||
title_el = card.find("p", class_="rns-post-title")
|
||||
raw_title = player_box.get("data-neo-player-title", "") if player_box else (title_el.text.strip() if title_el else ep_id)
|
||||
title = BeautifulSoup(raw_title, "html.parser").text.strip() if raw_title else ep_id
|
||||
|
||||
raw_subtitle = player_box.get("data-neo-player-subtitle", "") if player_box else ""
|
||||
authors = BeautifulSoup(raw_subtitle, "html.parser").text.strip() if raw_subtitle else "Radio Nowy Świat"
|
||||
|
||||
date_el = card.find("p", class_="rns-podcast-details-date")
|
||||
pub_date = parse_polish_date(date_el.text.strip()) if date_el else ""
|
||||
|
||||
time_el = card.find("p", class_="rns-podcast-details-long")
|
||||
duration_str = time_el.text.strip() if time_el else "0:00"
|
||||
duration_secs = 0
|
||||
if ":" in duration_str:
|
||||
parts = duration_str.split(":")
|
||||
if len(parts) == 3: duration_secs = int(parts[0])*3600 + int(parts[1])*60 + int(parts[2])
|
||||
elif len(parts) == 2: duration_secs = int(parts[0])*60 + int(parts[1])
|
||||
|
||||
img_url = player_box.get("data-neo-player-img", "") if player_box else ""
|
||||
if img_url and not img_url.startswith("http"): img_url = f"{BASE_URL}/{img_url.lstrip('/')}"
|
||||
|
||||
desc_el = card.find("p", class_="rns-post-card-desc")
|
||||
description = desc_el.get_text(separator="\n").strip() if desc_el else ""
|
||||
|
||||
ep = self.db.query(Episode).filter_by(station="rns", ep_id=ep_id).first()
|
||||
if not ep:
|
||||
ep = Episode(
|
||||
station="rns", program_slug=program_slug, ep_id=ep_id,
|
||||
title=title, authors=authors, url=audio_url, image=img_url,
|
||||
pub_date=pub_date, duration_secs=duration_secs, description=description,
|
||||
is_broken=False
|
||||
)
|
||||
self.db.add(ep)
|
||||
new_found += 1
|
||||
else:
|
||||
if not ep.url and audio_url:
|
||||
ep.url = audio_url
|
||||
ep.is_broken = False
|
||||
# Celowo nie zwiększamy new_found, aby Daily Catchup nie wchodził w tryb głębokiego archiwum.
|
||||
# Łataniem starych dziur zajmie się faza Backfill.
|
||||
|
||||
if cards and not valid_cards:
|
||||
invalid_response("rns", "HTML_SCHEMA_CHANGED", "Znaleziono karty podcastów bez oczekiwanych identyfikatorów", url, content=html)
|
||||
|
||||
if missing_player_count:
|
||||
warning = f"{missing_player_count} odcinków bez dostępnego playera audio; pomijam ich URL-e."
|
||||
if self.warning_callback:
|
||||
self.warning_callback(warning)
|
||||
else:
|
||||
self.logger(f"RNŚ: WARNING: {warning}")
|
||||
|
||||
self.db.commit()
|
||||
|
||||
page_links = soup.find_all("a", href=re.compile(r"page=\d+"))
|
||||
max_page = page
|
||||
for link in page_links:
|
||||
href = link.get("href")
|
||||
m = re.search(r"page=(\d+)", href)
|
||||
if m:
|
||||
max_page = max(max_page, int(m.group(1)))
|
||||
|
||||
return new_found, len(cards), max_page
|
||||
|
||||
def _compute_next_catchup(self, prog: Program) -> float:
|
||||
"""Wylicza kiedy najwcześniej warto znowu sprawdzać tę audycję."""
|
||||
episodes = self.db.query(Episode).filter_by(
|
||||
station="rns", program_slug=prog.slug
|
||||
).order_by(Episode.pub_date.desc()).limit(20).all()
|
||||
|
||||
if len(episodes) < 3:
|
||||
# Za mało danych – sprawdzamy przy każdym uruchomieniu
|
||||
return 0.0
|
||||
|
||||
# Wylicz średnią przerwę między odcinkami w dniach
|
||||
import re
|
||||
dates = []
|
||||
for ep in episodes:
|
||||
if ep.pub_date and re.match(r'\d{4}-\d{2}-\d{2}', ep.pub_date):
|
||||
try:
|
||||
import datetime
|
||||
dates.append(datetime.date.fromisoformat(ep.pub_date[:10]))
|
||||
except ValueError:
|
||||
pass
|
||||
|
||||
if len(dates) < 3:
|
||||
return 0.0
|
||||
|
||||
dates.sort(reverse=True)
|
||||
gaps = [(dates[i] - dates[i+1]).days for i in range(len(dates)-1)]
|
||||
avg_gap = sum(gaps) / len(gaps)
|
||||
|
||||
if avg_gap < 2:
|
||||
# Codziennie lub częściej → zawsze sprawdzamy (brak cooldownu)
|
||||
return 0.0
|
||||
elif avg_gap <= 8:
|
||||
# Tygodniowo → cooldown = 70% cyklu (żeby sprawdzić przed następnym odcinkiem)
|
||||
cooldown_days = avg_gap * 0.7
|
||||
else:
|
||||
# Rzadziej niż tygodniowo → max 7 dni
|
||||
cooldown_days = 7.0
|
||||
|
||||
return time.time() + cooldown_days * 86400
|
||||
|
||||
def phase_1_catchup(self, specific_program_slug=None):
|
||||
self.logger("RNŚ: Phase 1 (Daily Catchup)...")
|
||||
query = self.db.query(Program).filter_by(station="rns")
|
||||
if specific_program_slug:
|
||||
query = query.filter_by(slug=specific_program_slug)
|
||||
else:
|
||||
query = query.order_by(Program.last_catchup.asc())
|
||||
|
||||
catchup_limit = max(1, self.hard_limit - self.backfill_limit)
|
||||
skipped = 0
|
||||
|
||||
timeout_program = None
|
||||
for prog in query.all():
|
||||
if self.requests_made >= catchup_limit:
|
||||
self.logger(f"RNŚ: [{self.requests_made}/{self.hard_limit}] Limit Catchup osiągnięty. Zostawiam resztę dla Backfill.")
|
||||
break
|
||||
|
||||
# Pomiń jeśli za wcześnie (adaptive cooldown)
|
||||
if not specific_program_slug and prog.next_catchup_after and time.time() < prog.next_catchup_after:
|
||||
skipped += 1
|
||||
continue
|
||||
|
||||
is_first_sync = (prog.last_catchup == 0.0)
|
||||
|
||||
page = 1
|
||||
catchup_succeeded = False
|
||||
while True:
|
||||
if self.requests_made >= catchup_limit: break
|
||||
self.logger(f"RNŚ: [{self.requests_made}/{self.hard_limit}] Sprawdzam {prog.slug} (strona {page})")
|
||||
try:
|
||||
new_found, total_cards, max_page = self.fetch_program_page(prog.slug, page)
|
||||
except ScraperError as exc:
|
||||
if exc.code != "NETWORK_TIMEOUT":
|
||||
raise
|
||||
if timeout_program:
|
||||
raise ScraperError(
|
||||
"rns", "NETWORK_TIMEOUT",
|
||||
f"Timeout także dla kolejnej audycji {prog.slug}; poprzednia: {timeout_program}",
|
||||
exc.url,
|
||||
) from exc
|
||||
timeout_program = prog.slug
|
||||
self.logger(f"RNŚ: Pomijam {prog.slug} po trzech timeoutach i sprawdzam następną audycję.")
|
||||
break
|
||||
|
||||
if timeout_program:
|
||||
warning = f"Poprzednia audycja {timeout_program} miała trzy timeouty; kolejna audycja odpowiada poprawnie."
|
||||
if self.warning_callback:
|
||||
self.warning_callback(warning)
|
||||
else:
|
||||
self.logger(f"RNŚ: WARNING: {warning}")
|
||||
timeout_program = None
|
||||
|
||||
if max_page > prog.total_pages:
|
||||
prog.total_pages = max_page
|
||||
|
||||
if total_cards == 0 or new_found == 0:
|
||||
catchup_succeeded = True
|
||||
break
|
||||
|
||||
if is_first_sync:
|
||||
catchup_succeeded = True
|
||||
break
|
||||
|
||||
if page >= max_page:
|
||||
catchup_succeeded = True
|
||||
break
|
||||
|
||||
page += 1
|
||||
time.sleep(0.5)
|
||||
|
||||
if catchup_succeeded:
|
||||
prog.last_catchup = time.time()
|
||||
prog.next_catchup_after = self._compute_next_catchup(prog)
|
||||
self.db.commit()
|
||||
|
||||
if skipped:
|
||||
self.logger(f"RNŚ: Pominięto {skipped} audycji (cooldown adaptacyjny).")
|
||||
|
||||
def phase_2_backfill(self, specific_program_slug=None):
|
||||
backfill_count = 0
|
||||
if specific_program_slug:
|
||||
prog = self.db.query(Program).filter_by(station="rns", slug=specific_program_slug).first()
|
||||
if prog:
|
||||
page = prog.backfill_page
|
||||
self.logger(f"RNŚ: [{self.requests_made}/{self.hard_limit}] Backfill manualny dla {prog.slug} (od strony {page})")
|
||||
while self.requests_made < self.hard_limit and backfill_count < self.backfill_limit:
|
||||
new_found, total, max_p = self.fetch_program_page(prog.slug, page)
|
||||
backfill_count += 1
|
||||
if max_p > prog.total_pages:
|
||||
prog.total_pages = max_p
|
||||
self.logger(f"RNŚ: [{self.requests_made}/{self.hard_limit}] Backfill {prog.slug} strona {page} z {prog.total_pages or '?'} ({new_found} nowych)")
|
||||
if total == 0 or page >= max_p:
|
||||
prog.backfill_complete = True
|
||||
self.logger(f"RNŚ: Backfill {prog.slug} – zakończono.")
|
||||
break
|
||||
prog.backfill_page = page + 1
|
||||
page += 1
|
||||
self.db.commit()
|
||||
time.sleep(0.5)
|
||||
return
|
||||
|
||||
while self.requests_made < self.hard_limit and backfill_count < self.backfill_limit:
|
||||
prog = self.db.query(Program).filter_by(station="rns", backfill_complete=False).order_by(Program.backfill_page.asc()).first()
|
||||
if not prog:
|
||||
self.logger("RNŚ: Wszystkie audycje w pełni uzupełnione!")
|
||||
break
|
||||
|
||||
page = prog.backfill_page
|
||||
self.logger(f"RNŚ: [{self.requests_made}/{self.hard_limit}] Backfill dla {prog.slug} (strona {page} z {prog.total_pages or '?'})")
|
||||
new_found, total_cards, max_p = self.fetch_program_page(prog.slug, page)
|
||||
backfill_count += 1
|
||||
|
||||
if max_p > prog.total_pages:
|
||||
prog.total_pages = max_p
|
||||
|
||||
if total_cards == 0 or page >= max_p:
|
||||
prog.backfill_complete = True
|
||||
else:
|
||||
prog.backfill_page = page + 1
|
||||
self.db.commit()
|
||||
time.sleep(0.5)
|
||||
|
||||
def run_full_sync(self, specific_program_slug=None):
|
||||
if not self.ensure_auth():
|
||||
self.logger("RNŚ: Błąd logowania przed startem!")
|
||||
raise Exception("RNS Login Failed")
|
||||
|
||||
self.update_programs()
|
||||
self.phase_1_catchup(specific_program_slug)
|
||||
if self.requests_made < self.hard_limit:
|
||||
self.phase_2_backfill(specific_program_slug)
|
||||
|
||||
self.logger(f"RNŚ: Koniec. Wysłano zapytania: {self.requests_made}/{self.hard_limit} (w tym backfill ograniczony do {self.backfill_limit}).")
|
||||
@@ -0,0 +1,45 @@
|
||||
import base64
|
||||
import hashlib
|
||||
import os
|
||||
|
||||
from cryptography.fernet import Fernet, InvalidToken
|
||||
|
||||
_ENCRYPTED_PREFIX = "enc:v1:"
|
||||
|
||||
|
||||
def get_session_secret() -> str:
|
||||
secret = os.environ.get("RADIOSYNC_SECRET_KEY")
|
||||
if not secret:
|
||||
raise RuntimeError("RADIOSYNC_SECRET_KEY must be set")
|
||||
return secret
|
||||
|
||||
|
||||
def _get_fernet() -> Fernet:
|
||||
configured_key = os.environ.get("RADIOSYNC_CONFIG_KEY")
|
||||
if configured_key:
|
||||
try:
|
||||
return Fernet(configured_key.encode("ascii"))
|
||||
except Exception as exc:
|
||||
raise RuntimeError("RADIOSYNC_CONFIG_KEY must be a valid Fernet key") from exc
|
||||
|
||||
digest = hashlib.sha256(get_session_secret().encode("utf-8")).digest()
|
||||
return Fernet(base64.urlsafe_b64encode(digest))
|
||||
|
||||
|
||||
def encrypt_config_value(value: str) -> str:
|
||||
if value is None:
|
||||
return None
|
||||
if value.startswith(_ENCRYPTED_PREFIX):
|
||||
return value
|
||||
token = _get_fernet().encrypt(value.encode("utf-8")).decode("ascii")
|
||||
return f"{_ENCRYPTED_PREFIX}{token}"
|
||||
|
||||
|
||||
def decrypt_config_value(value: str) -> str:
|
||||
if value is None or not value.startswith(_ENCRYPTED_PREFIX):
|
||||
return value
|
||||
try:
|
||||
token = value[len(_ENCRYPTED_PREFIX):].encode("ascii")
|
||||
return _get_fernet().decrypt(token).decode("utf-8")
|
||||
except InvalidToken as exc:
|
||||
raise RuntimeError("Cannot decrypt config value with the configured key") from exc
|
||||
Reference in New Issue
Block a user