'i(tdZddlZddlmZmZmZmZddlmZddlm Z ddl m Z ddl m Z Gdd ZdS) zs Dependency Manager for Python Scheduler Handles job chaining and execution order with configurable error handling N)ListDictAnyOptional)datetime)logger)SchedulerDatabase) JobExecutorc eZdZdZdedefdZddedede ee ffd Z dd ed ed e defdZ ddedeedededef dZddededee ee ffdZdede ee ffdZdS)DependencyManagerz-Manages job dependencies and execution chainsdbexecutorcJ||_||_tjddS)NzDependencyManager initialized)r rrinfo)selfr rs 0c:\PYTHON\Scheduler\daemon\dependency_manager.py__init__zDependencyManager.__init__s'   344444 schedulergroup_id triggered_byreturnc |j|}|s!tjd|ddd|ddS|j|}|s#tjd|dddddStjd |dd t|d |d }d }g}tj }t|dD]\} } | d} | d} tjd| dt|d| d| d |j | d|d||} | sltjd| d| | | d ddd|dkrtjdn|dkrtjd n͌|| | }|j| }| | | | ||d!|d"d#|d$vrRtjd%| d&||dkrtjd'n |dkrtjd n| }tj |z }t%d(|D}t'd)|D}tjd|dd*|d+d,t)d-|Ddt||||dt|t|t)d.|Dt)d/|D||d0 S)1z Execute a job group with dependency chain Args: group_id: Job group ID to execute triggered_by: Who/what triggered execution Returns: Dictionary with execution results z Job group not foundF)successerrorz Job group 'namez' has no memberszNo jobs in groupz!Starting execution of job group 'z' (z jobs)error_handlingNjob_idjob_namezExecuting job /z: 'z' (ID )zgroup:)r rrparent_execution_idzFailed to queue job ''failedzFailed to queue)r r! execution_idstatusr stop_on_errorz-Stopping group execution due to queue failureskip_remainingz Skipping remaining jobs in groupduration_seconds exit_code)r r!r'r(durationr,) completedJob 'z' ended with status: z%Stopping group execution due to errorc3.K|]}|ddkVdSr(r.N.0rs r z6DependencyManager.execute_job_group..|s+RR1AhK;6RRRRRRrc3*K|]}|ddvVdS)r(r&timeoutNr2r3s rr6z6DependencyManager.execute_job_group..}s,YY!8(==YYYYYYrz' finished in z.1fzs - Completed: c32K|]}|ddkdVdSr(r.rNr2r3s rr6z6DependencyManager.execute_job_group..s1$`$`1Qx[T_E_E_QE_E_E_E_$`$`rc32K|]}|ddkdVdSr;r2r3s rr6z6DependencyManager.execute_job_group..s1!]!]!H+Q\B\B\!B\B\B\B\!]!]rc3.K|]}|ddv dVdS)r(r8rNr2r3s rr6z6DependencyManager.execute_job_group..s1ddQq{Nc?c?cq?c?c?c?cddr) rr group_name total_jobs executed_jobscompleted_jobs failed_jobsr+ executions)r get_job_grouprrget_job_group_memberswarningrlenrnow enumerater execute_jobappend_wait_for_completion get_executionget total_secondsallanysum)rrrgroupmembersrr$execution_resultsgroup_start_timeimemberr r!r' final_status executiongroup_duration all_completed any_faileds rexecute_job_groupz#DependencyManager.execute_job_groups%%h//  L:h::: ; ; ; :h:::  '//99  NHvHHH I I I +   ^f ^^#g,,^^^___/0"#<>>"7A..= /= /IAvH%Fj)H KWWWS\\WWhWWfWWW X X X =445eFm55!$7 5L   @X@@@AAA!(($ ($(&. **"_44K OPPPE#'777K BCCCE 44\8LLL--l;;I  $ $ $ ,&%MM*<==&]];77 &&   =00 RXRRLRRSSS!_44K GHHHE#'777K BCCCE#/  #,..+;;JJLLRR@QRRRRR YYGXYYYYY  |%-||~W||!$$`$`0A$`$`$`!`!`||cfgxcycy|| } } }% -g,, !233!!]!]->!]!]!]]]dd*;ddddd .+   r?r'r! poll_intervalctjd|d|d |j|}|dvrtjd|d||St j|O)a: Wait for job execution to complete Args: execution_id: Execution ID to monitor job_name: Job name for logging poll_interval: How often to check status (seconds) Returns: Final status ('completed', 'failed', 'timeout', 'cancelled') zWaiting for execution z ('z') to completeT)r.r&r9 cancelledz Execution z finished with status: )rdebugr get_execution_statustimesleep)rr'r!r`r(s rrLz&DependencyManager._wait_for_completions  WlWWxWWWXXX &W11,??FHHH W,WWvWWXXX J} % % % &rNr)rjob_ids descriptionrcZ|D]1}|j|}|std|d2|j|||}t |dD]!\}}|j|||"t jd|dt|d|S)aF Create a job group (chain) from a list of job IDs Args: name: Group name job_ids: List of job IDs in execution order description: Group description error_handling: 'stop_on_error', 'continue', or 'skip_remaining' Returns: Group ID Job ID rrzCreated job chain 'z' with z jobs) r get_job ValueErrorcreate_job_grouprIadd_job_to_grouprrrG) rrrgrhrr jobrorders rcreate_job_chainz"DependencyManager.create_job_chains  ? ?F'//&))C ? !=6!=!=!=>>> ?7++D+~NN'w22 > >ME6 G $ $Xvu = = = = J$JJs7||JJJKKKr limitc |j5}|d||f}d|D}dddn #1swxYwYi}|D]6}|ddd}||vrg||<|||7g} t |dD]\}} t | d } | d } | d } d }| D]}|d r ||d z }td | D}| | d| d|t| || d| S)z Get execution summary for a job group Args: group_id: Job group ID limit: Number of recent executions to retrieve Returns: List of execution summaries a  SELECT e.id, e.job_id, j.name as job_name, e.status, e.start_time, e.end_time, e.duration_seconds, e.exit_code, e.triggered_by FROM job_executions e JOIN jobs j ON e.job_id = j.id WHERE e.group_id = ? ORDER BY e.start_time DESC LIMIT ? c,g|]}t|Sr2)dict)r4rows r zADependencyManager.get_group_execution_summary..sAAA$s))AAArN start_timeT)reversec|dS)Nryr2)xs rz?DependencyManager.get_group_execution_summary..s 1\?r)keyrr+c3.K|]}|ddkVdSr1r2)r4exs rr6z@DependencyManager.get_group_execution_summary..s+YY8 ;YYYYYYrend_time)ryrtotal_duration jobs_executedr\rC) r _get_connectionexecutefetchallrKsorteditemsrPrG)rrrsconncursorrCbatchesr start_key summaries batch_execsbatch_execs_sorted first_exec last_execrr\s rget_group_execution_summaryz-DependencyManager.get_group_execution_summarysW $ $ & & B$\\# E"!$$F$BAv/@/@AAAJ' B B B B B B B B B B B B B B B, * *B<("-I''%' " I  % %b ) ) ) ) &,W]]__d&K&K&K   "I{!' 9R9R!S!S!S +A.J*2.IN( = =()="b);&<@@C@@TRVWZ\_W_R`Ma@@@@D% 3% 4S>% % % % % % rr )rretypingrrrrrlogurur core.databaser daemon.job_executorr r r2rrrs  ,,,,,,,,,,,,++++++++++++] ] ] ] ] ] ] ] ] ] r