#!/usr/bin/env python3
"""
Noctua Labs - Global Swarm Cloud Hub (Enterprise SaaS Aggregator)
Secure WebSockets/TCP router for customer Axiom Daemons.
Validates MASTER_LICENSE_KEY and routes events to authorized Customer Dashboards 
or the God-Mode Dev Control Center.
"""

import asyncio
import json
import logging
import uuid
import os

logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] [NoctuaCloudHub] %(message)s")

MASTER_ADMIN_KEY = os.environ.get("NOCTUA_MASTER_ADMIN_KEY", "NOCTUA_GOD_MODE_777")

class NoctuaCloudHub:
    def __init__(self, host='0.0.0.0', port=7443):
        self.host = host
        self.port = port
        
        # map: license_key -> set of active writer streams (dashboards)
        self.dashboard_subscribers = {} 
        
        # set of god-mode writer streams
        self.god_mode_subscribers = set()
        
        # map: server_id -> license_key
        self.active_edge_nodes = {}

    async def handle_client(self, reader, writer):
        addr = writer.get_extra_info('peername')
        logging.info(f"New connection from {addr}")
        
        # First message must be authentication
        try:
            auth_line = await reader.readline()
            if not auth_line:
                writer.close()
                return
                
            auth_data = json.loads(auth_line.decode('utf-8'))
            node_type = auth_data.get('type') # 'DAEMON' or 'DASHBOARD'
            license_key = auth_data.get('license_key')
            
            if not license_key:
                writer.write(json.dumps({"error": "Missing license key"}).encode() + b'\n')
                await writer.drain()
                writer.close()
                return

            if node_type == 'DASHBOARD':
                await self.handle_dashboard(reader, writer, license_key)
            elif node_type == 'DAEMON':
                await self.handle_daemon(reader, writer, license_key, addr)
            else:
                writer.close()
                
        except Exception as e:
            logging.error(f"Error handling client {addr}: {e}")
            writer.close()

    async def handle_dashboard(self, reader, writer, license_key):
        if license_key == MASTER_ADMIN_KEY:
            logging.info("God-Mode Dashboard Connected.")
            self.god_mode_subscribers.add(writer)
            try:
                while True:
                    data = await reader.readline()
                    if not data: break
            finally:
                self.god_mode_subscribers.discard(writer)
                writer.close()
        else:
            logging.info(f"Customer Dashboard Connected for key: {license_key[:8]}...")
            if license_key not in self.dashboard_subscribers:
                self.dashboard_subscribers[license_key] = set()
            self.dashboard_subscribers[license_key].add(writer)
            try:
                while True:
                    data = await reader.readline()
                    if not data: break
            finally:
                self.dashboard_subscribers[license_key].discard(writer)
                writer.close()

    async def handle_daemon(self, reader, writer, license_key, addr):
        node_id = str(uuid.uuid4())
        self.active_edge_nodes[node_id] = license_key
        logging.info(f"Edge Node {node_id} online for key: {license_key[:8]}...")
        
        try:
            while True:
                line = await reader.readline()
                if not line: break
                
                # Route to God Mode
                for sub in list(self.god_mode_subscribers):
                    try:
                        sub.write(line)
                        await sub.drain()
                    except Exception:
                        self.god_mode_subscribers.discard(sub)
                        
                # Route to specific Customer Dashboard
                if license_key in self.dashboard_subscribers:
                    for sub in list(self.dashboard_subscribers[license_key]):
                        try:
                            sub.write(line)
                            await sub.drain()
                        except Exception:
                            self.dashboard_subscribers[license_key].discard(sub)
        finally:
            del self.active_edge_nodes[node_id]
            writer.close()

    async def start(self):
        server = await asyncio.start_server(self.handle_client, self.host, self.port)
        logging.info(f"Noctua Cloud Hub listening on {self.host}:{self.port}")
        async with server:
            await server.serve_forever()

if __name__ == "__main__":
    hub = NoctuaCloudHub()
    try:
        asyncio.run(hub.start())
    except KeyboardInterrupt:
        pass
