impl postgres pool conn
This commit is contained in:
@@ -0,0 +1,27 @@
|
||||
use r2d2;
|
||||
|
||||
use diesel::pg::PgConnection;
|
||||
use diesel::r2d2::ConnectionManager;
|
||||
use lazy_static::lazy_static;
|
||||
|
||||
type Pool = r2d2::Pool<ConnectionManager<PgConnection>>;
|
||||
pub type DbConnection = r2d2::PooledConnection<ConnectionManager<PgConnection>>;
|
||||
|
||||
lazy_static! {
|
||||
static ref POOL: Pool = {
|
||||
let database_url = "postgres://rust:rust@localhost:5433/dispatcher";
|
||||
let manager = ConnectionManager::<PgConnection>::new(database_url);
|
||||
Pool::new(manager).expect("failed to create db pool")
|
||||
};
|
||||
}
|
||||
|
||||
pub fn init_database_pool() {
|
||||
lazy_static::initialize(&POOL);
|
||||
}
|
||||
|
||||
pub fn get_connection() -> Result<DbConnection, String> {
|
||||
match POOL.get() {
|
||||
Ok(conn) => Ok(conn),
|
||||
Err(e) => Err(format!("unable to get pool connection: {}", e)),
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,3 @@
|
||||
mod client;
|
||||
|
||||
pub use client::{get_connection, init_database_pool};
|
||||
@@ -0,0 +1,3 @@
|
||||
mod storer;
|
||||
|
||||
pub use storer::StorerError;
|
||||
@@ -0,0 +1,13 @@
|
||||
use thiserror::Error;
|
||||
|
||||
#[derive(Error, Debug)]
|
||||
pub enum StorerError {
|
||||
#[error("unable to insert `{0}`")]
|
||||
InsertError(String),
|
||||
#[error("unable to save `{0}`")]
|
||||
SaveError(String),
|
||||
#[error("unable to delete `{0}`")]
|
||||
DeleteError(String),
|
||||
#[error("unknown saver error")]
|
||||
Unknown,
|
||||
}
|
||||
@@ -1,4 +1,7 @@
|
||||
mod database;
|
||||
mod error;
|
||||
mod message;
|
||||
mod model;
|
||||
mod queue;
|
||||
mod worker;
|
||||
|
||||
@@ -6,6 +9,7 @@ use std::sync::mpsc::{channel, Sender};
|
||||
use std::thread;
|
||||
use std::time;
|
||||
|
||||
use database::{get_connection, init_database_pool};
|
||||
use message::Message;
|
||||
use queue::QueueController;
|
||||
|
||||
@@ -41,6 +45,8 @@ impl<T: std::fmt::Debug + std::marker::Send + 'static> Dispatch<T> {
|
||||
}
|
||||
|
||||
fn main() {
|
||||
init_database_pool();
|
||||
|
||||
let dispatch: Dispatch<Message> = Dispatch::new();
|
||||
|
||||
let wait = time::Duration::from_secs(2);
|
||||
|
||||
@@ -0,0 +1,79 @@
|
||||
use crate::database::get_connection;
|
||||
use std::time::Duration;
|
||||
use std::{thread, time};
|
||||
|
||||
use super::{Runner, RunnerStatus, Storer};
|
||||
use crate::error::StorerError;
|
||||
|
||||
#[derive(Debug)]
|
||||
pub enum JobAction {
|
||||
MegaportDeploy,
|
||||
MegaportUndeploy,
|
||||
MegaportCheckDeploy,
|
||||
MegaportCheckUndeploy,
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct Job {
|
||||
id: i32,
|
||||
action: JobAction,
|
||||
status: RunnerStatus,
|
||||
}
|
||||
|
||||
impl Job {
|
||||
pub fn new(id: i32, action: JobAction) -> Self {
|
||||
Self {
|
||||
id: id,
|
||||
action: action,
|
||||
status: RunnerStatus::Pending,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn set_status(&mut self, status: RunnerStatus) {
|
||||
self.status = status;
|
||||
}
|
||||
}
|
||||
|
||||
impl Runner for Job {
|
||||
fn run(&mut self) {
|
||||
self.set_status(RunnerStatus::Running);
|
||||
|
||||
match self.action {
|
||||
JobAction::MegaportDeploy => {
|
||||
println!("[job({:?})] deploying megaport...", self);
|
||||
}
|
||||
JobAction::MegaportUndeploy => {
|
||||
println!("[job({:?})] undeploying megaport...", self);
|
||||
}
|
||||
JobAction::MegaportCheckDeploy => {
|
||||
println!("[job({:?})] ckecking megaport deployment...", self);
|
||||
}
|
||||
JobAction::MegaportCheckUndeploy => {
|
||||
println!("[job({:?})] ckecking megaport undeployment...", self);
|
||||
}
|
||||
}
|
||||
|
||||
let wait = time::Duration::from_millis(1);
|
||||
thread::sleep(wait);
|
||||
|
||||
self.set_status(RunnerStatus::Success);
|
||||
}
|
||||
|
||||
fn set_status(&mut self, status: RunnerStatus) {
|
||||
self.status = status;
|
||||
}
|
||||
}
|
||||
|
||||
impl Storer for Job {
|
||||
fn insert(&self) -> Result<(), StorerError> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn save(&self) -> Result<(), StorerError> {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn delete(&self) -> Result<(), StorerError> {
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,5 @@
|
||||
mod job;
|
||||
mod traits;
|
||||
|
||||
pub use job::{Job, JobAction};
|
||||
pub use traits::{Runner, RunnerStatus, Storer};
|
||||
@@ -0,0 +1,20 @@
|
||||
use crate::error::StorerError;
|
||||
|
||||
#[derive(Debug)]
|
||||
pub enum RunnerStatus {
|
||||
Pending,
|
||||
Running,
|
||||
Failed,
|
||||
Success,
|
||||
}
|
||||
|
||||
pub trait Runner {
|
||||
fn run(&mut self);
|
||||
fn set_status(&mut self, status: RunnerStatus);
|
||||
}
|
||||
|
||||
pub trait Storer {
|
||||
fn insert(&self) -> Result<(), StorerError>;
|
||||
fn save(&self) -> Result<(), StorerError>;
|
||||
fn delete(&self) -> Result<(), StorerError>;
|
||||
}
|
||||
+42
-87
@@ -1,84 +1,34 @@
|
||||
use std::collections::VecDeque;
|
||||
use std::sync::{Arc, Mutex, Condvar};
|
||||
use std::sync::{Arc, Condvar, Mutex};
|
||||
use std::thread;
|
||||
use std::time;
|
||||
use std::time::Duration;
|
||||
|
||||
use crate::message::Message;
|
||||
use crate::model::*;
|
||||
use crate::queue::QueueStatus;
|
||||
|
||||
struct Queue {
|
||||
status: Mutex<QueueStatus>,
|
||||
content: Mutex<VecDeque<Job>>,
|
||||
struct Queue<T: Runner> {
|
||||
content: Mutex<VecDeque<T>>,
|
||||
not_empty: Condvar,
|
||||
}
|
||||
|
||||
impl Queue {
|
||||
impl<T: Runner> Queue<T> {
|
||||
fn new() -> Self {
|
||||
Self {
|
||||
status: Mutex::new(QueueStatus::Pending),
|
||||
content: Mutex::new(VecDeque::new()),
|
||||
not_empty: Condvar::new(),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn set_status(&mut self, status: QueueStatus) {
|
||||
let mut s = self.status.lock().unwrap();
|
||||
*s = status;
|
||||
}
|
||||
|
||||
pub fn add_job(&self, job: Job) {
|
||||
pub fn add_runner(&self, runner: T) {
|
||||
let mut c = self.content.lock().unwrap();
|
||||
c.push_back(job);
|
||||
c.push_back(runner);
|
||||
|
||||
self.not_empty.notify_one();
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Copy, Clone)]
|
||||
pub enum JobAction {
|
||||
MegaportDeploy,
|
||||
MegaportUndeploy,
|
||||
MegaportCheckDeploy,
|
||||
MegaportCheckUndeploy
|
||||
}
|
||||
|
||||
pub struct Job {
|
||||
id: u32,
|
||||
action: JobAction,
|
||||
data: String,
|
||||
}
|
||||
|
||||
impl Job {
|
||||
fn new(id: u32, action: JobAction, data: &str) -> Self {
|
||||
Self {
|
||||
id: id,
|
||||
action: action,
|
||||
data: data.to_string()
|
||||
}
|
||||
}
|
||||
|
||||
fn run(&self) {
|
||||
let id = self.id.clone();
|
||||
match self.action {
|
||||
JobAction::MegaportDeploy => {
|
||||
println!("[job({})] deploying megaport...", id);
|
||||
},
|
||||
JobAction::MegaportUndeploy => {
|
||||
println!("[job({})] undeploying megaport...", id);
|
||||
},
|
||||
JobAction::MegaportCheckDeploy => {
|
||||
println!("[job({})] ckecking megaport deployment...", id);
|
||||
}
|
||||
JobAction::MegaportCheckUndeploy => {
|
||||
println!("[job({})] ckecking megaport undeployment...", id);
|
||||
}
|
||||
}
|
||||
let wait = time::Duration::from_millis(1);
|
||||
thread::sleep(wait);
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Copy, Clone, PartialEq)]
|
||||
pub enum WorkerStatus {
|
||||
Pending,
|
||||
@@ -112,11 +62,18 @@ pub struct Manager {
|
||||
}
|
||||
|
||||
impl Manager {
|
||||
fn new(name :&str) -> Self {
|
||||
Self { name: name.to_string(), workers: vec![] }
|
||||
fn new(name: &str) -> Self {
|
||||
Self {
|
||||
name: name.to_string(),
|
||||
workers: vec![],
|
||||
}
|
||||
}
|
||||
|
||||
fn launch_workers(&mut self, nb_workers: u32, shared_queue: Arc<Queue>) {
|
||||
fn launch_workers<T: Runner + std::marker::Send + 'static>(
|
||||
&mut self,
|
||||
nb_workers: u32,
|
||||
shared_queue: Arc<Queue<T>>,
|
||||
) {
|
||||
for i in 0..nb_workers {
|
||||
let shared_worker = Arc::new(Mutex::new(Worker::new()));
|
||||
|
||||
@@ -126,31 +83,29 @@ impl Manager {
|
||||
let queue = shared_queue.clone();
|
||||
let name = self.name.clone();
|
||||
|
||||
thread::spawn(move || {
|
||||
loop {
|
||||
let mut c = queue.content.lock().unwrap();
|
||||
thread::spawn(move || loop {
|
||||
let mut guard = queue.content.lock().unwrap();
|
||||
|
||||
let job = loop {
|
||||
if let Some(job) = c.pop_front() {
|
||||
break job;
|
||||
} else {
|
||||
c = queue.not_empty.wait(c).unwrap();
|
||||
}
|
||||
};
|
||||
let mut runner = loop {
|
||||
if let Some(r) = guard.pop_front() {
|
||||
break r;
|
||||
} else {
|
||||
guard = queue.not_empty.wait(guard).unwrap();
|
||||
}
|
||||
};
|
||||
|
||||
drop(c);
|
||||
drop(guard);
|
||||
|
||||
let mut guard = worker.lock().unwrap();
|
||||
guard.set_status(WorkerStatus::Running);
|
||||
drop(guard);
|
||||
let mut guard = worker.lock().unwrap();
|
||||
guard.set_status(WorkerStatus::Running);
|
||||
drop(guard);
|
||||
|
||||
println!("[worker({} - {})] launching job...", name, i);
|
||||
job.run();
|
||||
println!("[worker({} - {})] launching job...", name, i);
|
||||
runner.run();
|
||||
|
||||
let mut guard = worker.lock().unwrap();
|
||||
guard.set_status(WorkerStatus::Pending);
|
||||
drop(guard);
|
||||
}
|
||||
let mut guard = worker.lock().unwrap();
|
||||
guard.set_status(WorkerStatus::Pending);
|
||||
drop(guard);
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -158,11 +113,11 @@ impl Manager {
|
||||
fn healthcheck(&self, target: WorkerStatus) -> bool {
|
||||
for w in &self.workers {
|
||||
if w.lock().unwrap().get_status() != target {
|
||||
return false
|
||||
return false;
|
||||
}
|
||||
}
|
||||
true
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -171,18 +126,18 @@ fn test_manager() {
|
||||
use std::time::Duration;
|
||||
|
||||
let mut m = Manager::new("deploy");
|
||||
let queue_deploy = Arc::new(Queue::new());
|
||||
let queue_check = Arc::new(Queue::new());
|
||||
let queue_deploy = Arc::new(Queue::<Job>::new());
|
||||
let queue_check = Arc::new(Queue::<Job>::new());
|
||||
|
||||
m.launch_workers(5, queue_deploy.clone());
|
||||
|
||||
|
||||
assert_eq!(5, m.workers.len());
|
||||
assert!(!m.healthcheck(WorkerStatus::Failed));
|
||||
|
||||
for i in 0..500 {
|
||||
let j = Job::new(i, JobAction::MegaportDeploy, "");
|
||||
let j = Job::new(i, JobAction::MegaportDeploy);
|
||||
|
||||
queue_deploy.add_job(j);
|
||||
queue_deploy.add_runner(j);
|
||||
}
|
||||
|
||||
let wait = time::Duration::from_millis(200);
|
||||
|
||||
Reference in New Issue
Block a user