use std::collections::HashSet;
use std::collections::HashMap;
use std::sync::atomic::AtomicBool;
use crate::{ProofmanResult, ProofmanError};
#[derive(Clone, Debug, Eq, Hash, PartialEq)]
pub struct InstanceChunks {
pub chunks: Vec<usize>,
pub slow: bool,
}
#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
pub struct InstanceInfo {
pub airgroup_id: usize,
pub air_id: usize,
pub table: bool,
pub shared: bool,
pub n_chunks: usize,
pub weight: u64,
pub compressor_weight: u64, }
impl InstanceInfo {
pub fn new(
airgroup_id: usize,
air_id: usize,
table: bool,
shared: bool,
weight: u64,
compressor_weight: u64,
) -> Self {
Self { airgroup_id, air_id, table, shared, n_chunks: 0, weight, compressor_weight }
}
#[inline]
pub fn total_weight(&self) -> u64 {
self.weight + self.compressor_weight
}
#[inline]
pub fn has_compressor(&self) -> bool {
self.compressor_weight > 0
}
}
#[derive(Default)]
pub struct DistributionCtx {
pub n_partitions: usize, pub partition_mask: Vec<bool>,
pub n_processes: usize, pub process_id: usize,
pub n_instances: usize, pub instances: Vec<InstanceInfo>, pub instances_chunks: Vec<InstanceChunks>, pub instances_calculated: Vec<AtomicBool>, pub n_tables: usize, pub aux_tables: Vec<InstanceInfo>, pub aux_table_map: Vec<i32>,
pub partition_set: bool, pub instance_partition: Vec<i32>, pub worker_instances: Vec<usize>, pub partition_count: Vec<u32>, pub partition_weight: Vec<u64>, pub partition_compressor_count: Vec<u32>,
pub instance_process: Vec<(i32, usize)>, pub process_instances: Vec<usize>, pub skipped_process_instances: Vec<usize>, pub process_count: Vec<usize>, pub process_weight: Vec<u64>, pub process_compressor_count: Vec<u32>,
pub worker_index: i32,
pub assignation_done: bool, }
impl std::fmt::Debug for DistributionCtx {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let mut dbg = f.debug_struct("DistributionCtx");
dbg.field("=== STATIC PARAMS ===", &"");
dbg.field("n_partitions", &self.n_partitions)
.field("partition_mask", &self.partition_mask)
.field("n_processes", &self.n_processes)
.field("process_id", &self.process_id);
dbg.field("=== DYNAMIC PARAMS ===", &"");
dbg.field("n_instances", &self.n_instances)
.field("instances", &self.instances)
.field("n_tables", &self.n_tables)
.field("tables", &self.aux_tables)
.field("instance_partition", &self.instance_partition)
.field("worker_instances", &self.worker_instances)
.field("partition_count", &self.partition_count)
.field("partition_weight", &self.partition_weight)
.field("partition_compressor_count", &self.partition_compressor_count)
.field("instance_process", &self.instance_process)
.field("process_instances", &self.process_instances)
.field("process_count", &self.process_count)
.field("process_weight", &self.process_weight)
.field("process_compressor_count", &self.process_compressor_count)
.field("assignation_done", &self.assignation_done);
dbg.finish()
}
}
impl DistributionCtx {
pub fn new() -> Self {
DistributionCtx {
n_partitions: 0,
partition_mask: Vec::new(),
n_processes: 0,
process_id: 0,
n_instances: 0,
instances: Vec::new(),
instances_calculated: Vec::new(),
instances_chunks: Vec::new(),
n_tables: 0,
aux_tables: Vec::new(),
aux_table_map: Vec::new(),
instance_partition: Vec::new(),
worker_instances: Vec::new(),
partition_count: Vec::new(),
partition_weight: Vec::new(),
partition_compressor_count: Vec::new(),
instance_process: Vec::new(),
process_instances: Vec::new(),
skipped_process_instances: Vec::new(),
process_count: Vec::new(),
process_weight: Vec::new(),
process_compressor_count: Vec::new(),
worker_index: -1,
assignation_done: false,
partition_set: false,
}
}
pub fn setup_partitions(&mut self, n_partitions: usize, partition_ids: Vec<u32>) -> ProofmanResult<()> {
self.n_partitions = n_partitions;
self.partition_mask = vec![false; n_partitions];
for id in &partition_ids {
if *id < n_partitions as u32 {
self.partition_mask[*id as usize] = true;
} else {
return Err(ProofmanError::InvalidConfiguration(format!(
"Partition ID {} exceeds total partitions {}",
id, n_partitions
)));
}
}
self.partition_count = vec![0; n_partitions];
self.partition_weight = vec![0; n_partitions];
self.partition_compressor_count = vec![0; n_partitions];
self.partition_set = true;
Ok(())
}
pub fn setup_processes(&mut self, n_processes: usize, process_id: usize) -> ProofmanResult<()> {
if process_id >= n_processes {
return Err(ProofmanError::InvalidConfiguration(format!(
"Process rank {} exceeds total processes {}",
process_id, n_processes
)));
}
self.n_processes = n_processes;
self.process_id = process_id;
self.process_count = vec![0; n_processes];
self.process_weight = vec![0; n_processes];
self.process_compressor_count = vec![0; n_processes];
Ok(())
}
pub fn setup_worker_index(&mut self, worker_index: usize) {
self.worker_index = worker_index as i32;
}
pub fn reset_instances(&mut self) {
self.n_instances = 0;
self.instances.clear();
self.instances_chunks.clear();
self.instances_calculated.clear();
self.n_tables = 0;
self.aux_tables.clear();
self.aux_table_map.clear();
self.instance_partition.clear();
self.worker_instances.clear();
self.partition_count.fill(0);
self.partition_weight.fill(0);
self.partition_compressor_count.fill(0);
self.instance_process.clear();
self.process_instances.clear();
self.skipped_process_instances.clear();
self.process_count.fill(0);
self.process_weight.fill(0);
self.process_compressor_count.fill(0);
self.assignation_done = false;
self.partition_set = false;
}
#[inline]
pub fn validate_static_config(&self) -> Result<(), ProofmanError> {
if self.n_partitions == 0 {
return Err(ProofmanError::InvalidConfiguration(
"Partition configuration not set. Call setup_partitions() first.".to_string(),
));
}
if self.partition_mask.len() != self.n_partitions {
return Err(ProofmanError::InvalidConfiguration(
"Partition mask size mismatch with n_partitions".to_string(),
));
}
if self.n_processes == 0 {
return Err(ProofmanError::InvalidConfiguration(
"Process configuration not set. Call setup_processes() first.".to_string(),
));
}
if self.process_id >= self.n_processes {
return Err(ProofmanError::InvalidConfiguration(format!(
"Invalid process rank {} >= {}",
self.process_id, self.n_processes
)));
}
Ok(())
}
#[inline]
pub fn is_setup_partition_init(&self) -> bool {
self.partition_set
}
pub fn skip_instance(&mut self, instance_id: usize) {
if !self.skipped_process_instances.contains(&instance_id) {
self.skipped_process_instances.push(instance_id);
}
}
pub fn is_skipped_instance(&self, instance_id: usize) -> bool {
self.skipped_process_instances.contains(&instance_id)
}
#[inline]
pub fn is_my_process_instance(&self, instance_id: usize) -> ProofmanResult<bool> {
if instance_id >= self.instance_process.len() {
return Err(ProofmanError::OutOfBounds(format!(
"Instance index {} out of bounds (max: {})",
instance_id,
self.instance_process.len()
)));
}
Ok(self.instance_process[instance_id].0 == self.process_id as i32)
}
#[inline]
pub fn get_process_owner_instance(&self, instance_id: usize) -> ProofmanResult<i32> {
if instance_id >= self.instance_process.len() {
return Err(ProofmanError::OutOfBounds(format!(
"Instance index {} out of bounds (max: {})",
instance_id,
self.instance_process.len()
)));
}
let owner = self.instance_process[instance_id].0;
if owner == -1 {
return Err(ProofmanError::InvalidAssignation(format!(
"Instance {} is not owned by any process",
instance_id
)));
}
Ok(owner)
}
#[inline]
pub fn get_instance_info(&self, instance_id: usize) -> ProofmanResult<(usize, usize)> {
if instance_id >= self.instances.len() {
return Err(ProofmanError::OutOfBounds(format!(
"Instance index {} out of bounds (max: {})",
instance_id,
self.instances.len()
)));
}
Ok((self.instances[instance_id].airgroup_id, self.instances[instance_id].air_id))
}
#[inline]
pub fn get_table_info(&self, table_idx: usize) -> ProofmanResult<(usize, usize)> {
if self.assignation_done {
let instance_id = self.aux_table_map[table_idx] as usize;
self.get_instance_info(instance_id)
} else {
if table_idx >= self.aux_tables.len() {
return Err(ProofmanError::OutOfBounds(format!(
"Table index {} out of bounds (max: {})",
table_idx,
self.aux_tables.len()
)));
}
Ok((self.aux_tables[table_idx].airgroup_id, self.aux_tables[table_idx].air_id))
}
}
pub fn get_table_instance_idx(&self, table_idx: usize) -> ProofmanResult<usize> {
if self.assignation_done {
Ok(self.aux_table_map[table_idx] as usize)
} else {
Err(ProofmanError::InvalidAssignation("Table instances not yet assigned".into()))
}
}
#[inline]
pub fn get_instance_local_idx(&self, instance_id: usize) -> ProofmanResult<usize> {
if instance_id >= self.instance_process.len() {
return Err(ProofmanError::OutOfBounds(format!(
"Instance index {} out of bounds (max: {})",
instance_id,
self.instance_process.len()
)));
}
Ok(self.instance_process[instance_id].1)
}
#[inline]
pub fn get_instance_chunks(&self, instance_id: usize) -> ProofmanResult<usize> {
if instance_id >= self.instances.len() {
return Err(ProofmanError::OutOfBounds(format!(
"Instance index {} out of bounds (max: {})",
instance_id,
self.instances.len()
)));
}
Ok(self.instances_chunks[instance_id].chunks.len())
}
pub fn set_n_chunks(&mut self, instance_id: usize, n_chunks: usize) -> ProofmanResult<()> {
if instance_id >= self.instances.len() {
return Err(ProofmanError::OutOfBounds(format!(
"Instance index {} out of bounds (max: {})",
instance_id,
self.instances.len()
)));
}
let instance_info = &mut self.instances[instance_id];
instance_info.n_chunks = n_chunks;
Ok(())
}
#[inline]
pub fn is_my_worker_instance(&self, instance_id: usize) -> ProofmanResult<bool> {
if instance_id >= self.instance_process.len() {
return Err(ProofmanError::OutOfBounds(format!(
"Instance index {} out of bounds (max: {})",
instance_id,
self.instance_process.len()
)));
}
Ok(self.instance_process[instance_id].0 >= 0)
}
#[inline]
pub fn find_air_instance_id(&self, instance_id: usize) -> ProofmanResult<usize> {
let mut air_instance_id = 0;
let (airgroup_id, air_id) = self.get_instance_info(instance_id)?;
for idx in 0..instance_id {
let (instance_airgroup_id, instance_air_id) = self.get_instance_info(idx)?;
if (instance_airgroup_id, instance_air_id) == (airgroup_id, air_id) {
air_instance_id += 1;
}
}
Ok(air_instance_id)
}
#[inline]
pub fn find_process_instance(&self, airgroup_id: usize, air_id: usize) -> ProofmanResult<(bool, usize)> {
let mut matches = self
.process_instances
.iter()
.enumerate()
.filter(|&(_pos, &id)| {
let inst = &self.instances[id];
inst.airgroup_id == airgroup_id && inst.air_id == air_id
})
.map(|(pos, _)| pos);
match (matches.next(), matches.next()) {
(None, _) => Ok((false, 0)),
(Some(pos), None) => Ok((true, pos)),
(Some(_), Some(_)) => Err(ProofmanError::InvalidAssignation(format!(
"Multiple instances found for airgroup_id: {airgroup_id}, air_id: {air_id}"
))),
}
}
#[inline]
pub fn find_instance_id(&self, airgroup_id: usize, air_id: usize) -> ProofmanResult<(bool, usize)> {
let mut matches = self
.instances
.iter()
.enumerate()
.filter(|&(_gid, inst)| inst.airgroup_id == airgroup_id && inst.air_id == air_id)
.map(|(gid, _)| gid);
match (matches.next(), matches.next()) {
(None, _) => Ok((false, 0)),
(Some(gid), None) => Ok((true, gid)),
(Some(_), Some(_)) => Err(ProofmanError::InvalidAssignation(format!(
"Multiple instances found for airgroup_id: {airgroup_id}, air_id: {air_id}"
))),
}
}
#[inline]
pub fn find_process_table(&self, airgroup_id: usize, air_id: usize) -> ProofmanResult<(bool, usize)> {
if self.assignation_done {
self.find_process_instance(airgroup_id, air_id)
} else {
let mut matches = self
.aux_tables
.iter()
.enumerate()
.filter(|&(_pos, &info)| info.airgroup_id == airgroup_id && info.air_id == air_id)
.map(|(pos, _)| pos);
match (matches.next(), matches.next()) {
(None, _) => Ok((false, 0)),
(Some(pos), None) => Ok((true, pos)),
(Some(_), Some(_)) => Err(ProofmanError::InvalidAssignation(format!(
"Multiple tables found for airgroup_id: {airgroup_id}, air_id: {air_id}"
))),
}
}
}
#[inline]
pub fn find_by_air_instance_id(&self, airgroup_id: usize, air_id: usize, air_instance_id: usize) -> Option<usize> {
let mut count = 0;
for (instance_idx, instance) in self.instances.iter().enumerate() {
let (inst_airgroup_id, inst_air_id) = (instance.airgroup_id, instance.air_id);
if airgroup_id == inst_airgroup_id && air_id == inst_air_id {
if count == air_instance_id {
return Some(instance_idx);
}
count += 1;
}
}
None
}
#[inline]
fn least_loaded_partition(&self, has_compressor: bool) -> usize {
let mut best_idx = 0;
let mut best_key = (u64::MAX, u32::MAX, u32::MAX);
for (i, &weight) in self.partition_weight.iter().enumerate() {
let compressors = if has_compressor { self.partition_compressor_count[i] } else { 0 };
let key = (weight, compressors, self.partition_count[i]);
if key < best_key {
best_key = key;
best_idx = i;
}
}
best_idx
}
#[inline]
fn least_loaded_process(&self, has_compressor: bool, counts: &[usize]) -> usize {
let mut best_idx = 0;
let mut best_key = (u64::MAX, u32::MAX, usize::MAX);
for (i, &weight) in self.process_weight.iter().enumerate() {
let compressors = if has_compressor { self.process_compressor_count[i] } else { 0 };
let key = (weight, compressors, counts[i]);
if key < best_key {
best_key = key;
best_idx = i;
}
}
best_idx
}
#[inline]
pub fn add_instance(
&mut self,
airgroup_id: usize,
air_id: usize,
weight: u64,
compressor_weight: u64,
) -> ProofmanResult<usize> {
if self.assignation_done {
return Err(ProofmanError::InvalidAssignation("Instances already assigned".to_string()));
}
self.validate_static_config().expect("Static configuration invalid or incomplete");
let gid: usize = self.instances.len();
let instance = InstanceInfo::new(airgroup_id, air_id, false, false, weight, compressor_weight);
let (total_weight, has_compressor) = (instance.total_weight(), instance.has_compressor());
self.instances.push(instance);
self.instances_chunks.push(InstanceChunks { chunks: vec![], slow: false });
self.instances_calculated.push(AtomicBool::new(false));
self.n_instances += 1;
let partition_id = self.least_loaded_partition(has_compressor) as u32;
self.instance_partition.push(partition_id as i32);
self.partition_count[partition_id as usize] += 1;
self.partition_weight[partition_id as usize] += total_weight;
if has_compressor {
self.partition_compressor_count[partition_id as usize] += 1;
}
let mut local_idx = 0;
let mut owner = -1;
if self.partition_mask[partition_id as usize] {
self.worker_instances.push(gid);
let process_id = self.least_loaded_process(has_compressor, &self.process_count);
owner = process_id as i32;
local_idx = self.process_count[process_id];
self.process_count[process_id] += 1;
self.process_weight[process_id] += total_weight;
if has_compressor {
self.process_compressor_count[process_id] += 1;
}
if process_id == self.process_id {
self.process_instances.push(gid);
}
}
self.instance_process.push((owner, local_idx));
Ok(gid)
}
pub fn is_first_process(&self) -> bool {
self.partition_mask[0] && self.process_id == 0
}
#[inline]
pub fn add_instance_no_assign(
&mut self,
airgroup_id: usize,
air_id: usize,
weight: u64,
compressor_weight: u64,
) -> ProofmanResult<usize> {
if self.assignation_done {
return Err(ProofmanError::InvalidAssignation("Instances already assigned".to_string()));
}
self.validate_static_config().expect("Static configuration invalid or incomplete");
self.instances.push(InstanceInfo::new(airgroup_id, air_id, false, false, weight, compressor_weight));
self.instances_chunks.push(InstanceChunks { chunks: vec![], slow: false });
self.instances_calculated.push(AtomicBool::new(false));
self.instance_partition.push(-1);
self.instance_process.push((-1, 0_usize));
self.n_instances += 1;
Ok(self.n_instances - 1)
}
pub fn add_table(&mut self, airgroup_id: usize, air_id: usize, weight: u64) -> ProofmanResult<usize> {
if self.assignation_done {
return Err(ProofmanError::InvalidAssignation("Instances already assigned".to_string()));
}
self.validate_static_config().expect("Static configuration invalid or incomplete");
let lid = self.aux_tables.len();
self.aux_tables.push(InstanceInfo::new(airgroup_id, air_id, true, true, weight, 0));
self.aux_table_map.push(-1);
self.n_tables += 1;
Ok(lid)
}
pub fn set_chunks(&mut self, global_idx: usize, chunks: Vec<usize>, slow: bool) {
let instance_info = &mut self.instances_chunks[global_idx];
instance_info.chunks = chunks;
instance_info.slow = slow;
}
pub fn add_table_all(&mut self, airgroup_id: usize, air_id: usize, weight: u64) -> ProofmanResult<usize> {
if self.assignation_done {
return Err(ProofmanError::InvalidAssignation("Instances already assigned".to_string()));
}
self.validate_static_config().expect("Static configuration invalid or incomplete");
let lid = self.aux_tables.len();
self.aux_tables.push(InstanceInfo::new(airgroup_id, air_id, true, false, weight, 0));
self.aux_table_map.push(-1);
self.n_tables += 1;
Ok(lid)
}
pub fn assign_instances(&mut self) -> ProofmanResult<()> {
if self.assignation_done {
return Err(ProofmanError::InvalidAssignation("Instances already assigned".to_string()));
}
self.validate_static_config().expect("Static configuration invalid or incomplete");
let mut unassigned_instances = Vec::new();
for (gid, &partition_id) in self.instance_partition.iter().enumerate() {
if partition_id == -1 {
unassigned_instances.push((gid, self.instances[gid].total_weight()));
}
}
unassigned_instances.sort_by_key(|b| std::cmp::Reverse(b.1));
let mut instances_assigned_partition = vec![HashMap::<(usize, usize), usize>::new(); self.n_partitions];
let mut instances_assigned_process = vec![HashMap::<(usize, usize), usize>::new(); self.n_processes];
let mut local_process_count = self.process_count.clone();
for (gid, _) in &unassigned_instances {
let (airgroup_id, air_id) = self.get_instance_info(*gid)?;
let has_compressor = self.instances[*gid].has_compressor();
let min_weight_idx = self.least_loaded_partition(has_compressor);
if has_compressor {
self.partition_compressor_count[min_weight_idx] += 1;
}
*instances_assigned_partition[min_weight_idx].entry((airgroup_id, air_id)).or_insert(0) += 1;
self.partition_count[min_weight_idx] += 1;
self.partition_weight[min_weight_idx] += self.instances[*gid].total_weight();
if self.partition_mask[min_weight_idx] {
let min_weight_process_idx = self.least_loaded_process(has_compressor, &local_process_count);
if has_compressor {
self.process_compressor_count[min_weight_process_idx] += 1;
}
local_process_count[min_weight_process_idx] += 1;
self.process_weight[min_weight_process_idx] += self.instances[*gid].total_weight();
instances_assigned_process[min_weight_process_idx]
.entry((airgroup_id, air_id))
.and_modify(|c| *c += 1)
.or_insert(1);
}
}
unassigned_instances.sort_by_key(|&(idx, weight)| (self.instances_chunks[idx].slow, std::cmp::Reverse(weight)));
let partitions_chunks: &mut Vec<HashSet<usize>> = &mut (0..self.n_partitions).map(|_| HashSet::new()).collect();
let process_chunks: &mut Vec<HashSet<usize>> = &mut (0..self.n_processes).map(|_| HashSet::new()).collect();
for (gid, _) in &unassigned_instances {
let chunks = &self.instances_chunks[*gid].chunks;
let (airgroup_id, air_id) = self.get_instance_info(*gid)?;
let mut min_chunks = usize::MAX;
let mut min_chunks_idx = 0;
for partition_id in 0..self.n_partitions {
if instances_assigned_partition[partition_id].get(&(airgroup_id, air_id)).unwrap_or(&0) > &0 {
let mut new_chunks_added = 0;
for chunk in chunks {
if !partitions_chunks[partition_id].contains(chunk) {
new_chunks_added += 1;
}
}
if new_chunks_added < min_chunks {
min_chunks = new_chunks_added;
min_chunks_idx = partition_id;
}
}
}
if let Some(c) = instances_assigned_partition[min_chunks_idx].get_mut(&(airgroup_id, air_id)) {
*c -= 1;
}
self.instance_partition[*gid] = min_chunks_idx as i32;
for chunk in chunks {
partitions_chunks[min_chunks_idx].insert(*chunk);
}
if self.partition_mask[min_chunks_idx] {
self.worker_instances.push(*gid);
let mut min_chunks = usize::MAX;
let mut min_process_id = 0;
for process_id in 0..self.n_processes {
if instances_assigned_process[process_id].get(&(airgroup_id, air_id)).unwrap_or(&0) > &0 {
let mut new_chunks_added = 0;
for chunk in chunks {
if !process_chunks[process_id].contains(chunk) {
new_chunks_added += 1;
}
}
if new_chunks_added < min_chunks {
min_chunks = new_chunks_added;
min_process_id = process_id;
}
}
}
for chunk in chunks {
process_chunks[min_process_id].insert(*chunk);
}
if min_process_id == self.process_id {
self.process_instances.push(*gid);
}
self.instance_process[*gid].0 = min_process_id as i32;
self.instance_process[*gid].1 = self.process_count[min_process_id];
self.process_count[min_process_id] += 1;
if let Some(c) = instances_assigned_process[min_process_id].get_mut(&(airgroup_id, air_id)) {
*c -= 1;
}
}
}
self.n_tables = 0;
for (table_idx, table) in self.aux_tables.iter().enumerate() {
if table.shared {
let process_id = self.least_loaded_process(false, &self.process_count);
let gid = self.instances.len();
self.instances.push(*table);
self.instances_calculated.push(AtomicBool::new(false));
self.instances_chunks.push(InstanceChunks { chunks: vec![], slow: false });
self.n_instances += 1;
self.n_tables += 1;
self.instance_partition.push(-2); self.worker_instances.push(gid);
let lid = self.process_count[process_id];
self.process_count[process_id] += 1;
self.process_weight[process_id] += table.weight;
if process_id == self.process_id {
self.process_instances.push(gid);
}
self.aux_table_map[table_idx] = gid as i32;
self.instance_process.push((process_id as i32, lid));
} else {
for rank in 0..self.n_processes {
let gid = self.instances.len();
self.instances.push(InstanceInfo::new(
table.airgroup_id,
table.air_id,
true,
false,
table.weight,
0,
));
self.instances_chunks.push(InstanceChunks { chunks: vec![], slow: false });
self.n_instances += 1;
self.n_tables += 1;
self.instance_partition.push(-2); self.worker_instances.push(gid);
let lid = self.process_count[rank];
self.process_count[rank] += 1;
self.process_weight[rank] += table.weight;
if rank == self.process_id {
self.process_instances.push(gid);
self.aux_table_map[table_idx] = gid as i32;
}
self.instance_process.push((rank as i32, lid));
}
}
}
self.aux_tables.clear();
self.assignation_done = true;
Ok(())
}
pub fn load_balance_info_partition(&self) -> (f64, u64, u64, f64) {
let mut average_partition_weight = 0.0;
let mut max_partition_weight = 0;
let mut min_partition_weight = u64::MAX;
for i in 0..self.n_partitions {
average_partition_weight += self.partition_weight[i] as f64;
if self.partition_weight[i] > max_partition_weight {
max_partition_weight = self.partition_weight[i];
}
if self.partition_weight[i] < min_partition_weight {
min_partition_weight = self.partition_weight[i];
}
}
average_partition_weight /= self.n_partitions as f64;
let max_deviation = max_partition_weight as f64 / average_partition_weight;
(average_partition_weight, max_partition_weight, min_partition_weight, max_deviation)
}
pub fn load_balance_info_process(&self) -> (f64, u64, u64, f64) {
let mut average_process_weight = 0.0;
let mut max_process_weight = 0;
let mut min_process_weight = u64::MAX;
for i in 0..self.n_processes {
average_process_weight += self.process_weight[i] as f64;
if self.process_weight[i] > max_process_weight {
max_process_weight = self.process_weight[i];
}
if self.process_weight[i] < min_process_weight {
min_process_weight = self.process_weight[i];
}
}
average_process_weight /= self.n_processes as f64;
let max_deviation = max_process_weight as f64 / average_process_weight;
(average_process_weight, max_process_weight, min_process_weight, max_deviation)
}
}
#[cfg(test)]
mod tests {
use super::DistributionCtx;
const HEAVY_PLAIN: (usize, u64, u64) = (0, 5000, 0);
const COMPRESSOR: (usize, u64, u64) = (1, 1000, 100);
const LIGHT_PLAIN: (usize, u64, u64) = (2, 100, 0);
fn ctx(n_partitions: usize) -> DistributionCtx {
let mut dctx = DistributionCtx::new();
dctx.setup_processes(1, 0).unwrap();
dctx.setup_partitions(n_partitions, (0..n_partitions as u32).collect()).unwrap();
dctx
}
#[test]
fn compressor_air_follows_cost_not_compressor_count() {
let mut dctx = ctx(2);
for (air_id, weight, compressor_weight) in [HEAVY_PLAIN, COMPRESSOR, COMPRESSOR] {
dctx.add_instance_no_assign(0, air_id, weight, compressor_weight).unwrap();
}
dctx.assign_instances().unwrap();
assert_eq!(dctx.instance_partition, vec![0, 1, 1]);
assert_eq!(dctx.partition_weight, vec![5000, 2200]);
assert_eq!(dctx.partition_compressor_count, vec![0, 2]);
}
#[test]
fn immediate_assignation_balances_cost_instead_of_round_robin() {
let mut dctx = ctx(2);
for (air_id, weight, compressor_weight) in [HEAVY_PLAIN, LIGHT_PLAIN, LIGHT_PLAIN] {
dctx.add_instance(0, air_id, weight, compressor_weight).unwrap();
}
assert_eq!(dctx.instance_partition, vec![0, 1, 1]);
assert_eq!(dctx.partition_weight, vec![5000, 200]);
}
#[test]
fn first_instance_goes_to_partition_zero() {
let mut dctx = ctx(4);
dctx.add_instance(0, LIGHT_PLAIN.0, LIGHT_PLAIN.1, LIGHT_PLAIN.2).unwrap();
assert_eq!(dctx.instance_partition[0], 0);
assert_eq!(dctx.instance_process[0], (0, 0));
}
#[test]
fn compressor_weight_counts_towards_the_partition_load() {
let mut dctx = ctx(1);
dctx.add_instance(0, COMPRESSOR.0, COMPRESSOR.1, COMPRESSOR.2).unwrap();
assert_eq!(dctx.partition_weight, vec![1100]);
assert_eq!(dctx.process_weight, vec![1100]);
assert_eq!(dctx.partition_compressor_count, vec![1]);
}
}