From 05eb6dd545dc9863693ec124951d3521f8fdf6ff Mon Sep 17 00:00:00 2001 From: mhdzumair Date: Sat, 18 Jan 2025 00:09:08 +0530 Subject: [PATCH] Refactor DLHD Live Sports Events scraping and storing the data Reworked the caching logic to handle a missing cache key for "events" type gracefully. Integrated a new DLHD schedule service for improved event metadata and stream handling, leveraging fuzzy matching and IPTV-org data retrieval for enhanced TV channel mappings. Updated scraper configuration, enhanced genre handling, and added support for dynamic event descriptions. --- api/main.py | 23 +- db/crud.py | 33 +-- mediafusion_scrapy/spiders/dlhd.py | 301 ++++++++++++++++++-------- resources/json/scraper_config.json | 95 ++++---- resources/templates/manifest.json.j2 | 2 +- scrapers/dlhd.py | 311 +++++++++++++++++++++++++++ 6 files changed, 605 insertions(+), 160 deletions(-) create mode 100644 scrapers/dlhd.py diff --git a/api/main.py b/api/main.py index ab40251..1a4cf5e 100644 --- a/api/main.py +++ b/api/main.py @@ -505,11 +505,12 @@ async def search_meta( ) @wrappers.auth_required async def get_meta( + response: Response, catalog_type: Literal["movie", "series", "tv", "events"], meta_id: str, user_data: schemas.UserData = Depends(get_user_data), ): - cache_key = f"{catalog_type}_{meta_id}_meta" + cache_key = f"{catalog_type}_{meta_id}_meta" if catalog_type != "events" else None if catalog_type in ["movie", "series"]: cache_key += "_" + "_".join( @@ -517,13 +518,16 @@ async def get_meta( ) # Try retrieving the cached data - cached_data = await REDIS_ASYNC_CLIENT.get(cache_key) - if cached_data: - try: - meta_data = schemas.MetaItem.model_validate_json(cached_data) - return await update_rpdb_poster(meta_data, user_data, catalog_type) - except ValidationError: - pass + if cache_key: + cached_data = await REDIS_ASYNC_CLIENT.get(cache_key) + if cached_data: + try: + meta_data = schemas.MetaItem.model_validate_json(cached_data) + return await update_rpdb_poster(meta_data, user_data, catalog_type) + except ValidationError: + pass + else: + response.headers.update(const.NO_CACHE_HEADERS) if catalog_type == "movie": if meta_id.startswith("dl"): @@ -541,7 +545,8 @@ async def get_meta( # Cache the data with a TTL of 30 minutes # If the data is not found, cached the empty data to avoid db query. - await REDIS_ASYNC_CLIENT.set(cache_key, json.dumps(data, default=str), ex=1800) + if cache_key: + await REDIS_ASYNC_CLIENT.set(cache_key, json.dumps(data, default=str), ex=1800) if not data: raise HTTPException(status_code=404, detail="Meta ID not found.") diff --git a/db/crud.py b/db/crud.py index beaf5e0..b660d8d 100644 --- a/db/crud.py +++ b/db/crud.py @@ -30,6 +30,7 @@ from db.models import ( ) from db.redis_database import REDIS_ASYNC_CLIENT from db.schemas import Stream, TorrentStreamsList +from scrapers.dlhd import dlhd_schedule_service from scrapers.mdblist import initialize_mdblist_scraper from scrapers.scraper_tasks import run_scrapers, meta_fetcher from streaming_providers.cache_helpers import store_cached_info_hashes @@ -1436,31 +1437,9 @@ async def save_events_data(metadata: dict) -> str: async def get_events_meta_list(genre=None, skip=0, limit=25) -> list[schemas.Meta]: - if genre: - key_pattern = f"events:genre:{genre}" - else: - key_pattern = "events:all" - - # Fetch event keys sorted by timestamp in descending order - events_keys = await REDIS_ASYNC_CLIENT.zrevrange( - key_pattern, skip, skip + limit - 1 + return await dlhd_schedule_service.get_scheduled_events( + genre=genre, skip=skip, limit=limit ) - events = [] - - # Iterate over event keys, fetching and decoding JSON data - for key in events_keys: - events_json = await REDIS_ASYNC_CLIENT.get(key) - if events_json: - meta_data = schemas.Meta.model_validate_json(events_json) - meta_data.poster = ( - f"{settings.poster_host_url}/poster/events/{meta_data.id}.jpg" - ) - events.append(meta_data) - else: - # Cleanup: Remove expired or missing event key from the index - await REDIS_ASYNC_CLIENT.zrem(key_pattern, key) - - return events async def get_event_meta(meta_id: str) -> dict: @@ -1469,7 +1448,11 @@ async def get_event_meta(meta_id: str) -> dict: if not events_json: return {} - event_data = MediaFusionTVMetaData.model_validate_json(events_json) + event_data = MediaFusionEventsMetaData.model_validate_json(events_json) + if event_data.event_start_timestamp: + # Update description with localized time + event_data.description = f"🎬 {event_data.title} - ⏰ {dlhd_schedule_service.format_event_time(event_data.event_start_timestamp)}" + return { "meta": { "_id": meta_id, diff --git a/mediafusion_scrapy/spiders/dlhd.py b/mediafusion_scrapy/spiders/dlhd.py index 3d3a4b1..5a05c3e 100644 --- a/mediafusion_scrapy/spiders/dlhd.py +++ b/mediafusion_scrapy/spiders/dlhd.py @@ -1,111 +1,234 @@ import logging -import random -from datetime import datetime, timedelta +from datetime import datetime -import pytz import scrapy -from dateutil import parser as date_parser +from thefuzz import fuzz from utils.config import config_manager -from utils.runtime_const import SPORTS_ARTIFACTS -class DaddyLiveHDSpider(scrapy.Spider): +class DaddyLiveHDChannelsSpider(scrapy.Spider): name = "dlhd" - start_urls = [config_manager.get_scraper_config(name, "schedule_url")] - - # The number of hours to consider the event as starting within next hours from now. - start_within_next_hours = 1 - started_within_hours_ago = 6 - - category_map = config_manager.get_scraper_config(name, "category_mapping") - m3u8_base_url = config_manager.get_scraper_config(name, "m3u8_base_url") - referer = config_manager.get_scraper_config(name, "referer") - gmt = pytz.timezone("Etc/GMT") + iptv_org_api = "https://iptv-org.github.io/api/channels.json" custom_settings = { "ITEM_PIPELINES": { "mediafusion_scrapy.pipelines.LiveStreamResolverPipeline": 100, - "mediafusion_scrapy.pipelines.LiveEventStorePipeline": 300, + "mediafusion_scrapy.pipelines.TVStorePipeline": 300, }, - "DUPEFILTER_DEBUG": True, } + def __init__(self, *args, **kwargs): + super().__init__(*args, **kwargs) + self.start_urls = [config_manager.get_scraper_config(self.name, "channels_url")] + self.m3u8_base_url = config_manager.get_scraper_config( + self.name, "m3u8_base_url" + ) + self.referer = config_manager.get_scraper_config(self.name, "referer") + self.category_map = config_manager.get_scraper_config( + self.name, "category_mapping" + ) + self.channels_data = {} # Will store IPTV-org channels data + self.min_match_ratio = 85 # Minimum ratio for fuzzy matching + def start_requests(self): - yield scrapy.Request(self.start_urls[0], self.parse) + # First fetch IPTV-org channels data + yield scrapy.Request( + self.iptv_org_api, + callback=self.parse_iptv_org_data, + headers={"Accept": "application/json"}, + ) - def parse(self, response, **kwargs): - data = response.json() - current_time = datetime.now(tz=self.gmt) - for date_section, sports in data.items(): - date_str = date_section.split(" - ")[0] - event_date = date_parser.parse(date_str).date() - for sport, events in sports.items(): - for event in events: - time = datetime.strptime(event["time"], "%H:%M").time() - datetime_obj = datetime.combine(event_date, time) - # Make the datetime object timezone aware (UK GMT) - aware_datetime = self.gmt.localize(datetime_obj) + def parse_iptv_org_data(self, response): + """Parse IPTV-org channels data and store it.""" + channels = response.json() + current_time = datetime.now() - # Check if event starts within the specified time range - time_difference = aware_datetime - current_time - if not ( - timedelta(hours=-self.started_within_hours_ago) - <= time_difference - <= timedelta(hours=self.start_within_next_hours) - ): - logging.warning( - "Skipping event %s as it doesn't start within the specified time range. %s", - event["event"], - time_difference, - ) - continue + # Filter out closed channels and store active ones + for channel in channels: + # Skip closed or NSFW channels + if channel.get("is_nsfw") or channel.get("closed"): + continue - # Convert to UNIX timestamp - event_start_timestamp = int(aware_datetime.timestamp()) - if event_start_timestamp != 0: - event_start_time = datetime.fromtimestamp( - event_start_timestamp - ).strftime("%I:%M%p GMT") - description = f'{event["event"]} - {event_start_time}' - else: - description = event["event"] - category = self.category_map.get(sport, "Other Sports") + # Store both the main name and alternate names for matching + channel_names = [channel["name"].lower()] + if channel.get("alt_names"): + channel_names.extend(name.lower() for name in channel["alt_names"]) - item = { - "stream_source": "DaddyLiveHD (1.dlhd.sx)", - "genres": [category], - "poster": random.choice(SPORTS_ARTIFACTS[category]["poster"]), - "background": random.choice( - SPORTS_ARTIFACTS[category]["background"] - ), - "logo": random.choice(SPORTS_ARTIFACTS[category]["logo"]), - "is_add_title_to_poster": True, - "title": event["event"], - "description": description, - "channels": event["channels"], - "event_start_timestamp": event_start_timestamp, - "streams": [], - } - yield from self.parse_channels(item) - - def parse_channels(self, item): - for channel in item["channels"]: - item_copy = item.copy() - m3u8_url = self.m3u8_base_url.format(channel_id=channel["channel_id"]) - item_copy.update( - { - "stream_name": channel["channel_name"], - "stream_url": m3u8_url, - "stream_headers": { - "Referer": self.referer, - "Origin": self.referer.rstrip("/"), - }, - "response_headers": { - "Content-Type": "application/vnd.apple.mpegurl", - }, - "channel_id": channel["channel_id"], + for name in channel_names: + self.channels_data[name] = { + "id": channel["id"], + "country": channel.get("country"), + "languages": channel.get("languages", []), + "categories": channel.get("categories", []), + "logo": channel.get("logo"), + "website": channel.get("website"), + "broadcast_area": channel.get("broadcast_area", []), + "network": channel.get("network"), } - ) - yield item_copy + # Now proceed to scrape channels + yield scrapy.Request( + self.start_urls[0], + callback=self.parse_channels, + headers={"Referer": self.referer}, + ) + + def find_matching_channel(self, title): + """Find matching channel in IPTV-org data using fuzzy matching.""" + title_lower = title.lower() + + # First try direct matching + if title_lower in self.channels_data: + return self.channels_data[title_lower] + + # Try fuzzy matching + best_match = None + best_ratio = 0 + + for channel_name in self.channels_data.keys(): + ratio = fuzz.ratio(title_lower, channel_name) + if ratio > best_ratio and ratio >= self.min_match_ratio: + best_ratio = ratio + best_match = self.channels_data[channel_name] + + return best_match + + def parse_channels(self, response): + """Parse the main channels page.""" + for channel_div in response.css("div.grid-item"): + channel_link = channel_div.css("a") + if not channel_link: + continue + + title = channel_link.css("strong::text").get() + if not title: + continue + + # Skip adult content + if "18+" in title: + logging.info(f"Skipping adult content channel: {title}") + continue + + # Extract channel ID from the href + href = channel_link.attrib.get("href", "") + channel_id = href.split("stream-")[-1].split(".")[0] + if not channel_id.isdigit(): + continue + + # Try to find matching channel in IPTV-org data + iptv_data = self.find_matching_channel(title) + + # Determine category and other metadata + if iptv_data: + category = self.map_iptv_category(iptv_data["categories"]) + languages = iptv_data["languages"] + country = iptv_data["country"] + logo = iptv_data.get("logo") + else: + category = self.determine_category(title) + languages = [] + country = None + logo = None + + # Build stream URL + stream_url = self.m3u8_base_url.format(channel_id=channel_id) + + # Create item with required fields for LiveStreamResolverPipeline + item = { + "title": title, + "stream_name": title, + "stream_url": stream_url, + "stream_source": "DaddyLiveHD", + "stream_headers": { + "Referer": self.referer, + "Origin": self.referer.rstrip("/"), + }, + "response_headers": { + "Content-Type": "application/vnd.apple.mpegurl", + }, + "genres": [category], + "languages": languages, + "country": country, + "poster": logo, + "background": logo, + "logo": logo, + "streams": [], + } + + # Add additional metadata if available from IPTV-org + if iptv_data: + item.update( + { + "network": iptv_data.get("network"), + "broadcast_area": iptv_data.get("broadcast_area"), + "website": iptv_data.get("website"), + "iptv_id": iptv_data.get("id"), + } + ) + + yield item + + def map_iptv_category(self, categories): + """Map IPTV-org categories to our category system.""" + category_mapping = { + "sports": "Sports", + "news": "News", + "entertainment": "Entertainment", + "movies": "Entertainment", + "series": "Entertainment", + "animation": "Entertainment", + "documentary": "Documentary", + "education": "Education", + "music": "Music", + "business": "News", + "legislative": "News", + } + + for category in categories: + if category in category_mapping: + return category_mapping[category] + return "Other" + + def determine_category(self, title): + """Fallback category determination based on channel name.""" + title_lower = title.lower() + + # Check for sports channels + sports_keywords = [ + "espn", + "sport", + "fox sports", + "bein", + "sky sports", + "tennis", + "golf", + "nba", + "nfl", + "mlb", + "nhl", + ] + if any(keyword in title_lower for keyword in sports_keywords): + return "Sports" + + # Check for news channels + news_keywords = ["news", "cnn", "bbc", "msnbc", "fox news"] + if any(keyword in title_lower for keyword in news_keywords): + return "News" + + # Check for entertainment channels + entertainment_keywords = [ + "hbo", + "showtime", + "movie", + "disney", + "netflix", + "comedy", + "cartoon", + "nick", + "mtv", + ] + if any(keyword in title_lower for keyword in entertainment_keywords): + return "Entertainment" + + return "Other" # Default category diff --git a/resources/json/scraper_config.json b/resources/json/scraper_config.json index 8b95197..bbb3aae 100644 --- a/resources/json/scraper_config.json +++ b/resources/json/scraper_config.json @@ -1,10 +1,14 @@ { "start_urls": { "nowsports": "https://neesport.me/", - "nowmetv": "https://neeplay.tv/", + "nowmetv": "https://neeplay.tv/", "tamilultra": "https://tamilultratv.co.in/", "tamilbulb": "https://tamilbulb.tv/", - "tgx": ["torrentgalaxy.to", "tgx.rs", "torrentgalaxy.mx"] + "tgx": [ + "torrentgalaxy.to", + "tgx.rs", + "torrentgalaxy.mx" + ] }, "tamil_blasters": { "homepage": "https://www.1tamilblasters.party", @@ -251,8 +255,21 @@ "homepage": "https://www.arab-torrents.com", "catalogs": { "arabic": { - "movies": ["41", "59", "88", "116", "114"], - "series": ["44", "57", "89", "71", "115", "113"] + "movies": [ + "41", + "59", + "88", + "116", + "114" + ], + "series": [ + "44", + "57", + "89", + "71", + "115", + "113" + ] } } }, @@ -263,38 +280,44 @@ "key_url": "https://key.keylocking.ru", "referer": "https://cookiewebplay.xyz/", "category_mapping": { - "Tv Shows": "Other Sports", - "Soccer": "Football", - "Cricket": "Cricket", - "Tennis": "Tennis", - "Motorsport": "Motor Sport", - "Boxing": "Boxing", - "MMA": "MMA", - "Golf": "Golf", - "Snooker": "Other Sports", - "Am. Football": "American Football", - "Athletics": "Athletics", - "Aussie rules": "Aussie Rules", - "Baseball": "Baseball", - "Basketball": "Basketball", - "Bowling": "Bowling", - "Cycling": "Cycling", - "Darts": "Dart", - "Floorball": "Floorball", - "Futsal": "Futsal", - "Gymnastics": "Gymnastics", - "Handball": "Handball", - "Horse Racing": "Horse Racing", - "Ice Hockey": "Hockey", - "Lacrosse": "Lacrosse", - "Netball": "Netball", - "Rugby League": "Rugby/AFL", - "Rugby Union": "Rugby/AFL", - "Squash": "Squash", - "Volleyball": "Volleyball", - "GAA": "GAA", - "Clubber": "Other Sports" - } + "Tv Shows": "Other Sports", + "Soccer": "Football", + "Cricket": "Cricket", + "Tennis": "Tennis", + "Tennis ATP - AUSTRALIAN OPEN": "Tennis", + "Tennis ATP DOUBLES - AUSTRALIAN OPEN": "Tennis", + "Tennis WTA - AUSTRALIAN OPEN": "Tennis", + "Tennis WTA DOUBLES - AUSTRALIAN OPEN": "Tennis", + "Tennis ATP - FRENCH OPEN": "Tennis", + "Tennis MIXED DOUBLES - AUSTRALIAN OPEN": "Tennis", + "Motorsport": "Motor Sport", + "Boxing": "Boxing", + "MMA": "MMA", + "Golf": "Golf", + "Snooker": "Other Sports", + "Am. Football": "American Football", + "Athletics": "Athletics", + "Aussie rules": "Aussie Rules", + "Baseball": "Baseball", + "Basketball": "Basketball", + "Bowling": "Bowling", + "Cycling": "Cycling", + "Darts": "Dart", + "Floorball": "Floorball", + "Futsal": "Futsal", + "Gymnastics": "Gymnastics", + "Handball": "Handball", + "Horse Racing": "Horse Racing", + "Ice Hockey": "Hockey", + "Lacrosse": "Lacrosse", + "Netball": "Netball", + "Rugby League": "Rugby/AFL", + "Rugby Union": "Rugby/AFL", + "Squash": "Squash", + "Volleyball": "Volleyball", + "GAA": "GAA", + "Clubber": "Other Sports" + } }, "sport_video": { "categories": { diff --git a/resources/templates/manifest.json.j2 b/resources/templates/manifest.json.j2 index ac40d54..c1e3eb7 100644 --- a/resources/templates/manifest.json.j2 +++ b/resources/templates/manifest.json.j2 @@ -87,7 +87,7 @@ "tgx_movie": {"type": "movie", "name": "TGx Movies"}, "tgx_series": {"type": "series", "name": "TGx Series"}, "other_sports": {"type": "movie", "name": "Other Sports", "genres": []}, - "live_sport_events": {"type": "events", "name": "Live Sport Events", "genres": ["Afl", "American Football", "Baseball", "Basketball", "Billiards", "Cricket", "Dart", "Fighting", "Football", "Golf", "Hockey", "Motor Sport", "Other Sports", "Rugby", "Tennis"]}, + "live_sport_events": {"type": "events", "name": "Live Sport Events", "genres": ["American Football", "Athletics", "Aussie Rules", "Baseball", "Basketball", "Bowling", "Boxing", "Cricket", "Cycling", "Dart", "Floorball", "Football", "Futsal", "GAA", "Golf", "Gymnastics", "Handball", "Hockey", "Horse Racing", "Lacrosse", "MMA", "Motor Sport", "Netball", "Other Sports", "Rugby/AFL", "Squash", "Tennis", "Volleyball"]}, "prowlarr_movies": {"type": "movie", "name": "Prowlarr Scraped Movies"}, "prowlarr_series": {"type": "series", "name": "Prowlarr Scraped Series"}, "arabic_movies": {"type": "movie", "name": "Arabic Movies"}, diff --git a/scrapers/dlhd.py b/scrapers/dlhd.py new file mode 100644 index 0000000..227f3f5 --- /dev/null +++ b/scrapers/dlhd.py @@ -0,0 +1,311 @@ +import logging +import random +from datetime import datetime, timedelta +from typing import Optional, List + +import httpx +import humanize +import pytz +from dateutil import parser as date_parser + +from db.config import settings +from db.models import MediaFusionTVMetaData, TVStreams, MediaFusionEventsMetaData +from db.redis_database import REDIS_ASYNC_CLIENT +from utils import crypto +from utils.config import config_manager +from utils.runtime_const import SPORTS_ARTIFACTS + + +class DLHDScheduleService: + def __init__( + self, start_within_next_hours: int = 1, started_within_hours_ago: int = 6 + ): + self.name = "dlhd" + self.schedule_url = config_manager.get_scraper_config(self.name, "schedule_url") + self.m3u8_base_url = config_manager.get_scraper_config( + self.name, "m3u8_base_url" + ) + self.referer = config_manager.get_scraper_config(self.name, "referer") + self.gmt = pytz.timezone("Etc/GMT") + self.start_within_next_hours = start_within_next_hours + self.started_within_hours_ago = started_within_hours_ago + self.category_map = config_manager.get_scraper_config( + self.name, "category_mapping" + ) + + def format_event_time(self, event_timestamp: int) -> str: + """Format event time in specified timezone""" + event_time = datetime.fromtimestamp(event_timestamp) + return humanize.naturaltime(event_time) + + async def fetch_and_parse_schedule(self) -> dict: + """Fetch and parse the schedule data""" + logging.info("Fetching fresh schedule data from DLHD") + try: + async with httpx.AsyncClient(proxy=settings.requests_proxy_url) as client: + response = await client.get(self.schedule_url) + if response.status_code == 200: + return response.json() + logging.error(f"Failed to fetch schedule: {response.status_code}") + except Exception as e: + logging.error(f"Error fetching schedule: {e}") + return {} + + def create_event_id(self, title: str) -> str: + """Create a unique event ID""" + return f"mfdlhd{crypto.get_text_hash(title)}" + + def is_event_in_timewindow(self, event_datetime: datetime) -> bool: + """Check if event is within the configured time window""" + current_time = datetime.now(tz=self.gmt) + time_difference = event_datetime - current_time + + # For future events + if time_difference > timedelta(0): + return time_difference <= timedelta(hours=self.start_within_next_hours) + + # For past events + return time_difference >= timedelta(hours=-self.started_within_hours_ago) + + def create_event_metadata( + self, event: dict, date_str: str, mapped_category: str + ) -> Optional[MediaFusionEventsMetaData]: + """Create event metadata from schedule data""" + event_date = date_parser.parse(date_str.split(" - ")[0]).date() + time = datetime.strptime(event["time"], "%H:%M").time() + aware_datetime = datetime.combine(event_date, time).replace(tzinfo=self.gmt) + + if not self.is_event_in_timewindow(aware_datetime): + return None + + event_start_timestamp = int(aware_datetime.timestamp()) + meta_id = self.create_event_id(event["event"]) + + return MediaFusionEventsMetaData( + id=meta_id, + title=event["event"], + description="", + genres=[mapped_category], + event_start_timestamp=event_start_timestamp, + poster=random.choice(SPORTS_ARTIFACTS[mapped_category]["poster"]), + background=random.choice(SPORTS_ARTIFACTS[mapped_category]["background"]), + logo=random.choice(SPORTS_ARTIFACTS[mapped_category]["logo"]), + is_add_title_to_poster=True, + streams=[], + ) + + async def find_tv_channel( + self, channel_name: str + ) -> Optional[MediaFusionTVMetaData]: + """Find TV channel in database""" + return await MediaFusionTVMetaData.find_one({"title": channel_name}) + + def create_stream( + self, channel_id: str, channel_name: str, meta_id: str + ) -> TVStreams: + """Create a new stream for a channel""" + return TVStreams( + meta_id=meta_id, + name=channel_name, + source="DaddyLiveHD", + url=self.m3u8_base_url.format(channel_id=channel_id), + behaviorHints={ + "notWebReady": True, + "proxyHeaders": { + "request": { + "Referer": self.referer, + "Origin": self.referer.rstrip("/"), + }, + "response": { + "Content-Type": "application/vnd.apple.mpegurl", + }, + }, + }, + ) + + async def get_event_streams(self, event: dict, meta_id: str) -> list[TVStreams]: + """Get all available streams for an event""" + streams = [] + found_streams = 0 + created_streams = 0 + + async def process_channel(channel): + nonlocal found_streams, created_streams + tv_channel = await self.find_tv_channel(channel["channel_name"]) + if tv_channel: + channel_streams = await TVStreams.find( + { + "meta_id": tv_channel.id, + "namespaces": "mediafusion", + "is_working": True, + } + ).to_list(None) + found_streams += len(channel_streams) + streams.extend(channel_streams) + if not channel_streams: + stream = self.create_stream( + channel["channel_id"], channel["channel_name"], meta_id + ) + created_streams += 1 + streams.append(stream) + else: + stream = self.create_stream( + channel["channel_id"], channel["channel_name"], meta_id + ) + created_streams += 1 + streams.append(stream) + + # Process main channels + if "channels" in event: + for channel in event["channels"]: + await process_channel(channel) + + # Process SD channels + if "channels2" in event: + for channel in event["channels2"]: + stream = self.create_stream( + channel["channel_id"], channel["channel_name"], meta_id + ) + created_streams += 1 + streams.append(stream) + + logging.debug( + "Event %s: Found %s existing streams, created %s new streams", + meta_id, + found_streams, + created_streams, + ) + return streams + + async def cache_event(self, event: MediaFusionEventsMetaData): + """Cache single event in Redis with appropriate TTL""" + event_key = f"event:{event.id}" + event_json = event.model_dump_json(exclude_none=True) + + # Calculate TTL: time until event start + buffer period after start + current_time = datetime.now(tz=self.gmt) + event_time = datetime.fromtimestamp(event.event_start_timestamp, self.gmt) + + if event_time > current_time: + # Future event: TTL = time until event + buffer after start + ttl = int( + (event_time - current_time).total_seconds() + + self.started_within_hours_ago * 3600 + ) + else: + # Past event: TTL = remaining time in the window + ttl = int( + ( + event_time + + timedelta(hours=self.started_within_hours_ago) + - current_time + ).total_seconds() + ) + + # Ensure minimum TTL of 1 hour + ttl = max(3600, ttl) + + await REDIS_ASYNC_CLIENT.set(event_key, event_json, ex=ttl) + await REDIS_ASYNC_CLIENT.zadd( + "events:all", {event_key: event.event_start_timestamp} + ) + for genre in event.genres: + await REDIS_ASYNC_CLIENT.zadd( + f"events:genre:{genre}", {event_key: event.event_start_timestamp} + ) + + async def parse_and_cache_schedule(self) -> List[MediaFusionEventsMetaData]: + """Parse schedule and cache it in Redis""" + schedule_data = await self.fetch_and_parse_schedule() + if not schedule_data: + return [] + + events_metadata = [] + processed_count = 0 + skipped_count = 0 + + for date_section, sports in schedule_data.items(): + for sport, events in sports.items(): + mapped_category = self.category_map.get(sport, "Other Sports") + for event in events: + metadata = self.create_event_metadata( + event, date_section, mapped_category + ) + if not metadata: + skipped_count += 1 + continue + + metadata.streams = await self.get_event_streams(event, metadata.id) + await self.cache_event(metadata) + + metadata.poster = ( + f"{settings.poster_host_url}/poster/events/{metadata.id}.jpg" + ) + events_metadata.append(metadata) + processed_count += 1 + + logging.info( + f"Processed {processed_count} events, skipped {skipped_count} events outside time window" + ) + return events_metadata + + async def get_scheduled_events( + self, + force_refresh: bool = False, + genre: Optional[str] = None, + skip: int = 0, + limit: int = 25, + ) -> List[MediaFusionEventsMetaData]: + """Get scheduled events with pagination and genre filtering""" + events = [] + current_time = datetime.now(tz=self.gmt) + min_time = int( + (current_time - timedelta(hours=self.started_within_hours_ago)).timestamp() + ) + max_time = int( + (current_time + timedelta(hours=self.start_within_next_hours)).timestamp() + ) + + if not force_refresh: + cache_key = f"events:all" if not genre else f"events:genre:{genre}" + event_keys = await REDIS_ASYNC_CLIENT.zrevrangebyscore( + cache_key, max_time, min_time, start=skip, num=limit + ) + + if event_keys: + logging.info(f"Found {len(event_keys)} events in cache") + for event_key in event_keys: + event_json = await REDIS_ASYNC_CLIENT.get(event_key) + if event_json: + event = MediaFusionEventsMetaData.model_validate_json( + event_json + ) + event.poster = ( + f"{settings.poster_host_url}/poster/events/{event.id}.jpg" + ) + event.description = f"🎬 {event.title} - ⏰ {self.format_event_time(event.event_start_timestamp)}" + events.append(event) + else: + await REDIS_ASYNC_CLIENT.zrem(cache_key, event_key) + + if len(events) == limit and not force_refresh: + return events + + all_events = await self.parse_and_cache_schedule() + filtered_events = [ + event for event in all_events if not genre or genre in event.genres + ] + filtered_events.sort(key=lambda x: x.event_start_timestamp, reverse=True) + + start_idx = min(skip, len(filtered_events)) + end_idx = min(start_idx + limit, len(filtered_events)) + events = filtered_events[start_idx:end_idx] + + for event in events: + event.poster = f"{settings.poster_host_url}/poster/events/{event.id}.jpg" + event.description = f"🎬 {event.title} - ⏰ {self.format_event_time(event.event_start_timestamp)}" + + return events + + +dlhd_schedule_service = DLHDScheduleService()