diff --git a/src/baseline.py b/src/baseline.py index 666b4ac..92024c8 100644 --- a/src/baseline.py +++ b/src/baseline.py @@ -24,7 +24,7 @@ def baseline_step(baseline_state, baseline_cores_available, baseline_running_job new_jobs_cores: List of cores per node required per job baseline_next_empty_slot: Index of next empty slot in baseline queue next_job_id: Next available job ID - metrics: Dictionary to update with baseline job metrics + metrics: MetricsTracker object to update with baseline job metrics env_print: Print function for logging Returns: @@ -38,8 +38,8 @@ def baseline_step(baseline_state, baseline_cores_available, baseline_running_job job_queue_2d, new_jobs_count, new_jobs_durations, new_jobs_nodes, new_jobs_cores, baseline_next_empty_slot ) - metrics['baseline_jobs_submitted'] += len(new_baseline_jobs) - metrics['baseline_jobs_rejected_queue_full'] += (new_jobs_count - len(new_baseline_jobs)) + metrics.baseline_jobs_submitted += len(new_baseline_jobs) + metrics.baseline_jobs_rejected_queue_full += (new_jobs_count - len(new_baseline_jobs)) _, baseline_next_empty_slot, _, next_job_id = assign_jobs_to_available_nodes( job_queue_2d, baseline_state['nodes'], baseline_cores_available, @@ -52,8 +52,8 @@ def baseline_step(baseline_state, baseline_cores_available, baseline_running_job num_unprocessed_jobs = np.sum(job_queue_2d[:, 0] > 0) # Track baseline max queue size - if num_unprocessed_jobs > metrics['baseline_max_queue_size_reached']: - metrics['baseline_max_queue_size_reached'] = num_unprocessed_jobs + if num_unprocessed_jobs > metrics.baseline_max_queue_size_reached: + metrics.baseline_max_queue_size_reached = num_unprocessed_jobs baseline_state['job_queue'] = job_queue_2d.flatten() diff --git a/src/environment.py b/src/environment.py index 6fbc2cc..b421783 100644 --- a/src/environment.py +++ b/src/environment.py @@ -306,29 +306,11 @@ def step(self, action): # Assign jobs to available nodes self.env_print(f"[4] Assigning jobs to available nodes...") - # Create metrics dict for job assignment - job_metrics = { - 'jobs_completed': self.metrics.jobs_completed, - 'total_job_wait_time': self.metrics.total_job_wait_time, - 'jobs_dropped': self.metrics.jobs_dropped, - 'dropped_this_episode': self.metrics.dropped_this_episode, - 'baseline_jobs_completed': self.metrics.baseline_jobs_completed, - 'baseline_total_job_wait_time': self.metrics.baseline_total_job_wait_time, - 'baseline_jobs_dropped': self.metrics.baseline_jobs_dropped, - 'baseline_dropped_this_episode': self.metrics.baseline_dropped_this_episode, - } - num_launched_jobs, self.next_empty_slot, num_dropped_this_step, self.next_job_id = assign_jobs_to_available_nodes( job_queue_2d, self.state['nodes'], self.cores_available, self.running_jobs, - self.next_empty_slot, self.next_job_id, job_metrics, is_baseline=False + self.next_empty_slot, self.next_job_id, self.metrics, is_baseline=False ) - # Update metrics from job_metrics dict - self.metrics.jobs_completed = job_metrics['jobs_completed'] - self.metrics.total_job_wait_time = job_metrics['total_job_wait_time'] - self.metrics.jobs_dropped = job_metrics['jobs_dropped'] - self.metrics.dropped_this_episode = job_metrics['dropped_this_episode'] - self.env_print(f" {num_launched_jobs} jobs launched") # Calculate node utilization stats @@ -353,31 +335,12 @@ def step(self, action): self.env_print(f"[5] Calculating reward...") # Baseline step - baseline_metrics = { - 'baseline_jobs_submitted': self.metrics.baseline_jobs_submitted, - 'baseline_jobs_rejected_queue_full': self.metrics.baseline_jobs_rejected_queue_full, - 'baseline_jobs_completed': self.metrics.baseline_jobs_completed, - 'baseline_total_job_wait_time': self.metrics.baseline_total_job_wait_time, - 'baseline_jobs_dropped': self.metrics.baseline_jobs_dropped, - 'baseline_dropped_this_episode': self.metrics.baseline_dropped_this_episode, - 'baseline_max_queue_size_reached': self.metrics.baseline_max_queue_size_reached, - } - baseline_cost, baseline_cost_off, self.baseline_next_empty_slot, self.next_job_id = baseline_step( self.baseline_state, self.baseline_cores_available, self.baseline_running_jobs, current_price, new_jobs_count, new_jobs_durations, new_jobs_nodes, new_jobs_cores, - self.baseline_next_empty_slot, self.next_job_id, baseline_metrics, self.env_print + self.baseline_next_empty_slot, self.next_job_id, self.metrics, self.env_print ) - # Update metrics from baseline_metrics dict - self.metrics.baseline_jobs_submitted = baseline_metrics['baseline_jobs_submitted'] - self.metrics.baseline_jobs_rejected_queue_full = baseline_metrics['baseline_jobs_rejected_queue_full'] - self.metrics.baseline_jobs_completed = baseline_metrics['baseline_jobs_completed'] - self.metrics.baseline_total_job_wait_time = baseline_metrics['baseline_total_job_wait_time'] - self.metrics.baseline_jobs_dropped = baseline_metrics['baseline_jobs_dropped'] - self.metrics.baseline_dropped_this_episode = baseline_metrics['baseline_dropped_this_episode'] - self.metrics.baseline_max_queue_size_reached = baseline_metrics['baseline_max_queue_size_reached'] - self.metrics.baseline_cost += baseline_cost self.metrics.baseline_cost_off += baseline_cost_off diff --git a/src/job_management.py b/src/job_management.py index 0036036..0fc0e34 100644 --- a/src/job_management.py +++ b/src/job_management.py @@ -111,7 +111,7 @@ def assign_jobs_to_available_nodes(job_queue_2d, nodes, cores_available, running running_jobs: Dictionary of currently running jobs next_empty_slot: Index of next empty slot in queue next_job_id: Next available job ID - metrics: Dictionary to update with job completion metrics + metrics: MetricsTracker object to update with job completion metrics is_baseline: Whether this is baseline simulation Returns: @@ -154,11 +154,11 @@ def assign_jobs_to_available_nodes(job_queue_2d, nodes, cores_available, running # Track job completion and wait time if is_baseline: - metrics['baseline_jobs_completed'] += 1 - metrics['baseline_total_job_wait_time'] += job_age + metrics.baseline_jobs_completed += 1 + metrics.baseline_total_job_wait_time += job_age else: - metrics['jobs_completed'] += 1 - metrics['total_job_wait_time'] += job_age + metrics.jobs_completed += 1 + metrics.total_job_wait_time += job_age num_processed_jobs += 1 continue @@ -176,11 +176,11 @@ def assign_jobs_to_available_nodes(job_queue_2d, nodes, cores_available, running num_dropped += 1 if is_baseline: - metrics['baseline_jobs_dropped'] += 1 - metrics['baseline_dropped_this_episode'] += 1 + metrics.baseline_jobs_dropped += 1 + metrics.baseline_dropped_this_episode += 1 else: - metrics['jobs_dropped'] += 1 - metrics['dropped_this_episode'] += 1 + metrics.jobs_dropped += 1 + metrics.dropped_this_episode += 1 else: job_queue_2d[job_idx][1] = new_age