use crate::grpc::manager::ManagerClient;
use block_list::BlockList;
use dragonfly_api::manager::v2::{ListSchedulersResponse, Scheduler};
use dragonfly_client_config::dfdaemon::Config;
use dragonfly_client_core::Result;
use dragonfly_client_util::shutdown;
use local::Local;
use regex::Regex;
use remote::Remote;
use serde::{Deserialize, Serialize};
use std::path::PathBuf;
use std::sync::Arc;
use tokio::sync::{mpsc, Mutex, RwLock};
use tracing::{debug, error, info, instrument};
pub mod block_list;
pub mod local;
pub mod remote;
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
#[serde(default)]
pub struct SchedulerClusterClientConfig {
#[serde(alias = "blockList")]
pub block_list: Option<SchedulerClusterConfigBlockList>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
#[serde(default)]
pub struct SchedulerClusterSeedClientConfig {
#[serde(alias = "blockList")]
pub block_list: Option<SchedulerClusterConfigBlockList>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
#[serde(default)]
pub struct SchedulerClusterConfigBlockList {
pub task: Option<SchedulerClusterConfigTaskBlockList>,
#[serde(alias = "persistentTask")]
pub persistent_task: Option<SchedulerClusterConfigPersistentTaskBlockList>,
#[serde(alias = "persistentCacheTask")]
pub persistent_cache_task: Option<SchedulerClusterConfigPersistentCacheTaskBlockList>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
#[serde(default)]
pub struct SchedulerClusterConfigTaskBlockList {
pub download: Option<SchedulerClusterConfigDownloadBlockList>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
#[serde(default)]
pub struct SchedulerClusterConfigPersistentTaskBlockList {
pub download: Option<SchedulerClusterConfigDownloadBlockList>,
pub upload: Option<SchedulerClusterConfigUploadBlockList>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
#[serde(default)]
pub struct SchedulerClusterConfigPersistentCacheTaskBlockList {
pub download: Option<SchedulerClusterConfigDownloadBlockList>,
pub upload: Option<SchedulerClusterConfigUploadBlockList>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
#[serde(default)]
pub struct SchedulerClusterConfigDownloadBlockList {
pub applications: Option<Vec<String>>,
#[serde(with = "serde_regex")]
pub urls: Vec<Regex>,
pub tags: Option<Vec<String>>,
pub priorities: Option<Vec<i32>>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
#[serde(default)]
pub struct SchedulerClusterConfigUploadBlockList {
pub applications: Option<Vec<String>>,
#[serde(with = "serde_regex")]
pub urls: Vec<Regex>,
pub tags: Option<Vec<String>>,
}
#[derive(Default)]
pub struct Data {
pub schedulers: ListSchedulersResponse,
pub available_schedulers: Vec<Scheduler>,
pub available_scheduler_cluster_id: Option<u64>,
pub client_config: Option<SchedulerClusterClientConfig>,
pub seed_client_config: Option<SchedulerClusterSeedClientConfig>,
}
enum Backend {
Remote(Remote),
Local(Local),
}
pub struct Dynconfig {
pub data: Arc<RwLock<Data>>,
pub block_list: Arc<BlockList>,
config: Arc<Config>,
backend: Backend,
mutex: Mutex<()>,
shutdown: shutdown::Shutdown,
_shutdown_complete: mpsc::UnboundedSender<()>,
}
impl Dynconfig {
pub async fn new(
config: Arc<Config>,
dynconfig_path: PathBuf,
shutdown: shutdown::Shutdown,
shutdown_complete_tx: mpsc::UnboundedSender<()>,
) -> Result<Self> {
let data = Arc::new(RwLock::new(Data::default()));
let dc = Dynconfig {
block_list: Arc::new(BlockList::new(config.clone(), data.clone())),
data,
config: config.clone(),
backend: Self::backend(config, dynconfig_path.clone()).await?,
mutex: Mutex::new(()),
shutdown,
_shutdown_complete: shutdown_complete_tx,
};
dc.refresh().await?;
Ok(dc)
}
async fn backend(config: Arc<Config>, dynconfig_path: PathBuf) -> Result<Backend> {
match config.manager.addr {
Some(ref addr) => {
let manager_client = ManagerClient::new(&config.manager, addr.clone())
.await
.inspect_err(|err| {
error!("initialize manager client failed: {}", err);
})?;
info!("refresh dynamic configuration from manager {}", addr);
Ok(Backend::Remote(Remote::new(
config.clone(),
Arc::new(manager_client),
)))
}
None => {
info!(
"manager address is not configured, load dynamic configuration from {}",
dynconfig_path.display()
);
let local = Local::new(config.clone(), dynconfig_path);
local.generate_default().await?;
Ok(Backend::Local(local))
}
}
}
pub async fn run(&self) {
let mut shutdown = self.shutdown.clone();
let mut interval = tokio::time::interval(self.config.dynconfig.refresh_interval);
loop {
tokio::select! {
_ = interval.tick() => {
match self.refresh().await {
Err(err) => error!("refresh dynconfig failed: {}", err),
Ok(_) => debug!("refresh dynconfig success"),
}
}
_ = shutdown.recv() => {
info!("dynconfig server shutting down");
return
}
}
}
}
#[instrument(skip_all)]
pub async fn refresh(&self) -> Result<()> {
let Ok(_guard) = self.mutex.try_lock() else {
debug!("refresh is already running");
return Ok(());
};
let data = match &self.backend {
Backend::Remote(remote) => remote.refresh().await?,
Backend::Local(local) => local.refresh().await?,
};
*self.data.write().await = data;
Ok(())
}
}