
    4 j1                     l    d dl Z d dlZd dlZd dlmZ d dlZd dlmZm	Z	m
Z
mZ d dlmZmZ  G d d      Zy)    N)partial)DictOptionalListTuple)NonPickledSyncManagerformat_secondsc                       e Zd ZdZdej
                  j                  dededdfdZ	deddfd	Z
d
edeeeeef         fdZd
ededdfdZd
eddfdZ	 dd
ededeeeeef         dedef
dZdefdZdefdZy)WorkerInsightsz
    Worker insights class for profiling the worker start up time, waiting time and working time. When worker init and
    exit functions are provided it will time those as well.
    ctxn_jobsuse_dillreturnNc                     || _         || _        || _        d| _        d| _        d| _        d| _        d| _        d| _        d| _	        d| _
        d| _        d| _        d| _        y)z
        Parameter class for worker insights.

        :param ctx: Multiprocessing context
        :param n_jobs: Number of workers
        :param use_dill: Whether dill is used as serialization library
        FN)r   r   r   insights_enabledinsights_managerinsights_manager_lockworker_start_up_timeworker_init_timeworker_n_completed_tasksworker_waiting_timeworker_working_timeworker_exit_timemax_task_durationmax_task_args)selfr   r   r   s       A/home/agent/.local/lib/python3.12/site-packages/mpire/insights.py__init__zWorkerInsights.__init__   s       !& !%%)" %)! !% )-% $(  $(  !% "& "    enable_insightsc                    |r,t        | j                        | _        | j                  j                          | j                  j                         | _        | j                  j                  t        j                  | j                  d      | _        | j                  j                  t        j                  | j                  d      | _        | j                  j                  t        j                  | j                  d      | _        | j                  j                  t        j                  | j                  d      | _        | j                  j                  t        j                  | j                  d      | _        | j                  j                  t        j                  | j                  d      | _        | j                  j                  t        j                  | j                  dz  d      | _        | j                  j'                  dg| j                  z  dz        | _        || _        yd| _        d| _        d| _        d| _        d| _        d| _        d| _        d| _        d| _        d| _        || _        y)zs
        Resets the insights containers

        :param enable_insights: Whether to enable worker insights
        F)lock    N)r   r   r   startr   Lockr   Arrayctypesc_doubler   r   r   c_intr   r   r   r   r   listr   r   )r   r    s     r   reset_insightszWorkerInsights.reset_insights=   s     %:$--$HD!!!''))-D&(,vZ_(`D%$(HHNN6??DKKV[N$\D!,0HHNN6<<[`N,aD)'+xx~~foot{{Y^~'_D$'+xx~~foot{{Y^~'_D$$(HHNN6??DKKV[N$\D!%)XX^^FOOT[[ST_[`^%aD"!%!6!6!;!;RD4;;<NQR<R!SD !0 %)D!)-D&(,D%$(D!,0D)'+D$'+D$$(D!%)D"!%D /r   	worker_idc           
          | j                   re| j                  5  | j                  At        t	        | j                  |dz  |dz   dz   | j
                  |dz  |dz   dz               ndcddd       S y# 1 sw Y   yxY w)z`
        Initialize insights for a specific worker

        :param worker_id: worker ID
        Nr#      )r   r   r   r+   zipr   r   r-   s     r   get_max_task_duration_listz)WorkerInsights.get_max_task_duration_list`   s        ++ I  11= S!7!7	AyST}XYFY!Z!%!3!3IM9q=TUBU!VX YCGI I !I Is   AA33A<
start_timec                 f    | j                   r%t        j                         |z
  | j                  |<   yy)zp
        Update start up time

        :param worker_id: Worker ID
        :param start_time: Timestamp
        N)r   timer   )r   r-   r3   s      r   update_start_up_timez#WorkerInsights.update_start_up_timen   s-       3799;3KD%%i0 !r   c                 L    | j                   r| j                  |xx   dz  cc<   yy)zn
        Increment the number of completed tasks for this worker

        :param worker_id: Worker ID
        r/   N)r   r   r1   s     r   update_n_completed_tasksz'WorkerInsights.update_n_completed_tasksx   s(       )))494 !r   max_task_duration_last_updatedmax_task_duration_listforce_updatec                    t        j                          }| j                  r\|s||z
  dkD  rRt        | \  }}|| j                  |dz  |dz   dz   | j                  5  || j
                  |dz  |dz   dz   ddd       |}|S # 1 sw Y   xY w)a  
        Update synced containers with new top 5 max task duration + args. Updates every 2 seconds.

        :param worker_id: Worker ID
        :param max_task_duration_last_updated: Last updated timestamp
        :param max_task_duration_list: Local worker insights container that holds (task duration, task args) tuples,
            sorted for heapq
        :param force_update: Whether to force the update
        :return: Last updated timestamp
           r#   r/   N)r5   r   r0   r   r   r   )r   r-   r9   r:   r;   nowtask_durations	task_argss           r   update_task_insightsz#WorkerInsights.update_task_insights   s     iik  ls=[7[_`6`(+-C(D%NIJXD""9q=IMQ3FG++ TJS""9q=IMQ3FGT-0*--	T Ts   A??Bc                    d }d }| j                   si S t        t        d      } || j                        dd ddd   }g g }}|D ]k  }| j                  |   dk(  r nW| j                  |   d	k(  r*|j                   || j                  |                |j                  | j                  |          m t        | j                        }t        | j                        }	t        | j                        }
t        | j                        }t        | j                        }||	z   |
z   |z   |z   }t        t        | j                        t        t        || j                              t        t        || j                              t        t        || j                              t        t        || j                              t        t        || j                               ||       ||	       ||
       ||       ||      ||
      } ||      |d<   d|fd|	fd|
fd|fd|ffD ]H  \  }} |t!        | d| d            \  }}||dz   z  || d<    ||      || d<    ||      || d<   J |S )zt
        Creates insights from the raw insight data

        :return: dictionary containing worker insights
        c                 T    t        t        t        |             | j                        S )z
            argsort, as to not be dependent on numpy, by
            https://stackoverflow.com/questions/3382352/equivalent-of-numpy-argsort-in-basic-python/3382369#3382369
            )key)sortedrangelen__getitem__)seqs    r   argsortz,WorkerInsights.get_insights.<locals>.argsort   s    
 %C/s??r   c                     t        |       t        |       z  t        fd| D              t        |       z  }t        j                  |      }|fS )za
            Calculates mean and standard deviation, as to not be dependent on numpy
            c              3   <   K   | ]  }t        |z
  d         yw)r=   N)pow).0x_means     r   	<genexpr>z@WorkerInsights.get_insights.<locals>.mean_std.<locals>.<genexpr>   s     6Qs1u9a(6s   )sumrG   mathsqrt)rI   _var_stdrP   s      @r   mean_stdz-WorkerInsights.get_insights.<locals>.mean_std   sG     Hs3x'E6#66SAD99T?D$;r   T)with_millisecondsNr   r$   )n_completed_tasksstart_up_time	init_timewaiting_timeworking_time	exit_timetotal_start_up_timetotal_init_timetotal_waiting_timetotal_working_timetotal_exit_timetop_5_max_task_durationstop_5_max_task_args
total_timestart_upinitwaitingworkingexitworker__timeg:0yE>_ratio
_time_mean	_time_std)r   r   r	   r   r   appendrR   r   r   r   r   r   dictr+   r   mapgetattr)r   rJ   rW   format_seconds_func
sorted_idxrf   rg   idxra   rb   rc   rd   re   rh   insightsparttotalmeanstds                      r   get_insightszWorkerInsights.get_insights   s   	@	 $$I%nM T334RS9$B$?
8:B"5  	@C%%c*a/!!#&",$++,?@V@VWZ@[,\]&&t'9'9#'>?	@ "$";";<d334 !9!9: !9!9:d334(?:=OORddgvv
$t/L/L*M&*3/BDD]D]+^&_"&s+>@U@U'V"W%)#.A4C[C[*\%]%)#.A4C[C[*\%]"&s+>@U@U'V"W,?@S,T(;O(L+>?Q+R+>?Q+R(;O(L1I,?A "5Z!@ ()<=#_5&(:;&(:;#_5	7 	DKD%
 !e/D!EFID#(-d1B(CHvV_%,?,EHvZ()+>s+CHvY'(	D r   c                 6   | j                   sy| j                         }dddt        |d          g}dD ]P  }|j                  d|j	                  dd	       d
|d| d    d|| d    d|| d    d|| d   dz  dd       R |d   dk  r|j                  ddg       |j                  g d       t        | j                        D ]i  }d| d|d   |    g}dD ]2  }|j                  |j	                  dd	       d|| d   |    d       4 |j                  dj                  |             k |j                  g d        t        t        |d!   |d"         d#$      D ]!  \  }\  }}|j                  | d%| d|        # d&j                  |      S )'zs
        Formats the worker insights_str and returns a string

        :return: worker insights_str string
        zSNo profiling stats available. Try to run a function first with insights enabled ...zWorkerPool insights-------------------z!Total number of tasks completed: r[   )ri   rj   rk   rl   rm   zTotal _ z time: total_ro   z	s (mean: rq   z, std: rr   z	, ratio: rp   g      Y@z.2fz%)working_ratiog?r$   z+Efficiency warning: working ratio is < 80%!)r$   zStats per workerz----------------zWorker zTasks completed: z: sz - )r$   zTop 5 longest tasksr   rf   rg   r/   )r%   z. Time: 
)r   r   rR   rs   replaceextendrF   r   join	enumerater0   )	r   rz   insights_strr{   r-   
worker_strtask_idxdurationargss	            r   get_insights_stringz"WorkerInsights.get_insights_string   sK    $$h$$&--;CI\@]<^;_`b
 G 	TD&c3)?(@SYZ^Y__dQeHfGg h))1TF*2E)F(Gwx[_Z``iXjOkNl m**2dV6?*Cd*J3)Or!S T	T O$s*!N!P Q 	 1 	2 t{{+ 	8I#I;/-h7J.KI.V-WXZJJ g!!T\\#s%;$<Bx4&PU?WXa?b>ccd"efg

: 67	8 	 4 	5 +4CA[8\8@AV8W5Y`a+c 	J&H&x8*HXJc$ HI	J yy&&r   )F)__name__
__module____qualname____doc__multiprocessingcontextBaseContextintboolr   r,   r   r   r   floatstrr2   r6   r8   rA   r   r   r    r   r   r   r      s   
)"O33?? )" )"X\ )"ae )"V!0d !0t !0FIC IHT%PUWZPZJ[E\<] ILc Lu L L:# :$ : 38.c .SX .5=d5PSCT>U5V.+/.<A..Id IV,'S ,'r   r   )r(   rS   multiprocessing.contextr   	functoolsr   r5   typingr   r   r   r   mpire.utilsr   r	   r   r   r   r   <module>r      s)         . . =D' D'r   