import os, re, sqlite3, hashlib, time, asyncio, json, shutil, threading
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
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'))
LOG = APP_DIR / 'data' / 'sync.log'
pwd = CryptContext(schemes=['bcrypt'], deprecated='auto')
app = FastAPI(title='VOD Creator Smart Cache V2')
templates = Jinja2Templates(directory=str(APP_DIR / 'app' / 'templates'))

# Lock para impedir sincronizações concorrentes (duplo clique)
_sync_lock = threading.Lock()


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);
    ''')
    # V3 migration for databases created by V2.
    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'
    }
    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_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('')):
    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}
    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.post('/sync')
async def sync(request: Request):
    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. Aguarde alguns minutos e tente novamente.'), 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:
        log('SYNC ERROR: ' + str(e))
        return RedirectResponse('/?error=' + quote(str(e)), 302)
    finally:
        _sync_lock.release()


@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, ''
    # XtreamUI clássico: status em STATUS_xxx
    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
    # XUI.one / outras variações: success: true/false
    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):
    name = (name or '').strip() or 'Sem categoria'
    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():
                return int(row.get('id') or row.get('category_id'))
        resp = await xui_call('create_category', {'category_name': name, 'category_type': kind}, 'POST')
        d = xui_data(resp)
        return int(d.get('id') or d.get('category_id'))
    except Exception as e:
        log(f'CATEGORY ERROR ({name}/{kind}): {e}')
        return None


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 source_value(url):
    # XUI expects stream_source as a JSON array string in Admin API.
    return json.dumps([url], ensure_ascii=False)


async def _run_xui_sync(mode):
    """Executa a sincronização com o XUI e retorna resumo com contagem de erros."""
    if not xui_base():
        return 'XUI não sincronizado: configure URL + Código de Acesso + Chave API.'
    if mode == 'vod':
        c = db(); rows = c.execute("SELECT * FROM items WHERE kind IN ('movie','series') ORDER BY id").fetchall(); c.close()
        movie_ok, movie_err = 0, 0
        series_ok, series_err = 0, 0
        # Movies
        for r in rows:
            if r['kind'] != 'movie':
                continue
            try:
                cat = await xui_category_id(r['group_title'], 'movie')
                payload = {'stream_display_name': r['name'], 'stream_source': source_value(local_stream_url(r['id'])), 'stream_icon': r['tvg_logo'] or '',
                           'target_container': ext(r['source_url']).lstrip('.') or 'mp4', 'direct_source': '1'}
                if cat:
                    payload['category_id'] = json.dumps([cat])
                if r['xui_id']:
                    await xui_call('edit_movie', dict(payload, id=r['xui_id']), 'POST')
                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()
                movie_ok += 1
            except Exception as e:
                log(f'MOVIE ERROR id={r["id"]} {r["name"]}: {e}')
                movie_err += 1
        # Series: one XUI series per normalized show title, episodes underneath.
        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():
            existing = next((r['xui_series_id'] for r in eps if r['xui_series_id']), None)
            try:
                cat = await xui_category_id(group, 'series')
            except Exception:
                cat = None
            payload = {'name': stitle, 'series_name': stitle}
            if cat:
                payload['category_id'] = json.dumps([cat])
            if existing:
                sid = int(existing)
                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:
                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': source_value(local_stream_url(r['id'])), 'target_container': ext(r['source_url']).lstrip('.') or 'mp4',
                      'direct_source': '1'}
                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')
                    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()
                    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']
        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_err:
            parts.append(f'{series_err} episódios com erro (veja o log)')
        return 'XUI: ' + ', '.join(parts) + '.'
    # Channels are sent as live channels using the same source mirror.
    c = db(); rows = c.execute("SELECT * FROM items WHERE kind='channel' ORDER BY id").fetchall(); c.close()
    ok, err = 0, 0
    for r in rows:
        try:
            cat = await xui_category_id(r['group_title'], 'live')
            payload = {'stream_display_name': r['name'], 'stream_source': source_value(local_stream_url(r['id'])), '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):
    return await _run_xui_sync(mode)


async def clear_xui_kind(kind):
    if not xui_base():
        return
    c = db()
    rows = c.execute("SELECT xui_id,xui_series_id FROM items WHERE kind=?", (kind,)).fetchall()
    c.close()
    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


@app.post('/clear/movies')
async def clear_movies(request: Request):
    if not auth(request):
        return RedirectResponse('/login', 302)
    await clear_xui_kind('movie')
    c = db(); rows = c.execute("SELECT * FROM items WHERE kind='movie'").fetchall()
    for r in rows:
        if r['local_path']:
            try:
                Path(r['local_path']).unlink(missing_ok=True)
            except Exception:
                pass
    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):
    if not auth(request):
        return RedirectResponse('/login', 302)
    await clear_xui_kind('series')
    c = db(); rows = c.execute("SELECT * FROM items WHERE kind='series'").fetchall()
    for r in rows:
        if r['local_path']:
            try:
                Path(r['local_path']).unlink(missing_ok=True)
            except Exception:
                pass
    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
        # 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:
        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):
    # 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('#'):
            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 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)
