import os import html import requests from openai import OpenAI, AuthenticationError from pydantic import BaseModel, Field from celery import shared_task from django.utils import timezone from ads.models import Ad, AdEvaluation, NotificationLog # Pydantic model for strict OpenAI structured output format class AdEvaluationResult(BaseModel): is_flagged: bool = Field(description="True if the ad matches the criteria in the prompt, otherwise False") reason: str = Field(description="A concise explanation in Persian describing why the ad matches or doesn't match the criteria") confidence: float = Field(description="Confidence score between 0.0 and 1.0") class SingleBatchAdResult(BaseModel): evaluation_id: str = Field(description="The exact UUID evaluation_id provided for the ad") is_flagged: bool = Field(description="True if the ad matches the target criteria, otherwise False") reason: str = Field(description="A concise explanation in Persian describing why the ad matches or doesn't match") confidence: float = Field(description="Confidence score between 0.0 and 1.0") class BatchAdEvaluationResult(BaseModel): results: list[SingleBatchAdResult] = Field(description="List of evaluation results corresponding to each ad in the batch") @shared_task(bind=True, max_retries=3, default_retry_delay=10) def evaluate_ad_with_ai(self, evaluation_id): """ Asynchronously evaluates a single ad using OpenAI's gpt-4o-mini model against the specific CrawlTask search criteria. """ try: evaluation = AdEvaluation.objects.get(pk=evaluation_id) except AdEvaluation.DoesNotExist: return ad = evaluation.ad task = evaluation.crawl_task api_key = os.getenv('OPENAI_API_KEY') # Fallback if OpenAI API Key is missing if not api_key or 'placeholder' in api_key or 'your_openai_api_key' in api_key or 'your-openai-api-key-here' in api_key: evaluation.is_flagged = False evaluation.reason = "AI Evaluation skipped: OpenAI API key is not configured." evaluation.confidence = 0.0 evaluation.save() return try: # Support OpenRouter keys seamlessly client_kwargs = {"api_key": api_key} if api_key.startswith("sk-or-"): client_kwargs["base_url"] = "https://openrouter.ai/api/v1" client = OpenAI(**client_kwargs) system_instruction = ( "You are an expert Iranian market analyst. Your job is to read listing descriptions " "and decide if they match specific target criteria. You MUST reply using the structured JSON response format " "with Persian strings." ) user_content = ( f"User Search Criteria: {task.detection_prompt}\n\n" f"Ad Title: {ad.title}\n" f"Ad Price: {ad.price or 'Not specified'}\n" f"Ad Category: {ad.category or 'Not specified'}\n" f"Ad Description:\n{ad.description}" ) default_model = "openrouter/free" if api_key.startswith("sk-or-") else "gpt-4o-mini" model_name = os.getenv("OPENAI_MODEL", default_model) completion = client.beta.chat.completions.parse( model=model_name, messages=[ {"role": "system", "content": system_instruction}, {"role": "user", "content": user_content} ], response_format=AdEvaluationResult, max_tokens=300, timeout=25 ) result = completion.choices[0].message.parsed evaluation.is_flagged = result.is_flagged evaluation.reason = result.reason evaluation.confidence = result.confidence evaluation.save() # Update run stats on successful AI flagging if result.is_flagged: from crawler.models import CrawlRun from django.db.models import F latest_run = CrawlRun.objects.filter(crawl_task=task).order_by('-started_at').first() if latest_run: latest_run.ads_flagged_count = F('ads_flagged_count') + 1 latest_run.save(update_fields=['ads_flagged_count']) # Send telegram channel notification if configured if task.telegram_channel_id: send_telegram_notification.delay(evaluation.id) except AuthenticationError as auth_err: evaluation.is_flagged = False evaluation.reason = "AI Evaluation failed: Invalid or incorrect OpenAI API key." evaluation.confidence = 0.0 evaluation.save() except Exception as e: # Retry in case of API rate limits or network issues try: self.retry(exc=e) except self.MaxRetriesExceededError: evaluation.is_flagged = False evaluation.reason = f"AI Evaluation failed: {str(e)}" evaluation.confidence = 0.0 evaluation.save() @shared_task(bind=True, max_retries=3, default_retry_delay=10) def evaluate_ad_batch_with_ai(self, evaluation_ids): """ Asynchronously evaluates a batch of ads in a single LLM API call. """ if not evaluation_ids: return evaluations = list(AdEvaluation.objects.filter(id__in=evaluation_ids).select_related('ad', 'crawl_task')) if not evaluations: return api_key = os.getenv('OPENAI_API_KEY') if not api_key or any(k in api_key for k in ['placeholder', 'your_openai_api_key', 'your-openai-api-key-here']): for ev in evaluations: ev.is_flagged = False ev.reason = "AI Evaluation skipped: OpenAI API key is not configured." ev.confidence = 0.0 ev.save() return task = evaluations[0].crawl_task try: client_kwargs = {"api_key": api_key} if api_key.startswith("sk-or-"): client_kwargs["base_url"] = "https://openrouter.ai/api/v1" client = OpenAI(**client_kwargs) system_instruction = ( "You are an expert Iranian market analyst. Your job is to evaluate a batch of ad listings " "against the target search criteria. You MUST reply with a structured JSON object containing " "evaluations for each ad ID provided, using Persian text for the reasons." ) ads_text = [] for index, ev in enumerate(evaluations, 1): ads_text.append( f"--- Ad #{index} ---\n" f"Evaluation ID: {ev.id}\n" f"Title: {ev.ad.title}\n" f"Price: {ev.ad.price or 'Not specified'}\n" f"Category: {ev.ad.category or 'Not specified'}\n" f"Description: {ev.ad.description[:600]}\n" ) user_content = ( f"Target Search Criteria: {task.detection_prompt}\n\n" "Evaluate each of the following ads:\n" + "\n".join(ads_text) ) default_model = "openrouter/free" if api_key.startswith("sk-or-") else "gpt-4o-mini" model_name = os.getenv("OPENAI_MODEL", default_model) completion = client.beta.chat.completions.parse( model=model_name, messages=[ {"role": "system", "content": system_instruction}, {"role": "user", "content": user_content} ], response_format=BatchAdEvaluationResult, max_tokens=1500, timeout=45 ) parsed = completion.choices[0].message.parsed res_map = {str(res.evaluation_id): res for res in parsed.results} flagged_count = 0 for ev in evaluations: res = res_map.get(str(ev.id)) if res: ev.is_flagged = res.is_flagged ev.reason = res.reason ev.confidence = res.confidence else: ev.is_flagged = False ev.reason = "Evaluation result omitted in batch response." ev.confidence = 0.0 ev.save() if ev.is_flagged: flagged_count += 1 if task.telegram_channel_id: send_telegram_notification.delay(ev.id) if flagged_count > 0: from crawler.models import CrawlRun from django.db.models import F latest_run = CrawlRun.objects.filter(crawl_task=task).order_by('-started_at').first() if latest_run: latest_run.ads_flagged_count = F('ads_flagged_count') + flagged_count latest_run.save(update_fields=['ads_flagged_count']) except Exception as e: # Fallback to individual evaluations if batch call fails for ev in evaluations: evaluate_ad_with_ai.delay(ev.id) @shared_task(bind=True, max_retries=3, default_retry_delay=15) def send_telegram_notification(self, evaluation_id): """ Sends an HTML formatted alert message to the target Telegram Channel notifying them of a flagged ad. """ try: evaluation = AdEvaluation.objects.get(pk=evaluation_id) except AdEvaluation.DoesNotExist: return ad = evaluation.ad task = evaluation.crawl_task token = os.getenv('TELEGRAM_BOT_TOKEN') channel = task.telegram_channel_id if not token or not channel: return # Escape HTML to prevent telegram parsing errors title_esc = html.escape(ad.title) price_esc = html.escape(ad.price or 'مشخص نشده') cat_esc = html.escape(ad.category or 'مشخص نشده') reason_esc = html.escape(evaluation.reason or '') message_html = ( f"🔔 آگهی پرچم‌گذاری شده دیوار\n\n" f"📌 عنوان: {title_esc}\n" f"💰 قیمت: {price_esc}\n" f"🗂 دسته‌بندی: {cat_esc}\n\n" f"🤖 علت انتخاب AI:\n{reason_esc}\n\n" f"🔗 مشاهده آگهی در دیوار" ) url = f"https://api.telegram.org/bot{token}/sendMessage" payload = { 'chat_id': channel, 'text': message_html, 'parse_mode': 'HTML' } try: res = requests.post(url, json=payload, timeout=10) if res.status_code == 200: NotificationLog.objects.create( evaluation=evaluation, channel_id=channel, status='SENT' ) else: raise Exception(f"Telegram API responded with code {res.status_code}: {res.text}") except Exception as e: try: self.retry(exc=e) except self.MaxRetriesExceededError: NotificationLog.objects.create( evaluation=evaluation, channel_id=channel, status='FAILED', error_message=str(e) )