You can not select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
251 lines
9.1 KiB
251 lines
9.1 KiB
import re
|
|
import json
|
|
import requests
|
|
import urllib.parse
|
|
from urllib.parse import urlparse
|
|
from django.utils import timezone
|
|
from celery import shared_task
|
|
from crawler.models import CrawlRun, CrawlTask
|
|
from ads.models import Ad, AdEvaluation
|
|
|
|
def parse_divar_url(divar_url):
|
|
"""
|
|
Parses a public Divar search link to extract city, category path,
|
|
and query parameters.
|
|
"""
|
|
parsed = urlparse(divar_url)
|
|
path_parts = [p for p in parsed.path.split('/') if p]
|
|
|
|
city = 'tehran'
|
|
category = ''
|
|
|
|
if len(path_parts) >= 2 and path_parts[0] == 's':
|
|
city = path_parts[1]
|
|
if len(path_parts) > 2:
|
|
category = '/'.join(path_parts[2:])
|
|
elif len(path_parts) >= 1:
|
|
city = path_parts[0]
|
|
if len(path_parts) > 1:
|
|
category = '/'.join(path_parts[1:])
|
|
|
|
query_params = dict(urllib.parse.parse_qsl(parsed.query))
|
|
return city, category, query_params
|
|
|
|
def extract_widgets_from_dict(d, found_ads):
|
|
"""
|
|
Recursively searches a JSON structure for dictionaries matching
|
|
the properties of an ad widget, extracting token, title, description,
|
|
price, and image details.
|
|
"""
|
|
if isinstance(d, dict):
|
|
# Retrieve token
|
|
token = d.get('token')
|
|
if not token:
|
|
token = d.get('action', {}).get('payload', {}).get('token')
|
|
if not token:
|
|
token = d.get('data', {}).get('token')
|
|
|
|
# Retrieve title
|
|
title = d.get('title') or d.get('data', {}).get('title')
|
|
|
|
# Verify it represents an ad widget (token length >= 6 and has a title)
|
|
if token and title and isinstance(token, str) and len(token) >= 6:
|
|
desc = (
|
|
d.get('description') or
|
|
d.get('data', {}).get('description') or
|
|
d.get('data', {}).get('top_description_text') or
|
|
d.get('data', {}).get('middle_description_text') or
|
|
d.get('data', {}).get('bottom_description_text') or
|
|
''
|
|
)
|
|
|
|
price = (
|
|
d.get('price') or
|
|
d.get('data', {}).get('price_text') or
|
|
d.get('data', {}).get('price') or
|
|
''
|
|
)
|
|
|
|
image_url = ''
|
|
image_data = d.get('image') or d.get('data', {}).get('image') or d.get('data', {}).get('image_url')
|
|
if isinstance(image_data, str):
|
|
image_url = image_data
|
|
elif isinstance(image_data, dict):
|
|
image_url = image_data.get('url') or image_data.get('src') or ''
|
|
|
|
ad_info = {
|
|
'token': token,
|
|
'title': title,
|
|
'description': desc,
|
|
'price': price,
|
|
'image_url': image_url,
|
|
'category': d.get('category') or d.get('data', {}).get('category') or ''
|
|
}
|
|
found_ads[token] = ad_info
|
|
|
|
for v in d.values():
|
|
extract_widgets_from_dict(v, found_ads)
|
|
|
|
elif isinstance(d, list):
|
|
for item in d:
|
|
extract_widgets_from_dict(item, found_ads)
|
|
|
|
@shared_task
|
|
def run_crawl_pipeline(run_id, force=False):
|
|
"""
|
|
Background Celery task that executes a crawling run.
|
|
Parses Divar listing data, registers new Ads, and queues AI evaluations.
|
|
"""
|
|
try:
|
|
run = CrawlRun.objects.get(pk=run_id)
|
|
except CrawlRun.DoesNotExist:
|
|
return
|
|
|
|
task = run.crawl_task
|
|
|
|
# 1. Verification checks
|
|
if not task.is_active:
|
|
run.status = 'FAILED'
|
|
run.finished_at = timezone.now()
|
|
run.error_log = "CrawlTask is inactive."
|
|
run.save()
|
|
return
|
|
|
|
if not force:
|
|
now = timezone.localtime(timezone.now()).time()
|
|
if not (task.start_hour <= now <= task.end_hour):
|
|
run.status = 'SUCCESS'
|
|
run.finished_at = timezone.now()
|
|
run.error_log = f"Skipped: current time {now.strftime('%H:%M')} is outside allowed window {task.start_hour.strftime('%H:%M')} to {task.end_hour.strftime('%H:%M')}"
|
|
run.save()
|
|
return
|
|
|
|
# Update run status to RUNNING
|
|
run.status = 'RUNNING'
|
|
run.save()
|
|
|
|
headers = {
|
|
'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/115.0.0.0 Safari/537.36'
|
|
}
|
|
|
|
city, category, query_params = parse_divar_url(task.divar_url)
|
|
found_ads = {}
|
|
|
|
# 2. Try fetching from Divar public web search API
|
|
api_url = f"https://api.divar.ir/v8/web-search/{city}/{category}"
|
|
try:
|
|
res = requests.get(api_url, params=query_params, headers=headers, timeout=15)
|
|
if res.status_code == 200:
|
|
extract_widgets_from_dict(res.json(), found_ads)
|
|
except Exception as api_err:
|
|
# Log error but proceed to HTML scraper fallback
|
|
pass
|
|
|
|
# 3. Fallback: Parse HTML state if API failed/returned nothing
|
|
if not found_ads:
|
|
try:
|
|
res = requests.get(task.divar_url, headers=headers, timeout=15)
|
|
if res.status_code == 200:
|
|
match = re.search(r'window\.__PRELOADED_STATE__\s*=\s*(\{[\s\S]*?\});', res.text)
|
|
if not match:
|
|
match = re.search(r'window\.__INITIAL_STATE__\s*=\s*(\{[\s\S]*?\});', res.text)
|
|
if match:
|
|
state_data = json.loads(match.group(1))
|
|
extract_widgets_from_dict(state_data, found_ads)
|
|
except Exception as html_err:
|
|
run.status = 'FAILED'
|
|
run.finished_at = timezone.now()
|
|
run.error_log = f"Scraping failed. API and HTML fallbacks both failed.\nAPI Error: {str(api_err) if 'api_err' in locals() else 'None'}\nHTML Error: {str(html_err)}"
|
|
run.save()
|
|
return
|
|
|
|
# 4. Save listings & trigger AI evaluations
|
|
ads_fetched = len(found_ads)
|
|
new_eval_ids = []
|
|
|
|
for token, ad_info in found_ads.items():
|
|
# Check if run was cancelled/stopped by user during loop
|
|
run.refresh_from_db()
|
|
if run.status != 'RUNNING':
|
|
return
|
|
|
|
try:
|
|
# Create or get Ad metadata
|
|
ad, created = Ad.objects.get_or_create(
|
|
divar_token=token,
|
|
defaults={
|
|
'title': ad_info['title'],
|
|
'description': ad_info['description'],
|
|
'price': ad_info['price'],
|
|
'category': ad_info['category'],
|
|
'images': [ad_info['image_url']] if ad_info['image_url'] else [],
|
|
'url': f"https://divar.ir/v/{token}"
|
|
}
|
|
)
|
|
|
|
# Check if an evaluation already exists for this crawler on this ad
|
|
eval_exists = AdEvaluation.objects.filter(ad=ad, crawl_task=task).exists()
|
|
if not eval_exists:
|
|
eval_obj = AdEvaluation.objects.create(ad=ad, crawl_task=task)
|
|
new_eval_ids.append(str(eval_obj.id))
|
|
except Exception:
|
|
# If a single ad save fails, continue processing other ads
|
|
pass
|
|
|
|
# Batch process AI evaluations (5 ads per API request)
|
|
if new_eval_ids:
|
|
from ads.tasks import evaluate_ad_batch_with_ai
|
|
batch_size = 5
|
|
for i in range(0, len(new_eval_ids), batch_size):
|
|
batch = new_eval_ids[i:i + batch_size]
|
|
evaluate_ad_batch_with_ai.delay(batch)
|
|
|
|
# Update runs counters and status
|
|
run.refresh_from_db()
|
|
if run.status == 'RUNNING':
|
|
run.status = 'SUCCESS'
|
|
run.finished_at = timezone.now()
|
|
run.ads_fetched_count = ads_fetched
|
|
run.ads_evaluated_count = len(new_eval_ids)
|
|
run.save()
|
|
|
|
|
|
@shared_task
|
|
def check_and_schedule_crawl_tasks():
|
|
"""
|
|
Periodic master task executed every minute by Celery Beat.
|
|
Evaluates all active CrawlTasks and dispatches run_crawl_pipeline if:
|
|
- Current time falls within start_hour and end_hour window
|
|
- Time elapsed since the last run is >= interval_minutes (or has never run)
|
|
"""
|
|
now_dt = timezone.now()
|
|
now_time = timezone.localtime(now_dt).time()
|
|
|
|
active_tasks = CrawlTask.objects.filter(is_active=True)
|
|
for task in active_tasks:
|
|
# Check start_hour and end_hour window (handles overnight windows like 22:00 to 06:00)
|
|
in_window = False
|
|
if task.start_hour <= task.end_hour:
|
|
in_window = task.start_hour <= now_time <= task.end_hour
|
|
else:
|
|
in_window = now_time >= task.start_hour or now_time <= task.end_hour
|
|
|
|
if not in_window:
|
|
continue
|
|
|
|
# Check last run
|
|
last_run = CrawlRun.objects.filter(crawl_task=task).order_by('-started_at').first()
|
|
if last_run:
|
|
# If a run is currently in progress, skip scheduling another
|
|
if last_run.status == 'RUNNING':
|
|
continue
|
|
|
|
elapsed_minutes = (now_dt - last_run.started_at).total_seconds() / 60.0
|
|
if elapsed_minutes < task.interval_minutes:
|
|
continue
|
|
|
|
# Create new run and trigger background pipeline
|
|
run = CrawlRun.objects.create(crawl_task=task, status='RUNNING')
|
|
run_crawl_pipeline.delay(str(run.id))
|
|
|
|
|