o
    ;j\                     @  s  d Z ddlmZ ddlZddlZddlmZ ddlZddlmZ ddl	m
Z
 ejejejedZejde eejed dd	lmZmZmZmZmZmZ dd
lmZ ddlmZ ddlmZ ddlm Z  G dd dZ!d7ddZ"d8ddZ#d9ddZ$d:ddd;d#d$Z%e&d%krddl'Z'e'j(d&d'Z)e)j*d(e+d)d*d+ e)j*d,e, -d-d+ e)j*d.ddd/gd0 e)j*d1e.dd2 e)/ Z0e1d3e0j2 d4e0j3 d5e0j4 d6 e%e0j2e0j3e0j4e0j5d dS dS )<zfSync calls from Mcube source DB into raw_calls via shared ingest core (same pipeline as webhook/poll).    )annotationsN)datetime)load_dotenv)
DictCursorzdashboard-backendz.env)check_min_duration_for_ingesteffective_min_duration_sgroup_allowedmaybe_queue_after_ingestnormalize_sync_callupsert_raw_call)Config)DatabaseHandler)purge_unprocessed_if_below_min)evaluate_usage_allocationc                   @  s&   e Zd Zdd ZdddZdd ZdS )	_ConfigWrapperc                 C  s
   || _ d S N)_cfg)selfcfg r   sync_calls.py__init__!      
z_ConfigWrapper.__init__Nc                 C  s&   t | j|rt| j||S t||S r   )hasattrr   getattrosgetenv)r   keydefaultr   r   r   get$   s   &z_ConfigWrapper.getc                 C  s
   |  |S r   )r   )r   r   r   r   r   __getitem__'   r   z_ConfigWrapper.__getitem__r   )__name__
__module____qualname__r   r   r    r   r   r   r   r       s    
r   returndictc                
   C  s^   t dt ddtt dt ddt dt dd	t d
t ddt ddddS )NSYNC_SOURCE_DB_HOSTDB_HOSTz	127.0.0.1SYNC_SOURCE_DB_PORTDB_PORT3306SYNC_SOURCE_DB_USERDB_USERadminSYNC_SOURCE_DB_PASSWORDDB_PASSWORD SYNC_SOURCE_DB_NAME	mcube_cl1utf8mb4)hostportuserpassworddatabasecharset)r   r   intr   r   r   r   _source_config+   s   
r;   
db_handlerr   c              	   C  sF   | j dt| j dpd| j d| j d| j ddtdd	S )
Nr'   r)   i  r,   r/   DB_NAMEr3   T)r4   r5   r6   r7   r8   r9   cursorclass
autocommit)configr   r:   r   )r<   r   r   r   _dest_config6   s   



rA   bidstrc                 C  sX   |    | |p
i }tdt|dpddkr | ||d< | |d|d< |S )Nr   min_call_duration_smin_call_duration_effective_atallowed_groupnames_allowed_groupnames)%ensure_business_pipeline_config_tableget_pipeline_configmaxr:   r    ensure_min_duration_effective_at_decode_allowed_groupnames)r<   rB   r   r   r   r   _load_bid_cfgC   s   rM   callhistoryd   )	max_callsdate_strsource_tablerP   r:   c             
   C  s~  t |  } ttt }t|| }t|}|d}tj	d i t
 dti}tj	d i t|}	d}
d}d}zz| }|	 }|  d| }|d| d||f | pZg }tdt| d| d	| d
 |D ]}t|| d}|d|}|d }|d dkr|d7 }qmt||||d |dkd\}}}|rt|| |||| |d7 }qm|durtdtt||d< t|dpd|\}}|s|d7 }qmt|| }|d rtd|  d  nt|| | |
d7 }
t|| ||||dr|d7 }qmtd|
 d| d| d |
W W |  |	  S  ty5 } ztd|  W Y d}~W |  |	  dS d}~ww |  |	  w )!z;Sync ANSWER calls for one date through unified ingest core.rE   r>   r   _z
            SELECT
                callid, bid, agentname, groupname, starttime, endtime,
                dialstatus, direction, filename, emp_phone, clicktocalldid
            FROM `z`
            WHERE DATE(starttime) = %s
              AND dialstatus = 'ANSWER'
            ORDER BY starttime ASC
            LIMIT %s
            zFound z ANSWER calls for z (limit ))rB   _duration_probe_rowcallidcall_statusANSWER   fileurl)min_duration_seffective_atrecording_urlprobe_audioNduration_seconds	groupnamer0   blockedzUsage limit exhausted for BID z; stopping sync.)r<   r[   r\   zSync complete: z ingested, z queued for STT, z skippedzError syncing calls: r   )rC   stripr   r   r   rM   r   r   pymysqlconnectr;   r   rA   cursorexecutefetchallprintlenr
   popr   r   rJ   r:   roundr   r   r   r	   close	Exception)rB   rQ   rR   rP   r<   r   r[   r\   source_conn	dest_conn	processedqueuedskipped
source_curdest_curtablecallscall
normalized	probe_rowcall_idskip_min
min_reasonprobed_audioallowedrS   usageexcr   r   r   
sync_callsL   s   


 


r   __main__z(Sync Mcube calls via unified ingest core)descriptionz--bidSYNC_BID7987)r   z--datez%Y-%m-%dz--tablecallarchive)r   choicesz--limit)typer   zSyncing BID z on z from z...)r$   r%   )r<   r   r$   r%   )r<   r   rB   rC   r$   r%   )rN   )
rB   rC   rQ   rC   rR   rC   rP   r:   r$   r:   )6__doc__
__future__r   r   sysr   rc   dotenvr   pymysql.cursorsr   pathjoindirnameabspath__file___BACKENDinsertcall_ingest_corer   r   r   r	   r
   r   r@   r   r<   r   min_duration_utilr   usage_allocation_utilr   r   r;   rA   rM   r   r!   argparseArgumentParserparseradd_argumentr   nowstrftimer:   
parse_argsargsrh   rB   dateru   limitr   r   r   r   <module>   s@    


	]"