Job tracker
cutana.job_tracker
¶
Job tracker for Cutana - manages overall job progress and resource utilization.
This module handles: - Overall job progress tracking across multiple processes - Resource monitoring and constraint checking - Integration with ProcessStatusReader/Writer for file operations - Status aggregation and completion detection - Error recording and reporting
JobTracker(tracking_file='job_tracking.json', progress_dir=None, session_id=None)
¶
Tracks overall job progress and manages system resources.
This class coordinates overall job status by: - Managing job-level state (total sources, completion counts) - Monitoring system resources (CPU, memory) - Delegating process-level operations to ProcessStatusReader/Writer - Aggregating status from multiple processes - Detecting job completion
Initialize the job tracker.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
tracking_file
|
str
|
Path to file for persisting job-level tracking data |
'job_tracking.json'
|
progress_dir
|
str
|
Directory for progress files (default: system temp) |
None
|
session_id
|
str
|
Session ID for progress file isolation (auto-generated if None) |
None
|
start_job(total_sources)
¶
Start tracking a new job.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
total_sources
|
int
|
Total number of sources to process |
required |
register_process(process_id, sources_assigned)
¶
Register a new process and create its progress file.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
process_id
|
str
|
Unique identifier for the process |
required |
sources_assigned
|
int
|
Number of sources assigned to this process |
required |
Returns:
| Type | Description |
|---|---|
bool
|
True if registration was successful |
update_process_progress(process_id, progress_update)
¶
Update progress for a specific process.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
process_id
|
str
|
Process identifier |
required |
progress_update
|
Dict[str, Any]
|
Dictionary with progress information |
required |
Returns:
| Type | Description |
|---|---|
bool
|
True if update was successful |
complete_process(process_id, completed_count, failed_count)
¶
Mark a process as completed and update overall counters.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
process_id
|
str
|
Process identifier |
required |
completed_count
|
int
|
Final count of successfully processed sources |
required |
failed_count
|
int
|
Count of failed sources |
required |
Returns:
| Type | Description |
|---|---|
bool
|
True if completion was recorded successfully |
calculate_smoothed_eta(completed_batches, total_batches, start_time)
¶
Calculate estimated time to completion with exponential smoothing.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
completed_batches
|
int
|
Number of completed batches |
required |
total_batches
|
int
|
Total number of batches |
required |
start_time
|
float
|
Workflow start time |
required |
Returns:
| Type | Description |
|---|---|
Optional[float]
|
Smoothed estimated seconds to completion, or None if not calculable |
report_process_progress(process_id, completed_sources, total_sources=None)
¶
Simplified progress reporting method for worker processes.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
process_id
|
str
|
Process identifier |
required |
completed_sources
|
int
|
Number of sources completed so far |
required |
total_sources
|
int
|
Total sources (optional, will read from file if not provided) |
None
|
Returns:
| Type | Description |
|---|---|
bool
|
True if report was successful |
update_process_stage(process_id, stage)
¶
Update the current processing stage for a process.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
process_id
|
str
|
Process identifier |
required |
stage
|
str
|
Current processing stage |
required |
Returns:
| Type | Description |
|---|---|
bool
|
True if update was successful |
has_process_progress_file(process_id)
¶
Check if a process has a progress file.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
process_id
|
str
|
Process identifier |
required |
Returns:
| Type | Description |
|---|---|
bool
|
True if progress file exists |
get_process_start_time(process_id)
¶
Get the start time for a process.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
process_id
|
str
|
Process identifier |
required |
Returns:
| Type | Description |
|---|---|
Optional[float]
|
Process start time or None if not found |
record_error(error_info)
¶
Record an error for tracking.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
error_info
|
Dict[str, Any]
|
Dictionary containing error information |
required |
get_process_details()
¶
Get detailed information about active processes.
Returns:
| Type | Description |
|---|---|
Dict[str, Dict[str, Any]]
|
Dictionary mapping process IDs to their details |
check_completion_status()
¶
Check completion status by aggregating all progress files.
Returns:
| Type | Description |
|---|---|
Dict[str, Any]
|
Dictionary containing aggregated completion status |
get_status()
¶
Get job status information (job-level only, no system resources).
Returns:
| Type | Description |
|---|---|
Dict[str, Any]
|
Dictionary containing job progress and process information |
cleanup_stale_processes(timeout=1800)
¶
Clean up processes that haven't updated in a while.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
timeout
|
int
|
Timeout in seconds for considering a process stale |
1800
|
Returns:
| Type | Description |
|---|---|
List[str]
|
List of cleaned up process IDs |
cleanup_all_progress_files()
¶
Clean up all progress files for this job.
Returns:
| Type | Description |
|---|---|
int
|
Number of files cleaned up |
get_sources_assigned_to_process(process_id)
¶
Get the number of sources assigned to a process.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
process_id
|
str
|
Process identifier |
required |
Returns:
| Type | Description |
|---|---|
int
|
Number of sources assigned to the process |