Migrate to TamilMV and TamilBlasters scrapy spiders & add scrapy pipeline functionality to parse and store torrent data

This commit is contained in:
mhdzumair
2024-04-13 21:56:37 +05:30
parent d10e8890d4
commit a7d05cc488
12 changed files with 1053 additions and 1625 deletions
-3
View File
@@ -26,8 +26,6 @@ pillow = "*"
httpx = "*"
uvloop = { version = "*", markers = "sys_platform != 'win32'" }
thefuzz = "*"
playwright = "==1.38.0"
playwright-stealth = "*"
pikpakapi = {git = "git+https://github.com/mhdzumair/PikPakAPI.git"}
parse-torrent-title = {git = "git+https://github.com/mhdzumair/parse-torrent-title"}
diskcache = "*"
@@ -38,7 +36,6 @@ gunicorn = "*"
scrapy = "*"
aioqbt = {git = "git+https://github.com/mhdzumair/aioqbt.git"}
aiowebdav = "*"
scrapy-playwright = "*"
m3u-ipytv = "*"
python-multipart = "*"
Generated
+743 -703
View File
File diff suppressed because it is too large Load Diff
+1 -1
View File
@@ -8,7 +8,7 @@ from scrapers import prowlarr # noqa: F401
from mediafusion_scrapy import task # noqa: F401
from utils import torrent
from utils import validation_helper # noqa: F401
from scrapers import tamil_blasters, tamilmv, tv # noqa: F401
from scrapers import tv # noqa: F401
async def async_setup():
+6 -6
View File
@@ -430,12 +430,12 @@ async def save_movie_metadata(metadata: dict, is_imdb: bool = True):
)
if not matching_stream:
existing_movie.streams.append(new_stream)
await existing_movie.save(link_rule=WriteRules.WRITE)
logging.info(
"Updated movie %s. total streams: %s",
existing_movie.title,
len(existing_movie.streams),
)
logging.info(
"Updated movie %s. Total streams: %d",
existing_movie.title,
len(existing_movie.streams),
)
await existing_movie.save(link_rule=WriteRules.WRITE)
else:
# If the movie doesn't exist, create a new one
movie_data = MediaFusionMovieMetaData(
+87 -21
View File
@@ -1,15 +1,17 @@
import asyncio
import logging
import os
import re
from datetime import datetime
from uuid import uuid4
import redis.asyncio as redis_async
import scrapy
from beanie import WriteRules
from itemadapter import ItemAdapter
from scrapy import signals
from scrapy.exceptions import DropItem
from scrapy.http.request import NO_CALLBACK
from scrapy.utils.defer import maybe_deferred_to_future
from db import crud
from db.config import settings
@@ -468,9 +470,8 @@ class FormulaParserPipeline:
torrent_data["episodes"] = episodes
class FormulaStorePipeline:
class QueueBasedPipeline:
def __init__(self):
self.redis = redis_async.Redis.from_url(settings.redis_url)
self.queue = asyncio.Queue()
self.processing_task = None
@@ -479,8 +480,6 @@ class FormulaStorePipeline:
async def close(self):
await self.queue.join()
await self.redis.aclose()
logging.info("Closed pipeline")
self.processing_task.cancel()
@classmethod
@@ -506,6 +505,19 @@ class FormulaStorePipeline:
finally:
self.queue.task_done()
async def parse_item(self, item, spider):
raise NotImplementedError
class FormulaStorePipeline(QueueBasedPipeline):
def __init__(self):
super().__init__()
self.redis = redis_async.Redis.from_url(settings.redis_url)
async def close(self):
await super().close()
await self.redis.aclose()
async def parse_item(self, item, spider):
if "title" not in item:
logging.warning(f"title not found in item: {item}")
@@ -603,8 +615,8 @@ class FormulaStorePipeline:
) # Ensure episodes are ordered by episode number
class TVStorePipeline:
async def process_item(self, item, spider):
class TVStorePipeline(QueueBasedPipeline):
async def parse_item(self, item, spider):
if "title" not in item:
logging.warning(f"title not found in item: {item}")
raise DropItem(f"title not found in item: {item}")
@@ -614,19 +626,42 @@ class TVStorePipeline:
return item
class TorrentFileParserPipeline:
def process_item(self, item, spider):
class TorrentDownloadAndParsePipeline:
async def process_item(self, item, spider):
adapter = ItemAdapter(item)
if "torrent_file_path" not in adapter:
raise DropItem(f"torrent_file_path not found in item: {item}")
torrent_link = adapter.get("torrent_link")
with open(item["torrent_file_path"], "rb") as torrent_file:
torrent_metadata = torrent.extract_torrent_metadata(
torrent_file.read(), item.get("is_parse_ptn", True)
if not torrent_link:
raise DropItem(f"No torrent link found in item: {item}")
response = await maybe_deferred_to_future(
spider.crawler.engine.download(
scrapy.Request(torrent_link, callback=NO_CALLBACK)
)
os.remove(item["torrent_file_path"])
)
if response.status != 200:
spider.logger.error(
f"Failed to download torrent file: {response.url} with status {response.status}"
)
return item
# Validate the content-type of the response
if "application/x-bittorrent" not in response.headers.get(
"Content-Type", b""
).decode("utf-8", "ignore"):
spider.logger.warning(
f"Unexpected Content-Type for {response.url}: {response.headers.get('Content-Type')}"
)
return item
torrent_metadata = torrent.extract_torrent_metadata(
response.body, item.get("is_parse_ptn", True)
)
if not torrent_metadata:
raise DropItem(f"Failed to extract torrent metadata: {item}")
item.update(torrent_metadata)
return item
@@ -714,16 +749,47 @@ class SportVideoParserPipeline:
return None
class MovieStorePipeline:
async def process_item(self, item, spider):
class MovieStorePipeline(QueueBasedPipeline):
async def parse_item(self, item, spider):
if "title" not in item:
raise DropItem(f"title not found in item: {item}")
if item.get("type") != "movie":
return item
await crud.save_movie_metadata(item, item.get("is_imdb", True))
return item
class LiveEventStorePipeline:
class SeriesStorePipeline(QueueBasedPipeline):
async def parse_item(self, item, spider):
if "title" not in item:
raise DropItem(f"title not found in item: {item}")
if item.get("type") != "series":
return item
await crud.save_series_metadata(item)
return item
class LiveEventStorePipeline(QueueBasedPipeline):
def __init__(self):
super().__init__()
self.redis = redis_async.Redis.from_url(settings.redis_url)
async def close(self):
await super().close()
await self.redis.aclose()
async def parse_item(self, item, spider):
if "title" not in item:
raise DropItem(f"name not found in item: {item}")
await crud.save_events_data(self.redis, item)
return item
class RedisCacheURLPipeline:
def __init__(self):
self.redis = redis_async.Redis.from_url(settings.redis_url)
@@ -737,8 +803,8 @@ class LiveEventStorePipeline:
return p
async def process_item(self, item, spider):
if "title" not in item:
raise DropItem(f"name not found in item: {item}")
if "webpage_url" not in item:
raise DropItem(f"webpage_url not found in item: {item}")
await crud.save_events_data(self.redis, item)
await self.redis.sadd(item["scraped_url_key"], item["webpage_url"])
return item
+1 -1
View File
@@ -94,4 +94,4 @@ HTTPCACHE_STORAGE = "scrapy.extensions.httpcache.FilesystemCacheStorage"
REQUEST_FINGERPRINTER_IMPLEMENTATION = "2.7"
TWISTED_REACTOR = "twisted.internet.asyncioreactor.AsyncioSelectorReactor"
FEED_EXPORT_ENCODING = "utf-8"
LOG_LEVEL = "INFO"
LOG_LEVEL = "DEBUG"
+203
View File
@@ -0,0 +1,203 @@
import re
import redis
import scrapy
from db.config import settings
from scrapers.helpers import get_scraper_config
class CommonTamilSpider(scrapy.Spider):
name = None
source = None
custom_settings = {
"ITEM_PIPELINES": {
"mediafusion_scrapy.pipelines.TorrentDownloadAndParsePipeline": 100,
"mediafusion_scrapy.pipelines.MovieStorePipeline": 200,
"mediafusion_scrapy.pipelines.SeriesStorePipeline": 300,
"mediafusion_scrapy.pipelines.RedisCacheURLPipeline": 400,
},
"RETRY_HTTP_CODES": [500, 502, 503, 504, 522, 524, 408, 429],
"RETRY_TIMES": 5,
}
def __init__(
self,
pages: int = 1,
start_page: int = 1,
search_keyword: str = None,
scrap_catalog_id: str = "all",
*args,
**kwargs,
):
super(CommonTamilSpider, self).__init__(*args, **kwargs)
self.pages = pages
self.start_page = start_page
self.search_keyword = search_keyword
if scrap_catalog_id != "all" and "_" not in scrap_catalog_id:
self.logger.error(
f"Invalid catalog ID: {scrap_catalog_id}. Expected format: <language>_<video_type>"
)
return
self.scrap_catalog_id = scrap_catalog_id
print(f"Scraping catalog ID: {self.scrap_catalog_id}")
self.redis = redis.Redis(
connection_pool=redis.ConnectionPool.from_url(settings.redis_url)
)
self.scraped_urls_key = f"{self.name}_scraped_urls"
self.catalogs = get_scraper_config(self.name, "catalogs")
self.homepage = get_scraper_config(self.name, "homepage")
self.supported_search_forums = get_scraper_config(
self.name, "supported_search_forums"
)
def __del__(self):
self.redis.close()
def generate_forum_data(self):
data = []
link_prefix = f"{self.homepage}/index.php?/forums/forum/"
for language, catalog in self.catalogs.items():
for video_type, forum_ids in catalog.items():
scrap_links = (
[link_prefix + str(forum_id) for forum_id in forum_ids]
if isinstance(forum_ids, list)
else [link_prefix + str(forum_ids)]
)
for link in scrap_links:
for page in range(self.start_page, self.start_page + self.pages):
paginated_link = f"{link}/page/{page}/"
data.append((paginated_link, language, video_type))
return data
def start_requests(self):
if self.search_keyword:
# Construct the search URL and initiate search
search_url = f"{self.homepage}/index.php?/search/&q={self.search_keyword}&type=forums_topic&search_and_or=or&search_in=titles&sortby=relevancy"
yield scrapy.Request(search_url, self.parse_search_results)
else:
if self.scrap_catalog_id == "all":
for url, language, video_type in self.generate_forum_data():
yield scrapy.Request(
url,
self.parse_page_results,
meta={
"item": {
"language": language,
"video_type": video_type,
"source": self.source,
}
},
)
else:
language, video_type = self.scrap_catalog_id.split("_")
forum_ids = self.catalogs.get(language, {}).get(video_type)
if not forum_ids:
self.logger.error(f"Invalid catalog ID: {self.scrap_catalog_id}")
return
forum_links = (
[
f"{self.homepage}/index.php?/forums/forum/{forum_id}"
for forum_id in forum_ids
]
if isinstance(forum_ids, list)
else [f"{self.homepage}/index.php?/forums/forum/{forum_ids}"]
)
for forum_link in forum_links:
yield scrapy.Request(
forum_link,
self.parse_page_results,
meta={
"item": {
"language": language,
"video_type": video_type,
"source": self.source,
}
},
)
def parse_search_results(self, response):
movies = response.css("li[data-role='activityItem']")
for movie in movies:
movie_page_link = movie.css("a[data-linktype='link']::attr(href)").get()
if not movie_page_link:
continue
if not self.check_scraped_urls(movie_page_link):
yield response.follow(
movie_page_link,
self.parse_movie_page,
meta={
"item": {
"source": self.source,
"webpage_url": movie_page_link,
}
},
)
# Handling pagination
next_page_link = response.css("a[rel='next']::attr(href)").get()
if next_page_link:
yield response.follow(next_page_link, self.parse_search_results)
def parse_page_results(self, response):
movies = response.css("li[data-rowid]")
item = response.meta["item"].copy()
for movie in movies:
movie_page_link = movie.css("a::attr(href)").get()
if not movie_page_link:
continue
if not self.check_scraped_urls(movie_page_link):
item["webpage_url"] = movie_page_link
yield response.follow(
movie_page_link, self.parse_movie_page, meta={"item": item}
)
def check_scraped_urls(self, page_link):
if self.redis.sismember(self.scraped_urls_key, page_link):
self.logger.info(f"Skipping already scraped URL: {page_link}")
return True
return False
def parse_movie_page(self, response):
item = response.meta["item"].copy()
poster = response.css("div[data-commenttype='forums'] img::attr(src)").get()
created_at = response.css("time::attr(datetime)").get()
torrent_links = response.css("a[data-fileext='torrent']::attr(href)").getall()
if self.search_keyword:
forum_link = response.css("a[href*='forums/forum/']::attr(href)").get()
if not forum_link:
return
forum_id = re.search(r"forums/forum/([^/]+)/", forum_link).group(1)
language, video_type = self.supported_search_forums.get(
forum_id, {"language": None, "video_type": None}
).values()
if not language:
self.logger.error(f"Unsupported forum {forum_id}")
return
item.update({"language": language, "video_type": video_type})
if not torrent_links:
self.logger.warning(f"No torrents found for {response.url}")
return
for torrent_link in torrent_links:
item.update(
{
"catalog": f"{item['language']}_{item['video_type']}",
"type": "series" if item["video_type"] == "series" else "movie",
"poster": poster,
"created_at": created_at,
"scrap_language": item["language"].title(),
"torrent_link": torrent_link,
"scraped_url_key": self.scraped_urls_key,
}
)
yield item
@@ -0,0 +1,6 @@
from mediafusion_scrapy.spiders.common import CommonTamilSpider
class TamilBlastersSpider(CommonTamilSpider):
name = "tamil_blasters"
source = "TamilBlasters"
+6
View File
@@ -0,0 +1,6 @@
from mediafusion_scrapy.spiders.common import CommonTamilSpider
class TamilMVSpider(CommonTamilSpider):
name = "tamilmv"
source = "TamilMV"
-344
View File
@@ -1,344 +0,0 @@
#!/usr/bin/env python3
import argparse
import asyncio
import logging
import math
import random
import re
from multiprocessing import Process
import dramatiq
from bs4 import BeautifulSoup
from dateutil.parser import parse as dateparser
from playwright.async_api import async_playwright
from playwright_stealth import stealth_async
from db import database
from db.config import settings
from scrapers.helpers import (
get_page_content,
get_scraper_session,
download_and_save_torrent,
get_scraper_config,
)
from utils.wrappers import worker_rate_limit
HOMEPAGE = get_scraper_config("tamil_blasters", "homepage")
TAMIL_BLASTER_CATALOGS = get_scraper_config("tamil_blasters", "catalogs")
async def get_search_results(page, keyword, page_number=1):
search_link = f"{HOMEPAGE}/index.php?/search/&q={keyword}&type=forums_topic&page={page_number}&search_and_or=or&search_in=titles&sortby=relevancy"
# Get page content and initialize BeautifulSoup
page_content = await get_page_content(page, search_link)
soup = BeautifulSoup(page_content, "html.parser")
return soup
async def process_movie(
movie,
scraper=None,
page=None,
keyword=None,
language=None,
media_type=None,
supported_forums=None,
):
if keyword:
movie_link = movie.find("a", {"data-linktype": "link"})
forum_link = movie.find("a", href=re.compile(r"forums/forum/")).get("href")
forum_id = re.search(r"forums/forum/([^/]+)/", forum_link)[1]
if forum_id not in supported_forums:
logging.error(f"Unsupported forum {forum_id}")
return
# Extracting language and media_type from supported_forums
language = supported_forums[forum_id]["language"]
media_type = supported_forums[forum_id]["media_type"]
else:
movie_link = movie.find("a")
if not movie_link:
logging.error(f"Movie link not found")
return
page_link = movie_link.get("href")
try:
if scraper: # If using the scraper
response = scraper.get(page_link)
movie_page_content = response.content
else: # If using playwright
movie_page_content = await get_page_content(page, page_link)
movie_page = BeautifulSoup(movie_page_content, "html.parser")
# Extracting other details
poster_element = movie_page.select_one(
"div[data-commenttype='forums'] img[data-src]"
)
poster = poster_element.get("data-src") if poster_element else None
datetime_element = movie_page.select_one("time")
created_at = (
dateparser(datetime_element.get("datetime")) if datetime_element else None
)
# Define metadata
metadata = {
"catalog": f"{language}_{media_type}",
"poster": poster,
"created_at": created_at,
"scrap_language": language.title(),
"source": "TamilBlasters",
}
# Extracting torrent details
torrent_elements = movie_page.select("a[data-fileext='torrent']")
if not torrent_elements:
logging.error(f"No torrents found for {page_link}")
return
for torrent_element in torrent_elements:
try:
await download_and_save_torrent(
torrent_element,
scraper=scraper,
page=page,
metadata=metadata.copy(),
media_type=media_type,
page_link=page_link,
)
except Exception as e:
logging.error(
f"Error processing torrent {page_link}: {e}",
exc_info=True,
stack_info=True,
)
return True
except Exception as e:
logging.error(
f"Error processing movie {page_link}: {e}", exc_info=True, stack_info=True
)
return False
async def scrap_page(url, language, media_type):
scraper = get_scraper_session()
response = scraper.get(url)
if response.status_code == 403:
logging.error(
"Cloudflare validation required. Run with --scrap-with-playwright"
)
return
response.raise_for_status()
tamil_blasters = BeautifulSoup(response.content, "html.parser")
movies = tamil_blasters.select("li[data-rowid]")
for movie in movies:
await process_movie(
movie, scraper=scraper, language=language, media_type=media_type
)
async def scrap_page_with_playwright(url, language, media_type):
async with async_playwright() as p:
# Launch a new browser session
browser = await p.firefox.launch(
headless=False,
proxy={"server": settings.scraper_proxy_url}
if settings.scraper_proxy_url
else None,
)
page = await browser.new_page()
await stealth_async(page)
await asyncio.sleep(2)
page_content = await get_page_content(page, url)
tamil_blasters = BeautifulSoup(page_content, "html.parser")
movies = tamil_blasters.select("li[data-rowid]")
for movie in movies:
await process_movie(
movie, page=page, language=language, media_type=media_type
)
await browser.close()
async def scrap_search_keyword(keyword):
supported_forums = {
TAMIL_BLASTER_CATALOGS[language][media_type]: {
"language": language,
"media_type": media_type,
}
for language in TAMIL_BLASTER_CATALOGS
for media_type in TAMIL_BLASTER_CATALOGS[language]
}
async with async_playwright() as p:
# Launch a new browser session
browser = await p.firefox.launch(
headless=False,
proxy={"server": settings.scraper_proxy_url}
if settings.scraper_proxy_url
else None,
)
page = await browser.new_page()
await stealth_async(page)
await asyncio.sleep(2)
soup = await get_search_results(page, keyword)
results_element = soup.find("div", {"data-role": "resultsArea"})
results_count = int(re.search(r"\d+", results_element.find("p").text).group())
logging.info(f"Found {results_count} results for {keyword}")
movies = results_element.select("li[data-role='activityItem']")
if results_count > 25:
number_of_pages = math.ceil(results_count / 25)
logging.info(f"Found {number_of_pages} pages for {keyword}")
for page_number in range(2, number_of_pages + 1):
soup = await get_search_results(page, keyword, page_number)
movies.extend(soup.select("li[data-role='activityItem']"))
await asyncio.sleep(random.randint(2, 5))
for movie in movies:
await process_movie(
movie, page=page, keyword=keyword, supported_forums=supported_forums
)
await browser.close()
async def run_scraper(
language: str = None,
video_type: str = None,
pages: int = None,
start_page: int = None,
search_keyword: str = None,
scrap_with_playwright: bool = None,
):
if search_keyword:
await scrap_search_keyword(search_keyword)
return
try:
scrap_link_prefix = f"{HOMEPAGE}/index.php?/forums/forum/{TAMIL_BLASTER_CATALOGS[language][video_type]}"
except KeyError:
logging.error(f"Unsupported language or video type: {language}_{video_type}")
return
for page in range(start_page, pages + start_page):
scrap_link = f"{scrap_link_prefix}/page/{page}/"
logging.info(f"Scrap page: {page}")
if scrap_with_playwright is True:
await scrap_page_with_playwright(scrap_link, language, video_type)
else:
await scrap_page(scrap_link, language, video_type)
logging.info(f"Scrap completed for : {language}_{video_type}")
async def run_schedule_scrape(
pages: int = 1, start_page: int = 1, scrap_with_playwright: bool = False
):
await database.init()
async with asyncio.TaskGroup() as tg:
for language in TAMIL_BLASTER_CATALOGS:
for video_type in TAMIL_BLASTER_CATALOGS[language]:
tg.create_task(
run_scraper(
language,
video_type,
pages=pages,
start_page=start_page,
scrap_with_playwright=scrap_with_playwright,
)
)
def run_schedule_scrape_sync(pages, start_page, scrap_with_playwright):
asyncio.run(run_schedule_scrape(pages, start_page, scrap_with_playwright))
@dramatiq.actor(priority=5, time_limit=60 * 60 * 1000)
@worker_rate_limit(limit=1)
def run_tamil_blasters_scraper(pages: int = 1, start_page: int = 1):
# Use a separate process to run the scraper
process = Process(target=run_schedule_scrape_sync, args=(pages, start_page, False))
try:
process.start()
process.join()
except Exception as e:
logging.error(f"Error running tamilblasters scraper: {e}")
if __name__ == "__main__":
parser = argparse.ArgumentParser(
description="Scrap Movie metadata from TamilBlasters"
)
parser.add_argument(
"--all", action="store_true", help="scrap all type of movies & series"
)
parser.add_argument(
"-l",
"--language",
help="scrap movie language",
default="tamil",
choices=["tamil", "malayalam", "telugu", "hindi", "kannada", "english"],
)
parser.add_argument(
"-t",
"--video-type",
help="scrap movie video type",
default="hdrip",
choices=["hdrip", "tcrip", "dubbed", "series", "old"],
)
parser.add_argument(
"-p", "--pages", type=int, default=1, help="number of scrap pages"
)
parser.add_argument(
"-s", "--start-pages", type=int, default=1, help="page number to start scrap."
)
parser.add_argument(
"-k",
"--search-keyword",
help="search keyword to scrap movies & series. ex: 'bigg boss'",
default=None,
)
parser.add_argument(
"--scrap-with-playwright", action="store_true", help="scrap with playwright"
)
parser.add_argument(
"--proxy-url",
help="proxy url to scrap. ex: socks5://127.0.0.1:1080",
default=None,
)
args = parser.parse_args()
logging.basicConfig(
format="%(levelname)s::%(asctime)s - %(message)s",
datefmt="%d-%b-%y %H:%M:%S",
level=logging.INFO,
)
if args.all:
asyncio.run(
run_schedule_scrape(
args.pages, args.start_pages, args.scrap_with_playwright
)
)
else:
asyncio.run(
run_scraper(
args.language,
args.video_type,
args.pages,
args.start_pages,
args.search_keyword,
args.scrap_with_playwright,
)
)
-326
View File
@@ -1,326 +0,0 @@
#!/usr/bin/env python3
import argparse
import asyncio
import logging
import math
import random
import re
from multiprocessing import Process
import dramatiq
from bs4 import BeautifulSoup
from dateutil.parser import parse as dateparser
from playwright.async_api import async_playwright
from playwright_stealth import stealth_async
from db import database
from db.config import settings
from scrapers.helpers import (
get_page_content,
get_scraper_session,
download_and_save_torrent,
get_scraper_config,
)
from utils.wrappers import worker_rate_limit
HOMEPAGE = get_scraper_config("tamilmv", "homepage")
TAMIL_MV_CATALOGS = get_scraper_config("tamilmv", "catalogs")
SUPPORTED_SEARCH_FORUMS = get_scraper_config("tamilmv", "supported_search_forums")
async def process_movie(
movie,
scraper=None,
page=None,
keyword=None,
language=None,
media_type=None,
supported_forums=None,
):
if keyword:
movie_link = movie.find("a", {"data-linktype": "link"})
forum_link = movie.find("a", href=re.compile(r"forums/forum/")).get("href")
forum_id = re.search(r"forums/forum/([^/]+)/", forum_link)[1]
if forum_id not in supported_forums:
logging.error(f"Unsupported forum {forum_id}")
return
# Extracting language and media_type from supported_forums
language = supported_forums[forum_id]["language"]
media_type = supported_forums[forum_id]["media_type"]
else:
movie_link = movie.find("a")
if not movie_link:
logging.error(f"Movie link not found")
return
page_link = movie_link.get("href")
try:
if scraper: # If using the scraper
response = scraper.get(page_link)
movie_page_content = response.content
else: # If using playwright
movie_page_content = await get_page_content(page, page_link)
movie_page = BeautifulSoup(movie_page_content, "html.parser")
# Extracting other details
poster_element = movie_page.select_one("div[data-commenttype='forums'] img")
poster = poster_element.get("src") if poster_element else None
datetime_element = movie_page.select_one("time")
created_at = (
dateparser(datetime_element.get("datetime")) if datetime_element else None
)
# Define metadata
metadata = {
"catalog": f"{language}_{media_type}",
"poster": poster,
"created_at": created_at,
"scrap_language": language.title(),
"source": "TamilMV",
}
# Extracting torrent details
torrent_elements = movie_page.select("a[data-fileext='torrent']")
if not torrent_elements:
logging.error(f"No torrents found for {page_link}")
return
for torrent_element in torrent_elements:
await download_and_save_torrent(
torrent_element,
scraper=scraper,
page=page,
metadata=metadata.copy(),
media_type=media_type,
page_link=page_link,
)
return True
except Exception as e:
logging.error(
f"Error processing movie {page_link}: {e}", exc_info=True, stack_info=True
)
return False
async def scrap_page(url, language, media_type):
scraper = get_scraper_session()
response = scraper.get(url)
if response.status_code == 403:
logging.error(
"Cloudflare validation required. Run with --scrap-with-playwright"
)
return
response.raise_for_status()
tamil_blasters = BeautifulSoup(response.content, "html.parser")
movies = tamil_blasters.select("li[data-rowid]")
for movie in movies:
await process_movie(
movie, scraper=scraper, language=language, media_type=media_type
)
async def scrap_page_with_playwright(url, language, media_type):
async with async_playwright() as p:
# Launch a new browser session
browser = await p.firefox.launch(
headless=False,
proxy={"server": settings.scraper_proxy_url}
if settings.scraper_proxy_url
else None,
)
page = await browser.new_page()
await stealth_async(page)
await asyncio.sleep(2)
page_content = await get_page_content(page, url)
tamil_blasters = BeautifulSoup(page_content, "html.parser")
movies = tamil_blasters.select("li[data-rowid]")
for movie in movies:
await process_movie(
movie, page=page, language=language, media_type=media_type
)
await browser.close()
async def get_search_results(scraper, keyword, page_number=1):
search_link = f"{HOMEPAGE}/index.php?/search/&q={keyword}&type=forums_topic&page={page_number}&search_and_or=or&search_in=titles&sortby=relevancy"
# Get page content and initialize BeautifulSoup
response = scraper.get(search_link)
response.raise_for_status()
page_content = response.content
soup = BeautifulSoup(page_content, "html.parser")
return soup
async def scrap_search_keyword(keyword):
scraper = get_scraper_session()
soup = await get_search_results(scraper, keyword)
results_element = soup.find("div", {"data-role": "resultsArea"})
results_count = int(re.search(r"\d+", results_element.find("p").text).group())
logging.info(f"Found {results_count} results for {keyword}")
movies = results_element.select("li[data-role='activityItem']")
if results_count > 25:
number_of_pages = math.ceil(results_count / 25)
logging.info(f"Found {number_of_pages} pages for {keyword}")
for page_number in range(2, number_of_pages + 1):
soup = await get_search_results(scraper, keyword, page_number)
movies.extend(soup.select("li[data-role='activityItem']"))
await asyncio.sleep(random.randint(2, 5))
for movie in movies:
await process_movie(
movie,
scraper=scraper,
keyword=keyword,
supported_forums=SUPPORTED_SEARCH_FORUMS,
)
async def run_scraper(
language: str = None,
video_type: str = None,
pages: int = None,
start_page: int = None,
search_keyword: str = None,
scrap_with_playwright: bool = None,
):
if search_keyword:
await scrap_search_keyword(search_keyword)
return
link_prefix = f"{HOMEPAGE}/index.php?/forums/forum/"
try:
forum_ids = TAMIL_MV_CATALOGS[language][video_type]
scrap_links = (
[link_prefix + link for link in forum_ids]
if isinstance(forum_ids, list)
else [link_prefix + forum_ids]
)
except KeyError:
logging.error(f"Unsupported language or video type: {language}_{video_type}")
return
for scrap_link_prefix in scrap_links:
for page in range(start_page, pages + start_page):
scrap_link = f"{scrap_link_prefix}/page/{page}/"
logging.info(f"Scrap page: {scrap_link}")
if scrap_with_playwright is True:
await scrap_page_with_playwright(scrap_link, language, video_type)
else:
await scrap_page(scrap_link, language, video_type)
logging.info(f"Scrap completed for : {language}_{video_type}")
async def run_schedule_scrape(
pages: int = 1,
start_page: int = 1,
scrap_with_playwright: bool = None,
):
await database.init()
async with asyncio.TaskGroup() as tg:
for language in TAMIL_MV_CATALOGS:
for video_type in TAMIL_MV_CATALOGS[language]:
tg.create_task(
run_scraper(
language,
video_type,
pages=pages,
start_page=start_page,
scrap_with_playwright=scrap_with_playwright,
)
)
def run_schedule_scrape_sync(pages, start_page, scrap_with_playwright):
asyncio.run(run_schedule_scrape(pages, start_page, scrap_with_playwright))
@dramatiq.actor(priority=5, time_limit=60 * 60 * 1000)
@worker_rate_limit(limit=1)
def run_tamilmv_scraper(pages: int = 1, start_page: int = 1):
# Use a separate process to run the scraper
process = Process(target=run_schedule_scrape_sync, args=(pages, start_page, False))
try:
process.start()
process.join()
except Exception as e:
logging.error(f"Error running tamilmv scraper: {e}")
if __name__ == "__main__":
parser = argparse.ArgumentParser(description="Scrap Movie metadata from TamilMV")
parser.add_argument(
"--all", action="store_true", help="scrap all type of movies & series"
)
parser.add_argument(
"-l",
"--language",
help="scrap movie language",
default="tamil",
choices=["tamil", "malayalam", "telugu", "hindi", "kannada", "english"],
)
parser.add_argument(
"-t",
"--video-type",
help="scrap movie video type",
default="hdrip",
choices=["hdrip", "tcrip", "dubbed", "series"],
)
parser.add_argument(
"-p", "--pages", type=int, default=1, help="number of scrap pages"
)
parser.add_argument(
"-s", "--start-pages", type=int, default=1, help="page number to start scrap."
)
parser.add_argument(
"-k",
"--search-keyword",
help="search keyword to scrap movies & series. ex: 'bigg boss'",
default=None,
)
parser.add_argument(
"--scrap-with-playwright", action="store_true", help="scrap with playwright"
)
parser.add_argument(
"--proxy-url",
help="proxy url to scrap. ex: socks5://127.0.0.1:1080",
default=None,
)
args = parser.parse_args()
logging.basicConfig(
format="%(levelname)s::%(asctime)s - %(message)s",
datefmt="%d-%b-%y %H:%M:%S",
level=logging.INFO,
)
if args.all:
asyncio.run(
run_schedule_scrape(
args.pages, args.start_pages, args.scrap_with_playwright
)
)
else:
asyncio.run(
run_scraper(
args.language,
args.video_type,
args.pages,
args.start_pages,
args.search_keyword,
args.scrap_with_playwright,
)
)
-220
View File
@@ -1,220 +0,0 @@
import argparse
import asyncio
import json
import logging
from urllib.parse import urlparse, urljoin
import requests
from playwright.async_api import async_playwright
logging.basicConfig(
format="%(levelname)s::%(asctime)s - %(message)s", level=logging.INFO
)
BASE_URL = "https://tamilultra.in"
MEDIAFUSION_URL = "http://127.0.0.1:8000"
async def scrape_tv_channels(page):
# Scrape channel metadata
channels_data = []
channel_elements = await page.query_selector_all("article.item.movies")
# First, store all channel information in a list
channel_info_list = []
for channel_element in channel_elements:
title_element = await channel_element.query_selector("h3 > a")
title = (
(await title_element.text_content())
.replace("\u2013", "-")
.split("-")[0]
.strip()
.title()
if title_element
else "No Title"
)
poster_element = await channel_element.query_selector(".poster > img")
poster_url = (
await poster_element.get_attribute("src")
if poster_element
else "No Poster URL"
)
stream_page_link_element = await channel_element.query_selector(".poster > a")
stream_page_url = (
await stream_page_link_element.get_attribute("href")
if stream_page_link_element
else "No Stream Page URL"
)
channel_info_list.append((title, poster_url, stream_page_url))
# Then, navigate to each channel's stream page and capture the M3U8 URLs
for title, poster_url, stream_page_url in channel_info_list:
# Navigate to the stream page
await page.goto(stream_page_url)
m3u8_url_data = []
# Scrape genre tags
genre_elements = await page.query_selector_all(".sgeneros a[rel='tag']")
genres = [await genre.text_content() for genre in genre_elements]
genres = [genre.title() for genre in genres]
# Query for player option elements and click to load M3U8 URL
player_option_elements = await page.query_selector_all(
"#playeroptionsul > li.dooplay_player_option"
)
for player_option_element in player_option_elements:
# Click the player option element
await player_option_element.click()
# Wait for the iframe to load and get its 'src' attribute
iframe_element = await page.wait_for_selector("iframe.metaframe.rptss")
iframe_src = await iframe_element.get_attribute("src")
m3u8_url_part = iframe_src.replace("/player.php?", "")
# Check if iframe_src is a valid URL
parsed_src = urlparse(m3u8_url_part)
behavior_hints = {}
if parsed_src.scheme and parsed_src.netloc:
# If it's a valid URL, use it directly
m3u8_url = m3u8_url_part
if "jio.tamilultra.in" in m3u8_url:
behavior_hints = {
"notWebReady": True,
"proxyHeaders": {
"request": {
"User-Agent": "Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/118.0.0.0 Safari/537.36",
"Referer": BASE_URL + "/",
}
},
}
else:
# Otherwise, join with the BASE_URL
m3u8_url = urljoin(BASE_URL, m3u8_url_part)
behavior_hints = {
"is_redirect": True,
"proxyHeaders": {
"request": {
"User-Agent": "Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/118.0.0.0 Safari/537.36",
"Referer": urljoin(BASE_URL, iframe_src),
}
},
}
m3u8_url_data.append((m3u8_url, behavior_hints))
channels_data.append(
{
"title": title.replace("Hd", "").strip(),
"poster": poster_url,
"genres": genres,
"country": "India", # Set default country
"tv_language": genres[0],
"streams": [
{
"name": f"{title} - {index}",
"url": m3u8_url,
"source": "TamilUltra",
"behaviorHints": behavior_hints,
}
for index, (m3u8_url, behavior_hints) in enumerate(m3u8_url_data, 1)
],
}
)
logging.info("Scraped %s", title)
return channels_data
async def scrape_category(category_url, page):
# Navigate to the category page
await page.goto(category_url)
# Scrape channels from the current page
channels_data = await scrape_tv_channels(
page
) # Assuming scrape_tv_channels accepts a page argument
# Try to find a pagination control and collect all page URLs
pagination_links = await page.query_selector_all("div.pagination a.inactive")
page_urls = [category_url] # Include the first page
for link in pagination_links:
page_url = await link.get_attribute("href")
if page_url:
page_urls.append(page_url)
logging.info("found %d pages", len(page_urls))
# Iterate over each page URL and scrape channels
for page_url in page_urls:
if page_url != category_url: # We already scraped the first page
await page.goto(page_url)
# Scrape channels from this page
channels_data.extend(await scrape_tv_channels(page))
return channels_data
async def scrape_all_categories():
async with async_playwright() as p:
browser = await p.firefox.launch(headless=True)
page = await browser.new_page()
# Extract category URLs
await page.goto(BASE_URL)
category_elements = await page.query_selector_all(".main-header a")
category_urls = [
urljoin(BASE_URL, await element.get_attribute("href"))
for element in category_elements
]
# Scrape channels from each category
all_channels_data = []
for category_url in category_urls:
logging.info("Scraping %s", category_url)
all_channels_data.extend(await scrape_category(category_url, page))
await browser.close()
# remove duplicates
unique_channels = {channel["title"]: channel for channel in all_channels_data}
unique_channels_data = list(unique_channels.values())
logging.info("found %d channels", len(unique_channels_data))
with open("tamilultra.json", "w") as file:
json.dump({"channels": unique_channels_data}, file, indent=4)
logging.info(
"Done scraping TamilUltra. Manually verify the data & add it via /tv-metadata endpoint"
)
def main(is_scraping: bool = True):
if is_scraping:
asyncio.run(scrape_all_categories())
return
with open("tamilultra.json") as file:
channels = json.load(file)["channels"]
for channel in channels:
logging.info("Adding %s", channel["title"])
response = requests.post(f"{MEDIAFUSION_URL}/tv-metadata", json=channel)
try:
response.raise_for_status()
except requests.HTTPError as err:
logging.info("Response data: %s", response.text)
continue
logging.info("Response data: %s", response.json())
if __name__ == "__main__":
parser = argparse.ArgumentParser(description="Scrape TamilUltra Live TV")
parser.add_argument(
"--no-scrape",
action="store_true",
help="Don't scrape TamilUltra. Use this option to add the data to MediaFusion",
)
args = parser.parse_args()
main(not args.no_scrape)