o 9iC@sdZddlZddlZddlZddlZddlZddlmZddlmZddl m Z m Z m Z ddl mZddlZddlmZddlmZGd d d ZdS) z Job Executor for Python Scheduler Handles job execution in background threads with timeout support Based on DataHub web_panel/utils/job_manager.py pattern N)datetime)Path)OptionalDictAny)logger)SchedulerDatabase)SchedulerConfigc@seZdZdZd'dedefddZdefdd Zd(d ed e d e efddZ dede e e ffddZde e e fd efddZde e d e fddZdede d efddZd)deded efdd Zdede e e ffd!d"Zd efd#d$Zded efd%d&ZdS)* JobExecutorz8Executes jobs in background threads with proper trackingNdbconfigcCs0||_||_||_i|_t|_tddS)NzJobExecutor initialized) r r notification_manager active_jobs threadingLocklockrinfo)selfr r r r*c:\PYTHON\Scheduler\daemon\job_executor.py__init__s  zJobExecutor.__init__pidc Csbzt|}|jdd}|D] }ztd|jd||Wqtjy.Yqwztd|jd||Wn tjyMYnwtj |g|dd\}}|D] }zt d|jd||Wq\tjy|Yq\wWd Stjyt d |d Yd St y}zt d |d|WYd }~d Sd }~ww) z| Terminate process and all its children recursively Args: pid: Process ID to terminate T) recursivezKilling child process : zKilling parent process timeoutzForce killing process zProcess z already terminatedzError killing process tree N) psutilProcesschildrenrdebugrnamekillZ NoSuchProcessZ wait_procswarning Exceptionerror) rrparentrchildZgonealiveperrr_kill_process_tree!s>     $zJobExecutor._kill_process_tree schedulerjob_id triggered_byreturncKs|j|}|std|ddS|ds'td|dd|ddS|jjd|d |td |}td |dd |d t j |j ||fdd}| |S)a. Execute a job asynchronously Args: job_id: ID of job to execute triggered_by: Who/what triggered execution **kwargs: Additional execution parameters Returns: Execution ID for tracking, or None if job not found/disabled zJob z not foundNenabledJob 'r!z' (ID z) is disabled, skippingZqueued)r-statusr. start_timez Queued job 'z' (execution ID )F)targetargsdaemonr) r get_jobrr%r#create_executionrnowrrThread_run_jobstart)rr-r.kwargsjob execution_idthreadrrr execute_jobKs. zJobExecutor.execute_jobr@r?c Csf|d}|||}z|jj|dtt|dtd|d|d||}t}| dd}|rdtd |d t |d d d [}| d| d|d| d|d| d| dd| dt |trv|nd|d| d|dd| dd| d| dWdn1swYtd |dt |trddl} | |}t |dd d }tj||dr|dnd|tjd d!} Wdn1swY|j | |j|<Wdn1swY|d"} z| j| d#} | dkrd$nd%} Wn.tjyGtd |d&| d'td(|d)| jd|| jd*} d+} Ynw|j|j|dWdn 1s^wYnt |d d d }| d| d|d| d|d| d| dd| dd|d| d|dd| dd|tj||tj|dr|dndd d,} |j | |j|<Wdn 1swY|d"} z| j| d#} | dkrd$nd%} Wn.tjy"td |d&| d'td(|d)| jd|| jd*} d+} Ynw|j|j|dWdn 1s9wYWdn 1sIwYt}||}||\}}|jj|| ||| ||d-| d$krtd |d.|d/d0n| d+krtd |d&| d0n td |d1| |j r|j j!||| || d2| d%kr|d3dkr|"||WdSWdSWdSt#y2}z`t$d4|d5||jj|d%tt|d6z4t |dd d "}| ddd| d7|d| ddWdn 1swYWn YWYd}~dSWYd}~dSd}~ww)8z Internal: Execute job in subprocess Args: execution_id: Execution record ID job: Job definition dict r!running)r2r3log_filezStarting job 'z ' (execution r4 show_consoleFr1zW' configured to show console window (debug mode) - output will be visible in CMD windowwutf-8encodingz=== Job Execution Log === zJob:  zExecution ID: z Start Time: z%Y-%m-%d %H:%M:%Sz Command:  zWorking Directory: working_directoryz2==================================================z z6NOTE: Output displayed in console window (debug mode) z6Log file not captured - see console window for output Nz:' has show_console enabled - log can be monitored from GUIraT)cwdstdoutstderrtexttimeout_secondsr completedfailedz' timed out after z secondsz"Terminating process tree for job 'z' (PID: r)rOrPrNrQ)r2end_timeduration_seconds exit_codestdout_previewstderr_previewz' completed successfully in z.1fsz' failed with exit code )r@job_namer2durationrX retry_countzError executing job 'z': )r2rV error_messagezEXECUTION ERROR: )%_create_log_file_pathr update_executionrr:strrr_build_commandgetopenwritestrftime isinstancejoinshlexsplit subprocessPopenSTDOUTrrwaitTimeoutExpiredr#rr+popflush total_seconds_read_log_previewsr%r Znotify_job_result _handle_retryr$ exception)rr@r?r\ log_file_pathcmdr3rErDrjprocessrrXr2rVr]rYrZr*rrrr<us    &           )     zJobExecutor._run_jobcCs|d}|d}|dkr|g}n*|dkr!||d}||g}n|dkr3||d}|d|g}ntd||d r`zt|d }||W|Stjy_||d Y|Sw|S) znBuild command list for subprocess Returns: list: Command and arguments as a list job_typeexecutable_pathbatchscriptrL python_modulez-mzUnknown job type: arguments)_get_python_exerd ValueErrorjsonloadsextendJSONDecodeErrorappend)rr?rzr{rxZ python_exer6rrrrc/s*   zJobExecutor._build_commandrLcCsddl}|rt|ddd}|rt|S|jdi}|rL|D]&\}}||ddrK|d}|rKt|dd}|rKt|Sq%|jS) z3Get Python executable, preferring venv if availablerNvenvZScriptsz python.exeprojectsroot) sysrexistsrbr rditems startswith executable)rrLrZ venv_pythonrZ project_nameZproject_configZ venv_pathrrrrTs   zJobExecutor._get_python_exer\cCsXt|jj}|jdddtd}|dddd}|d|d}||S)z#Create log file path with timestampT)parentsexist_okz %Y%m%d_%H%M%SrK_/z.log) rr execution_log_dirmkdirrr:rglowerreplace)rr@r\log_dir timestampZsanitized_nameZ log_filenamerrrr`ls z!JobExecutor._create_log_file_pathrw preview_linesc Csz.t|ddd }|}Wdn1swY|r(d|| dnd}|dfWStyI}ztd|WYd}~dSd}~ww)z+Read last N lines from log file for previewrrGrHNrzCould not read log preview: )NN)re readlinesrir$rr#)rrwrflinesZpreviewr*rrrrtws  zJobExecutor._read_log_previewsc sjjddd}tdd|D}|dkrJdtdd d d |d d dd d fdd}tj|dddSdS)zHandle job retry logic id)limitr-css |] }|ddkrdVqdS)r2rTNr).0exrrr sz,JobExecutor._handle_retry..r^retry_delay_secondszRetrying job 'r!z' in z seconds (attempt rrr4cs tjddddS)Nrretry)r.)timesleeprBrr? retry_delayrrr retry_jobs z,JobExecutor._handle_retry..retry_jobT)r5r7N)r get_recent_executionssumrrrr;r=)rr@r?Zrecent_executions failed_countrrrrrus 4zJobExecutor._handle_retrycCs8|jt|jWdS1swYdS)z(Get list of currently running executionsN)rlistrkeys)rrrrget_active_executionss $z!JobExecutor.get_active_executionsc Cs$|j||jvr WddS|j|}|durz>|z|jddWntjy8|Ynw|jj |dt d|j |dt d|WWddSty}zt d |d |WYd}~WddSd}~wwWddS1swYdS) z Stop a running execution Args: execution_id: Execution ID to stop Returns: True if stopped successfully NFrr cancelled)r2rVzStopped execution TzError stopping execution r)rrpoll terminaterorlrpr"r rarr:rqrrr$r%)rr@ryr*rrrstop_executionsB      ##zJobExecutor.stop_execution)N)r,)r)__name__ __module__ __qualname____doc__rr rintr+rbrrBrrr<rrcrrr`tuplertrurboolrrrrrr s**;% r )rrlrrrrrpathlibrtypingrrrlogururr core.databaser core.configr r rrrrs