import asyncio
import json
import logging
from src.outreach.service import OutreachService

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

async def handle_client(reader: asyncio.StreamReader, writer: asyncio.StreamWriter):
    addr = writer.get_extra_info('peername')
    logger.info(f"Accepted connection from {addr}")

    try:
        data = await reader.readline()
        if not data:
            return

        message = data.decode('utf-8').strip()
        logger.info(f"Received message: {message}")

        try:
            payload = json.loads(message)
            run_id = payload.get("run_id")
            user_id = payload.get("user_id")

            if not run_id or not user_id:
                response = {"success": False, "error": "Missing run_id or user_id"}
            else:
                service = OutreachService()
                logger.info(f"Starting Woodpecker sync for run_id={run_id}, user_id={user_id}")
                response = await service.sync_approved_emails_to_woodpecker(run_id, user_id)
                logger.info(f"Sync complete for run_id={run_id}: {response}")

        except json.JSONDecodeError:
            logger.error("Invalid JSON received")
            response = {"success": False, "error": "Invalid JSON"}
        except Exception as e:
            logger.exception("Error processing sync task")
            response = {"success": False, "error": str(e)}

        writer.write((json.dumps(response) + "\n").encode('utf-8'))
        await writer.drain()

    except Exception as e:
        logger.error(f"Connection error with {addr}: {e}")
    finally:
        logger.info(f"Closing connection from {addr}")
        writer.close()
        try:
            await writer.wait_closed()
        except Exception:
            pass

async def main():
    server = await asyncio.start_server(handle_client, '0.0.0.0', 8888)
    addr = server.sockets[0].getsockname()
    logger.info(f'Woodpecker sync TCP worker serving on {addr}')

    async with server:
        await server.serve_forever()

if __name__ == '__main__':
    try:
        asyncio.run(main())
    except KeyboardInterrupt:
        logger.info("TCP worker shutdown requested")
