use std::{collections::HashMap, sync::Arc, time::Duration};
use tokio::time::sleep;
use rpki::ca::{
idexchange::{CaHandle, ParentHandle},
provisioning::{ResourceClassName, RevocationRequest},
};
use crate::{
commons::{actor::Actor, api::Timestamp, bgp::BgpAnalyser, KrillResult},
constants::{
SCHEDULER_INTERVAL_RENEW_MINS, SCHEDULER_INTERVAL_REPUBLISH_MINS, SCHEDULER_RESYNC_REPO_CAS_THRESHOLD,
SCHEDULER_USE_JITTER_CAS_THRESHOLD,
},
daemon::{
ca::CaManager,
config::Config,
mq::{in_hours, in_minutes, now, Task, TaskQueue},
},
pubd::RepositoryManager,
};
#[cfg(feature = "multi-user")]
use crate::daemon::auth::common::session::LoginSessionCache;
pub struct Scheduler {
tasks: Arc<TaskQueue>,
ca_manager: Arc<CaManager>,
repo_manager: Arc<RepositoryManager>,
bgp_analyser: Arc<BgpAnalyser>,
#[cfg(feature = "multi-user")]
login_session_cache: Arc<LoginSessionCache>,
config: Arc<Config>,
system_actor: Actor,
started: Timestamp,
}
impl Scheduler {
pub fn build(
tasks: Arc<TaskQueue>,
ca_manager: Arc<CaManager>,
repo_manager: Arc<RepositoryManager>,
bgp_analyser: Arc<BgpAnalyser>,
#[cfg(feature = "multi-user")] login_session_cache: Arc<LoginSessionCache>,
config: Arc<Config>,
system_actor: Actor,
) -> Self {
Scheduler {
tasks,
ca_manager,
repo_manager,
bgp_analyser,
#[cfg(feature = "multi-user")]
login_session_cache,
config,
system_actor,
started: Timestamp::now(),
}
}
pub async fn run(&self) {
loop {
while let Some(evt) = self.tasks.pop(now()) {
if let Err(e) = match evt {
Task::QueueStartTasks => self.queue_start_tasks().await,
Task::SyncRepo { ca } => self.sync_repo(ca).await,
Task::SyncParent { ca, parent } => self.sync_parent(ca, parent).await,
Task::SuspendChildrenIfNeeded { ca } => self.suspend_children_if_needed(ca).await,
Task::RepublishIfNeeded => self.republish_if_needed().await,
Task::RenewObjectsIfNeeded => self.renew_objects_if_needed().await,
Task::RefreshAnnouncementsInfo => self.announcements_refresh().await,
#[cfg(feature = "multi-user")]
Task::SweepLoginCache => self.sweep_login_cache(),
Task::UpdateSnapshots => self.update_snapshots(),
Task::RrdpUpdateIfNeeded => self.update_rrdp_if_needed(),
Task::ResourceClassRemoved {
ca,
parent,
rcn,
revocation_requests,
} => self.resource_class_removed(ca, parent, rcn, revocation_requests).await,
Task::UnexpectedKey {
ca,
rcn,
revocation_request,
} => self.unexpected_key(ca, rcn, revocation_request).await,
} {
error!("Fatal error in scheduler: {}", e);
return;
}
}
sleep(Duration::from_millis(500)).await;
}
}
async fn queue_start_tasks(&self) -> KrillResult<()> {
let ca_list = self.ca_manager.ca_list(&self.system_actor)?;
let cas = ca_list.cas();
debug!("Adding tasks at start up");
let too_many_cas_resync_parent = cas.len() >= SCHEDULER_USE_JITTER_CAS_THRESHOLD;
let too_many_cas_resync_repo = cas.len() >= SCHEDULER_RESYNC_REPO_CAS_THRESHOLD;
for summary in cas {
let ca = self.ca_manager.get_ca(summary.handle()).await?;
let too_many_parents = ca.nr_parents() >= self.config.ca_refresh_parents_batch_size;
let use_parent_sync_jitter = too_many_cas_resync_parent || too_many_parents;
debug!(
"Adding tasks for CA {}, using jitter: {}",
ca.handle(),
use_parent_sync_jitter
);
if !too_many_cas_resync_parent && too_many_parents {
debug!(
"Will force jitter for sync between CA {} and parents. Nr of parents ({}) exceeds batch size ({})",
ca.handle(),
ca.nr_parents(),
self.config.ca_refresh_parents_batch_size
)
}
for parent in ca.parents() {
self.tasks.sync_parent(
ca.handle().clone(),
parent.clone(),
self.config.ca_refresh_start_up(use_parent_sync_jitter),
);
}
if !too_many_cas_resync_repo {
self.tasks.sync_repo(ca.handle().clone(), now());
}
if self.config.suspend_child_after_inactive_seconds().is_some() {
self.tasks.suspend_children(ca.handle().clone(), now())
}
}
self.tasks.republish_if_needed(now());
self.tasks.renew_if_needed(now());
self.tasks.refresh_announcements_info(now());
#[cfg(feature = "multi-user")]
self.tasks.sweep_login_cache(in_minutes(1));
self.tasks.update_snapshots(in_hours(24));
Ok(())
}
async fn sync_repo(&self, ca: CaHandle) -> KrillResult<()> {
debug!("Synchronize CA {} with repository", ca);
if let Err(e) = self
.ca_manager
.cas_repo_sync_single(self.repo_manager.as_ref(), &ca)
.await
{
let next = self.config.requeue_remote_failed();
error!(
"Failed to publish for '{}'. Will reschedule to: '{}'. Error: {}",
ca, next, e
);
self.tasks.sync_repo(ca, next);
}
Ok(())
}
async fn sync_parent(&self, ca: CaHandle, parent: ParentHandle) -> KrillResult<()> {
info!("Synchronize CA '{}' with its parent '{}'", ca, parent);
if let Err(e) = self.ca_manager.ca_sync_parent(&ca, &parent, &self.system_actor).await {
let next = self.config.requeue_remote_failed();
error!(
"Failed to synchronize CA '{}' with its parent '{}'. Will reschedule to: '{}'. Error: {}",
ca, parent, next, e
);
self.tasks.sync_parent(ca, parent, next);
} else {
let next = self.config.ca_refresh_next();
self.tasks.sync_parent(ca, parent, next);
}
Ok(())
}
async fn suspend_children_if_needed(&self, ca_handle: CaHandle) -> KrillResult<()> {
debug!("Verify if CA '{}' has children that need to be suspended", ca_handle);
self.ca_manager
.ca_suspend_inactive_children(&ca_handle, self.started, &self.system_actor)
.await;
self.tasks.suspend_children(ca_handle, in_hours(1));
Ok(())
}
async fn republish_if_needed(&self) -> KrillResult<()> {
let cas = self.ca_manager.republish_all(false).await?;
for ca in cas {
info!("Re-issued MFT and CRL for CA: {}", ca);
self.tasks.sync_repo(ca, now());
}
self.tasks
.republish_if_needed(in_minutes(SCHEDULER_INTERVAL_REPUBLISH_MINS));
Ok(())
}
async fn announcements_refresh(&self) -> KrillResult<()> {
if let Err(e) = self.bgp_analyser.update().await {
error!("Failed to update BGP announcements: {}", e)
}
self.tasks.refresh_announcements_info(in_minutes(10));
Ok(())
}
async fn renew_objects_if_needed(&self) -> KrillResult<()> {
self.ca_manager.renew_objects_all(&self.system_actor).await?;
self.tasks.renew_if_needed(in_minutes(SCHEDULER_INTERVAL_RENEW_MINS));
Ok(())
}
#[cfg(feature = "multi-user")]
fn sweep_login_cache(&self) -> KrillResult<()> {
if let Err(e) = self.login_session_cache.sweep() {
error!("Background sweep of session decryption cache failed: {}", e);
}
self.tasks.sweep_login_cache(in_minutes(1));
Ok(())
}
fn update_snapshots(&self) -> KrillResult<()> {
if let Err(e) = self.repo_manager.update_snapshots() {
error!("Could not update snapshots on disk! Error: {}", e);
}
self.tasks.update_snapshots(in_hours(24));
Ok(())
}
fn update_rrdp_if_needed(&self) -> KrillResult<()> {
match self.repo_manager.update_rrdp_if_needed() {
Err(e) => {
error!("Could not update RRDP deltas! Error: {}", e);
self.tasks.update_rrdp_if_needed(in_hours(1));
}
Ok(None) => {
}
Ok(Some(later_time)) => {
self.tasks.update_rrdp_if_needed(later_time.into());
}
}
Ok(())
}
async fn resource_class_removed(
&self,
ca: CaHandle,
parent: ParentHandle,
rcn: ResourceClassName,
revocation_requests: Vec<RevocationRequest>,
) -> KrillResult<()> {
info!(
"Trigger send revoke requests for removed RC for '{}' under '{}'",
ca, parent
);
let requests = HashMap::from([(rcn, revocation_requests)]);
if self
.ca_manager
.send_revoke_requests(&ca, &parent, requests)
.await
.is_err()
{
warn!(
"Could not revoke key for removed resource class. This is not \
an issue, because typically the parent will revoke our keys pro-actively, \
just before removing the resource class entitlements."
);
}
Ok(())
}
async fn unexpected_key(
&self,
ca: CaHandle,
rcn: ResourceClassName,
revocation_request: RevocationRequest,
) -> KrillResult<()> {
info!(
"Trigger sending revocation requests for unexpected key with id '{}' in RC '{}'",
revocation_request.key(),
rcn
);
if let Err(e) = self
.ca_manager
.send_revoke_unexpected_key(&ca, rcn, revocation_request)
.await
{
error!("Could not revoke unexpected surplus key at parent: {}", e);
}
Ok(())
}
}