import os, re, sqlite3, hashlib, time, asyncio, json, shutil
from pathlib import Path
from urllib.parse import urlparse, quote
import aiohttp
from fastapi import FastAPI, Request, Form
from fastapi.responses import HTMLResponse, RedirectResponse, StreamingResponse, PlainTextResponse, JSONResponse
from fastapi.templating import Jinja2Templates
from passlib.context import CryptContext

APP_DIR = Path(os.getenv('VOD_APP_DIR', '/opt/vodcreator'))
DB = APP_DIR / 'data' / 'vodcreator.db'
PORT = int(os.getenv('VOD_PORT', '8090'))
pwd = CryptContext(schemes=['bcrypt'], deprecated='auto')
app = FastAPI(title='VOD Creator Smart Cache V2')
templates = Jinja2Templates(directory=str(APP_DIR / 'app' / 'templates'))


def db():
    DB.parent.mkdir(parents=True, exist_ok=True)
    c = sqlite3.connect(DB)
    c.row_factory = sqlite3.Row
    return c


def init_db():
    c = db()
    c.executescript('''
    CREATE TABLE IF NOT EXISTS settings(k TEXT PRIMARY KEY, v TEXT NOT NULL);
    CREATE TABLE IF NOT EXISTS users(id INTEGER PRIMARY KEY, username TEXT UNIQUE, password_hash TEXT NOT NULL);
    CREATE TABLE IF NOT EXISTS items(
      id INTEGER PRIMARY KEY, name TEXT NOT NULL, group_title TEXT, tvg_id TEXT, tvg_logo TEXT,
      source_url TEXT NOT NULL UNIQUE, kind TEXT NOT NULL, is_live INTEGER DEFAULT 0,
      local_path TEXT, status TEXT NOT NULL DEFAULT 'remote', size INTEGER DEFAULT 0,
      last_access INTEGER DEFAULT 0, created_at INTEGER NOT NULL, updated_at INTEGER NOT NULL
    );
    CREATE TABLE IF NOT EXISTS channel_sources(id INTEGER PRIMARY KEY, name TEXT UNIQUE, m3u_url TEXT NOT NULL, enabled INTEGER DEFAULT 1);
    CREATE TABLE IF NOT EXISTS channel_categories(id INTEGER PRIMARY KEY, name TEXT UNIQUE, enabled INTEGER DEFAULT 1);
    CREATE INDEX IF NOT EXISTS idx_items_access ON items(last_access);
    ''')
    if not c.execute("SELECT 1 FROM users WHERE username='admin'").fetchone():
        c.execute("INSERT INTO users(username,password_hash) VALUES(?,?)", ('admin', pwd.hash('admin')))
    defaults = {
      'm3u_url':'', 'base_url':'', 'storage_root':'', 'expire_days':'7', 'max_gb':'',
      'xui_url':'', 'xui_api':'', 'xui_user':'', 'xui_password':'', 'xui_token':'',
      'xui_category_mode':'group', 'channel_primary':'', 'channel_backup':'',
      'auto_sync':'0'
    }
    for k,v in defaults.items(): c.execute('INSERT OR IGNORE INTO settings(k,v) VALUES(?,?)',(k,v))
    c.commit(); c.close()


def settings():
    c=db(); rows=c.execute('SELECT k,v FROM settings').fetchall(); c.close(); return {r['k']:r['v'] for r in rows}

def set_setting(k,v):
    c=db(); c.execute('INSERT INTO settings(k,v) VALUES(?,?) ON CONFLICT(k) DO UPDATE SET v=excluded.v',(k,v)); c.commit(); c.close()


def storage_root():
    s=settings(); root=(s.get('storage_root') or '').strip()
    if not root: root='/opt/vodcreator/cache'
    p=Path(root).expanduser().resolve(); p.mkdir(parents=True, exist_ok=True); return p

def film_dir(): p=storage_root()/'filmes'; p.mkdir(parents=True,exist_ok=True); return p
def series_dir(): p=storage_root()/'series'; p.mkdir(parents=True,exist_ok=True); return p

def parse_m3u(text):
    items=[]; current={}
    for raw in text.splitlines():
        line=raw.strip()
        if not line: continue
        if line.startswith('#EXTINF:'):
            attrs={m.group(1):m.group(2) for m in re.finditer(r'([\w-]+)="([^"]*)"',line)}
            title=line.split(',',1)[1].strip() if ',' in line else attrs.get('tvg-name','Sem título')
            current={'name':title,'group_title':attrs.get('group-title',''),'tvg_id':attrs.get('tvg-id',''),'tvg_logo':attrs.get('tvg-logo','')}
        elif not line.startswith('#') and current:
            name=current['name']; group=current['group_title']
            # Channel detection is intentionally conservative: groups commonly used for live TV or stream URLs without VOD markers.
            live = bool(re.search(r'(?i)\b(canais?|tv|ao vivo|live|news|esportes?|sports|radio|rádio|dazn|sportv|espn)\b', group+' '+name))
            series = bool(re.search(r'(?i)(s\d{1,2}e\d{1,3}|\bseason\s*\d+|\btemporada\s*\d+|\bepis[oó]dio\s*\d+)', name))
            kind='channel' if live else ('series' if series else 'movie')
            items.append({**current,'url':line,'kind':kind,'is_live':1 if live else 0}); current={}
    return items

def slug(s):
    s=re.sub(r'[^\w\s.-]+','',s,flags=re.UNICODE).strip(); return (s[:180] or 'video')

def is_hls(url): return '.m3u8' in urlparse(url).path.lower() or '.m3u8' in url.lower()

def ext(url):
    p=urlparse(url).path.lower(); m=re.search(r'\.([a-z0-9]{2,5})$',p); return '.'+m.group(1) if m else '.mp4'

def item_path(row):
    root=film_dir() if row['kind']=='movie' else series_dir()
    return root/f"{row['id']}_{slug(row['name'])}{'.m3u8' if is_hls(row['source_url']) else ext(row['source_url'])}"

def auth(request): return request.cookies.get('vod_session')=='ok'
def render(request,t,**ctx): return templates.TemplateResponse(t, {'request':request, **ctx})

@app.on_event('startup')
async def startup():
    init_db(); storage_root(); asyncio.create_task(cleaner_loop())

@app.get('/',response_class=HTMLResponse)
async def home(request:Request):
    if not auth(request): return RedirectResponse('/login',302)
    c=db();
    total=c.execute('SELECT COUNT(*) n FROM items WHERE kind<>"channel"').fetchone()['n']
    channels=c.execute('SELECT COUNT(*) n FROM items WHERE kind="channel"').fetchone()['n']
    movies=c.execute('SELECT COUNT(*) n FROM items WHERE kind="movie"').fetchone()['n']
    series=c.execute('SELECT COUNT(*) n FROM items WHERE kind="series"').fetchone()['n']
    cached=c.execute('SELECT COUNT(*) n FROM items WHERE status="complete"').fetchone()['n']
    active=c.execute('SELECT COUNT(*) n FROM items WHERE status="downloading"').fetchone()['n']
    size=c.execute('SELECT COALESCE(SUM(size),0) n FROM items WHERE status="complete"').fetchone()['n']; c.close()
    return render(request,'dashboard.html',total=total,channels=channels,movies=movies,series=series,cached=cached,active=active,size_gb=round(size/1024**3,2),s=settings(),storage=str(storage_root()))

@app.get('/login',response_class=HTMLResponse)
async def login_page(request:Request): return render(request,'login.html',error='')
@app.post('/login')
async def login(request:Request,username:str=Form(...),password:str=Form(...)):
    c=db(); u=c.execute('SELECT * FROM users WHERE username=?',(username,)).fetchone(); c.close()
    if u and pwd.verify(password,u['password_hash']):
        r=RedirectResponse('/',302); r.set_cookie('vod_session','ok',httponly=True,samesite='lax'); return r
    return render(request,'login.html',error='Usuário ou senha inválidos.')
@app.get('/logout')
async def logout(): r=RedirectResponse('/login',302); r.delete_cookie('vod_session'); return r

@app.get('/settings',response_class=HTMLResponse)
async def settings_page(request:Request):
    if not auth(request): return RedirectResponse('/login',302)
    return render(request,'settings.html',s=settings(),storage=str(storage_root()))

@app.post('/settings')
async def save_settings(request:Request, m3u_url:str=Form(''),base_url:str=Form(''),storage_root_:str=Form('',alias='storage_root'),
    xui_url:str=Form(''),xui_api:str=Form(''),xui_user:str=Form(''),xui_password:str=Form(''),xui_token:str=Form(''),
    expire_days:str=Form('7'),max_gb:str=Form(''),channel_primary:str=Form(''),channel_backup:str=Form('')):
    if not auth(request): return RedirectResponse('/login',302)
    vals={'m3u_url':m3u_url,'base_url':base_url,'storage_root':storage_root_,'xui_url':xui_url,'xui_api':xui_api,'xui_user':xui_user,'xui_password':xui_password,'xui_token':xui_token,'expire_days':expire_days,'max_gb':max_gb,'channel_primary':channel_primary,'channel_backup':channel_backup}
    for k,v in vals.items(): set_setting(k,v.strip())
    storage_root(); film_dir(); series_dir(); return RedirectResponse('/settings?saved=1',302)

@app.post('/password')
async def change_password(request:Request,password:str=Form(...)):
    if not auth(request): return RedirectResponse('/login',302)
    if len(password)<8: return RedirectResponse('/settings?error=senha_min_8',302)
    c=db(); c.execute("UPDATE users SET password_hash=? WHERE username='admin'",(pwd.hash(password),)); c.commit(); c.close(); return RedirectResponse('/settings?saved=password',302)

async def fetch_text(url):
    if not url: raise RuntimeError('URL não configurada.')
    timeout=aiohttp.ClientTimeout(total=180)
    async with aiohttp.ClientSession(timeout=timeout) as session:
        async with session.get(url) as r: r.raise_for_status(); return await r.text(errors='ignore')

async def sync_items(url, only_channels=False, clear_channels=False):
    text=await fetch_text(url); parsed=parse_m3u(text); now=int(time.time()); c=db(); count=0
    if clear_channels:
        c.execute("DELETE FROM items WHERE kind='channel'")
    for it in parsed:
        if only_channels and it['kind']!='channel': continue
        if not only_channels and it['kind']=='channel': continue
        c.execute('''INSERT INTO items(name,group_title,tvg_id,tvg_logo,source_url,kind,is_live,status,created_at,updated_at)
                     VALUES(?,?,?,?,?,?,?,?,?,?) ON CONFLICT(source_url) DO UPDATE SET name=excluded.name,group_title=excluded.group_title,tvg_id=excluded.tvg_id,tvg_logo=excluded.tvg_logo,kind=excluded.kind,is_live=excluded.is_live,updated_at=excluded.updated_at''',
                  (it['name'],it['group_title'],it['tvg_id'],it['tvg_logo'],it['url'],it['kind'],it['is_live'],'remote',now,now)); count+=1
    c.commit(); c.close(); return count

@app.post('/sync')
async def sync(request:Request):
    if not auth(request): return RedirectResponse('/login',302)
    try:
        n=await sync_items(settings().get('m3u_url',''),False,False)
        x=await xui_sync('vod')
        return RedirectResponse('/?sync='+str(n)+'&xui='+quote(x),302)
    except Exception as e: return RedirectResponse('/?error='+quote(str(e)),302)

@app.get('/channels',response_class=HTMLResponse)
async def channels_page(request:Request):
    if not auth(request): return RedirectResponse('/login',302)
    c=db(); groups=[r[0] for r in c.execute("SELECT DISTINCT group_title FROM items WHERE kind='channel' AND group_title<>'' ORDER BY group_title").fetchall()]; n=c.execute("SELECT COUNT(*) n FROM items WHERE kind='channel'").fetchone()['n']; c.close()
    return render(request,'channels.html',s=settings(),groups=groups,count=n)

@app.post('/channels/sync')
async def channels_sync(request:Request, primary:str=Form(''),backup:str=Form(''),clear: str=Form('')):
    if not auth(request): return RedirectResponse('/login',302)
    if primary: set_setting('channel_primary',primary.strip())
    if backup: set_setting('channel_backup',backup.strip())
    s=settings();
    try:
        n=await sync_items(s.get('channel_primary',''),True,clear=='1') if s.get('channel_primary') else 0
        # Backup is stored for the player/proxy; it is not duplicated into the XUI catalog.
        if s.get('channel_backup'):
            await sync_items(s.get('channel_backup',''),True,False)
        x=await xui_sync('channels')
        return RedirectResponse('/channels?ok='+quote(f'{n} canais processados; {x}'),302)
    except Exception as e: return RedirectResponse('/channels?error='+quote(str(e)),302)

@app.post('/channels/clear')
async def channels_clear(request:Request):
    if not auth(request): return RedirectResponse('/login',302)
    c=db(); c.execute("DELETE FROM items WHERE kind='channel'"); c.commit(); c.close(); return RedirectResponse('/channels?ok=Lista de canais limpa do catálogo local. Confirme a remoção no XUI pelo botão de sincronização quando a API estiver configurada.',302)

@app.get('/api/m3u')
async def mirror_m3u(request:Request,kind:str='vod'):
    c=db(); rows=c.execute("SELECT * FROM items WHERE kind<>'channel' ORDER BY id" if kind=='vod' else "SELECT * FROM items WHERE kind='channel' ORDER BY id").fetchall(); c.close()
    s=settings(); base=(s.get('base_url') or str(request.base_url).rstrip('/')).rstrip('/'); out=['#EXTM3U']
    for r in rows:
        out.append(f'#EXTINF:-1 tvg-id="{r["tvg_id"] or ""}" tvg-logo="{r["tvg_logo"] or ""}" group-title="{r["group_title"] or ""}",{r["name"]}')
        out.append(f'{base}/stream/{r["id"]}')
    return PlainTextResponse('\n'.join(out),media_type='audio/x-mpegurl')

async def xui_sync(mode):
    s=settings(); endpoint=s.get('xui_api','').strip()
    if not endpoint: return 'XUI não sincronizado: configure o endpoint/API.'
    # Safe generic adapter: sends normalized catalog to an explicitly configured endpoint.
    # It never edits XUI's database directly and never guesses undocumented endpoints.
    c=db();
    rows=c.execute("SELECT name,group_title,tvg_id,tvg_logo,source_url,kind FROM items WHERE kind<>? ORDER BY id",('channel',) if mode=='vod' else ('movie',)).fetchall(); c.close()
    payload={'mode':mode,'items':[dict(r) for r in rows]}
    headers={'Content-Type':'application/json'}
    if s.get('xui_token'): headers['Authorization']='Bearer '+s['xui_token']
    auth=aiohttp.BasicAuth(s['xui_user'],s['xui_password']) if s.get('xui_user') else None
    timeout=aiohttp.ClientTimeout(total=120)
    async with aiohttp.ClientSession(timeout=timeout,auth=auth) as session:
        async with session.post(endpoint,json=payload,headers=headers) as r:
            text=await r.text()
            if r.status>=400: raise RuntimeError(f'API XUI respondeu HTTP {r.status}: {text[:300]}')
            return f'API XUI OK ({r.status})'

async def hls_manifest(url):
    timeout=aiohttp.ClientTimeout(total=60)
    async with aiohttp.ClientSession(timeout=timeout) as session:
        async with session.get(url) as r:
            r.raise_for_status(); return await r.text(), r.headers.get('Content-Type','application/vnd.apple.mpegurl')

async def proxy_hls(item):
    # For HLS, return a local playlist whose segment requests are proxied through /hls-segment.
    text,ctype=await hls_manifest(item['source_url']); base=item['source_url']
    lines=[]
    for line in text.splitlines():
        if line and not line.startswith('#'):
            from urllib.parse import urljoin
            target=urljoin(base,line.strip())
            lines.append(f'/hls-segment/{item["id"]}?u={quote(target,safe="")}')
        else: lines.append(line)
    return PlainTextResponse('\n'.join(lines)+'\n',media_type=ctype)

@app.get('/stream/{item_id}')
async def stream(request:Request,item_id:int):
    c=db(); item=c.execute('SELECT * FROM items WHERE id=?',(item_id,)).fetchone(); c.close()
    if not item: return PlainTextResponse('Conteúdo não encontrado',404)
    if item['kind']=='channel':
        # Live channels are proxied; they are not written to disk.
        return await proxy_live(item)
    now=int(time.time()); c=db(); c.execute('UPDATE items SET last_access=?,updated_at=? WHERE id=?',(now,now,item_id)); c.commit(); c.close()
    if is_hls(item['source_url']): return await proxy_hls(item)
    path=Path(item['local_path']) if item['local_path'] else item_path(item)
    if item['status']=='complete' and path.exists(): return StreamingResponse(file_iter(path),media_type='video/mp4')
    asyncio.create_task(download_file(item, path))
    return await proxy_origin(item['source_url'])

async def proxy_live(item):
    if is_hls(item['source_url']): return await proxy_hls(item)
    return await proxy_origin(item['source_url'])

@app.get('/hls-segment/{item_id}')
async def hls_segment(request:Request,item_id:int,u:str):
    c=db(); item=c.execute('SELECT * FROM items WHERE id=?',(item_id,)).fetchone(); c.close()
    if not item: return PlainTextResponse('Not found',404)
    # Segment-level proxy. VOD segments can be cached later without pretending a live channel is a VOD file.
    return await proxy_origin(u)

async def proxy_origin(url):
    session=aiohttp.ClientSession(timeout=aiohttp.ClientTimeout(total=None,connect=30,sock_read=120))
    try: resp=await session.get(url)
    except Exception as e: await session.close(); return PlainTextResponse(f'Erro na origem: {e}',502)
    async def gen():
        try:
            async for chunk in resp.content.iter_chunked(1024*1024): yield chunk
        finally: await resp.release(); await session.close()
    return StreamingResponse(gen(),media_type=resp.headers.get('Content-Type','video/mp4'),status_code=resp.status)

async def download_file(item,path):
    tmp=path.with_suffix(path.suffix+'.part'); path.parent.mkdir(parents=True,exist_ok=True); now=int(time.time())
    c=db(); c.execute('UPDATE items SET status="downloading",local_path=?,updated_at=? WHERE id=?',(str(path),now,item['id'])); c.commit(); c.close(); total=0
    try:
        async with aiohttp.ClientSession(timeout=aiohttp.ClientTimeout(total=None,connect=30,sock_read=120)) as session:
            async with session.get(item['source_url']) as r:
                r.raise_for_status()
                with open(tmp,'wb') as f:
                    async for chunk in r.content.iter_chunked(1024*1024): f.write(chunk); total+=len(chunk)
        os.replace(tmp,path); c=db(); c.execute('UPDATE items SET status="complete",size=?,last_access=?,updated_at=? WHERE id=?',(total,int(time.time()),int(time.time()),item['id'])); c.commit(); c.close()
    except Exception:
        try: tmp.unlink(missing_ok=True)
        except: pass
        c=db(); c.execute('UPDATE items SET status="remote",updated_at=? WHERE id=?',(int(time.time()),item['id'])); c.commit(); c.close()

async def file_iter(path,chunk=1024*1024):
    with open(path,'rb') as f:
        while True:
            b=f.read(chunk)
            if not b: break
            yield b

async def cleaner_loop():
    while True:
        try: await asyncio.sleep(3600); await cleanup_cache()
        except asyncio.CancelledError: break
        except Exception: pass

async def cleanup_cache():
    s=settings(); days=max(1,int(s.get('expire_days','7') or 7)); cutoff=int(time.time())-days*86400
    c=db(); rows=c.execute("SELECT * FROM items WHERE status='complete' AND kind<>'channel' AND last_access>0 AND last_access<?",(cutoff,)).fetchall();
    for r in rows:
        p=Path(r['local_path']) if r['local_path'] else None
        if p and p.exists():
            try: p.unlink()
            except: continue
        c.execute("UPDATE items SET status='remote',local_path=NULL,size=0 WHERE id=?",(r['id'],))
    c.commit(); c.close()
    maxgb=s.get('max_gb','').strip()
    if maxgb:
        try: limit=int(float(maxgb)*1024**3)
        except: limit=0
        if limit>0:
            c=db(); rows=c.execute("SELECT * FROM items WHERE status='complete' AND kind<>'channel' ORDER BY last_access ASC").fetchall(); total=sum(r['size'] or 0 for r in rows)
            for r in rows:
                if total<=limit: break
                p=Path(r['local_path']) if r['local_path'] else None
                try:
                    if p and p.exists(): p.unlink()
                    total-=r['size'] or 0; c.execute("UPDATE items SET status='remote',local_path=NULL,size=0 WHERE id=?",(r['id'],))
                except: pass
            c.commit(); c.close()

@app.get('/api/status')
async def api_status():
    c=db(); q=lambda sql: c.execute(sql).fetchone()['n']; data={'catalog_vod':q("SELECT COUNT(*) n FROM items WHERE kind<>'channel'"),'channels':q("SELECT COUNT(*) n FROM items WHERE kind='channel'"),'cached':q("SELECT COUNT(*) n FROM items WHERE status='complete'"),'downloading':q("SELECT COUNT(*) n FROM items WHERE status='downloading'"),'cache_gb':round(q("SELECT COALESCE(SUM(size),0) n FROM items WHERE status='complete'")/1024**3,2),'storage':str(storage_root())}; c.close(); return JSONResponse(data)
