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.
397 lines
15 KiB
397 lines
15 KiB
import os
|
|
import re
|
|
import json
|
|
import html
|
|
import logging
|
|
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
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
# 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")
|
|
|
|
|
|
def extract_and_parse_json(text: str):
|
|
"""
|
|
Cleans raw LLM text (removing markdown code fences, preambles, thinking process)
|
|
and extracts a valid JSON object or list.
|
|
"""
|
|
if not text:
|
|
raise ValueError("Empty response from LLM")
|
|
|
|
text = text.strip()
|
|
|
|
# 1. Check for markdown code fences
|
|
fence_match = re.search(r'```(?:json)?\s*([\s\S]*?)\s*```', text, re.IGNORECASE)
|
|
if fence_match:
|
|
text = fence_match.group(1).strip()
|
|
|
|
# 2. Try direct json.loads
|
|
try:
|
|
return json.loads(text)
|
|
except json.JSONDecodeError:
|
|
pass
|
|
|
|
# 3. Find outermost '{' ... '}' or '[' ... ']'
|
|
start_dict = text.find('{')
|
|
end_dict = text.rfind('}')
|
|
start_list = text.find('[')
|
|
end_list = text.rfind(']')
|
|
|
|
if start_dict != -1 and end_dict != -1 and end_dict > start_dict:
|
|
try:
|
|
return json.loads(text[start_dict:end_dict + 1])
|
|
except json.JSONDecodeError:
|
|
pass
|
|
|
|
if start_list != -1 and end_list != -1 and end_list > start_list:
|
|
try:
|
|
return json.loads(text[start_list:end_list + 1])
|
|
except json.JSONDecodeError:
|
|
pass
|
|
|
|
raise ValueError(f"Could not extract valid JSON from LLM response: {text[:200]}")
|
|
|
|
|
|
def call_llm_with_structured_fallback(client, model_name, system_instruction, user_content, pydantic_cls, max_tokens, timeout):
|
|
"""
|
|
Attempts to call LLM using OpenAI beta parse first, then falls back to standard completions
|
|
with json_object response format and robust markdown-stripping JSON parsing.
|
|
"""
|
|
# 1. First attempt: beta structured parse (ideal for official OpenAI models / compatible endpoints)
|
|
try:
|
|
completion = client.beta.chat.completions.parse(
|
|
model=model_name,
|
|
messages=[
|
|
{"role": "system", "content": system_instruction},
|
|
{"role": "user", "content": user_content}
|
|
],
|
|
response_format=pydantic_cls,
|
|
max_tokens=max_tokens,
|
|
timeout=timeout
|
|
)
|
|
if completion.choices and completion.choices[0].message.parsed:
|
|
return completion.choices[0].message.parsed
|
|
except Exception:
|
|
# Fallback to chat completions if beta parse fails (e.g. OpenRouter returning markdown wrappers or non-strict JSON)
|
|
pass
|
|
|
|
# 2. Second attempt: standard chat completions with json_object response format
|
|
system_instruction_json = (
|
|
f"{system_instruction}\n"
|
|
"CRITICAL REQUIREMENT: You MUST reply strictly with a valid JSON object. "
|
|
"Do NOT wrap the JSON in markdown code blocks like ```json ... ```. "
|
|
"Do NOT include any introduction, thinking process, or explanatory text before or after the JSON."
|
|
)
|
|
|
|
try:
|
|
completion = client.chat.completions.create(
|
|
model=model_name,
|
|
messages=[
|
|
{"role": "system", "content": system_instruction_json},
|
|
{"role": "user", "content": user_content}
|
|
],
|
|
response_format={"type": "json_object"},
|
|
max_tokens=max_tokens,
|
|
timeout=timeout
|
|
)
|
|
raw_text = completion.choices[0].message.content or ""
|
|
json_data = extract_and_parse_json(raw_text)
|
|
return pydantic_cls.model_validate(json_data)
|
|
except Exception:
|
|
pass
|
|
|
|
# 3. Third attempt: standard chat completions without response_format constraint
|
|
completion = client.chat.completions.create(
|
|
model=model_name,
|
|
messages=[
|
|
{"role": "system", "content": system_instruction_json},
|
|
{"role": "user", "content": user_content}
|
|
],
|
|
max_tokens=max_tokens,
|
|
timeout=timeout
|
|
)
|
|
raw_text = completion.choices[0].message.content or ""
|
|
json_data = extract_and_parse_json(raw_text)
|
|
return pydantic_cls.model_validate(json_data)
|
|
|
|
|
|
@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 / OpenRouter 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 any(k in api_key for k in ['placeholder', 'your_openai_api_key', 'your-openai-api-key-here']):
|
|
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)
|
|
|
|
result = call_llm_with_structured_fallback(
|
|
client=client,
|
|
model_name=model_name,
|
|
system_instruction=system_instruction,
|
|
user_content=user_content,
|
|
pydantic_cls=AdEvaluationResult,
|
|
max_tokens=2500,
|
|
timeout=35
|
|
)
|
|
|
|
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:\n{ev.ad.description[:2500]}\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)
|
|
|
|
parsed = call_llm_with_structured_fallback(
|
|
client=client,
|
|
model_name=model_name,
|
|
system_instruction=system_instruction,
|
|
user_content=user_content,
|
|
pydantic_cls=BatchAdEvaluationResult,
|
|
max_tokens=8000,
|
|
timeout=90
|
|
)
|
|
|
|
res_map = {}
|
|
for idx, res in enumerate(parsed.results):
|
|
eval_id_str = str(res.evaluation_id).strip()
|
|
res_map[eval_id_str] = res
|
|
res_map[str(idx)] = res
|
|
res_map[f"Ad #{idx + 1}"] = res
|
|
res_map[f"#{idx + 1}"] = res
|
|
|
|
flagged_count = 0
|
|
for idx, ev in enumerate(evaluations):
|
|
ev_id_str = str(ev.id).strip()
|
|
res = res_map.get(ev_id_str)
|
|
if not res and idx < len(parsed.results):
|
|
res = parsed.results[idx]
|
|
|
|
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"🔔 <b>آگهی پرچمگذاری شده دیوار</b>\n\n"
|
|
f"📌 <b>عنوان:</b> {title_esc}\n"
|
|
f"💰 <b>قیمت:</b> {price_esc}\n"
|
|
f"🗂 <b>دستهبندی:</b> {cat_esc}\n\n"
|
|
f"🤖 <b>علت انتخاب AI:</b>\n{reason_esc}\n\n"
|
|
f"🔗 <a href='{ad.url}'>مشاهده آگهی در دیوار</a>"
|
|
)
|
|
|
|
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)
|
|
)
|