import time
import json
import asyncio
import re
import aiohttp
import backoff
from aiohttp import ClientSession
from urllib.parse import quote, urlencode

# Configuration
CLIENT_ID = '1000.4F1T5L8L5KN9MZYRV2DOYWY3VLOJJU'
CLIENT_SECRET = 'eb827f5d29e12e73d43882a809244ffbb482f4aec6'
REFRESH_TOKEN = '1000.a6646c7ffa21f40b0b8f5d6a82d6fe89.9fe3f5ce4a376693c3c52a7572bf868c'
REDIRECT_URI = 'https://google.com'
BASE_URL = 'https://www.zohoapis.com/crm/v2/'
TIMEOUT = 300  # seconds

def rate_limited(max_per_minute):
    min_interval = 60.0 / float(max_per_minute)
    semaphore = asyncio.Semaphore(max_per_minute)

    def decorator(func):
        async def wrapper(*args, **kwargs):
            async with semaphore:
                await asyncio.sleep(min_interval)
                return await func(*args, **kwargs)
        return wrapper
    return decorator

class ZohoCRMClient:
    def __init__(self):
        self.access_token = '1000.80ed64118fad197b4430b59e2e5974fe.ff3205bd58587a64c56de17cdfc06739'
        self.token_expiry = None
        self.headers = {
            'Authorization': f'Zoho-oauthtoken {self.access_token}',
            'Content-Type': 'application/json'
        }

    async def get_access_token(self, session):
        url = 'https://accounts.zoho.com/oauth/v2/token'
        data = {
            'refresh_token': REFRESH_TOKEN,
            'client_id': CLIENT_ID,
            'client_secret': CLIENT_SECRET,
            'redirect_uri': REDIRECT_URI,
            'grant_type': 'refresh_token',
        }
        headers = {
            'Content-Type': 'application/x-www-form-urlencoded'
        }
        async with session.post(url, data=data, headers=headers) as response:
            if response.status == 200:
                response_data = await response.json()
                self.access_token = response_data.get('access_token')
                self.headers['Authorization'] = f'Zoho-oauthtoken {self.access_token}'
                expires_in = response_data.get('expires_in', 3600)
                self.token_expiry = time.time() + expires_in
            else:
                response_text = await response.text()
                print("Failed to refresh access token:", response_text)
                return None

    def needs_token_refresh(self):
        """Determine if the access token needs to be refreshed."""
        buffer = 300  # 5 minutes buffer
        return self.token_expiry is None or (time.time() + buffer) >= self.token_expiry

    @rate_limited(20)  # 20 requests per minute
    @backoff.on_exception(backoff.expo,
                          (aiohttp.ClientError, asyncio.TimeoutError),
                          max_time=60)  # Retry for up to 60 seconds
    async def make_api_call_with_retry(self, session, endpoint, method='GET', data=None, params=None):
        try:
            complete_url = BASE_URL + endpoint
            if params:
                complete_url += '?' + urlencode(params)

            print(f"Making API call: {complete_url}")
            print(f"Method: {method}")
            print(f"Headers: {self.headers}")
            print(f"Data: {data}")

            async with session.request(method, complete_url, headers=self.headers, json=data, timeout=TIMEOUT) as response:
                response_text = await response.text()
                if response.status == 429:
                    retry_after = response.headers.get('Retry-After')
                    if retry_after:
                        await asyncio.sleep(int(retry_after))
                    else:
                        await asyncio.sleep(1)
                    raise aiohttp.ClientError
                elif response.status >= 400:
                    print(f"Error response ({response.status}): {response_text}")
                response.raise_for_status()
                return await response.json()
        except aiohttp.ClientError as e:
            print(f"HTTP Client Error during API call: {e}")
            raise
        except asyncio.TimeoutError as e:
            print(f"Timeout Error: {str(e)}")
            raise

    async def find_contact_by_email(self, session, email):
        search_criteria = f"(Email:equals:{email})"
        try:
            response = await self.make_api_call_with_retry(session, f'Contacts/search?criteria={quote(search_criteria)}')
            if response:
                return response.get('data', [])
        except Exception as e:
            print(f"Error searching contact by email {email}: {e}")
        return []

    @rate_limited(20)
    async def update_contact(self, session, contact_id, lead_details):
        print(f"Updating contact {contact_id}...")
        url = f"{BASE_URL}Contacts/{contact_id}"
        update_data = {
            "data": [
                {
                    "Email": lead_details.get("Email", ""),
                    "Lead_Created_Time": lead_details.get("Created_Time", ""),
                    "Consent": lead_details.get("Consent", False),
                    "Funded": lead_details.get("Funded", False),
                    "Link_Bank_Screen": lead_details.get("Link_Bank_Screen", False),
                    "FA": lead_details.get("FA", False),
                    "Underwriting": lead_details.get("Underwriting", False),
                    "Plaid": lead_details.get("Plaid", False),
                    "SignedContract": lead_details.get("SignedContract", False),
                    "Remarketing_Closed": lead_details.get("Remarketing_Closed", False),
                    "GiggleID": lead_details.get("GiggleID", "").split("-")[0] if lead_details.get("GiggleID") else "",
                    "Funded_Amount": lead_details.get("Funded_Amount", ""),
                    "Funded_Date": lead_details.get("Funded_Date", ""),
                    "Renewal": lead_details.get("Renewal", False),
                    "GiggleID_marketing": lead_details.get("GiggleID", ""),
                    "Partner_Account_Name": lead_details.get("Partner_Account_Name","")
                }
            ]
        }
        self.headers['Authorization'] = f'Zoho-oauthtoken {self.access_token}'
        async with session.put(url, json=update_data, headers=self.headers) as response:
            if response.status == 200:
                response_data = await response.json()
                print(f"Response from updating contact {contact_id}: {response_data}")
                return response_data
            else:
                print(f"Error updating contact {contact_id}: {response.status}, {await response.text()}")
                return None

async def is_valid_email(email):
    pattern = re.compile(r'^[a-zA-Z0-9._%+-]+@[a-zA-Z0-9.-]+\.[a-zA-Z]{2,}$')
    return pattern.match(email) is not None

async def process_single_lead(client, session, lead):
    email = lead.get("Email")
    if not email or not await is_valid_email(email):
        print(f"Invalid or missing email for lead: {lead}. Skipping...")
        return

    try:
        contacts = await client.find_contact_by_email(session, email)
        for contact in contacts:
            await client.update_contact(session, contact['id'], lead)
    except Exception as e:
        print(f"Error processing lead with email {email}: {e}")

async def fetch_and_process_leads(client, session, module_name, cvid, page, per_page, fields):
    params = {
        'cvid': cvid,
        'page': page,
        'per_page': per_page,
        'fields': fields
    }
    response = await client.make_api_call_with_retry(session, module_name, params=params)
    leads = response.get('data', [])

    tasks = [process_single_lead(client, session, lead) for lead in leads]
    await asyncio.gather(*tasks)
    
    return len(leads)

async def process_crm_data():
    client = ZohoCRMClient()
    async with aiohttp.ClientSession() as session:
        if client.needs_token_refresh():
            await client.get_access_token(session)

        module_name = 'Leads'
        cvid = '4702815000363008122'
        #cvid = '4702815000202077003'
        
        leads_already_processed = 0
        per_page = 200
        fields = 'Partner_Account_Name,Email,Lead_Created_Time,Consent,Funded,Link_Bank_Screen,FA,Underwriting,Plaid,SignedContract,Remarketing_Closed,GiggleID,Funded_Amount,Funded_Date,Renewal,GiggleID_marketing'
        
        start_page = (leads_already_processed // per_page) + 1

        page = start_page
        total_leads_processed = 0

        while True:
            if client.needs_token_refresh():
                await client.get_access_token(session)
            num_leads_processed = await fetch_and_process_leads(client, session, module_name, cvid, page, per_page, fields)
            if num_leads_processed == 0:
                break
            total_leads_processed += num_leads_processed
            print(f"Processed {total_leads_processed} leads so far...")
            if total_leads_processed == -1:
                break;
            page += 1

            await asyncio.sleep(1)

        print(f"Total Leads Processed: {total_leads_processed}")

if __name__ == '__main__':
    asyncio.run(process_crm_data())
