leads-validator / main.py
vibhansh's picture
Update main.py
09531df verified
Raw
History Blame Contribute Delete
14.8 kB
import os
import time
import json
import random
import asyncio
import requests
import phonenumbers
from datetime import datetime, timedelta
from bs4 import BeautifulSoup
from fastapi import FastAPI, Depends, HTTPException, Security, BackgroundTasks
from fastapi.security import HTTPBearer, HTTPAuthorizationCredentials
from pydantic import BaseModel
from huggingface_hub import HfApi
app = FastAPI(title="Landsea Lead Validator")
security = HTTPBearer()
HF_TOKEN = os.getenv("HF_TOKEN")
SPACE_URL = os.getenv("SPACE_URL", "http://localhost:7860")
DATASET_REPO = os.getenv("DATASET_REPO", "vibhansh/landsea-validated-leads")
CACHE_FILE = "/tmp/validation_cache.json"
# Cache Configuration (TTL in Days)
CACHE_TTL_DAYS = 30
hf_api = HfApi(token=HF_TOKEN)
stats = {"total_validated": 0, "gold_leads": 0, "exclusive_leads": 0, "cache_pruned_entries": 0}
class LeadInput(BaseModel):
name: str
phone: str
city: str
category: str
google_maps_url: str
website: str = ""
address: str = ""
source: str = "leads-sync-engine"
def verify_token(credentials: HTTPAuthorizationCredentials = Security(security)):
if not HF_TOKEN:
raise HTTPException(status_code=500, detail="HF_TOKEN not configured")
if credentials.credentials != HF_TOKEN:
raise HTTPException(status_code=401, detail="Invalid Bearer Token")
return credentials.credentials
# ── CACHE AUTO-PRUNING ENGINE ────────────────────────────────────────────────
def prune_expired_cache():
"""Scans the local filesystem cache and removes expired entries automatically."""
global stats
if not os.path.exists(CACHE_FILE):
return {"status": "skipped", "message": "No cache file found to prune."}
try:
with open(CACHE_FILE, "r") as f:
cache = json.load(f)
now = datetime.now().timestamp()
# Filter entries whose expiration timestamp is still in the future
cleaned_cache = {k: v for k, v in cache.items() if v.get("expires_at", 0) > now}
pruned_count = len(cache) - len(cleaned_cache)
if pruned_count > 0:
with open(CACHE_FILE, "w") as f:
json.dump(cleaned_cache, f)
stats["cache_pruned_entries"] += pruned_count
print(f"[CACHE-PRUNE] Successfully removed {pruned_count} expired entries from local file storage.")
return {"status": "success", "pruned_entries": pruned_count}
return {"status": "skipped", "message": "All cache records are healthy and active."}
except Exception as e:
print(f"[CACHE-PRUNE] Error during automated execution: {e}")
return {"status": "error", "message": str(e)}
# ── VALIDATION SUB-ENGINES ───────────────────────────────────────────────────
def validate_phone_format(phone: str) -> int:
try:
parsed = phonenumbers.parse(phone, "IN")
if phonenumbers.is_valid_number(parsed):
if str(parsed.national_number)[0] in ['6', '7', '8', '9']:
return 20
except phonenumbers.phonenumberutil.NumberParseException:
pass
return 0
def check_whatsapp(phone: str) -> int:
time.sleep(random.uniform(1, 3))
try:
clean_phone = "".join(filter(str.isdigit, phone))
if not clean_phone.startswith("91") and len(clean_phone) == 10:
clean_phone = "91" + clean_phone
headers = {"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64)"}
response = requests.head(f"https://wa.me/{clean_phone}", headers=headers, allow_redirects=False, timeout=5)
if response.status_code in [301, 302] and "web.whatsapp.com" in response.headers.get("Location", ""):
return 30
except requests.RequestException:
pass
return 0
def check_justdial_and_website(name: str, city: str, category: str) -> dict:
time.sleep(random.uniform(2, 4))
query = f"{name} {city} {category}".replace(" ", "+")
url = f"https://www.justdial.com/{city}/{query}"
headers = {"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36"}
try:
response = requests.get(url, headers=headers, timeout=8)
soup = BeautifulSoup(response.text, 'html.parser')
page_text = soup.get_text().lower()
if "we could not find any results" in page_text or "no results found" in page_text:
return {"jd_points": 25, "exclusive": True}
if "website" in page_text or soup.find("a", href=lambda href: href and "http" in href and "justdial" not in href):
return {"jd_points": 5, "exclusive": False}
else:
return {"jd_points": 10, "exclusive": False}
except requests.RequestException:
return {"jd_points": 0, "exclusive": False}
def save_to_dataset(lead_data: dict):
timestamp = datetime.now().strftime("%Y%m%d_%H%M%S")
phone = lead_data.get("lead_info", {}).get("phone", "unknown")
filename = f"leads/{phone}_{timestamp}.json"
local_path = f"/tmp/{phone}_{timestamp}.json"
with open(local_path, "w") as f:
json.dump(lead_data, f)
try:
hf_api.upload_file(
path_or_fileobj=local_path,
path_in_repo=filename,
repo_id=DATASET_REPO,
repo_type="dataset"
)
except Exception as e:
print(f"Dataset upload failed: {e}")
finally:
if os.path.exists(local_path):
os.remove(local_path)
# ── ENDPOINTS ────────────────────────────────────────────────────────────────
@app.get("/")
def home():
return {
"app": "Landsea Lead Validator API",
"status": "Running - 10X Ultra UX Layer Active",
"cache_file_exists": os.path.exists(CACHE_FILE),
"docs_url": "/docs"
}
@app.post("/validate")
def validate_lead(lead: LeadInput, token: str = Depends(verify_token)):
global stats
# Improved Backend UI Execution Trail array
execution_trail = []
execution_trail.append(f"🏁 System handshake successful. Initializing process sequence for lead: {lead.name}")
cache_key = f"{lead.phone}_{lead.city}"
now_timestamp = datetime.now().timestamp()
# Load and clean cache dictionary on request
if os.path.exists(CACHE_FILE):
with open(CACHE_FILE, "r") as f:
try:
cache = json.load(f)
except Exception:
cache = {}
else:
cache = {}
# Check Cache with TTL Expiration validation logic
if cache_key in cache:
cached_entry = cache[cache_key]
if cached_entry.get("expires_at", 0) > now_timestamp:
time_left_sec = cached_entry.get("expires_at", 0) - now_timestamp
hours_left = round(time_left_sec / 3600, 2)
execution_trail.append(f"⚡ [CACHE HIT] Loaded instantly from localized memory storage. Entry is fresh. TTL remaining: {hours_left} Hours.")
# Inject descriptive tracking matrix logs inside the returned cache structure
result = cached_entry["result"]
result["backend_diagnostics"]["execution_logs"] = execution_trail
result["backend_diagnostics"]["retrieval_mode"] = "INTERNAL_CACHE_MEMORY"
return result
else:
execution_trail.append("⏳ [CACHE EXPIRED] Record matching key was discovered but it exceeded the 30-day TTL duration limit. Purging record and forcing fresh live execution sequence.")
cache.pop(cache_key, None)
execution_trail.append("🛰️ [CACHE MISS] Entry not found or expired. Launching live parsing matrix engines across the digital grid.")
score = 0
# 1. Phone Format Verification
execution_trail.append("📱 Scanning telecom network architecture parameters...")
phone_points = validate_phone_format(lead.phone)
score += phone_points
execution_trail.append(f"↳ Structure check completed. Added {phone_points} alignment points to the matrix metric.")
# 2. WhatsApp Network Check
execution_trail.append("💬 Establishing remote encryption test link to verification endpoint on wa.me API...")
wa_points = check_whatsapp(lead.phone)
score += wa_points
execution_trail.append(f"↳ Encryption handshake status logged. Added {wa_points} points.")
# 3. Justdial Scrape Check
execution_trail.append(f"🌐 Deploying BeautifulSoup crawler to cross-reference indexing logs on Justdial database...")
jd_result = check_justdial_and_website(lead.name, lead.city, lead.category)
score += jd_result["jd_points"]
is_exclusive = jd_result["exclusive"]
execution_trail.append(f"↳ Web directory validation returned score. Added {jd_result['jd_points']} points. Exclusivity Match: {is_exclusive}")
# 4. Digital Vulnerability & Social Matching
execution_trail.append("🔍 Auditing asset portfolio infrastructure for digital footprint vulnerability mapping...")
social_domains = ['facebook.com', 'instagram.com', 'linkedin.com']
weak_domains = ['blogspot', 'wordpress', 'sites.google', 'justdial', 'indiamart', 'wixsite']
is_social_only = False
is_vulnerable = False
if lead.website:
web_lower = lead.website.lower()
if any(domain in web_lower for domain in social_domains):
is_social_only, is_vulnerable = True, True
score += 15
execution_trail.append("↳ Vulnerability spotted: Asset relies purely on a corporate social network profile container. Added 15 points.")
elif any(domain in web_lower for domain in weak_domains):
is_vulnerable = True
score += 15
execution_trail.append("↳ Vulnerability spotted: Asset operates on un-optimized free hosting subdomain endpoints. Added 15 points.")
else:
score += 5
execution_trail.append("↳ Asset security clean: Custom unique top-level domain configuration mapped. Minimal digital vulnerability surface detected. Added 5 points.")
else:
is_vulnerable = True
score += 20
execution_trail.append("💥 Vulnerability Extreme: Business has absolute zero online ecosystem presence. High conversion potential. Added 20 points.")
# 5. Storefront Geolocation Structuring
execution_trail.append("📍 Deconstructing geographic address properties into string density matrix...")
storefront_type = "Unknown"
commercial_keywords = ["shop", "floor", "building", "tower", "complex", "mall", "plaza", "opp", "near", "phase"]
if lead.address:
addr_lower = lead.address.lower()
if any(kw in addr_lower for kw in commercial_keywords) or lead.address.count(',') >= 2:
storefront_type = "Verified Storefront"
score += 15
execution_trail.append("↳ High density address format verified. Profile matches established commercial physical operation unit. Added 15 points.")
else:
storefront_type = "Potential Home-Setup"
score += 5
execution_trail.append("↳ Low address structural density matched. Flagged profile as virtual asset/residential operation space. Added 5 points.")
# Capping parameters
score = min(score, 100)
# 6. Tiering assignment Logic
execution_trail.append("🥇 Compiling performance indexes and translating value vector into strategic ranking metrics...")
if score >= 90:
status = "GOLD_LEAD"
stats["gold_leads"] += 1
elif score >= 60:
status = "SILVER_LEAD"
else:
status = "BRONZE_LEAD"
execution_trail.append(f"↳ Lead classification processed successfully. Assigned Tier status: {status}")
if is_exclusive:
stats["exclusive_leads"] += 1
stats["total_validated"] += 1
# Compiling explicit results payload object
result = {
"lead_info": lead.model_dump(),
"validation": {
"score": score,
"tier": status,
"is_exclusive": is_exclusive,
"digital_vulnerability": is_vulnerable,
"is_social_only": is_social_only,
"storefront_type": storefront_type,
"timestamp": datetime.now().isoformat()
},
"backend_diagnostics": {
"retrieval_mode": "LIVE_NETWORK_SCAN",
"ttl_duration_days": CACHE_TTL_DAYS,
"execution_logs": execution_trail
}
}
# Calculate exact expiration timestamp parameter (Now + 30 Days)
expiry_timestamp = (datetime.now() + timedelta(days=CACHE_TTL_DAYS)).timestamp()
# Write into Cache Dictionary structure
cache[cache_key] = {
"expires_at": expiry_timestamp,
"result": result
}
with open(CACHE_FILE, "w") as f:
json.dump(cache, f)
execution_trail.append(f"💾 Storage pipeline synchronized. Verification logs secured into cache. Expires on: {datetime.fromtimestamp(expiry_timestamp).strftime('%Y-%m-%d %H:%M:%S')}")
save_to_dataset(result)
return result
@app.get("/health")
def health_check():
return {"status": "healthy", "timestamp": datetime.now().isoformat()}
@app.post("/prune-cache")
def manual_prune_endpoint(token: str = Depends(verify_token)):
"""Allows administrators to manually invoke the database cache pruning task pipeline."""
return prune_expired_cache()
@app.get("/stats")
def get_stats(token: str = Depends(verify_token)):
return stats
async def keep_awake_and_prune():
"""Loops every 5 minutes to stay awake and executes automatic pruning routines daily."""
loop_counter = 0
while True:
try:
requests.get(f"{SPACE_URL}/health", timeout=5)
# Every 288 ticks (which is exactly 24 Hours in 5-minute loops), auto-prune data
if loop_counter % 288 == 0:
print("[BACKGROUND-TASK] Executing daily cache validation and cleanup sweep cycle...")
prune_expired_cache()
except Exception:
pass
loop_counter += 1
await asyncio.sleep(300)
@app.on_event("startup")
async def startup_event():
# Run a prune task immediately upon container boot sequence setup
print("[STARTUP] Booting Landsea Core Layer Engine. Checking internal database files...")
prune_expired_cache()
asyncio.create_task(keep_awake_and_prune())
# --- END OF FILE ---