use std::sync::{Arc, mpsc};
use fisher_common::prelude::*;
use fisher_common::state::State;
use fisher_common::structs::HealthDetails;
use scheduler::{Scheduler, SchedulerInput};
#[cfg(test)] use scheduler::DebugDetails;
use timer::Timer;
use types::{Job, JobContext};
#[derive(Debug)]
pub struct Processor<S: ScriptsRepositoryTrait + 'static> {
input: mpsc::Sender<SchedulerInput<S>>,
timer: Timer,
wait: mpsc::Receiver<()>,
}
impl<S: ScriptsRepositoryTrait> Processor<S> {
pub fn new(max_threads: u16, hooks: Arc<S>, ctx: Arc<JobContext<S>>,
state: Arc<State>) -> Result<Self> {
let (input_send, input_recv) = mpsc::sync_channel(0);
let (wait_send, wait_recv) = mpsc::channel();
::std::thread::spawn(move || {
let inner = Scheduler::new(
max_threads, hooks, ctx, state,
);
input_send.send(inner.input()).unwrap();
inner.run().unwrap();
wait_send.send(()).unwrap();
});
let processor = Processor {
input: input_recv.recv()?,
timer: Timer::new(),
wait: wait_recv,
};
let api = processor.api();
processor.timer.add_task(30, move || {
let _ = api.cleanup();
})?;
Ok(processor)
}
pub fn stop(self) -> Result<()> {
self.timer.stop()?;
self.input.send(SchedulerInput::StopSignal)?;
self.wait.recv()?;
Ok(())
}
pub fn api(&self) -> ProcessorApi<S> {
ProcessorApi {
input: self.input.clone(),
}
}
}
#[derive(Debug, Clone)]
pub struct ProcessorApi<S: ScriptsRepositoryTrait> {
input: mpsc::Sender<SchedulerInput<S>>,
}
impl<S: ScriptsRepositoryTrait> ProcessorApi<S> {
#[cfg(test)]
pub fn debug_details(&self) -> Result<DebugDetails<S>> {
let (res_send, res_recv) = mpsc::channel();
self.input.send(SchedulerInput::DebugDetails(res_send))?;
Ok(res_recv.recv()?)
}
}
impl<S: ScriptsRepositoryTrait> ProcessorApiTrait<S> for ProcessorApi<S> {
fn queue(&self, job: Job<S>, priority: isize) -> Result<()> {
self.input.send(SchedulerInput::Job(job, priority))?;
Ok(())
}
fn health_details(&self) -> Result<HealthDetails> {
let (res_send, res_recv) = mpsc::channel();
self.input.send(SchedulerInput::HealthStatus(res_send))?;
Ok(res_recv.recv()?)
}
fn cleanup(&self) -> Result<()> {
self.input.send(SchedulerInput::Cleanup)?;
Ok(())
}
fn lock(&self) -> Result<()> {
self.input.send(SchedulerInput::Lock)?;
Ok(())
}
fn unlock(&self) -> Result<()> {
self.input.send(SchedulerInput::Unlock)?;
Ok(())
}
}