import os, re, sqlite3, time, asyncio, json, threading, html
from pathlib import Path
from urllib.parse import urlparse, quote, urljoin
import aiohttp
from fastapi import FastAPI, Request, Form
from fastapi.responses import HTMLResponse, RedirectResponse, StreamingResponse, PlainTextResponse, JSONResponse
from fastapi.templating import Jinja2Templates
import bcrypt as _bcrypt

APP_DIR = Path(os.getenv('VOD_APP_DIR', '/opt/vodcreator'))
DB = APP_DIR / 'data' / 'vodcreator.db'
PORT = int(os.getenv('VOD_PORT', '8090'))
LOG = APP_DIR / 'data' / 'sync.log'

class _Pwd:
    def hash(self, pw): return _bcrypt.hashpw(pw.encode(), _bcrypt.gensalt()).decode()
    def verify(self, pw, h): return _bcrypt.checkpw(pw.encode(), h.encode() if isinstance(h,str) else h)
pwd = _Pwd()

app = FastAPI(title='Sinc Brasil v6')
templates = Jinja2Templates(directory=str(APP_DIR / 'app' / 'templates'))

# Lock para impedir sincronizações concorrentes (duplo clique)
_sync_lock = threading.Lock()
# Estado de progresso da sincronização (usado pela barra de progresso)
_sync_state = {'running': False, 'total': 0, 'done': 0, 'label': '', 'error': ''}
# Flag de cancelamento da sincronização (botão Parar)
_sync_cancel = {'cancel': False}
# Watchdog: se a sincronização estiver rodando há mais de X segundos sem progresso,
# considera travada e encerra (ex.: crash do processo no meio de um item)
_sync_heartbeat = {'last_progress': time.time()}
_SYNC_STUCK_SECONDS = 120


def sync_canceled():
    """Retorna True se o usuário pediu para parar a sincronização."""
    _sync_heartbeat['last_progress'] = time.time()
    return _sync_cancel['cancel']


def log(msg):
    try:
        with open(LOG, 'a') as f:
            f.write(f"{time.strftime('%Y-%m-%d %H:%M:%S')} {msg}\n")
    except Exception:
        pass


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,
      xui_id INTEGER DEFAULT NULL,
      xui_series_id INTEGER DEFAULT 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);
    ''')
    cols = {r[1] for r in c.execute("PRAGMA table_info(items)").fetchall()}
    if 'xui_id' not in cols:
        c.execute("ALTER TABLE items ADD COLUMN xui_id INTEGER")
    if 'xui_series_id' not in cols:
        c.execute("ALTER TABLE items ADD COLUMN xui_series_id INTEGER")
    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_access_code': '', 'xui_api_key': '', 'xui_api': '', 'xui_user': '', 'xui_password': '', 'xui_token': '', 'xui_server_id': '',
        'xui_category_mode': 'group', 'channel_primary': '', 'channel_backup': '',
        'auto_sync': '0',
        # V4 — novos campos
        'movie_bouquet': '',           # bouquets para filmes (separados por ;)
        'series_bouquet': '',          # bouquets para séries (separados por ;)
        'bouquet_include_series_bouquet': '',  # nome do bouquet adulto p/ filtrar séries
    }
    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', '')}
            # V4: IMDB ID / TMDB ID podem vir como atributos tvg-name extra ou imdb-id/tmdb-id customizados
            for a in ('imdb-id', 'imdb_id', 'tmdb-id', 'tmdb_id'):
                if a in attrs and attrs[a]:
                    current[a] = attrs[a]
            # Também procura padrão ttNNNNNNN no título como fallback
            m = re.search(r'\b(tt\d{7,8})\b', title)
            if m:
                current.setdefault('imdb-id', m.group(1))
        elif not line.startswith('#') and current:
            name = current['name']; group = current['group_title']
            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()
    # V4: incluir o último erro da sincronização (para exibir ao final de um sync com falha)
    s = settings()
    ctx = dict(total=total, channels=channels, movies=movies, series=series, cached=cached, active=active, size_gb=round(size / 1024 ** 3, 2), s=s, storage=str(storage_root()), last_sync_error=s.get('last_sync_error', ''))
    if 'error' in dict(request.query_params):
        ctx['last_sync_error'] = request.query_params['error']
        set_setting('last_sync_error', request.query_params['error'])
    return render(request, 'dashboard.html', **ctx)


@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_access_code: str = Form(''), xui_api_key: str = Form(''), xui_api: str = Form(''),
                        xui_user: str = Form(''), xui_password: str = Form(''), xui_token: str = Form(''), xui_server_id: str = Form(''),
                        expire_days: str = Form('7'), max_gb: str = Form(''), channel_primary: str = Form(''), channel_backup: str = Form(''),
                        movie_bouquet: str = Form(''), series_bouquet: 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_access_code': xui_access_code, 'xui_api_key': xui_api_key, 'xui_api': xui_api, 'xui_user': xui_user, 'xui_password': xui_password, 'xui_token': xui_token, 'xui_server_id': xui_server_id, 'expire_days': expire_days, 'max_gb': max_gb, 'channel_primary': channel_primary, 'channel_backup': channel_backup, 'movie_bouquet': movie_bouquet, 'series_bouquet': series_bouquet}
    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:
        try:
            async with session.get(url) as r:
                if r.status >= 400:
                    raise RuntimeError(f'M3U retornou erro {r.status} (lista indisponível ou URL incorreta).')
                return await r.text(errors='ignore')
        except aiohttp.ClientConnectorError as e:
            raise RuntimeError(f'Não foi possível baixar a lista M3U: servidor inacessível ({url}). Confira o endereço.')
        except aiohttp.ServerTimeoutError:
            raise RuntimeError(f'Tempo esgotado ao baixar a lista M3U ({url}). O servidor da lista pode estar fora do ar ou lento.')


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.get('/api/sync-progress')
async def sync_progress(request: Request):
    """V5: endpoint da barra de progresso. Exige login.
    Watchdog: se running=True mas sem progresso há _SYNC_STUCK_SECONDS, considera travada."""
    if not auth(request):
        return JSONResponse({'running': False})
    if _sync_state.get('running') and (time.time() - _sync_heartbeat['last_progress']) > _SYNC_STUCK_SECONDS:
        log('SYNC STUCK (watchdog): encerrando sincronização travada')
        _sync_state.update(running=False, label='Sincronização travada e encerrada. Clique em Sincronizar para continuar.', done=_sync_state.get('total', 0))
        try:
            if _sync_lock.locked():
                _sync_lock.release()
        except Exception:
            pass
        _sync_cancel['cancel'] = False
    return JSONResponse(_sync_state)


@app.post('/sync')
async def sync(request: Request, mode: str = Form('all'), speed: str = Form('normal')):
    """V6: Sincronização flexível (tudo, só filmes ou só séries) com modo rápido/normal."""
    if not auth(request):
        return RedirectResponse('/login', 302)
    if not _sync_lock.acquire(blocking=False):
        return RedirectResponse('/?error=' + quote('Já existe uma sincronização em andamento.'), 302)
    try:
        s = settings(); url = s.get('m3u_url', '')
        if not url:
            raise RuntimeError('Configure a URL da M3U em Configurações.')
        
        # 1. Baixa e atualiza catálogo local (sempre lê tudo da M3U)
        _sync_state.update(running=True, total=0, done=0, label='Baixando lista M3U...')
        n = await sync_items(url, False, False)
        
        # 2. Sincroniza com XUI conforme o modo
        # speed='fast' pula itens que já têm xui_id (não chama edit_movie/edit_series)
        x = await xui_sync('vod', filter_kind=mode, fast_mode=(speed == 'fast'))
        
        set_setting('last_sync_error', '')
        return RedirectResponse(f'/?sync={n}&xui=' + quote(x), 302)
    except Exception as e:
        msg = str(e); log('SYNC ERROR: ' + msg); set_setting('last_sync_error', msg)
        return RedirectResponse('/?error=' + quote(msg), 302)
    finally:
        _sync_cancel['cancel'] = False
        _sync_state.update(running=False, done=_sync_state.get('total', 0), label='Sincronização concluída')
        if _sync_lock.locked(): _sync_lock.release()


@app.post('/sync/cancel')
@app.get('/sync/cancel')
async def sync_cancel(request: Request):
    """Para a sincronização em andamento (finaliza o item atual e encerra)."""
    if not auth(request):
        return RedirectResponse('/login', 302)
    _sync_cancel['cancel'] = True
    _sync_heartbeat['last_progress'] = time.time()
    # Se estiver travada (watchdog), libera o lock na hora para não bloquear uma nova sincronização
    try:
        if _sync_lock.locked():
            _sync_lock.release()
    except Exception:
        pass
    _sync_state.update(running=True, label='Parando a sincronização…')
    return RedirectResponse('/', 302)



@app.post('/xui/test')
async def xui_test(request: Request):
    if not auth(request):
        return RedirectResponse('/login', 302)
    try:
        detail = await xui_diagnose()
        return RedirectResponse('/settings?xui_ok=' + quote(detail), 302)
    except Exception as e:
        return RedirectResponse('/settings?xui_error=' + quote(str(e)), 302)


async def xui_diagnose():
    """Teste de conexão em etapas, mostrando onde o problema está."""
    s = settings()
    base = xui_base()
    key = (s.get('xui_api_key') or s.get('xui_token') or '').strip()
    if not (s.get('xui_url') or '').strip():
        raise RuntimeError('XUI: campo "URL do XUI" está vazio. Informe http://IP:PORTA do painel.')
    if not base:
        raise RuntimeError('XUI: informe "URL do XUI" + "Código de Acesso" (Management > Access Control > Access Codes > tipo Admin API).')
    if not key:
        raise RuntimeError('XUI: campo "Chave API" está vazio. Gere a chave no perfil do administrador (clique no ícone de refresh se estiver vazia).')
    timeout = aiohttp.ClientTimeout(total=30)
    session = aiohttp.ClientSession(timeout=timeout)
    try:
        # Etapa 1: alcança o painel?
        try:
            async with session.get(s.get('xui_url', '').strip().rstrip('/')) as r:
                if r.status >= 500:
                    raise RuntimeError(f'XUI: o painel respondeu com erro {r.status} em {s.get("xui_url")}. Verifique se o serviço do painel está rodando.')
                # 4xx na raiz é normal para painéis que não servem página pública; o importante é alcançar o host.
        except aiohttp.ClientConnectorError as e:
            raise RuntimeError(f'XUI: não foi possível alcançar {s.get("xui_url")} ({e}). Confira IP, porta e se o painel está ligado.')
        except aiohttp.ServerTimeoutError:
            raise RuntimeError(f'XUI: tempo esgotado ao conectar em {s.get("xui_url")}. Confira IP e porta (a API usa a porta do painel, ex.: 25000/8000).')
        # Etapa 2: access code existe e a rota da API responde?
        async with session.get(base, params={'api_key': key, 'action': 'user_info'}) as r:
            text = await r.text()
            if r.status == 404:
                raise RuntimeError(f'XUI: rota {base} retornou 404. O Código de Acesso está errado ou não é do tipo Admin API.')
            if r.status >= 400:
                raise RuntimeError(f'XUI HTTP {r.status}: {text[:300]}')
            try:
                data = json.loads(text)
            except Exception:
                raise RuntimeError(f'XUI: resposta não é JSON em {base} — {text[:300]}. Verifique a URL/porta do painel.')
            ok, err = parse_xui_status(data)
            if not ok:
                raise RuntimeError(f'XUI: o painel rejeitou a conexão — {err or "acesso negado"}. Confira se a Chave API foi gerada no perfil do administrador.')
    finally:
        await session.close()
    return f'Conexão XUI OK (painel em {s.get("xui_url")} + código {s.get("xui_access_code")} + API key válida).'


def parse_xui_status(data):
    """Normaliza o status de respostas do XtreamUI/XUI.one.
    Retorna (ok: bool, erro: str)."""
    if not isinstance(data, dict):
        return True, ''
    status = str(data.get('status', ''))
    if status.startswith('STATUS_'):
        if status in ('STATUS_SUCCESS', 'STATUS_SUCCESS_REPLACE', 'STATUS_SUCCESS_MULTI'):
            return True, ''
        return False, data.get('error') or data.get('message') or status
    if 'success' in data:
        if data['success'] in (True, 'true', 1, '1'):
            return True, ''
        return False, data.get('error') or data.get('message') or 'success=false'
    return True, ''


async def xui_call(action, params=None, method='GET'):
    s = settings(); base = xui_base()
    key = (s.get('xui_api_key') or s.get('xui_token') or '').strip()
    if not base or not key:
        raise RuntimeError('XUI: informe URL, Código de Acesso e Chave API.')
    params = dict(params or {}); params.update({'api_key': key, 'action': action})
    timeout = aiohttp.ClientTimeout(total=120)
    async with aiohttp.ClientSession(timeout=timeout) as session:
        last_err = None
        for attempt in range(3):
            try:
                if method.upper() == 'POST':
                    async with session.post(base, data=params) as r:
                        text = await r.text()
                else:
                    async with session.get(base, params=params) as r:
                        text = await r.text()
                if r.status >= 400:
                    raise RuntimeError(f'XUI HTTP {r.status}: {text[:400]}')
                try:
                    data = json.loads(text)
                except Exception:
                    raise RuntimeError(f'XUI resposta inválida: {text[:400]}')
                ok, err = parse_xui_status(data)
                if not ok:
                    raise RuntimeError(f'XUI: {err}')
                return data
            except RuntimeError:
                raise
            except Exception as e:
                last_err = e
                if attempt < 2:
                    await asyncio.sleep(2)
                else:
                    raise RuntimeError(f'XUI {action}: {last_err}')


def xui_data(resp):
    if isinstance(resp, dict):
        return resp.get('data', resp)
    return resp


def series_title(name):
    x = re.sub(r'(?i)[._ -]*S\d{1,2}E\d{1,3}.*$', '', name).strip(' .-_')
    x = re.sub(r'(?i)[._ -]*(?:season|temporada)\s*\d+.*$', '', x).strip(' .-_')
    return x or name


def season_episode(name):
    m = re.search(r'(?i)S(\d{1,2})E(\d{1,3})', name)
    if m:
        return int(m.group(1)), int(m.group(2))
    m = re.search(r'(?i)(?:season|temporada)\s*(\d+).*?(?:episode|epis[oó]dio)\s*(\d+)', name)
    return (int(m.group(1)), int(m.group(2))) if m else (1, 1)


_category_cache = {}


async def xui_category_id(name, kind):
    """Retorna o ID da categoria no XUI, criando se necessário. Cache em memória.
    V5.2: tenta tipos alternativos se a criação falhar e loga erros detalhados."""
    name = (name or '').strip() or 'Sem categoria'
    cache_key = (name.lower(), kind)
    if cache_key in _category_cache:
        return _category_cache[cache_key]
    try:
        data = xui_data(await xui_call('get_categories'))
        if isinstance(data, dict):
            data = list(data.values())
        for row in data or []:
            if str(row.get('category_name') or row.get('name') or '').strip().lower() == name.lower():
                cid = int(row.get('id') or row.get('category_id'))
                _category_cache[cache_key] = cid
                return cid
        cid = None
        for t in (kind, 'vod', 'movie', 'series', 'live'):
            try:
                resp = await xui_call('create_category', {'category_name': name, 'category_type': t}, 'POST')
                d = xui_data(resp)
                cid = int(d.get('id') or d.get('category_id'))
                break
            except Exception as e:
                log(f'CATEGORY CREATE ERROR {name} (type={t}): {e}')
                continue
        if cid:
            _category_cache[cache_key] = cid
            if cid is not None:
                log(f'CATEGORY CREATED {name} id={cid} type={t}')
        return cid
    except Exception as e:
        log(f'CATEGORY ERROR ({name}/{kind}): {e}')
        return None


_bouquet_cache = {}


async def xui_bouquet_ids(names):
    """Resolve nomes de bouquets para IDs no XUI. Retorna lista de IDs (vazia se nenhum encontrado)."""
    out = []
    names = [n.strip() for n in names if n.strip()]
    if not names:
        return out
    key = tuple(n.lower() for n in names)
    if key in _bouquet_cache:
        return _bouquet_cache[key]
    try:
        data = xui_data(await xui_call('get_bouquets'))
        if isinstance(data, dict):
            data = list(data.values())
        want = {n.lower() for n in names}
        for row in data or []:
            bname = str(row.get('bouquet_name') or row.get('name') or '').strip().lower()
            if bname in want:
                out.append(int(row.get('id') or row.get('bouquet_id')))
    except Exception as e:
        log(f'BOUQUET ERROR: {e}')
    _bouquet_cache[key] = out
    return out


def local_stream_url(item_id):
    s = settings(); base = (s.get('base_url') or '').strip().rstrip('/')
    if not base:
        base = 'http://127.0.0.1:8090'
    return f'{base}/stream/{item_id}'


def movie_payload_source(r):
    """V6: Corrige montagem da fonte para o XUI.one.
    Retorna uma string JSON que o XUI.one aceita como array de fontes.
    Se o link vier quebrado ou como '{', garante a URL original."""
    url = str(r['source_url']).strip()
    mirror = local_stream_url(r['id'])
    # XUI.one espera um array JSON de strings para fontes múltiplas
    sources = [url, mirror]
    return json.dumps(sources, ensure_ascii=False)


def imdb_fields(r):
    """Campos IMDB/TMDB para o XUI. Tenta imdb-id da M3U e, se não houver, procura tt no nome."""
    d = dict(r) if not isinstance(r, dict) else r
    fields = {}
    imdb = d.get('imdb-id') or d.get('imdb_id') or ''
    tmdb = d.get('tmdb-id') or d.get('tmdb_id') or ''
    if imdb:
        fields['tmdb_id'] = imdb  # XUI.one usa tmdb_id para buscar pôster/dados
        fields['imdb_id'] = imdb
    if tmdb:
        fields['tmdb_id'] = tmdb
    return fields


async def add_to_bouquets(xui_item_id, item_type, bouquet_ids):
    """Adiciona um item já criado aos bouquets informados (IDs numéricos).
    Retorna True se pelo menos um bouquet aceitou o item."""
    any_ok = False
    for bid in bouquet_ids:
        ok = False
        for action in ('add_to_bouquet', 'add_stream_bouquet'):
            try:
                await xui_call(action, {'bouquet_id': bid, 'stream_id': xui_item_id, 'type': item_type}, 'POST')
                ok = True
                any_ok = True
                break
            except Exception:
                continue
        if not ok:
            log(f'BOUQUET ADD ERROR item={xui_item_id} bouquet_id={bid}: falhou em todas as variações da ação')
    return any_ok


async def ensure_bouquets(names, kind):
    """Garante que os bouquets existem no XUI e retorna os IDs numéricos.
    Tenta as variações de type ('vod'/'series'/'movie') se a criação falhar."""
    ids = []
    try:
        data = xui_data(await xui_call('get_bouquets'))
        if isinstance(data, dict):
            data = list(data.values())
        existing = {str(row.get('bouquet_name') or row.get('name') or '').strip().lower(): int(row.get('id') or row.get('bouquet_id')) for row in data or []}
        for n in names:
            low = n.strip().lower()
            if not low:
                continue
            if low in existing:
                ids.append(existing[low])
                continue
            created = None
            for t in (kind, 'vod', 'movie', 'series', 'live'):
                try:
                    resp = await xui_call('create_bouquet', {'bouquet_name': n.strip(), 'type': t}, 'POST')
                    d = xui_data(resp)
                    created = int(d.get('id') or d.get('bouquet_id'))
                    break
                except Exception as e:
                    log(f'BOUQUET CREATE ERROR {n} (type={t}): {e}')
                    continue
            if created:
                ids.append(created)
                # atualiza o cache para os próximos itens acharem na hora
                _bouquet_cache[tuple(x.lower() for x in names)] = ids
                log(f'BOUQUET CREATED {n} id={created} type={t}')
            else:
                log(f'BOUQUET ALERT: não foi possível criar o bouquet "{n}" — verifique o log data/sync.log')
    except Exception as e:
        log(f'BOUQUET ERROR: {e}')
    return ids


async def _run_xui_sync(mode, filter_kind='all', fast_mode=False):
    """Executa a sincronização com o XUI com barra de progresso e itens idempotentes.
    V6: suporte a filtro (filmes/series) e modo rápido (pula edições)."""
    _sync_cancel['cancel'] = False
    if not xui_base():
        return 'XUI não sincronizado: configure URL + Código de Acesso + Chave API.'

    s = settings()
    movie_bouquets = [b for b in s.get('movie_bouquet', '').split(';') if b.strip()]
    series_bouquets = [b for b in s.get('series_bouquet', '').split(';') if b.strip()]
    
    movie_bouquet_ids = await ensure_bouquets(movie_bouquets, 'vod')
    series_bouquet_ids = await ensure_bouquets(series_bouquets, 'series')
    _bouquet_cache.clear()

    if mode == 'vod':
        c = db()
        if filter_kind == 'movie':
            rows = c.execute("SELECT * FROM items WHERE kind='movie' ORDER BY id").fetchall()
        elif filter_kind == 'series':
            rows = c.execute("SELECT * FROM items WHERE kind='series' ORDER BY id").fetchall()
        else:
            rows = c.execute("SELECT * FROM items WHERE kind IN ('movie','series') ORDER BY id").fetchall()
        c.close()
        
        _sync_state.update(running=True, total=len(rows), done=0, label='Sincronizando VOD...', error='')
        _sync_heartbeat['last_progress'] = time.time()
        movie_ok, movie_skip, movie_err = 0, 0, 0
        series_ok, series_skip, series_err = 0, 0, 0

        rows = [dict(r) for r in rows]
        for idx, r in enumerate(rows):
            _sync_state.update(done=idx + 1, label=f'{r["kind"].title()}: {r["name"][:60]}')
            if sync_canceled():
                _sync_state.update(running=False, label='Sincronização parada')
                return 'XUI: sincronização parada pelo usuário.'
            
            if r['kind'] == 'movie':
                if fast_mode and r['xui_id']:
                    movie_skip += 1; continue
                try:
                    cat = await xui_category_id(r['group_title'], 'movie')
                    payload = {'stream_display_name': r['name'], 'stream_source': movie_payload_source(r), 'stream_icon': r['tvg_logo'] or '',
                               'target_container': ext(r['source_url']).lstrip('.') or 'mp4', 'direct_source': '1'}
                    payload.update(imdb_fields(r))
                    if cat: payload['category_id'] = json.dumps([cat])
                    
                    if r['xui_id']:
                        await xui_call('edit_movie', dict(payload, id=r['xui_id']), 'POST')
                        movie_skip += 1
                    else:
                        resp = await xui_call('create_movie', payload, 'POST')
                        d = xui_data(resp)
                        xid = (d.get('id') or d.get('movie_id')) if isinstance(d, dict) else None
                        if xid:
                            c = db(); c.execute("UPDATE items SET xui_id=? WHERE id=?", (int(xid), r['id'])); c.commit(); c.close()
                            if movie_bouquet_ids: await add_to_bouquets(int(xid), 'movie', movie_bouquet_ids)
                        movie_ok += 1
                except Exception as e:
                    log(f'MOVIE ERROR id={r["id"]} {r["name"]}: {e}'); movie_err += 1
        
        # Séries
        groups = {}
        for r in rows:
            if r['kind'] == 'series':
                groups.setdefault((series_title(r['name']), r['group_title']), []).append(r)
        
        for (stitle, group), eps in groups.items():
            if sync_canceled():
                _sync_state.update(running=False, label='Sincronização parada')
                return 'XUI: sincronização parada pelo usuário.'
            
            existing = next((r['xui_series_id'] for r in eps if r['xui_series_id']), None)
            cat = await xui_category_id(group, 'series')
            payload = {'name': stitle, 'series_name': stitle}
            if cat: payload['category_id'] = json.dumps([cat])
            
            if existing:
                sid = int(existing)
                if not fast_mode:
                    try: await xui_call('edit_series', dict(payload, id=sid), 'POST')
                    except Exception: pass
            else:
                try:
                    resp = await xui_call('create_series', payload, 'POST')
                    d = xui_data(resp)
                    sid = int(d.get('id') or d.get('series_id'))
                except Exception as e:
                    log(f'SERIES ERROR {stitle}: {e}'); series_err += len(eps); continue
            
            for r in eps:
                if fast_mode and r['xui_id']:
                    series_skip += 1; continue
                season, episode = season_episode(r['name'])
                ep = {'series': sid, 'series_id': sid, 'season_num': season, 'episode_num': episode,
                      'stream_display_name': r['name'], 'episode_name': r['name'],
                      'stream_source': movie_payload_source(r), 'target_container': ext(r['source_url']).lstrip('.') or 'mp4',
                      'direct_source': '1'}
                ep.update(imdb_fields(r))
                if r['tvg_logo']: ep['stream_icon'] = r['tvg_logo']
                
                try:
                    if r['xui_id']:
                        await xui_call('edit_episode', dict(ep, id=r['xui_id']), 'POST')
                        series_skip += 1
                    else:
                        resp = await xui_call('create_episode', ep, 'POST')
                        d = xui_data(resp)
                        xid = (d.get('id') or d.get('episode_id')) if isinstance(d, dict) else None
                        if xid:
                            c = db(); c.execute("UPDATE items SET xui_id=?,xui_series_id=? WHERE id=?", (int(xid), sid, r['id'])); c.commit(); c.close()
                            if series_bouquet_ids: await add_to_bouquets(int(xid), 'series', series_bouquet_ids)
                        series_ok += 1
                except Exception as e:
                    log(f'EPISODE ERROR id={r["id"]} {r["name"]}: {e}'); series_err += 1
        parts = [f'{movie_ok} filmes enviados']
        parts.append(f'{len(movie_bouquet_ids)} bouquets de filmes usados')
        parts.append(f'{len(series_bouquet_ids)} bouquets de séries usados')
        if movie_skip:
            parts.append(f'{movie_skip} filmes já existiam (atualizados)')
        if movie_err:
            parts.append(f'{movie_err} filmes com erro (veja o log)')
        parts.append(f'{series_ok} episódios em {len(groups)} séries')
        if series_skip:
            parts.append(f'{series_skip} episódios já existiam')
        if series_err:
            parts.append(f'{series_err} episódios com erro')
        return 'XUI: ' + ', '.join(parts) + '.'

    # Channels
    c = db(); rows = [dict(r) for r in c.execute("SELECT * FROM items WHERE kind='channel' ORDER BY id").fetchall()]; c.close()
    _sync_state.update(running=True, total=len(rows), done=0, label='Sincronizando canais...')
    _sync_heartbeat['last_progress'] = time.time()
    ok, err = 0, 0
    for idx, r in enumerate(rows):
        _sync_state.update(done=idx + 1, label=r['name'][:60])
        if sync_canceled():
            _sync_state.update(running=False, label='Sincronização parada')
            return 'XUI: sincronização parada pelo usuário.'
        try:
            cat = await xui_category_id(r['group_title'], 'live')
            payload = {'stream_display_name': r['name'], 'stream_source': json.dumps([r['source_url']], ensure_ascii=False), 'stream_icon': r['tvg_logo'] or '', 'target_container': 'ts', 'direct_source': '1'}
            if cat:
                payload['category_id'] = json.dumps([cat])
            if r['xui_id']:
                await xui_call('edit_channel', dict(payload, id=r['xui_id']), 'POST')
            else:
                resp = await xui_call('create_channel', payload, 'POST')
                d = xui_data(resp)
                xid = (d.get('id') or d.get('stream_id')) if isinstance(d, dict) else None
                if xid:
                    c = db(); c.execute("UPDATE items SET xui_id=? WHERE id=?", (int(xid), r['id'])); c.commit(); c.close()
            ok += 1
        except Exception as e:
            log(f'CHANNEL ERROR id={r["id"]} {r["name"]}: {e}')
            err += 1
    parts = [f'{ok} canais sincronizados']
    if err:
        parts.append(f'{err} canais com erro')
    return 'XUI: ' + ', '.join(parts) + '.'


async def xui_sync(mode, filter_kind='all', fast_mode=False):
    return await _run_xui_sync(mode, filter_kind=filter_kind, fast_mode=fast_mode)


async def clear_xui_kind(kind, remove_from_xui=True):
    """Remove itens do XUI (se remove_from_xui=True) e apaga o catálogo local."""
    if remove_from_xui and not xui_base():
        remove_from_xui = False
    c = db()
    rows = c.execute("SELECT xui_id,xui_series_id FROM items WHERE kind=?", (kind,)).fetchall()
    c.close()
    if remove_from_xui:
        if kind == 'movie':
            for r in rows:
                if r['xui_id']:
                    try: await xui_call('delete_movie', {'id': int(r['xui_id'])}, 'POST')
                    except Exception: pass
        else:
            series_ids = {int(r['xui_series_id']) for r in rows if r['xui_series_id']}
            for r in rows:
                if r['xui_id']:
                    try: await xui_call('delete_episode', {'id': int(r['xui_id'])}, 'POST')
                    except Exception: pass
            for sid in series_ids:
                try: await xui_call('delete_series', {'id': sid}, 'POST')
                except Exception: pass
    # Sempre apaga o catálogo local (e os arquivos do HD quando solicitado pela rota)
    c = db(); c.execute("DELETE FROM items WHERE kind=?", (kind,)); c.commit(); c.close()


@app.post('/clear/movies')
async def clear_movies(request: Request, mode: str = Form('normal')):
    """V6: Limpeza rápida (só catálogo) ou normal (catálogo + HD + XUI)."""
    if not auth(request):
        return RedirectResponse('/login', 302)
    if mode == 'normal':
        await clear_xui_kind('movie', remove_from_xui=True)
        c = db(); rows = c.execute("SELECT * FROM items WHERE kind='movie'").fetchall(); c.close()
        for r in rows:
            if r['local_path']:
                try: Path(r['local_path']).unlink(missing_ok=True)
                except Exception: pass
    c = db(); c.execute("DELETE FROM items WHERE kind='movie'"); c.commit(); c.close()
    return RedirectResponse('/?cleared=movies', 302)


@app.post('/clear/series')
async def clear_series(request: Request, mode: str = Form('normal')):
    """V6: Limpeza rápida (só catálogo) ou normal (catálogo + HD + XUI)."""
    if not auth(request):
        return RedirectResponse('/login', 302)
    if mode == 'normal':
        await clear_xui_kind('series', remove_from_xui=True)
        c = db(); rows = c.execute("SELECT * FROM items WHERE kind='series'").fetchall(); c.close()
        for r in rows:
            if r['local_path']:
                try: Path(r['local_path']).unlink(missing_ok=True)
                except Exception: pass
    c = db(); c.execute("DELETE FROM items WHERE kind='series'"); c.commit(); c.close()
    return RedirectResponse('/?cleared=series', 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
        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:
        log('CHANNELS SYNC ERROR: ' + str(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()
    sql = "SELECT * FROM items WHERE kind<>'channel' ORDER BY id" if kind == 'vod' else "SELECT * FROM items WHERE kind='channel' ORDER BY id"
    rows = c.execute(sql).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')


def xui_base():
    s = settings()
    explicit = (s.get('xui_api') or '').strip().rstrip('/')
    if explicit:
        return explicit
    url = (s.get('xui_url') or '').strip().rstrip('/')
    access = (s.get('xui_access_code') or '').strip().strip('/')
    if not url or not access:
        return ''
    return f"{url}/{access}/"


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):
    text, ctype = await hls_manifest(item['source_url']); base = item['source_url']
    lines = []
    for line in text.splitlines():
        if line and not line.startswith('#'):
            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':
        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)
    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 Exception:
            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 Exception:
                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 Exception:
            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 Exception:
                    pass
            c.commit(); c.close()


@app.get('/api/status')
async def api_status():
    c = db()

    def q(sql):
        return 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)
