import asyncio
import httpx
import json
import logging
import websockets
import uuid

# Configuration
BASE_URL = "https://dev-strategistai.galaxiq.ai"
WS_URL = "wss://dev-strategistai.galaxiq.ai"
TENANT_ID = "org_d1e1c1ea-6c90-431f-9337-3375dcdfef5c"

logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s')
logger = logging.getLogger("verify_handover")

class HandoverTester:
    def __init__(self, thread_id):
        self.thread_id = thread_id
        self.admin_received = []
        self.user_received = []
        self.stop_event = asyncio.Event()

    async def listen_admin(self):
        url = f"{WS_URL}/ws/{TENANT_ID}/{self.thread_id}?role=admin&user_name=AdminUser"
        logger.info(f"Admin connecting to WS: {url}")
        try:
            async with websockets.connect(url) as ws:
                logger.info("Admin connected.")
                while not self.stop_event.is_set():
                    try:
                        msg = await asyncio.wait_for(ws.recv(), timeout=1.0)
                        data = json.loads(msg)
                        logger.info(f" [ADMIN RECEIVED]: {data.get('role')}: {data.get('text') or data.get('message')}")
                        self.admin_received.append(data)
                    except asyncio.TimeoutError:
                        continue
        except Exception as e:
            logger.error(f"Admin WS error: {e}")

    async def listen_user(self):
        url = f"{WS_URL}/ws/{TENANT_ID}/{self.thread_id}?role=user&user_name=EndUser"
        logger.info(f"User connecting to WS: {url}")
        try:
            async with websockets.connect(url) as ws:
                logger.info("User connected.")
                while not self.stop_event.is_set():
                    try:
                        msg = await asyncio.wait_for(ws.recv(), timeout=1.0)
                        data = json.loads(msg)
                        logger.info(f" [USER RECEIVED]: {data.get('role')}: {data.get('text') or data.get('message')}")
                        self.user_received.append(data)
                        
                        # If user receives an admin message, they can respond
                        if data.get('role') == 'admin' and "Hello from Admin" in data.get('text', ''):
                            await ws.send(json.dumps({
                                "role": "user",
                                "user_name": "EndUser",
                                "text": "Hello Admin! I see you now.",
                                "type": "chat_message"
                            }))
                            logger.info("User sent reply via WebSocket.")
                    except asyncio.TimeoutError:
                        continue
        except Exception as e:
            logger.error(f"User WS error: {e}")

async def run_test():
    thread_id = f"handover-{uuid.uuid4().hex[:8]}"
    tester = HandoverTester(thread_id)
    
    # 1. Start listeners
    admin_task = asyncio.create_task(tester.listen_admin())
    user_task = asyncio.create_task(tester.listen_user())
    
    await asyncio.sleep(3) # Wait for connections
    
    async with httpx.AsyncClient(timeout=30.0) as client:
        # 2. Activate Handover
        logger.info(f"Activating human takeover for thread {thread_id}...")
        takeover_res = await client.post(
            f"{BASE_URL}/chat/takeover",
            data={
                "tenant_id": TENANT_ID,
                "thread_id": thread_id,
                "status": "true",
                "user_id": "admin-123"
            }
        )
        logger.info(f"Takeover Response: {takeover_res.json()}")

        await asyncio.sleep(2)
        
        # 3. Admin sends a message via WebSocket (simulated by a separate small connection if needed, but we can do it in listen_admin if we had a handle)
        # For simplicity, we'll open a temporary sender connection for Admin
        logger.info("Admin sending message via WebSocket...")
        async with websockets.connect(f"{WS_URL}/ws/{TENANT_ID}/{thread_id}?role=admin") as ws:
            await ws.send(json.dumps({
                "role": "admin",
                "user_name": "AdminUser",
                "text": "Hello from Admin! How can I help you today?",
                "type": "chat_message"
            }))
        
        # 4. Wait for user to receive and reply (user listener handles reply)
        await asyncio.sleep(5)

    # 5. Cleanup
    tester.stop_event.set()
    await asyncio.gather(admin_task, user_task, return_exceptions=True)
    
    print("\n" + "="*50)
    print("HANDOVER TEST RESULTS")
    print("="*50)
    
    admin_saw_user = any(d.get('role') == 'user' and "Hello Admin" in d.get('text', '') for d in tester.admin_received)
    user_saw_admin = any(d.get('role') == 'admin' and "Hello from Admin" in d.get('text', '') for d in tester.user_received)
    
    if admin_saw_user and user_saw_admin:
        print("RESULT: SUCCESS - Real-time handover bi-directional communication verified!")
    else:
        print(f"RESULT: FAILURE - Admin saw user: {admin_saw_user}, User saw admin: {user_saw_admin}")
    
    print("="*50)

if __name__ == "__main__":
    asyncio.run(run_test())
