clean code
This commit is contained in:
+19
-14
@@ -18,7 +18,7 @@ impl<T> Queue<T> {
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Copy, Clone, PartialEq)]
|
||||
#[derive(Debug, Copy, Clone, PartialEq)]
|
||||
#[allow(dead_code)]
|
||||
pub enum WorkerStatus {
|
||||
Pending,
|
||||
@@ -27,12 +27,16 @@ pub enum WorkerStatus {
|
||||
}
|
||||
|
||||
struct Worker {
|
||||
id: u32,
|
||||
manager: String,
|
||||
status: WorkerStatus,
|
||||
}
|
||||
|
||||
impl Worker {
|
||||
fn new() -> Self {
|
||||
fn new(id: u32, manager: String) -> Self {
|
||||
Self {
|
||||
id: id,
|
||||
manager: manager,
|
||||
status: WorkerStatus::Pending,
|
||||
}
|
||||
}
|
||||
@@ -64,12 +68,11 @@ impl<T: Runner + std::marker::Send + 'static> Manager<T> {
|
||||
|
||||
pub fn launch_workers(&mut self, nb_workers: u32) {
|
||||
for i in 0..nb_workers {
|
||||
let shared_worker = Arc::new(Mutex::new(Worker::new()));
|
||||
let shared_worker = Arc::new(Mutex::new(Worker::new(i, self.name.clone())));
|
||||
self.workers.push(shared_worker.clone());
|
||||
|
||||
let worker = shared_worker.clone();
|
||||
let queue = self.queue.clone();
|
||||
let name = self.name.clone();
|
||||
|
||||
thread::spawn(move || loop {
|
||||
let mut guard = queue.content.lock().unwrap();
|
||||
@@ -84,17 +87,9 @@ impl<T: Runner + std::marker::Send + 'static> Manager<T> {
|
||||
|
||||
drop(guard);
|
||||
|
||||
let mut guard = worker.lock().unwrap();
|
||||
guard.set_status(WorkerStatus::Running);
|
||||
drop(guard);
|
||||
|
||||
println!("[worker({} - {})] launching job...", name, i);
|
||||
Manager::<T>::set_worker_status(&worker, WorkerStatus::Running);
|
||||
runner.run();
|
||||
println!("[worker({} - {})] job done", name, i);
|
||||
|
||||
let mut guard = worker.lock().unwrap();
|
||||
guard.set_status(WorkerStatus::Pending);
|
||||
drop(guard);
|
||||
Manager::<T>::set_worker_status(&worker, WorkerStatus::Pending);
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -114,4 +109,14 @@ impl<T: Runner + std::marker::Send + 'static> Manager<T> {
|
||||
}
|
||||
true
|
||||
}
|
||||
|
||||
fn set_worker_status(worker: &Arc<Mutex<Worker>>, status: WorkerStatus) {
|
||||
let mut guard = worker.lock().unwrap();
|
||||
guard.set_status(status);
|
||||
|
||||
println!(
|
||||
"[worker({} - {})] status: {:?}",
|
||||
guard.manager, guard.id, guard.status
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user