"""AI Content Studio — single-operator deployable MVP.
Run one uvicorn process; durable SQLite job state resumes at startup.
"""
from __future__ import annotations
import os, sys, json, sqlite3, time, io, zipfile, csv, hmac, hashlib, secrets, threading, uuid, re, mimetypes, traceback
from contextlib import asynccontextmanager
from pathlib import Path
from datetime import datetime, timezone
from difflib import SequenceMatcher
from typing import Literal
from fastapi import FastAPI, UploadFile, File, Form, HTTPException, Request, Response, Depends
from fastapi.responses import FileResponse, StreamingResponse, HTMLResponse, JSONResponse
from starlette.middleware.trustedhost import TrustedHostMiddleware
from fastapi.staticfiles import StaticFiles
from pydantic import BaseModel, Field
from PIL import Image, ImageOps
from dotenv import load_dotenv
from . import ai_client, graphics

load_dotenv()
ROOT=Path(__file__).resolve().parent.parent
DATA=Path(os.getenv('AI_STUDIO_DATA', str(ROOT/'data'))).resolve()
MEDIA=DATA/'media'
DB=DATA/'studio.sqlite3'
PUBLIC=ROOT/'app'/'static'
DATA.mkdir(parents=True,exist_ok=True); MEDIA.mkdir(parents=True,exist_ok=True)
PASSWORD=os.getenv('APP_PASSWORD','')
APP_ENV=os.getenv('APP_ENV','development').lower()
if APP_ENV=='production' and (not PASSWORD or len(PASSWORD)<12):
    raise RuntimeError('APP_PASSWORD wajib minimal 12 karakter saat APP_ENV=production')
SESSION_SECRET=os.getenv('SESSION_SECRET') or (PASSWORD if APP_ENV!='production' and PASSWORD else secrets.token_hex(24))
APP_DOMAIN=os.getenv('APP_DOMAIN','studio.kurniawangroupofficial.com').strip().lower()
if APP_ENV=='production' and (len(SESSION_SECRET)<32 or SESSION_SECRET==PASSWORD):
    raise RuntimeError('SESSION_SECRET wajib unik, berbeda dari APP_PASSWORD, minimal 32 karakter saat production')
shutdown=threading.Event()
_worker=None

def now():return datetime.now(timezone.utc).isoformat()
def uid():return uuid.uuid4().hex

def db():
    con=sqlite3.connect(str(DB),timeout=30,check_same_thread=False)
    con.row_factory=sqlite3.Row
    con.execute('PRAGMA journal_mode=WAL')
    con.execute('PRAGMA foreign_keys=ON')
    return con

def query_one(sql, params=()):
    with db() as con:
        r=con.execute(sql,params).fetchone()
        return dict(r) if r else None

def query_all(sql,params=()):
    with db() as con: return [dict(x) for x in con.execute(sql,params).fetchall()]

def execute(sql, params=()):
    with db() as con:
        con.execute(sql,params)
        con.commit()

def init_db():
    with db() as con:
        con.executescript('''
        CREATE TABLE IF NOT EXISTS brands(
           id TEXT PRIMARY KEY, name TEXT NOT NULL, industry TEXT DEFAULT '', audience TEXT DEFAULT '',
           primary_color TEXT DEFAULT '#7545C5', secondary_color TEXT DEFAULT '#FAF7FF',
           logo_path TEXT, reference_path TEXT, style_notes TEXT DEFAULT '',
           caption_prefix TEXT DEFAULT '', max_hashtags INTEGER DEFAULT 5, created_at TEXT NOT NULL
        );
        CREATE TABLE IF NOT EXISTS batches(
           id TEXT PRIMARY KEY, brand_id TEXT NOT NULL REFERENCES brands(id),
           requested INTEGER NOT NULL, theme TEXT DEFAULT '', mode TEXT NOT NULL, style_mode TEXT DEFAULT 'medium',
           manual_topics TEXT DEFAULT '', status TEXT NOT NULL DEFAULT 'queued',
           error TEXT, stop_requested INTEGER DEFAULT 0, created_at TEXT NOT NULL, updated_at TEXT NOT NULL
        );
        CREATE TABLE IF NOT EXISTS contents(
           id TEXT PRIMARY KEY, batch_id TEXT NOT NULL REFERENCES batches(id),
           position INTEGER NOT NULL, title TEXT NOT NULL, subtitle TEXT DEFAULT '', category TEXT DEFAULT '',
           visual_prompt TEXT DEFAULT '', caption TEXT DEFAULT '', image_path TEXT, asset_path TEXT,
           status TEXT NOT NULL DEFAULT 'queued', error TEXT, created_at TEXT NOT NULL,
           UNIQUE(batch_id,position)
        );
        CREATE INDEX IF NOT EXISTS idx_contents_batch ON contents(batch_id,position);
        CREATE INDEX IF NOT EXISTS idx_batches_status ON batches(status);
        ''')
        con.commit()


def check_auth(request:Request):
    if not PASSWORD:return
    cookie=request.cookies.get('acs_session','')
    try:
        exp,sig=cookie.split(':',1)
        expected=hmac.new(SESSION_SECRET.encode(),exp.encode(),hashlib.sha256).hexdigest()
        valid=hmac.compare_digest(sig,expected) and int(exp)>time.time()
    except (ValueError,AttributeError):valid=False
    if not valid:raise HTTPException(401,'Silakan login terlebih dahulu.')


def ensure_brand(bid):
    brand=query_one('SELECT * FROM brands WHERE id=?',(bid,))
    if not brand:raise HTTPException(404,'Brand tidak ditemukan.')
    return brand

@asynccontextmanager
async def lifespan(app):
    global _worker
    init_db();shutdown.clear()
    _worker=threading.Thread(target=worker_loop,daemon=True,name='studio-worker');_worker.start()
    yield
    shutdown.set()
    if _worker:_worker.join(timeout=3)

app=FastAPI(title='AI Content Studio',version='0.2.0',lifespan=lifespan,
    docs_url=None if APP_ENV=='production' else '/docs',
    redoc_url=None if APP_ENV=='production' else '/redoc',
    openapi_url=None if APP_ENV=='production' else '/openapi.json')
if APP_ENV=='production':
    app.add_middleware(TrustedHostMiddleware, allowed_hosts=[APP_DOMAIN,'localhost','127.0.0.1'])

@app.middleware('http')
async def origin_guard(request:Request, call_next):
    # Block cross-site cookie-authenticated writes. Normal browser requests send Origin.
    if APP_ENV=='production' and request.method in ('POST','PUT','PATCH','DELETE'):
        origin=request.headers.get('origin','').rstrip('/')
        if origin and origin != f'https://{APP_DOMAIN}':
            return JSONResponse({'detail':'Origin tidak diizinkan.'},status_code=403)
    return await call_next(request)

app.mount('/static',StaticFiles(directory=str(PUBLIC)),name='static')

@app.get('/',include_in_schema=False)
def home():return FileResponse(PUBLIC/'index.html')

@app.get('/api/health')
def health():return {'status':'ok','version':'0.1.0'}

@app.post('/api/login')
async def login(req:Request,response:Response):
    payload=await req.json()
    if not PASSWORD or not hmac.compare_digest(str(payload.get('password','')),PASSWORD):
        raise HTTPException(401,'Kata sandi salah atau belum dikonfigurasi.')
    exp=str(int(time.time()+60*60*24*3))
    sig=hmac.new(SESSION_SECRET.encode(),exp.encode(),hashlib.sha256).hexdigest()
    response.set_cookie('acs_session',f'{exp}:{sig}',max_age=60*60*24*3,secure=APP_ENV=='production',httponly=True,samesite='lax')
    return {'ok':True}

@app.post('/api/logout',dependencies=[Depends(check_auth)])
def logout(response:Response):
    response.delete_cookie('acs_session')
    return {'ok':True}

@app.get('/api/config',dependencies=[Depends(check_auth)])
def config():
    return {'ai_configured':ai_client.configured(),'text_model':ai_client.TEXT_MODEL,
            'image_model':ai_client.IMAGE_MODEL,'mode':APP_ENV,
            'local_templates':True,'worker_type':'SQLite persistent single-worker','version':'0.1.0'}

class BrandIn(BaseModel):
    name:str=Field(min_length=2,max_length=90)
    industry:str=Field(default='',max_length=180)
    audience:str=Field(default='',max_length=260)
    primary_color:str=Field(default='#7545C5',pattern=r'^#[0-9a-fA-F]{6}$')
    secondary_color:str=Field(default='#FAF7FF',pattern=r'^#[0-9a-fA-F]{6}$')
    caption_prefix:str=Field(default='',max_length=360)
    max_hashtags:int=Field(default=5,ge=0,le=15)

@app.get('/api/brands',dependencies=[Depends(check_auth)])
def brands():return query_all('SELECT * FROM brands ORDER BY created_at DESC')

@app.post('/api/brands',dependencies=[Depends(check_auth)])
def create_brand(item:BrandIn):
    b=uid()
    execute('''INSERT INTO brands(id,name,industry,audience,primary_color,secondary_color,caption_prefix,max_hashtags,created_at)
    VALUES(?,?,?,?,?,?,?,?,?)''',(b,item.name,item.industry,item.audience,item.primary_color,item.secondary_color,item.caption_prefix,item.max_hashtags,now()))
    return ensure_brand(b)

@app.put('/api/brands/{bid}',dependencies=[Depends(check_auth)])
def update_brand(bid:str,item:BrandIn):
    ensure_brand(bid)
    execute('''UPDATE brands SET name=?,industry=?,audience=?,primary_color=?,secondary_color=?,caption_prefix=?,max_hashtags=? WHERE id=?''',
      (item.name,item.industry,item.audience,item.primary_color,item.secondary_color,item.caption_prefix,item.max_hashtags,bid))
    return ensure_brand(bid)


def save_image(data:bytes,kind='reference'):
    if len(data)>12*1024*1024:raise HTTPException(413,'Gambar maksimal 12MB.')
    try:
        # Full decoding catches invalid files, decompression bombs, and fake extensions.
        Image.MAX_IMAGE_PIXELS=32_000_000
        im=Image.open(io.BytesIO(data)); im.load()
        if im.format not in ('PNG','JPEG','WEBP'):raise ValueError('format')
        im=ImageOps.exif_transpose(im).convert('RGBA' if kind=='logo' else 'RGB')
        im.thumbnail((1600,1800),Image.Resampling.LANCZOS)
    except Exception:raise HTTPException(400,'Hanya gambar JPG/PNG/WEBP yang valid.')
    name=f'{uid()}.png' if kind=='logo' else f'{uid()}.jpg'
    p=MEDIA/name
    im.save(p,'PNG' if kind=='logo' else 'JPEG',quality=90)
    return name,im

@app.post('/api/brands/{bid}/upload',dependencies=[Depends(check_auth)])
async def upload(bid:str,kind:Literal['logo','reference']=Form(...),file:UploadFile=File(...)):
    ensure_brand(bid)
    name,im=save_image(await file.read(),kind)
    column='logo_path' if kind=='logo' else 'reference_path'
    if kind=='reference':
        colors=graphics.palettes(im)
        execute('UPDATE brands SET reference_path=?,primary_color=?,secondary_color=? WHERE id=?',(name,*colors,bid))
    else:execute(f'UPDATE brands SET {column}=? WHERE id=?',(name,bid))
    return {'brand':ensure_brand(bid),'media_url':f'/media/{name}','detected_colors':graphics.palettes(im) if kind=='reference' else None}

@app.post('/api/brands/{bid}/analyze',dependencies=[Depends(check_auth)])
def analyze(bid:str):
    brand=ensure_brand(bid)
    if not brand.get('reference_path'):raise HTTPException(400,'Upload desain referensi terlebih dahulu.')
    if not ai_client.configured():raise HTTPException(422,'Vision AI memerlukan OPENAI_API_KEY. Warna dasar sudah dianalisis secara lokal.')
    p=MEDIA/brand['reference_path']
    try:analysis=ai_client.analyze_reference(p.read_bytes())
    except ai_client.ProviderError as e:raise HTTPException(502,str(e))
    notes=json.dumps(analysis,ensure_ascii=False)
    execute('UPDATE brands SET style_notes=? WHERE id=?',(notes,bid))
    return {'analysis':analysis,'brand':ensure_brand(bid)}

class BatchIn(BaseModel):
    brand_id:str
    count:int=Field(ge=1,le=90)
    theme:str=Field(default='',max_length=600)
    mode:Literal['ai','template']='ai'
    style_mode:Literal['strict','medium','creative']='medium'
    topics:str=Field(default='',max_length=25000)

@app.post('/api/batches',dependencies=[Depends(check_auth)])
def start_batch(item:BatchIn):
    ensure_brand(item.brand_id)
    if item.mode=='ai' and not ai_client.configured():
        raise HTTPException(422,'Generasi AI membutuhkan OPENAI_API_KEY. Untuk uji offline pilih Template Lokal.')
    if item.mode=='template':
        topics=[x.strip() for x in item.topics.splitlines() if x.strip()]
        if len(topics)!=item.count:raise HTTPException(422,f'Mode template membutuhkan tepat {item.count} judul, satu per baris (ditemukan {len(topics)}).')
        if len({x.casefold() for x in topics})!=len(topics):raise HTTPException(422,'Judul dalam mode template harus berbeda.')
        seen=[]
        for title in topics:
            if not unique_title(title,seen):
                raise HTTPException(422,f'Judul terlalu mirip dengan yang lain: "{title[:75]}". Buat variasi topik nyata, bukan hanya mengganti nomor.')
            seen.append(title)
    bid=uid()
    execute('''INSERT INTO batches(id,brand_id,requested,theme,mode,style_mode,manual_topics,status,created_at,updated_at)
      VALUES(?,?,?,?,?,?,?,'queued',?,?)''',(bid,item.brand_id,item.count,item.theme,item.mode,item.style_mode,item.topics,now(),now()))
    return batch_status(bid)

@app.get('/api/batches',dependencies=[Depends(check_auth)])
def batch_list():
    return [batch_status(b['id']) for b in query_all('SELECT id FROM batches ORDER BY created_at DESC LIMIT 100')]

def batch_status(bid):
    b=query_one('SELECT * FROM batches WHERE id=?',(bid,))
    if not b:raise HTTPException(404,'Batch tidak ditemukan.')
    rows=query_all('SELECT status,COUNT(*) AS n FROM contents WHERE batch_id=? GROUP BY status',(bid,))
    counts={x['status']:x['n'] for x in rows}
    b['completed']=counts.get('completed',0)
    b['failed']=counts.get('failed',0)
    b['progress']=round(100*b['completed']/b['requested'])
    b['planned']=sum(counts.values())
    return b

@app.get('/api/batches/{bid}',dependencies=[Depends(check_auth)])
def get_batch(bid:str):return batch_status(bid)

@app.post('/api/batches/{bid}/cancel',dependencies=[Depends(check_auth)])
def cancel_batch(bid:str):
    batch_status(bid)
    execute('UPDATE batches SET stop_requested=1,updated_at=? WHERE id=?',(now(),bid))
    return batch_status(bid)

@app.post('/api/batches/{bid}/retry',dependencies=[Depends(check_auth)])
def retry_batch(bid:str):
    b=batch_status(bid)
    if b['status'] not in ('failed','partial','cancelled'):
        raise HTTPException(409,'Batch masih berjalan atau sudah selesai.')
    execute("UPDATE contents SET status='queued',error=NULL WHERE batch_id=? AND status='failed'",(bid,))
    execute("UPDATE batches SET status='queued',error=NULL,stop_requested=0,updated_at=? WHERE id=?",(now(),bid))
    return batch_status(bid)

@app.get('/api/batches/{bid}/contents',dependencies=[Depends(check_auth)])
def batch_contents(bid:str):
    batch_status(bid)
    return query_all('SELECT * FROM contents WHERE batch_id=? ORDER BY position',(bid,))

@app.get('/api/contents/{cid}',dependencies=[Depends(check_auth)])
def get_content(cid:str):
    row=query_one('SELECT * FROM contents WHERE id=?',(cid,))
    if not row:raise HTTPException(404,'Konten tidak ditemukan.')
    return row

class ContentEdit(BaseModel):
    title:str=Field(min_length=3,max_length=160)
    subtitle:str=Field(default='',max_length=300)
    caption:str=Field(min_length=1,max_length=4000)

@app.put('/api/contents/{cid}',dependencies=[Depends(check_auth)])
def edit_content(cid:str,item:ContentEdit):
    content=get_content(cid)
    b=query_one('SELECT batches.* FROM batches JOIN contents ON batches.id=contents.batch_id WHERE contents.id=?',(cid,))
    brand=ensure_brand(b['brand_id'])
    visual=(MEDIA/content['asset_path']).read_bytes() if content['asset_path'] and (MEDIA/content['asset_path']).exists() else None
    try:
        result=graphics.render({'name':brand['name'],'primary':brand['primary_color'],'secondary':brand['secondary_color']},
        {'title':item.title,'subtitle':item.subtitle,'category':content['category']},visual_bytes=visual,
        logo_path=str(MEDIA/brand['logo_path']) if brand['logo_path'] else None)
    except ValueError as e:raise HTTPException(422,str(e))
    image_filename=content['image_path'] or f'{uid()}.jpg';(MEDIA/image_filename).write_bytes(result)
    execute('UPDATE contents SET title=?,subtitle=?,caption=?,image_path=?,status=? WHERE id=?',
    (item.title,item.subtitle,item.caption,image_filename,'completed',cid))
    return get_content(cid)

@app.get('/media/{name}',dependencies=[Depends(check_auth)])
def media(name:str):
    if not re.fullmatch(r'[0-9a-f]{32}\.(?:jpg|png)',name):raise HTTPException(404,'File tidak ditemukan')
    p=MEDIA/name
    if not p.exists():raise HTTPException(404,'File tidak ditemukan')
    return FileResponse(p,media_type='image/png' if name.endswith('.png') else 'image/jpeg')

@app.get('/api/contents/{cid}/download',dependencies=[Depends(check_auth)])
def download(cid:str):
    content=get_content(cid)
    if not content['image_path']:raise HTTPException(404,'Gambar belum tersedia.')
    return FileResponse(MEDIA/content['image_path'],media_type='image/jpeg',filename=f'content_{content["position"]:03d}.jpg')

@app.get('/api/batches/{bid}/download',dependencies=[Depends(check_auth)])
def download_zip(bid:str):
    b=batch_status(bid)
    cs=query_all("SELECT * FROM contents WHERE batch_id=? AND status='completed' ORDER BY position",(bid,))
    if not cs:raise HTTPException(409,'Belum ada konten selesai.')
    buff=io.BytesIO()
    with zipfile.ZipFile(buff,'w',zipfile.ZIP_DEFLATED,compresslevel=6) as z:
        table=io.StringIO();w=csv.writer(table);w.writerow(['position','title','category','caption','file'])
        allcaps=[]
        for c in cs:
            tag=f'{c["position"]:03d}'
            p=MEDIA/c['image_path']
            if not p.exists():continue
            z.write(p,f'images/content_{tag}.jpg')
            z.writestr(f'captions/caption_{tag}.txt',c['caption'])
            w.writerow([c['position'],c['title'],c['category'],c['caption'],f'images/content_{tag}.jpg'])
            allcaps.append(f'KONTEN {tag}: {c["title"]}\n\n{c["caption"]}\n\n'+'-'*40)
        z.writestr('content_overview.csv',table.getvalue().encode('utf-8-sig'))
        z.writestr('all_captions.txt','\n'.join(allcaps).encode('utf-8'))
        z.writestr('README.txt',f'AI Content Studio\nBatch: {bid}\nMode: {b["mode"]}\nComplete: {len(cs)}/{b["requested"]}\nFiles in images correspond to captions by number.\n')
    buff.seek(0)
    return StreamingResponse(buff,media_type='application/zip',headers={'Content-Disposition':f'attachment; filename="AI-Content-Studio-{bid[:8]}.zip"'})


def unique_title(title,existing):
    def terms(x):return set(re.findall(r'\w+',x.lower()))
    for old in existing:
        if title.strip().lower()==old.strip().lower():return False
        a,b=terms(title),terms(old)
        if a and b and len(a&b)/len(a|b)>0.82:return False
        if SequenceMatcher(None,title.lower(),old.lower()).ratio()>0.89:return False
    return True


def caption_for(item,brand,mode):
    prefix=(brand.get('caption_prefix') or '').strip()
    if mode=='template':
        text=f'{item["title"]} ✨\n\n{item["subtitle"]}\n\nMenurut kamu bagaimana? Bagikan pendapatmu di komentar! 👇'
    else:text=f'{item.get("caption_body",item["title"])}\n\n{item.get("cta","Bagaimana pendapatmu?")}'
    count=brand.get('max_hashtags',5)
    # Clean generated hashtags; brand rules cap globally.
    text=re.sub(r'(?<!\w)#[\w]+','',text).strip()
    tags=['#'+re.sub(r'[^A-Za-z0-9]','',k.title()) for k in re.findall(r'[\w]+',brand['name'])[:2]]
    tags += ['#Instagram',' #KontenKreatif','#BisnisIndonesia']
    tags=[x.strip() for x in tags if len(x.strip())>1][:count]
    return ('\n\n'.join([s for s in [prefix,text,' '.join(tags)] if s])).strip()


def ensure_content(batch,brand,item,index):
    row=query_one('SELECT * FROM contents WHERE batch_id=? AND position=?',(batch['id'],index))
    if row:return row
    title=item.get('title','').strip()
    other_titles=[x['title'] for x in query_all('SELECT title FROM contents WHERE batch_id=?',(batch['id'],))]
    if not unique_title(title,other_titles):raise ai_client.ProviderError(f'Ide terdeteksi terlalu mirip: {title}')
    # Scope: only exact and high-similarity titles, not universal semantic uniqueness.
    cap=caption_for(item,brand,batch['mode'])
    cid=uid()
    execute('''INSERT INTO contents(id,batch_id,position,title,subtitle,category,visual_prompt,caption,status,created_at)
      VALUES(?,?,?,?,?,?,?,?,?,?)''', (cid,batch['id'],index,title,item.get('subtitle',''),item.get('category','konten'),item.get('visual_prompt',''),cap,'queued',now()))
    return get_content(cid)


def plan_batch(batch,brand):
    if batch['mode']=='template':
        topics=[x.strip() for x in batch['manual_topics'].splitlines() if x.strip()]
        return [{'title':t,'subtitle':batch['theme'] or 'Ide menarik untuk komunitas dan bisnis kamu.','category':'TEMPLATE',
                 'visual_prompt':'','caption_body':'','cta':''} for t in topics]
    else:
        earlier=[c['title'] for c in query_all('SELECT contents.title FROM contents JOIN batches ON batches.id=contents.batch_id WHERE batches.brand_id=? ORDER BY contents.created_at DESC LIMIT 200',(brand['id'],))]
        approved=[]
        # Semantic-ish headline filtering plus regeneration of rejected suggestions.
        # This is NOT a full embedding-based semantic anti-duplicate engine.
        for attempt in range(4):
            missing=batch['requested']-len(approved)
            if not missing:break
            candidates=ai_client.generate_plan({'name':brand['name'],'industry':brand['industry'],'audience':brand['audience']},batch['theme'],missing,earlier+[x['title'] for x in approved])
            for item in candidates:
                if unique_title(item['title'],earlier+[x['title'] for x in approved]):
                    approved.append(item)
            if len(approved)==batch['requested']:break
        if len(approved)!=batch['requested']:
            raise ai_client.ProviderError(f'AI tidak berhasil membuat {batch["requested"]} judul yang cukup berbeda setelah 4 percobaan. Batch ini dihentikan untuk menghindari duplikasi.')
        return approved


def process_batch(batch):
    brand=ensure_brand(batch['brand_id'])
    current=query_all('SELECT * FROM contents WHERE batch_id=? ORDER BY position',(batch['id'],))
    if len(current)<batch['requested']:
        execute("UPDATE batches SET status='planning',updated_at=? WHERE id=?",(now(),batch['id']))
        items=plan_batch(batch,brand)
        for i,item in enumerate(items,1):
            if not query_one('SELECT id FROM contents WHERE batch_id=? AND position=?',(batch['id'],i)):
                ensure_content(batch,brand,item,i)
    execute("UPDATE batches SET status='generating',updated_at=? WHERE id=?",(now(),batch['id']))
    current=query_all('SELECT * FROM contents WHERE batch_id=? ORDER BY position',(batch['id'],))
    # Maintain a set of visual fingerprints for generated AI assets in this batch.
    fingerprints=[]
    if batch['mode']=='ai':
        for existing in current:
            if existing['asset_path'] and existing['status']=='completed':
                saved=MEDIA/existing['asset_path']
                if saved.exists():
                    try:fingerprints.append(graphics.perceptual_dhash(saved.read_bytes()))
                    except Exception:pass
    for item in current:
        if shutdown.is_set():return
        latest=query_one('SELECT stop_requested FROM batches WHERE id=?',(batch['id'],))
        if latest and latest['stop_requested']:
            execute("UPDATE batches SET status='cancelled',updated_at=? WHERE id=?",(now(),batch['id']))
            return
        if item['status']=='completed':continue
        if item['status']=='failed':continue  # retried explicitly by user
        try:
            execute("UPDATE contents SET status='generating',error=NULL WHERE id=?",(item['id'],))
            visual=None;filename=None
            if batch['mode']=='ai':
                if item['asset_path'] and (MEDIA/item['asset_path']).exists():
                    visual=(MEDIA/item['asset_path']).read_bytes()
                else:
                    # Maximum 3 attempts when asset looks virtually identical to another.
                    for attempt in range(3):
                        visual=ai_client.generate_visual(item['visual_prompt'],
                          {'name':brand['name'],'primary':brand['primary_color'],'secondary':brand['secondary_color']},
                          {'notes':brand['style_notes']},batch['style_mode'])
                        visual_hash=graphics.perceptual_dhash(visual)
                        if all(graphics.hash_distance(visual_hash,other)>4 for other in fingerprints):
                            break
                    else:
                        raise ai_client.ProviderError('Tiga aset visual tampak identik dengan konten lain; coba ulangi dengan konsep baru.')
                    filename=f'{uid()}.png'; (MEDIA/filename).write_bytes(visual)
                    execute('UPDATE contents SET asset_path=? WHERE id=?',(filename,item['id']))
                fingerprints.append(graphics.perceptual_dhash(visual))
            rendered=graphics.render({'name':brand['name'],'primary':brand['primary_color'],'secondary':brand['secondary_color']},
                item,visual_bytes=visual,logo_path=str(MEDIA/brand['logo_path']) if brand['logo_path'] else None,mode=batch['style_mode'])
            file_name=item['image_path'] or f'{uid()}.jpg';(MEDIA/file_name).write_bytes(rendered)
            execute("UPDATE contents SET status='completed',image_path=?,error=NULL WHERE id=?",(file_name,item['id']))
        except Exception as exc:
            execute("UPDATE contents SET status='failed',error=? WHERE id=?",(str(exc)[:400],item['id']))
        execute('UPDATE batches SET updated_at=? WHERE id=?',(now(),batch['id']))
    statuses=query_all('SELECT status,COUNT(*) n FROM contents WHERE batch_id=? GROUP BY status',(batch['id'],))
    sc={x['status']:x['n'] for x in statuses}
    status='completed' if sc.get('completed',0)==batch['requested'] else 'partial' if sc.get('completed',0) else 'failed'
    execute('UPDATE batches SET status=?,updated_at=? WHERE id=?',(status,now(),batch['id']))


def worker_loop():
    while not shutdown.is_set():
        todo=query_one("SELECT * FROM batches WHERE status IN ('queued','planning','generating') ORDER BY created_at LIMIT 1")
        if todo:
            try:process_batch(todo)
            except Exception as exc:
                print('Generation job error',str(exc),file=sys.stderr)
                execute('UPDATE batches SET status=?,error=?,updated_at=? WHERE id=?',('failed',str(exc)[:400],now(),todo['id']))
        else:shutdown.wait(0.35)
