import json
import logging
import os
import threading
import time
import urllib.error
import urllib.request
from datetime import datetime, timedelta, timezone
from db import audit, connect, now, one, rows
from engine import ingest_event, process
from scheduling import IST

log=logging.getLogger('routeflow.worker')

def run_once():
    with connect() as db:
        events=rows(db,"SELECT * FROM webhook_events WHERE state='pending' ORDER BY id LIMIT 20")
    for event in events:
        try:
            with connect() as db:
                ingest_event(db,json.loads(event['raw_body']))
                db.execute("UPDATE webhook_events SET state='done',error=NULL WHERE id=?",(event['id'],))
        except Exception as e:
            log.exception('Webhook processing failed: %s',event['id'])
            with connect() as db:
                db.execute("UPDATE webhook_events SET state='failed',attempts=attempts+1,error=? WHERE id=?",(str(e)[:500],event['id']))
    with connect() as db:
        messages=rows(db,"SELECT id FROM messages WHERE direction='in' AND state='pending' ORDER BY id LIMIT 50")
    for message in messages:
        try:
            with connect() as db:
                process(db,message['id'])
        except Exception:
            log.exception('Message processing failed: %s',message['id'])
            from engine import review
            with connect() as db:
                review(db,message['id'],'Processing failed; raw message retained for staff')
    with connect() as db:
        outgoing=rows(db,"SELECT o.*,c.last_inbound_at,c.takeover,c.number_id,k.phone,n.phone_number_id,n.token_env FROM outbox o JOIN conversations c ON c.id=o.conversation_id JOIN customers k ON k.id=c.customer_id JOIN whatsapp_numbers n ON n.id=c.number_id WHERE o.state='pending' AND o.next_attempt_at<=? ORDER BY o.id LIMIT 20",(now(),))
        automation=one(db,"SELECT value FROM settings WHERE business_id=1 AND key='automation'")['value']
    for out in outgoing:
        try:
            payload=json.loads(out['payload_json'])
            if (out['takeover'] or automation!='true') and not payload.get('manual'):
                with connect() as db:
                    db.execute("UPDATE outbox SET state='suppressed' WHERE id=?",(out['id'],))
                    db.execute("UPDATE messages SET state='suppressed' WHERE id=?",(out['message_id'],))
                continue
            if out['phone_number_id']=='local':
                external_id='local-'+str(out['id'])
                state='simulated'
            elif out['phone_number_id'].startswith('linked:'):
                from linked import bridge_request
                if payload.get('template'):raise ValueError('Meta templates only apply to Cloud API numbers')
                sent=bridge_request('/send',{'phone':out['phone'],'body':payload['text'],'idempotencyKey':'outbox-'+str(out['id'])})
                external_id=sent['id'];state='sent'
            else:
                token=os.environ.get(out['token_env'],'')
                version=os.environ.get('WHATSAPP_GRAPH_VERSION','')
                if not token or not version: raise ValueError('WhatsApp token / Graph API version not configured')
                if not payload.get('template') and (not out['last_inbound_at'] or datetime.now(timezone.utc)-datetime.fromisoformat(out['last_inbound_at'])>timedelta(hours=24)):
                    raise ValueError('24-hour reply window closed; an approved Meta template is required')
                data={'messaging_product':'whatsapp','to':out['phone'],'type':'text','text':{'body':payload['text']}}
                if payload.get('template'):
                    data={'messaging_product':'whatsapp','to':out['phone'],'type':'template','template':payload['template']}
                req=urllib.request.Request(f"https://graph.facebook.com/{version}/{out['phone_number_id']}/messages",data=json.dumps(data).encode(),headers={'Authorization':'Bearer '+token,'Content-Type':'application/json'})
                with urllib.request.urlopen(req,timeout=20) as res:
                    result=json.load(res)
                external_id=result['messages'][0]['id']
                state='sent'
            with connect() as db:
                db.execute('UPDATE outbox SET state=?,external_id=?,error=NULL WHERE id=?',(state,external_id,out['id']))
                db.execute('UPDATE messages SET state=? WHERE id=?',(state,out['message_id']))
                if out['phone_number_id'].startswith('linked:'):
                    key='linked:main:'+out['phone']+'@s.whatsapp.net:'+external_id
                    if not one(db,'SELECT id FROM messages WHERE external_id=?',(key,)):
                        db.execute('UPDATE messages SET external_id=? WHERE id=?',(key,out['message_id']))
                audit(db,'WhatsApp Reply '+state,'messages',out['message_id'])
        except Exception as e:
            # A timeout can mean Meta accepted the send. Hold for review instead of
            # blindly retrying and sending duplicate confirmations.
            transient=isinstance(e,urllib.error.HTTPError) and e.code==429
            attempts=out['attempts']+1
            with connect() as db:
                db.execute('UPDATE outbox SET state=?,attempts=?,error=?,next_attempt_at=? WHERE id=?',('pending' if transient and attempts<5 else 'failed',attempts,type(e).__name__+': send failed; check server configuration / Meta delivery logs',(datetime.now(timezone.utc)+timedelta(seconds=min(3600,30*2**attempts))).isoformat(),out['id']))
                db.execute("UPDATE messages SET state='send_failed' WHERE id=?",(out['message_id'],))
    with connect() as db:
        time_setting=one(db,"SELECT value FROM settings WHERE key='summary_time' AND business_id=1")['value']
        today=datetime.now(IST)
        target=(today.date()+timedelta(days=1)).isoformat()
        if today.strftime('%H:%M')>=time_setting:
            existing=one(db,"SELECT id FROM notifications WHERE category='Daily Summary' AND body LIKE ?",(target+'%',))
            if not existing:
                counts=rows(db,"SELECT r.name,COUNT(*) AS orders,COUNT(DISTINCT o.customer_id) AS customers FROM orders o JOIN routes r ON r.id=o.route_id WHERE o.delivery_date=? AND o.status!='Cancelled' GROUP BY r.id",(target,))
                body=target+' delivery summary\n'+'\n'.join(f"{r['name']}: {r['orders']} orders / {r['customers']} customers" for r in counts)
                details=rows(db,"SELECT r.name AS route,a.name AS area,COUNT(*) AS orders,COALESCE(v.name,'Unassigned') AS vehicle FROM orders o JOIN routes r ON r.id=o.route_id JOIN areas a ON a.id=o.area_id LEFT JOIN delivery_schedules s ON s.route_id=o.route_id AND s.date=o.delivery_date LEFT JOIN vehicles v ON v.id=s.vehicle_id WHERE o.delivery_date=? AND o.status!='Cancelled' GROUP BY r.id,a.id",(target,))
                body+='\n'+'\n'.join(f"{r['area']}: {r['orders']} orders · Vehicle: {r['vehicle']}" for r in details)
                db.execute("INSERT INTO notifications(business_id,category,body,created_at) VALUES(1,'Daily Summary',?,?)",(body,now()))
                owner=os.environ.get('OWNER_SUMMARY_PHONE','')
                template=os.environ.get('OWNER_SUMMARY_TEMPLATE','')
                number=one(db,'SELECT * FROM whatsapp_numbers WHERE phone_number_id=? AND active=1',(os.environ.get('OWNER_SUMMARY_PHONE_NUMBER_ID',''),))
                if owner and template and number:
                    from engine import phone,reply
                    owner=phone(owner)
                    db.execute("INSERT OR IGNORE INTO customers(business_id,name,phone,notes) VALUES(1,'Business owner',?,'Summary recipient')",(owner,))
                    customer=one(db,'SELECT id FROM customers WHERE business_id=1 AND phone=?',(owner,))
                    db.execute('INSERT OR IGNORE INTO conversations(business_id,customer_id,number_id) VALUES(1,?,?)',(customer['id'],number['id']))
                    conv=one(db,'SELECT id FROM conversations WHERE customer_id=? AND number_id=?',(customer['id'],number['id']))
                    mid=reply(db,conv['id'],body)
                    payload={'text':body,'manual':True,'template':{'name':template,'language':{'code':os.environ.get('OWNER_SUMMARY_LANGUAGE','en')},'components':[{'type':'body','parameters':[{'type':'text','text':body[:900]}]}]}}
                    db.execute('UPDATE outbox SET payload_json=? WHERE message_id=?',(json.dumps(payload),mid))
                    audit(db,'Daily Owner Summary Queued','messages',mid)

def start():
    def loop():
        while True:
            try: run_once()
            except Exception: log.exception('Worker cycle failed')
            time.sleep(2)
    thread=threading.Thread(target=loop,daemon=True,name='queue-worker')
    thread.start()
    return thread
