from fastapi import APIRouter, Query, HTTPException, Depends
from typing import Optional
from src.auth.dependencies import get_current_user
import json
import asyncio
from .dependencies import OutreachServiceDep
from .models import (
    EmailBatchResponse,
    EmailBatchListResponse,
    RunIdListResponse,
    EmailPatchRequest,
    RegenerateRequest,
    RegenerateResponseAPI,
    DeleteResponse,
    BatchActionResponse,
    LeadListResponse,
    PromptRequest,
    PromptResponse,
    PromptListResponse,
    ResolvedPromptResponse,
)

router = APIRouter()


# ──────────────────────────────────────────────
# GET /emails/runs — List all run IDs
# ──────────────────────────────────────────────

@router.get("/emails/runs", response_model=RunIdListResponse)
def list_run_ids(service: OutreachServiceDep, user_id: str = Depends(get_current_user)):
    """List all distinct run IDs that have generated emails, with counts."""
    try:
        return service.get_all_run_ids(user_id)
    except Exception as e:
        raise HTTPException(status_code=500, detail=str(e))


# ──────────────────────────────────────────────
# GET /emails/runs/{run_id} — Emails by run ID
# ──────────────────────────────────────────────

@router.get("/emails/runs/{run_id}", response_model=EmailBatchListResponse)
def get_emails_by_run(
    run_id: str,
    service: OutreachServiceDep,
    user_id: str = Depends(get_current_user),
    page: int = Query(1, ge=1, description="Page number"),
    limit: int = Query(20, ge=1, le=100, description="Items per page"),
    approval_status: Optional[str] = Query(None, description="Filter by approval status"),
):
    """Fetch paginated email documents for a specific run_id."""
    try:
        result = service.get_emails_by_run_id(run_id, user_id, page, limit, approval_status)
        if result.total == 0:
            raise HTTPException(status_code=404, detail=f"No emails found for run_id: {run_id}")
        return result
    except HTTPException:
        raise
    except Exception as e:
        raise HTTPException(status_code=500, detail=str(e))


# ──────────────────────────────────────────────
# GET /leads/runs/{run_id} — Leads by run ID
# ──────────────────────────────────────────────

@router.get("/leads/runs/{run_id}", response_model=LeadListResponse)
def get_leads_by_run(
    run_id: str,
    service: OutreachServiceDep,
    user_id: str = Depends(get_current_user),
    page: int = Query(1, ge=1, description="Page number"),
    limit: int = Query(20, ge=1, le=100, description="Items per page"),
):
    """Fetch paginated enriched leads from PostgreSQL for a specific run_id."""
    try:
        result = service.get_leads_by_run_id(run_id, user_id, page, limit)
        if result.total == 0:
            raise HTTPException(status_code=404, detail=f"No leads found for run_id: {run_id}")
        return result
    except HTTPException:
        raise
    except Exception as e:
        raise HTTPException(status_code=500, detail=str(e))

@router.get("/leads/runs/{run_id}/stats")
def get_lead_stats(
    run_id: str,
    service: OutreachServiceDep,
    user_id: str = Depends(get_current_user),
):
    """Get lead statistics for a specific run."""
    return service.get_lead_stats_by_run(run_id, user_id)


# ──────────────────────────────────────────────
# GET /emails/{email_id} — Single email by ID
# ──────────────────────────────────────────────

@router.get("/emails/{email_id}", response_model=EmailBatchResponse)
async def get_email(email_id: str, service: OutreachServiceDep, user_id: str = Depends(get_current_user)):
    """Fetch a single email document by its MongoDB _id."""
    email = service.get_email_by_id(email_id, user_id)
    if not email:
        raise HTTPException(status_code=404, detail="Email not found")
    return email


# ──────────────────────────────────────────────
# PATCH /emails/{email_id} — Partial update
# ──────────────────────────────────────────────

@router.patch("/emails/{email_id}", response_model=EmailBatchResponse)
async def patch_email(
    email_id: str,
    patch: EmailPatchRequest,
    service: OutreachServiceDep,
    user_id: str = Depends(get_current_user),
):
    """Partially update an email document (approval_status, generated_emails, metadata)."""
    updated = service.patch_email(email_id, patch, user_id)
    if not updated:
        raise HTTPException(status_code=404, detail="Email not found")
    return updated


# ──────────────────────────────────────────────
# DELETE /emails/{email_id} — Delete single
# ──────────────────────────────────────────────

@router.delete("/emails/{email_id}", response_model=DeleteResponse)
async def delete_email(email_id: str, service: OutreachServiceDep, user_id: str = Depends(get_current_user)):
    """Delete a single email document by its MongoDB _id."""
    result = service.delete_email(email_id, user_id)
    if result.deleted_count == 0:
        raise HTTPException(status_code=404, detail=result.message)
    return result


# ──────────────────────────────────────────────
# DELETE /emails/runs/{run_id} — Bulk delete
# ──────────────────────────────────────────────

@router.delete("/emails/runs/{run_id}", response_model=DeleteResponse)
async def delete_emails_by_run(run_id: str, service: OutreachServiceDep, user_id: str = Depends(get_current_user)):
    """Delete all email documents for a given run_id."""
    result = service.delete_emails_by_run_id(run_id, user_id)
    if result.deleted_count == 0:
        raise HTTPException(status_code=404, detail=result.message)
    return result


@router.post("/emails/runs/{run_id}/regenerate", response_model=RegenerateResponseAPI)
async def regenerate_emails(
    run_id: str,
    service: OutreachServiceDep,
    user_id: str = Depends(get_current_user),
    body: Optional[RegenerateRequest] = None,
):
    """
    Re-generate all emails for a run_id.
    Applies filters, resets status to 'processing', and starts generation.
    """
    try:
        result = service.regenerate_emails_for_run(run_id, user_id)
        
        if result.total_leads == 0:
            raise HTTPException(
                status_code=404,
                detail=f"No eligible leads found for regeneration in run_id: {run_id}",
            )

        # Return the 'started' status and the count of leads being processed
        return {
            "status": "started",
            "run_id": run_id,
            "total_leads_queued": result.total_leads,
            "message": "Email regeneration has been initiated for approved leads."
        }

    except HTTPException:
        raise
    except Exception as e:
        raise HTTPException(status_code=500, detail=f"Failed to start regeneration: {str(e)}")


# ──────────────────────────────────────────────
# POST /emails/{email_id}/approve or decline
# ──────────────────────────────────────────────

@router.post("/emails/{email_id}/approve", response_model=EmailBatchResponse)
def approve_email(email_id: str, service: OutreachServiceDep, user_id: str = Depends(get_current_user)):
    """Approve a single email."""
    updated = service.update_email_status(email_id, "approved", user_id)
    if not updated:
        raise HTTPException(status_code=404, detail="Email not found")
    return updated


@router.post("/emails/{email_id}/decline", response_model=EmailBatchResponse)
def decline_email(email_id: str, service: OutreachServiceDep, user_id: str = Depends(get_current_user)):
    """Decline a single email."""
    updated = service.update_email_status(email_id, "declined", user_id)
    if not updated:
        raise HTTPException(status_code=404, detail="Email not found")
    return updated


# ──────────────────────────────────────────────
# POST /emails/runs/{run_id}/approve or decline
# ──────────────────────────────────────────────

@router.post("/emails/runs/{run_id}/approve", response_model=BatchActionResponse)
def approve_run(run_id: str, service: OutreachServiceDep, user_id: str = Depends(get_current_user)):
    """Approve all emails in a run."""
    return service.update_run_status(run_id, "approved", user_id)


@router.post("/emails/runs/{run_id}/decline", response_model=BatchActionResponse)
def decline_run(run_id: str, service: OutreachServiceDep, user_id: str = Depends(get_current_user)):
    """Decline all emails in a run."""
    return service.update_run_status(run_id, "declined", user_id)


@router.get("/sync-status/{run_id}")
def get_sync_status(run_id: str, service: OutreachServiceDep, user_id: str = Depends(get_current_user)):
    return service.get_sync_status(run_id, user_id)

@router.post("/emails/runs/{run_id}/sync-woodpecker")
async def sync_woodpecker(
    run_id: str,
    user_id: str = Depends(get_current_user)
):
    """
    Sync all approved emails for a specific run to the linked Woodpecker campaign.
    Dispatches request to the standalone Woodpecker TCP worker.
    """
    payload = json.dumps({"run_id": run_id, "user_id": user_id}) + "\n"
    try:
        reader, writer = await asyncio.open_connection('woodpecker-worker', 8888)
        writer.write(payload.encode('utf-8'))
        await writer.drain()

        data = await reader.readline()
        writer.close()
        await writer.wait_closed()

        if not data:
            raise HTTPException(status_code=500, detail="Worker returned empty response")

        result = json.loads(data.decode('utf-8'))
        if not result.get("success"):
            raise HTTPException(status_code=400, detail=result.get("error") or result.get("message"))
        return result
    except ConnectionRefusedError:
        raise HTTPException(status_code=503, detail="Woodpecker sync worker is unreachable")
    except Exception as e:
        raise HTTPException(status_code=500, detail=f"Failed to communicate with worker: {str(e)}")


# ──────────────────────────────────────────────
# PROMPTS
# ──────────────────────────────────────────────

@router.get("/prompts", response_model=PromptListResponse)
def list_prompts(service: OutreachServiceDep, user_id: str = Depends(get_current_user)):
    """List all custom prompts for the user."""
    return service.list_prompts(user_id)


@router.post("/prompts", response_model=PromptResponse)
def upsert_prompt(
    prompt: PromptRequest,
    service: OutreachServiceDep,
    user_id: str = Depends(get_current_user)
):
    """Create or update a custom prompt (run-specific or universal)."""
    return service.upsert_prompt(user_id, prompt)


@router.delete("/prompts/{prompt_id}", response_model=DeleteResponse)
def delete_prompt(
    prompt_id: str,
    service: OutreachServiceDep,
    user_id: str = Depends(get_current_user)
):
    """Delete a custom prompt."""
    return service.delete_prompt(prompt_id, user_id)


@router.get("/resolved-prompt/{run_id}", response_model=ResolvedPromptResponse)
def get_resolved_prompt(
    run_id: str,
    service: OutreachServiceDep,
    user_id: str = Depends(get_current_user)
):
    """
    Expose the final prompt logic (Run-specific > Universal > Fallback).
    Useful for reviewing the prompt before triggering generation.
    """
    return service.get_resolved_prompt(user_id, run_id)