use reifydb_core::{Result, interceptor::StandardInterceptorBuilder, util::ioc::IocContainer};
use reifydb_engine::{StandardCommandTransaction, StandardEngine};
use reifydb_sub_api::{Subsystem, SubsystemFactory};
use super::{WorkerBuilder, WorkerConfig, WorkerSubsystem};
pub type WorkerPoolConfigurator = Box<dyn FnOnce(WorkerBuilder) -> WorkerBuilder + Send>;
pub struct WorkerSubsystemFactory {
configurator: Option<WorkerPoolConfigurator>,
}
impl WorkerSubsystemFactory {
pub fn new() -> Self {
Self {
configurator: None,
}
}
pub fn with_configurator<F>(configurator: F) -> Self
where
F: FnOnce(WorkerBuilder) -> WorkerBuilder + Send + 'static,
{
Self {
configurator: Some(Box::new(configurator)),
}
}
pub fn with_config(config: WorkerConfig) -> Self {
Self::with_configurator(move |_| {
WorkerBuilder::new()
.num_workers(config.num_workers)
.max_queue_size(config.max_queue_size)
.scheduler_interval(config.scheduler_interval)
.task_timeout_warning(config.task_timeout_warning)
})
}
}
impl Default for WorkerSubsystemFactory {
fn default() -> Self {
Self::new()
}
}
impl SubsystemFactory<StandardCommandTransaction> for WorkerSubsystemFactory {
fn provide_interceptors(
&self,
builder: StandardInterceptorBuilder<StandardCommandTransaction>,
_ioc: &IocContainer,
) -> StandardInterceptorBuilder<StandardCommandTransaction> {
builder
}
fn create(self: Box<Self>, ioc: &IocContainer) -> Result<Box<dyn Subsystem>> {
let builder = if let Some(configurator) = self.configurator {
configurator(WorkerBuilder::new())
} else {
WorkerBuilder::default()
};
let engine = ioc.resolve::<StandardEngine>()?;
let config = builder.build();
let subsystem = WorkerSubsystem::new(config, engine);
Ok(Box::new(subsystem))
}
}