#!/usr/bin/env python3
"""
Full Pipeline: Layer 3.5 Research → Layer 4 Video Generation
For the top MAKE_NOW premise from each of the 7 channels.
"""

import json
import sys
import time
import requests
from pathlib import Path
from datetime import datetime, timezone

sys.path.insert(0, str(Path(__file__).parent))

from video_generation import run_layer4_v4_pipeline

OUTPUT_DIR = Path(__file__).parent / "output"
RESEARCH_SERVICE_URL = "http://127.0.0.1:8100"
RESEARCH_POLL_INTERVAL = 15  # seconds
RESEARCH_MAX_POLL_TIME = 600  # 10 minutes

CHANNELS = [
    "why_you_do_that",
    "how_it_actually_works",
    "sixty_second_rabbit_hole",
    "designed_to_trick_you",
    "the_money_thing",
    "what_happens_next",
    "one_minute_history"
]

# Work log
WORKLOG = []

def log(msg: str):
    """Log to console and worklog."""
    timestamp = datetime.now().strftime("%H:%M:%S")
    entry = f"[{timestamp}] {msg}"
    print(entry)
    WORKLOG.append(entry)


def save_worklog():
    """Save worklog to file."""
    worklog_path = OUTPUT_DIR / "pipeline_worklog.md"
    with open(worklog_path, "w") as f:
        f.write("# Full Pipeline Work Log\n\n")
        f.write(f"**Generated:** {datetime.now().isoformat()}\n\n")
        f.write("```\n")
        for entry in WORKLOG:
            f.write(entry + "\n")
        f.write("```\n")
    print(f"\nWorklog saved: {worklog_path}")


def check_research_service() -> bool:
    """Check if research service is running."""
    try:
        response = requests.get(f"{RESEARCH_SERVICE_URL}/health", timeout=5)
        return response.status_code == 200
    except Exception:
        return False


def submit_research_job(topic: str, thesis: str = None) -> dict:
    """Submit async research job. Returns job info or error."""
    try:
        payload = {"topic": topic}
        if thesis:
            payload["thesis"] = thesis

        response = requests.post(
            f"{RESEARCH_SERVICE_URL}/research/async",
            json=payload,
            timeout=30
        )

        if response.status_code == 200:
            return {"success": True, **response.json()}
        else:
            return {"success": False, "error": f"HTTP {response.status_code}: {response.text[:100]}"}
    except Exception as e:
        return {"success": False, "error": str(e)}


def poll_research_job(job_id: str) -> dict:
    """Poll for research job completion."""
    start_time = time.time()

    while True:
        elapsed = time.time() - start_time
        if elapsed > RESEARCH_MAX_POLL_TIME:
            return {"success": False, "error": f"Timeout after {RESEARCH_MAX_POLL_TIME}s"}

        try:
            response = requests.get(f"{RESEARCH_SERVICE_URL}/research/{job_id}", timeout=30)
            if response.status_code != 200:
                return {"success": False, "error": f"Poll HTTP {response.status_code}"}

            data = response.json()
            status = data.get("status")

            if status == "completed":
                result = data.get("result", {})
                return {"success": True, **result}
            elif status == "failed":
                return {"success": False, "error": data.get("error", "Unknown error")}

            # Still running
            log(f"    ... still running ({int(elapsed)}s elapsed)")
            time.sleep(RESEARCH_POLL_INTERVAL)

        except Exception as e:
            return {"success": False, "error": str(e)}


def construct_research_query(opportunity: dict) -> str:
    """Build search query from opportunity."""
    premise = str(opportunity.get("premise", opportunity.get("suggested_title", "")))[:80]
    reveal = str(opportunity.get("core_reveal", ""))[:120]

    query = premise
    if reveal and reveal not in ("", "N/A", "None"):
        query += f" {reveal}"

    return query[:380]


def find_top_make_now(opportunities_data: dict) -> tuple:
    """Find the MAKE_NOW with highest weighted_score."""
    opportunities = opportunities_data.get("opportunities", [])
    make_now_opps = [
        (i, opp) for i, opp in enumerate(opportunities)
        if opp.get("verdict") == "MAKE_NOW"
    ]

    if not make_now_opps:
        return None, None

    make_now_opps.sort(key=lambda x: x[1].get("weighted_score", 0), reverse=True)
    return make_now_opps[0]


def run_research_for_channel(channel_id: str) -> dict:
    """Run Layer 3.5 research for top MAKE_NOW in channel."""

    log(f"\n{'='*60}")
    log(f"CHANNEL: {channel_id}")
    log(f"{'='*60}")

    # Load opportunities
    opp_file = OUTPUT_DIR / f"analysis_opportunities_{channel_id}.json"
    if not opp_file.exists():
        log(f"  ERROR: File not found: {opp_file}")
        return {"channel": channel_id, "error": "File not found"}

    with open(opp_file, "r") as f:
        data = json.load(f)

    # Find top MAKE_NOW
    idx, opportunity = find_top_make_now(data)
    if opportunity is None:
        log(f"  ERROR: No MAKE_NOW opportunities")
        return {"channel": channel_id, "error": "No MAKE_NOW"}

    premise = opportunity.get("premise", opportunity.get("suggested_title", "Unknown"))
    log(f"  Selected: {premise[:60]}...")
    log(f"  Index: {idx}, Score: {opportunity.get('weighted_score', 0)}")

    # Check if already researched
    if opportunity.get("research_completed"):
        log(f"  Already researched, skipping...")
        return {
            "channel": channel_id,
            "premise": premise,
            "research_status": "already_complete",
            "opportunity_index": idx
        }

    # Build query
    query = construct_research_query(opportunity)
    thesis = premise
    log(f"  Query: {query[:80]}...")

    # Submit research job
    log(f"  [3.5] Submitting research job...")
    submit_result = submit_research_job(query, thesis)

    if not submit_result.get("success"):
        log(f"  ERROR: {submit_result.get('error')}")
        opportunity["research_completed"] = False
        opportunity["research_error"] = submit_result.get("error")

        # Save updated file
        with open(opp_file, "w") as f:
            json.dump(data, f, indent=2)

        return {
            "channel": channel_id,
            "premise": premise,
            "research_status": "submit_failed",
            "error": submit_result.get("error")
        }

    job_id = submit_result.get("job_id")
    log(f"  Job ID: {job_id}")
    log(f"  Estimated: ~{submit_result.get('estimated_duration_seconds', 290)}s")

    # Poll for completion
    log(f"  [3.5] Polling for completion...")
    poll_result = poll_research_job(job_id)

    if not poll_result.get("success"):
        log(f"  ERROR: {poll_result.get('error')}")
        opportunity["research_completed"] = False
        opportunity["research_error"] = poll_result.get("error")

        with open(opp_file, "w") as f:
            json.dump(data, f, indent=2)

        return {
            "channel": channel_id,
            "premise": premise,
            "research_status": "poll_failed",
            "error": poll_result.get("error")
        }

    # Research succeeded - update opportunity
    log(f"  ✓ Research complete!")
    log(f"    Words: {poll_result.get('word_count', 0)}")
    log(f"    Sources: {poll_result.get('source_count', 0)}")
    log(f"    Duration: {poll_result.get('duration_seconds', 0)}s")

    opportunity["research_completed"] = True
    opportunity["research_report_path"] = poll_result.get("report_path", "")
    opportunity["research_report_content"] = poll_result.get("report_content", "")
    opportunity["research_word_count"] = poll_result.get("word_count", 0)
    opportunity["research_source_count"] = poll_result.get("source_count", 0)
    opportunity["research_duration_seconds"] = poll_result.get("duration_seconds", 0)
    opportunity["research_generated_at"] = poll_result.get("generated_at", datetime.now(timezone.utc).isoformat())
    opportunity["research_query"] = query

    # Save updated file
    with open(opp_file, "w") as f:
        json.dump(data, f, indent=2)
    log(f"  Saved: {opp_file}")

    return {
        "channel": channel_id,
        "premise": premise,
        "research_status": "completed",
        "word_count": poll_result.get("word_count", 0),
        "source_count": poll_result.get("source_count", 0),
        "duration_seconds": poll_result.get("duration_seconds", 0),
        "opportunity_index": idx
    }


def run_layer4_for_channel(channel_id: str) -> dict:
    """Run Layer 4 pipeline for top MAKE_NOW in channel."""

    log(f"\n  [4] Running Layer 4 pipeline...")

    # Load opportunities (with updated research)
    opp_file = OUTPUT_DIR / f"analysis_opportunities_{channel_id}.json"
    with open(opp_file, "r") as f:
        data = json.load(f)

    idx, opportunity = find_top_make_now(data)
    if opportunity is None:
        return {"error": "No MAKE_NOW found"}

    premise = opportunity.get("premise", opportunity.get("suggested_title", "Unknown"))

    # Check research status
    has_research = opportunity.get("research_completed", False)
    log(f"  Research available: {has_research}")
    if has_research:
        log(f"    Content: {len(opportunity.get('research_report_content', ''))} chars")

    # Run v4 pipeline
    try:
        result = run_layer4_v4_pipeline(
            opportunity=opportunity,
            channel_id=channel_id,
            target_duration=90,
            max_iterations=2
        )
    except Exception as e:
        log(f"  ERROR: {e}")
        return {"error": str(e)}

    # Extract results
    evaluation = result.get("evaluation", {})
    composite_score = evaluation.get("composite_score", 0)
    verdict = evaluation.get("verdict", "UNKNOWN")
    iteration_count = result.get("iteration_count", 1)
    production_ready = result.get("production_ready", False)
    critical_failures = evaluation.get("critical_failures", [])

    log(f"  ✓ Layer 4 complete!")
    log(f"    Iterations: {iteration_count}")
    log(f"    Score: {composite_score:.2f}")
    log(f"    Verdict: {verdict}")
    if critical_failures:
        log(f"    Critical failures: {len(critical_failures)}")

    # Save output
    output_file = OUTPUT_DIR / f"layer4_{channel_id}.json"
    output_data = {
        "channel": channel_id,
        "premise": premise,
        "opportunity_index": idx,
        "weighted_score": opportunity.get("weighted_score", 0),
        "has_research": has_research,
        "research_word_count": opportunity.get("research_word_count", 0),
        "script": result.get("script"),
        "evaluation": evaluation,
        "director_metadata": result.get("director_metadata"),
        "evaluator_metadata": result.get("evaluator_metadata"),
        "iteration_count": iteration_count,
        "all_iterations": result.get("all_iterations", []),
        "production_ready": production_ready,
        "generated_at": datetime.now(timezone.utc).isoformat()
    }

    with open(output_file, "w") as f:
        json.dump(output_data, f, indent=2)
    log(f"  Saved: {output_file}")

    return {
        "composite_score": composite_score,
        "verdict": verdict,
        "iteration_count": iteration_count,
        "production_ready": production_ready,
        "critical_failures": [cf.get("dimension") for cf in critical_failures]
    }


def main():
    log("="*70)
    log("FULL PIPELINE: Layer 3.5 Research → Layer 4 Video Generation")
    log("="*70)
    log(f"Started: {datetime.now().strftime('%Y-%m-%d %H:%M:%S')}")
    log(f"Channels: {len(CHANNELS)}")

    # Check research service
    if not check_research_service():
        log("ERROR: Research service not available at " + RESEARCH_SERVICE_URL)
        log("Run: cd /home/sietch6/research-sandbox/option1-gpt-researcher && python research_service.py")
        save_worklog()
        return

    log(f"Research service: {RESEARCH_SERVICE_URL} ✓")

    results = []

    for channel_id in CHANNELS:
        # Layer 3.5: Research
        research_result = run_research_for_channel(channel_id)

        if "error" in research_result and research_result.get("research_status") != "already_complete":
            results.append({
                "channel": channel_id,
                "premise": research_result.get("premise", "?"),
                "research_status": research_result.get("research_status", "error"),
                "error": research_result.get("error")
            })
            continue

        # Layer 4: Video Generation
        layer4_result = run_layer4_for_channel(channel_id)

        results.append({
            "channel": channel_id,
            "premise": research_result.get("premise", "?"),
            "research_status": research_result.get("research_status", "?"),
            "research_words": research_result.get("word_count", 0),
            "research_sources": research_result.get("source_count", 0),
            "l4_score": layer4_result.get("composite_score", 0),
            "l4_verdict": layer4_result.get("verdict", "?"),
            "l4_iterations": layer4_result.get("iteration_count", 0),
            "l4_failures": layer4_result.get("critical_failures", [])
        })

    # Summary
    log("\n")
    log("="*100)
    log("FINAL SUMMARY")
    log("="*100)
    log(f"{'Channel':<25} {'Premise':<28} {'Research':<12} {'L4 Score':>8} {'Verdict':<18}")
    log("-"*100)

    for r in results:
        if "error" in r:
            log(f"{r['channel']:<25} ERROR: {r.get('error', '?')[:60]}")
            continue

        premise_short = r['premise'][:26] + ".." if len(r['premise']) > 28 else r['premise']
        research_status = f"{r.get('research_words', 0)}w" if r.get('research_status') == 'completed' else r.get('research_status', '?')[:10]

        log(f"{r['channel']:<25} {premise_short:<28} {research_status:<12} {r.get('l4_score', 0):>7.2f} {r.get('l4_verdict', '?'):<18}")

    log("-"*100)

    # Stats
    successful = [r for r in results if "error" not in r]
    production_ready = [r for r in successful if r.get("l4_verdict") == "PRODUCTION_READY"]
    researched = [r for r in successful if r.get("research_status") == "completed"]

    log(f"\nTotal: {len(results)} channels")
    log(f"Researched: {len(researched)}")
    log(f"Layer 4 Complete: {len(successful)}")
    log(f"Production Ready: {len(production_ready)}")
    if successful:
        log(f"Average L4 Score: {sum(r.get('l4_score', 0) for r in successful) / len(successful):.2f}")

    # Save summary
    summary_file = OUTPUT_DIR / "pipeline_summary.json"
    with open(summary_file, "w") as f:
        json.dump({
            "generated_at": datetime.now(timezone.utc).isoformat(),
            "channels_processed": len(results),
            "researched": len(researched),
            "production_ready": len(production_ready),
            "results": results
        }, f, indent=2)
    log(f"\nSummary saved: {summary_file}")

    save_worklog()


if __name__ == "__main__":
    main()
