"""Authenticated local bridge + inert historical import + opt-in keyword flows."""
import json
import os
import secrets
import shutil
import subprocess
import urllib.error
import urllib.request
from datetime import datetime,timezone,timedelta
from db import ROOT,audit,now,one,rows
from engine import receive,reply

PROCESS=None
SECRET=None

def bridge_secret():
    global SECRET
    if SECRET is None:
        path=ROOT/'data'/'bridge-secret.txt'
        path.parent.mkdir(parents=True,exist_ok=True)
        if not path.exists():path.write_text(secrets.token_urlsafe(48),encoding='utf-8')
        SECRET=path.read_text(encoding='utf-8').strip()
    return SECRET

def start_bridge():
    global PROCESS
    if PROCESS and PROCESS.poll() is None:return
    try:
        bridge_request('/status')
        return
    except ValueError:pass
    if not (ROOT/'bridge'/'node_modules'/'@whiskeysockets'/'baileys').exists():raise ValueError('Linked-device dependencies missing: run npm install in bridge folder')
    node=shutil.which('node')
    if not node:raise ValueError('Node.js is required for linked devices')
    env={**os.environ,'ROUTEFLOW_DATA':str(ROOT/'data'),'ROUTEFLOW_BRIDGE_SECRET':bridge_secret(),'ROUTEFLOW_BACKEND':'http://127.0.0.1:'+os.environ.get('PORT','8765')}
    log=open(ROOT/'data'/'bridge.log','ab')
    try:PROCESS=subprocess.Popen([node,str(ROOT/'bridge'/'index.mjs')],cwd=ROOT/'bridge',env=env,stdout=log,stderr=log,creationflags=subprocess.CREATE_NO_WINDOW if os.name=='nt' else 0)
    finally:log.close()

def bridge_request(endpoint,data=None):
    req=urllib.request.Request('http://127.0.0.1:'+os.environ.get('ROUTEFLOW_BRIDGE_PORT','8766')+endpoint,data=json.dumps(data).encode() if data is not None else None,headers={'Content-Type':'application/json','X-Bridge-Secret':bridge_secret()})
    try:
        with urllib.request.urlopen(req,timeout=30) as res:return json.load(res)
    except urllib.error.HTTPError as e:
        try:reason=json.load(e).get('error','Linked-device request failed')
        finally:e.close()
        raise ValueError(reason)
    except urllib.error.URLError as e:raise ValueError('Linked-device service is starting or unavailable. Retry shortly; check data/bridge.log.') from e

def matched_flow(db,text,number_id=None):
    for flow in rows(db,'SELECT * FROM keyword_flows WHERE business_id=1 AND active=1 ORDER BY priority DESC,id'):
        if flow['number_id'] and flow['number_id']!=number_id:continue
        words=[w.strip().casefold() for w in flow['keywords'].split(',') if w.strip()]
        t=text.strip().casefold()
        if not words:continue
        yes= any(t==w for w in words) if flow['match_mode']=='exact' else all(w in t for w in words) if flow['match_mode']=='all' else any(w in t for w in words)
        if yes:return flow
    return None

def ensure_number(db):
    db.execute("INSERT OR IGNORE INTO whatsapp_numbers(business_id,label,phone_number_id,token_env) VALUES(1,'Linked WhatsApp Business','linked:main','LINKED_DEVICE')")
    return one(db,"SELECT id FROM whatsapp_numbers WHERE phone_number_id='linked:main'")['id']

def contact(db,entry):
    jid=str(entry['jid'])
    phone=entry.get('phone')
    customer_id=None
    if phone and not entry.get('isGroup'):
        import re
        if not re.fullmatch(r'\d{8,15}',str(phone)):raise ValueError('Invalid linked phone')
        db.execute('INSERT OR IGNORE INTO customers(business_id,name,phone) VALUES(1,?,?)',(entry.get('name') or phone,phone))
        customer_id=one(db,'SELECT id FROM customers WHERE business_id=1 AND phone=?',(phone,))['id']
        if entry.get('name') and entry['name']!=jid:
            db.execute('UPDATE customers SET name=? WHERE id=? AND (name=phone OR name LIKE \'%@%\')',(entry['name'],customer_id))
    db.execute("INSERT INTO linked_chats(session,jid,name,phone,is_group,customer_id,updated_at) VALUES('main',?,?,?,?,?,?) ON CONFLICT(session,jid) DO UPDATE SET name=CASE WHEN excluded.name=excluded.jid THEN linked_chats.name ELSE excluded.name END,phone=COALESCE(excluded.phone,linked_chats.phone),customer_id=COALESCE(excluded.customer_id,linked_chats.customer_id),updated_at=excluded.updated_at",(jid,entry.get('name') or jid,phone,int(bool(entry.get('isGroup'))),customer_id,now()))
    return one(db,"SELECT * FROM linked_chats WHERE session='main' AND jid=?",(jid,))

def ingest(db,event):
    number_id=ensure_number(db)
    if event['type']=='connected':
        audit(db,'WhatsApp Device Linked',detail={'session':'main'})
        return
    if event['type']=='contacts':
        for entry in event.get('contacts',[]):contact(db,entry)
        return
    if event['type']!='messages':raise ValueError('Unknown linked event')
    for m in event.get('messages',[]):
        chat=contact(db,m)
        external='linked:main:'+str(m['jid'])+':'+str(m['id'])
        if one(db,'SELECT id FROM linked_messages WHERE external_id=?',(external,)):continue
        stamp=datetime.fromtimestamp(float(m.get('timestamp',0)),timezone.utc).isoformat()
        db.execute('INSERT INTO linked_messages(chat_id,external_id,direction,kind,body,raw_json,is_history,created_at) VALUES(?,?,?,?,?,?,?,?)',(chat['id'],external,m['direction'],m['kind'],m.get('body',''),json.dumps(m.get('raw',{}),ensure_ascii=False),int(bool(m.get('history'))),stamp))
        if not chat['customer_id']:continue # Groups and unresolved LIDs are visible, never guessed as customers.
        c=one(db,'SELECT * FROM customers WHERE id=?',(chat['customer_id'],))
        if m['direction']=='out':
            db.execute('INSERT OR IGNORE INTO conversations(business_id,customer_id,number_id) VALUES(1,?,?)',(c['id'],number_id))
            conv=one(db,'SELECT * FROM conversations WHERE customer_id=? AND number_id=?',(c['id'],number_id))
            if not one(db,'SELECT id FROM messages WHERE external_id=?',(external,)):
                db.execute("INSERT INTO messages(conversation_id,external_id,direction,kind,body,raw_json,state,created_at) VALUES(?,?,'out',?,?,?,'imported',?)",(conv['id'],external,m['kind'],m.get('body',''),json.dumps(m.get('raw',{})),stamp))
            continue
        mid=receive(db,number_id,c['phone'],m.get('body',''),external,m['kind'],m.get('raw'),c['name'])
        conv=one(db,'SELECT * FROM conversations WHERE customer_id=? AND number_id=?',(c['id'],number_id))
        history=bool(m.get('history')) or datetime.now(timezone.utc)-datetime.fromisoformat(stamp)>timedelta(minutes=5)
        db.execute("UPDATE messages SET state=?,created_at=? WHERE id=?",('history' if history else 'linked_waiting',stamp,mid))
        if history:
            # Historical imports must not open the Cloud reply window or trigger flows.
            db.execute('UPDATE conversations SET last_inbound_at=(SELECT MAX(created_at) FROM messages WHERE conversation_id=? AND direction=\'in\' AND state!=\'history\') WHERE id=?',(conv['id'],conv['id']))
            continue
        if conv['takeover'] or one(db,"SELECT value FROM settings WHERE key='automation' AND business_id=1")['value']!='true':continue
        if conv['pending_json']:
            db.execute("UPDATE messages SET state='pending' WHERE id=?",(mid,));continue
        flow=matched_flow(db,m.get('body',''),number_id)
        if not flow:continue
        audit(db,'Keyword Flow Matched','messages',mid,{'flow_id':flow['id']})
        if flow['action']=='order':db.execute("UPDATE messages SET state='pending' WHERE id=?",(mid,))
        elif flow['action']=='reply':reply(db,conv['id'],flow['response'])
        else:
            from engine import review
            review(db,mid,'Keyword flow requests staff review: '+flow['name'])
