"""
Internal job routes — Caller Type 3 (scheduled cron jobs).

POST /internal/jobs/collect      trigger data collection from source
POST /internal/jobs/cleanup      purge expired tokens + stale sessions
POST /internal/jobs/push-digest  send batched push updates to subscribers
GET  /internal/jobs/health       liveness check (no auth — for LB health probes)

Security rules:
  • All POST routes require scope 'internal:job' (via gateway_auth → verify_cron_key)
  • Blocked at nginx/ALB for public internet — only reachable inside VPC
  • Max 1 concurrent execution per job type (Redis SETNX distributed lock)
  • Every execution emits a structured log: event/job/status/duration_ms/records

Rate limit policy:
  No per-minute rate limit — cron jobs run at low frequency by design.
  Concurrency enforced instead via Redis lock (1 parallel execution max).
"""

from __future__ import annotations

import time
from datetime import datetime, timedelta, timezone

import structlog
from fastapi import APIRouter, Depends, HTTPException, status
from sqlalchemy import delete, select
from sqlalchemy.ext.asyncio import AsyncSession

from app.core.redis_client import acquire_job_lock
from app.db.models import RefreshToken
from app.db.session import get_db
from app.dependencies.auth import CallerContext, require_scope
from app.schemas.auth import CronCallerContext

log = structlog.get_logger(__name__)
router = APIRouter(prefix="/internal/jobs", tags=["internal"])

_cron_only = require_scope("internal:job")

# Redis lock TTLs (expected_max_duration + 60s safety buffer)
_LOCK_TTL = {
    "collect": 300,      # 5 min max + 60s = 360s lock
    "cleanup": 120,      # 2 min max + 60s = 180s lock
    "push-digest": 180,  # 3 min max + 60s = 240s lock
}


def _job_log(job_name: str, status_str: str, start_ns: int, records: int = 0) -> None:
    """
    Emit the mandatory structured log per execution.

    Required format (from spec):
    { event: job_run, job: name, status: started|completed|failed,
      duration_ms: N, records_processed: N, timestamp: ISO8601 }
    """
    duration_ms = (time.monotonic_ns() - start_ns) // 1_000_000
    log.info(
        "job_run",
        job=job_name,
        status=status_str,
        duration_ms=duration_ms,
        records_processed=records,
        timestamp=datetime.now(timezone.utc).isoformat(),
    )


# ===========================================================================
# GET /internal/jobs/health
# ===========================================================================

@router.get("/health", status_code=200, include_in_schema=False)
async def health_check() -> dict:
    """
    Liveness probe — no auth required.
    Used by ALB / Kubernetes readiness probes.
    nginx/ALB should allow this endpoint through even with /internal/* blocked.
    """
    return {"status": "ok", "timestamp": datetime.now(timezone.utc).isoformat()}


# ===========================================================================
# POST /internal/jobs/collect
# ===========================================================================

@router.post("/collect", status_code=202)
async def job_collect(
    caller: CallerContext = Depends(_cron_only),
    db: AsyncSession = Depends(get_db),
) -> dict:
    """
    Trigger data collection from configured external sources.

    Redis lock prevents concurrent runs of the same collection job.
    Lock TTL = 5 min max runtime + 60 s safety buffer.
    """
    start_ns = time.monotonic_ns()
    _job_log("collect", "started", start_ns)

    try:
        async with acquire_job_lock("collect", ttl_seconds=_LOCK_TTL["collect"]):
            # --- placeholder: replace with actual collection logic ---
            records_processed = await _run_data_collection(db)
            # ---------------------------------------------------------

    except RuntimeError as exc:
        # Lock already held — another instance is running
        raise HTTPException(
            status_code=status.HTTP_409_CONFLICT,
            detail=str(exc),
        ) from exc
    except Exception as exc:
        _job_log("collect", "failed", start_ns)
        log.exception("job_collect_error")
        raise HTTPException(status_code=500, detail="Collection job failed") from exc

    _job_log("collect", "completed", start_ns, records=records_processed)
    return {"status": "completed", "records_processed": records_processed}


# ===========================================================================
# POST /internal/jobs/cleanup
# ===========================================================================

@router.post("/cleanup", status_code=202)
async def job_cleanup(
    caller: CallerContext = Depends(_cron_only),
    db: AsyncSession = Depends(get_db),
) -> dict:
    """
    Purge expired refresh tokens and stale sessions.

    SQL: DELETE FROM refresh_tokens WHERE expires_at < NOW() - INTERVAL '1 day'
    The 1-day grace period keeps recently-expired tokens available for
    theft-detection forensics before they are purged.
    """
    start_ns = time.monotonic_ns()
    _job_log("cleanup", "started", start_ns)

    try:
        async with acquire_job_lock("cleanup", ttl_seconds=_LOCK_TTL["cleanup"]):
            cutoff = datetime.now(timezone.utc) - timedelta(days=1)

            result = await db.execute(
                delete(RefreshToken).where(RefreshToken.expires_at < cutoff)
            )
            deleted = result.rowcount

    except RuntimeError as exc:
        raise HTTPException(status_code=409, detail=str(exc)) from exc
    except Exception as exc:
        _job_log("cleanup", "failed", start_ns)
        log.exception("job_cleanup_error")
        raise HTTPException(status_code=500, detail="Cleanup job failed") from exc

    _job_log("cleanup", "completed", start_ns, records=deleted)
    return {"status": "completed", "records_deleted": deleted}


# ===========================================================================
# POST /internal/jobs/push-digest
# ===========================================================================

@router.post("/push-digest", status_code=202)
async def job_push_digest(
    caller: CallerContext = Depends(_cron_only),
    db: AsyncSession = Depends(get_db),
) -> dict:
    """
    Send batched push updates to all active WebSocket subscribers.

    Lock prevents double-delivery if two cron triggers fire close together.
    """
    start_ns = time.monotonic_ns()
    _job_log("push-digest", "started", start_ns)

    try:
        async with acquire_job_lock("push-digest", ttl_seconds=_LOCK_TTL["push-digest"]):
            # --- placeholder: replace with actual push-digest logic ---
            records_sent = await _run_push_digest(db)
            # -----------------------------------------------------------

    except RuntimeError as exc:
        raise HTTPException(status_code=409, detail=str(exc)) from exc
    except Exception as exc:
        _job_log("push-digest", "failed", start_ns)
        log.exception("job_push_digest_error")
        raise HTTPException(status_code=500, detail="Push-digest job failed") from exc

    _job_log("push-digest", "completed", start_ns, records=records_sent)
    return {"status": "completed", "records_sent": records_sent}


# ---------------------------------------------------------------------------
# Stub implementations — replace with real business logic
# ---------------------------------------------------------------------------

async def _run_data_collection(db: AsyncSession) -> int:
    """Fetch data from external sources. Returns number of records ingested."""
    return 0


async def _run_push_digest(db: AsyncSession) -> int:
    """Batch and push pending events to WebSocket subscribers. Returns count sent."""
    return 0
