o 'i(@sfdZddlZddlmZmZmZmZddlmZddlm Z ddl m Z ddl m Z Gdd d ZdS) zs Dependency Manager for Python Scheduler Handles job chaining and execution order with configurable error handling N)ListDictAnyOptional)datetime)logger)SchedulerDatabase) JobExecutorc @seZdZdZdedefddZd deded e ee ffd d Z d!d edede d efddZ  d"dedeededed ef ddZd#deded ee ee ffddZded e ee ffddZdS)$DependencyManagerz-Manages job dependencies and execution chainsdbexecutorcCs||_||_tddS)NzDependencyManager initialized)r r rinfo)selfr r r0c:\PYTHON\Scheduler\daemon\dependency_manager.py__init__szDependencyManager.__init__ schedulergroup_id triggered_byreturnc Cs|j|}|std|ddd|ddS|j|}|s2td|dddddStd |dd t|d |d }d }g}t }t |dD]\} } | d} | d} td| dt|d| d| d |j j | d|d||d} | std| d| | | d ddd|dkrtdnR|dkrtd nGqS|| | }|j| }| | | | ||d!|d"d#|d$vrtd%| d&||dkrtd'n|dkrtd n| }qSt |}td(d)|D}td*d)|D}td|dd+|d,d-td.d)|Ddt||||dt|t|td/d)|Dtd0d)|D||d1 S)2z 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:)rrrparent_execution_idzFailed to queue job ''failedzFailed to queue)rr execution_idstatusr stop_on_errorz-Stopping group execution due to queue failureskip_remainingz Skipping remaining jobs in groupduration_seconds exit_code)rrr#r$durationr() completedJob 'z' ended with status: z%Stopping group execution due to errorcs|] }|ddkVqdSr$r*Nr.0rrrr |z6DependencyManager.execute_job_group..css|] }|ddvVqdS)r$r"timeoutNrr.rrrr1}r2z' finished in z.1fzs - Completed: cs |] }|ddkrdVqdSr$r*rNrr.rrrr1csr5r6rr.rrrr1r7css |] }|ddvrdVqdS)r$r3rNrr.rrrr1r7) rr group_name total_jobs executed_jobscompleted_jobs failed_jobsr' executions)r get_job_grouprrget_job_group_memberswarningr lenrnow enumerater execute_jobappend_wait_for_completion get_executionget total_secondsallanysum)rrrgroupmembersrr Zexecution_resultsZgroup_start_timeimemberrrr#Z final_status executionZgroup_duration all_completedZ any_failedrrrexecute_job_groups    (        z#DependencyManager.execute_job_group?r#r poll_intervalcCsTtd|d|d |j|}|dvr$td|d||St|q )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"r4 cancelledz Execution z finished with status: )rdebugr get_execution_statustimesleep)rr#rrUr$rrrrFs   z&DependencyManager._wait_for_completionNr%rjob_ids descriptionrc Cs~|D]}|j|}|std|dq|j|||}t|dD] \}}|j|||q"td|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_grouprCadd_job_to_grouprr rA) rrr[r\rrjobrorderrrrcreate_job_chains z"DependencyManager.create_job_chain limitc Cs"|j}|d||f}dd|D}Wdn1s!wYi}|D]}|ddd}||vrszADependencyManager.get_group_execution_summary..N start_timeT)reversecSs|dS)Nrjr)xrrrsz?DependencyManager.get_group_execution_summary..)keyrr'csr,r-r)r/exrrrr1r2z@DependencyManager.get_group_execution_summary..end_time)rjrrtotal_durationZ jobs_executedrRr=) r _get_connectionexecutefetchallrEsorteditemsrJrA)rrrfconncursorr=ZbatchesrqZ start_keyZ summariesZ batch_execsZbatch_execs_sortedZ first_execZ last_execrsrRrrrget_group_execution_summarys@    z-DependencyManager.get_group_execution_summarycCs|j|}|sddgdS|j|}g}g}|s|d|D]'}|j|d}|s9|d|ddq!|dsH|d |d d q!t|d k||t|d S)z Validate a job group configuration Args: group_id: Job group ID Returns: Validation results FzGroup not found)validerrorszGroup has no membersrr]renabledr+rz ' is disabledr)r|r}warningsr9)r r>r?rEr^rA)rrrMrNr}rrPrbrrrvalidate_job_groups,    z$DependencyManager.validate_job_group)r)rT)Nr%)re)__name__ __module__ __qualname____doc__rr rintstrrrrSfloatrFrrdr{rrrrrr s* x   $Br )rrYtypingrrrrrlogurur core.databaserdaemon.job_executorr r rrrrs