U
    uie                     @   s  d dl Z d dlZd dlZd dlZd dlmZmZmZmZ d dl	m
Z
 d dlmZmZ d dlmZmZ d dlmZmZmZ d dlmZmZmZ d dlmZmZmZ ed	eZd
ddhZdddhZ e!dddZ"e#dddZ$e!e%dddZ&e!e'e'dddZ(e!dddZ)e!e!dd d!Z*d"d# Z+e!e!e!d$d%d&Z,ej-d'd(gd)d*d+d) Z.ej-d,d(gd-d*d.d- Z/ej-d/d(gd0d*d1d0 Z0ej-d2d(gd3d*d4d3 Z1ej-d5d(gd6d*d7d6 Z2dS )8    N)	Blueprintrequestjsonifycurrent_app)bulk_db)require_bulk_api_keyget_access_token_for_alias)parse_publish_at_iso_to_tscreator_can_post_now)safe_filenamevideo_duration_seconds_ffprobebuild_public_media_url)process_one_scheduled_jobresolve_access_token_for_jobrefresh_submitted_job_status)query_creator_infoupload_video_direct_postfetch_post_statusZbulk_apicompletefailed	cancelled	scheduled
submitting	submitted)
table_namec              	   C   sb   |  d| d |  }g }|D ]:}z||d  W q" tk
rZ   ||d  Y q"X q"|S )z
    Devuelve lista de columnas reales del schema para inserts robustos.
    Funciona con row_factory sqlite3.Row o tuplas.
    zPRAGMA table_info()name   )executefetchallappend	Exception)curr   rowscolsr r&   5/var/www/html/luxverbi-app/app/blueprints/bulk_api.py_table_columns   s    r(   )jobc                    s|   |   }t|d}t|  fdd| D }d| }ddgt| }d| d| d}||t|	  |S )	z^
    Inserta en bulk_jobs filtrando por columnas reales (evita 'N values for M columns').
    	bulk_jobsc                    s   i | ]\}}| kr||qS r&   r&   ).0kvcolsetr&   r'   
<dictcomp>3   s       z$_insert_bulk_job.<locals>.<dictcomp>z, ?zINSERT INTO bulk_jobs (z
) VALUES (r   )
cursorr(   setitemsjoinkeyslenr   tuplevalues)connr)   r"   r$   filteredcolnamesplaceholderssqlr&   r.   r'   _insert_bulk_job+   s    
r?   )errreturnc                    s   sdS t   dr"dS dr0dS dddddd	d
dddddddddddddg}t fdd|D rvdS dsdrdS ddddd d!d"g}tfd#d|D rd$S d$S )%u  
    Decide si el fallo merece reintento (requeue) en lugar de marcar el job como failed.

    Filosofía:
      - Reintentar todo lo que sea probablemente transitorio: 5xx, 429, internal_error,
        timeouts, errores de red, "cannot_post_now", etc.
      - Marcar failed solo para errores claramente permanentes (fichero inexistente, privacy
        inválida, duración excedida, etc.)
    Tcannot_post_nowZmissing_publish_idzhttp=429ztoo many requestsZ
rate_limitZthrottlezhttp=500zhttp=502zhttp=503zhttp=504Zinternal_errorzservice unavailablezbad gatewayzgateway timeouttimeoutz	timed outZreadtimeoutZconnecttimeout
connectionzconnection resetZtemporarilyz	try againc                 3   s   | ]}| kV  qd S Nr&   )r+   m)er&   r'   	<genexpr>[   s     z&_is_transient_error.<locals>.<genexpr>direct_post_failedcreator_info_failedZfile_not_found_on_serverduration_exceeds_maxinvalid_privacy_levelbranded_cannot_be_self_onlymissing_privacy_levelmissing_videomissing_filenamec                 3   s   | ]}  |V  qd S rE   )
startswith)r+   p)r@   r&   r'   rH   l   s     F)strlowerrQ   any)r@   Ztransient_markersZpermanent_prefixesr&   )rG   r@   r'   _is_transient_error<   sV    


               
	rV   )r@   attempts_nextrA   c                 C   sV   d}| p
d drdnd}t||dtd|d   }t|td	d
 }t|| S )u   
    Backoff exponencial con jitter, con cap. Evita martillear la API.

    - cannot_post_now: arranca más alto.
    - cap: 6h (el sistema sigue intentando, pero sin bucle agresivo).
    i`T   rB   <         r   r   g        g?)rQ   minmaxintrandomuniform)r@   rW   capbaseexpZjitterr&   r&   r'   _backoff_secondsr   s
    rd   )rA   c                 C   s   z| d pd  }W n tk
r,   d}Y nX |r<d| S z| d pHd  }W n tk
rh   d}Y nX |rd|dd  S dS )	z
    Clave para agrupar 'cooldown' por cuenta.
    Preferimos account_alias; si no existe (legacy), degradamos a un hash parcial del access_token.
    account_aliasrX   zalias:access_tokenztoken:N   unknown)stripr!   )rowaliastokr&   r&   r'   _job_account_key   s    


rm   )valuerA   c                 C   s   | pd   S )NrX   )ri   rT   )rn   r&   r&   r'   _normalize_job_status   s    ro   c                 C   s8   | sd S | d | d | d | d | d | d | d dS )	Nidstatuspublish_at_tsupdated_at_ts
publish_idlast_statuslast_fail_reason)job_idrq   rr   rs   rt   ru   rv   r&   )rj   r&   r&   r'   _serialize_job_match   s    rx   )db_pathre   original_filenamec          
      C   s   t | }| }|d||f | }|  d}d}|D ]:}t|d }	|	tkrb|dkrb|}q<|	tkr<|dkr<|}q<||t||dkt	|t	|dS )z
    Busca trabajos previos del mismo video para una cuenta concreta.

    Regla de negocio:
      - Un job activo (scheduled/submitting/submitted) bloquea nuevas altas.
      - Un job terminal (complete/failed/cancelled) NO bloquea re-subidas.
    a  
        SELECT id, status, publish_at_ts, updated_at_ts, publish_id, last_status, last_fail_reason
        FROM bulk_jobs
        WHERE account_alias=? AND original_filename=?
        ORDER BY
          CASE
            WHEN LOWER(COALESCE(status, '')) IN ('scheduled', 'submitting', 'submitted') THEN 0
            ELSE 1
          END,
          publish_at_ts DESC,
          updated_at_ts DESC
        Nrq   )re   rz   Zmatches_totalZcan_scheduleactive_matchlatest_terminal_match)
r   r2   r   r   closero   ACTIVE_JOB_STATUSESTERMINAL_JOB_STATUSESr7   rx   )
ry   re   rz   r:   r"   r#   r{   r|   rj   rq   r&   r&   r'   *_lookup_existing_jobs_by_original_filename   s.    r   z/api/bulk/publishPOSTapi_bulk_publish)methodsendpointc                  C   s  t   tjdpd } tjdp(d }| r~zt| }W n> tk
r| } z tddt|ddf W Y S d}~X Y nX |stdd	d
dfS dtj	krtddd
dfS tj	d }|j
stddd
dfS tjdpd }tjdpd }|stddd
dfS tjdddk}tjdddk}tjdddk}tjdddk}	tjdddk}
tjdddk}tjdddk}|r|dkrtddd
dfS zt|}W n@ tk
r } z tddt|ddf W Y S d}~X Y nX t|\}}|s&tdd|ddfS |dp4g }||krTtdd |d!dfS |d"d#krhd}|d$d#kr|d}|d%d#krd}| }| }| }tjd& }t|j
}tt  d'| }tj||}|| |d(}t|}t|tr\|d)kr\|d)kr\||kr\zt| W n tk
rD   Y nX tdd*||d+dfS t|tjd, tjd- tjd. d/}zVt||||||||
|d0|d1}|d2pi d3}td#||||| pdd4d5fW S  tk
r } z tdd6t|ddf W Y S d}~X Y nX dS )7u   
    Publish inmediato (server-to-server). Aquí SÍ es correcto generar video_url firmado,
    porque se publica en el momento (la ventana de TTL es suficiente).
    Soporta: account_alias (recomendado) o access_token.
    re   rX   rf   Falias_token_failedokerrordetails  NZ%missing_access_token_or_account_aliasr   r   videorO   rP   titleprivacy_levelrN   allow_comment1
allow_duetallow_stitchcommercial_toggle0brand_organic_togglebrand_content_toggleis_aigc	SELF_ONLYrM   rJ   rB   )r   r   reasonprivacy_level_optionsrL   )r   r   optionscomment_disabledTduet_disabledstitch_disabled
UPLOAD_DIR_max_video_post_duration_secr   rK   )r   r   durr]   PUBLIC_BASE_URLMEDIA_SIGNING_SECRETMEDIA_TOKEN_TTL_SECONDS)stored_filenamepublic_base_urlsigning_secretttl_secondsPULL_FROM_URL)rf   captionr   disable_commentdisable_duetdisable_stitchr   r   r   mode	video_urldatart   )r   rt   r   	init_respr   re      rI   )r   r   formgetri   r   r!   r   rS   filesfilenamer   r
   r   configr   r^   timeospathr5   saver   
isinstanceremover   r   )re   rf   rG   filer   r   r   r   r   r   r   r   r   creator_infocan_post_nowr   r   r   r   r   
upload_diroriginal_fnr   	save_pathmax_durr   r   r   rt   r&   r&   r'   r      s    .

.




*	
z/api/bulk/statusapi_bulk_statusc               
   C   s  t   tjdpd } tjdp(d }| r~zt| }W n> tk
r| } z tddt|ddf W Y S d}~X Y nX tjd	pd }|r|stdd
ddfS zt	||}td|ddfW S  tk
r } ztdt|ddf W Y S d}~X Y nX dS )zJ
    Status (server-to-server). Soporta account_alias o access_token.
    re   rX   rf   Fr   r   r   Nrt   Z"missing_access_token_or_publish_idr   T)r   r   r   )
r   r   r   r   ri   r   r!   r   rS   r   )re   rf   rG   rt   r   r&   r&   r'   r   R  s     .
z/api/bulk/find_existingapi_bulk_find_existingc                  C   s   t   tjdpd } tjdp4tjdp4d }| sPtddddfS |sftdd	ddfS t|}ttj	d
 | |d}tddi|dfS )u  
    Lookup server-to-server para deduplicaciÃ³n por alias + nombre de archivo original.

    El scheduler externo puede usar `can_schedule`:
      - True  => no hay job activo bloqueando.
      - False => ya existe uno en curso para ese video/cuenta.
    re   rX   rz   r   FZaccount_alias_requiredr   r   Zmissing_original_filenameBULK_DB_PATHre   rz   r   Tr   )
r   r   r   r   ri   r   r   r   r   r   )re   Zoriginal_filename_rawrz   lookupr&   r&   r'   r   m  s    	 z/api/bulk/scheduleapi_bulk_schedulec                  C   s  t   tjdpd } | s0tddddfS zt| }W n> tk
rz } z tddt|ddf W Y S d	}~X Y nX tjd
pd }|stddddfS zt	|}W n> tk
r } z tddt|ddf W Y S d	}~X Y nX dtj
krtddddfS tj
d }|js4tddddfS tjdpDd }tjdpZd }|sxtddddfS tjdpd  }|dkrtddddfS tjdddk}	tjdddk}
tjdddk}tjdddk}tjdddk}tjdddk}tjdddk}|rN|d krNtdd!ddfS tjd" }t|j}ttjd# | |d$}d	}|dkr|d% }n|d&kr|d% p|d' }|rt|d( }td)d)|tkrd*nd+|d, |d( |d- || |d.	d/fS t j}| d0| }tj||}|| tt }||||d1d2d	|| |||	rRd3nd2|
r^d3nd2|rjd3nd2|rvd3nd2|rd3nd2|rd3nd2|rd3nd2||dd	d	d	|d4}ttjd# }t|| |  |  td)|||| d5d/fS )6uM  
    Schedule (server-to-server). Para jobs programados, account_alias es requerido
    para que el cron/worker pueda refrescar tokens con refresh_token.

    Importante:
      - NO generamos video_url firmado aquí, porque caduca (TTL) antes de publish.
      - El worker debe generar video_url justo en el momento de publicar.
    re   rX   FZ)account_alias_required_for_scheduled_jobsr   r   r   r   N
publish_atmissing_publish_atZinvalid_publish_atr   rO   rP   r   r   rN   dedupe_modeactive_only>   allnoner   Zinvalid_dedupe_moder   r   r   r   r   r   r   r   r   r   rM   r   r   r   r{   r   r|   rq   TZalready_scheduledZalready_existsrw   rr   )	r   Zskippedr   rw   rq   rr   rz   re   r   r   r   r   r   r   )rp   created_at_tsrr   next_attempt_tsrq   attempts
last_errorrf   re   r   r   r   r   r   r   r   r   r   r   rz   r   rt   ru   rv   rs   )r   rw   rr   r   re   ) r   r   r   r   ri   r   r   r!   rS   r	   r   r   rT   r   r   r   r   ro   r~   uuiduuid4hexr   r   r5   r   r^   r   r   r?   commitr}   )re   rf   rG   r   rr   r   r   r   r   r   r   r   r   r   r   r   r   r   Zexisting_lookupZexisting_matchZexisting_statusrw   r   r   now_tsr)   r:   r&   r&   r'   r     s    
..









z/api/bulk/process_dueapi_bulk_process_duec                  C   s  t   ttjddpd} tt }ttjddp:d}tt	j
d }| }dddddd}|d|| f | }td	t| d
 g }g }t }	|D ]4}
t|
}||	kr||
 q|	| ||
 q|r.|dkr.|| }|D ]}|d|||d |f q|d  t|7  < |  |D ]}
|
d }|d||f |jdkr`q2|  z4t|
t	j
d t	j
d t	j
d t	j
d d\}}}W nB tk
r } z"d}dt|j dt| }W 5 d}~X Y nX |r0|}t|
}|d|||||f |jdkr$|d  d7  < |  q2|p8d}t|
d d }t|r|t|| }|d|||||f |jdkr|d  d7  < n0|d||||f |jdkr|d   d7  < |  q2|d!| f | }|D ]}
|
d }t|
\}}|sq|d"krB|d#kr&d$nd%}|d&|||||f n|d'||||f |jdkrr|d(  d7  < |  q|  t d)||d*d+fS ),u  
    Cron endpoint:
      - Auth: X-Api-Key
      - Procesa jobs due (scheduled) y refresca jobs submitted.

    Mejoras:
      - Reintento con backoff para errores transitorios (incl. 5xx/internal_error).
      - Cooldown por cuenta para no disparar 2 publicaciones seguidas en segundos.
      - Protección anti-atasco: si el worker lanza excepción, requeue en lugar de dejar "submitting".
    max_jobs10ZBULK_MIN_SECONDS_BETWEEN_POSTSZ180r   r   )scheduled_submittedscheduled_failedscheduled_requeuedsubmitted_updatedscheduled_deferred_by_cooldownz
        SELECT * FROM bulk_jobs
        WHERE status='scheduled' AND next_attempt_ts <= ?
        ORDER BY publish_at_ts ASC
        LIMIT ?
        zDEBUG: Found z	 due jobsz
                UPDATE bulk_jobs
                SET next_attempt_ts=?, updated_at_ts=?
                WHERE id=? AND status='scheduled' AND next_attempt_ts <= ?
                rp   r   z
            UPDATE bulk_jobs
            SET status='submitting', updated_at_ts=?
            WHERE id=? AND status='scheduled'
            r   r   r   r   r   )r   r   r   r   Fzworker_exception: z: Na  
                UPDATE bulk_jobs
                SET status='submitted',
                    publish_id=?,
                    attempts=attempts+1,
                    last_error=NULL,
                    last_status=NULL,
                    last_fail_reason=NULL,
                    access_token=?,
                    next_attempt_ts=?,
                    updated_at_ts=?
                WHERE id=? AND status='submitting'
                r   Zunknown_errorr   a  
                UPDATE bulk_jobs
                SET status='scheduled',
                    attempts=?,
                    last_error=?,
                    next_attempt_ts=?,
                    updated_at_ts=?
                WHERE id=? AND status='submitting'
                r   z
                UPDATE bulk_jobs
                SET status='failed',
                    attempts=?,
                    last_error=?,
                    updated_at_ts=?
                WHERE id=? AND status='submitting'
                r   z}
        SELECT * FROM bulk_jobs
        WHERE status='submitted'
        ORDER BY updated_at_ts ASC
        LIMIT ?
        )PUBLISH_COMPLETEZFAILEDZSEND_TO_USER_INBOXr   r   r   z
                UPDATE bulk_jobs
                SET status=?,
                    last_status=?,
                    last_fail_reason=?,
                    updated_at_ts=?
                WHERE id=? AND status='submitted'
                z
                UPDATE bulk_jobs
                SET last_status=?,
                    last_fail_reason=?,
                    updated_at_ts=?
                WHERE id=? AND status='submitted'
                r   T)r   	processedr   r   )!r   r^   r   r   r   r   r   environr   r   r   r2   r   r   printr7   r3   rm   r    addr   rowcountr   r!   type__name__rS   r   rV   rd   r   r}   r   )r   r   Z
cooldown_sr:   r"   r   ZdueZpickeddeferredseenrj   keyZbump_tsr%   rw   r   msgZ
_init_resprG   rt   Zused_access_tokenr@   rW   Znext_tssubsrq   Zfail_reasonZ
new_statusr&   r&   r'   r     s    		


,

	
	

)3r   r   r   r_   flaskr   r   r   r   db.bulkr   services.authr   r   services.creatorr	   r
   services.mediar   r   r   Zservices.bulk_workerr   r   r   tiktok_clientr   r   r   r   bpr   r~   rS   r(   dictr?   boolrV   r^   rd   rm   ro   rx   r   router   r   r   r   r   r&   r&   r&   r'   <module>   s@   


6/
}


 