diff --git a/src/daemon/ca/server.rs b/src/daemon/ca/server.rs index 71db753b..50efe5d9 100644 --- a/src/daemon/ca/server.rs +++ b/src/daemon/ca/server.rs @@ -33,7 +33,7 @@ use crate::daemon::ca::{ RtaContentRequest, RtaPrepareRequest, StatusStore, }; use crate::daemon::config::Config; -use crate::daemon::mq::EventQueueListener; +use crate::daemon::mq::MessageQueue; //------------ CaServer ------------------------------------------------------ @@ -111,16 +111,13 @@ pub struct CaServer { ca_store: Arc>, locks: Arc, status_store: Arc>, + mq: Arc, } impl CaServer { /// Builds a new CaServer. Will return an error if the TA store cannot be /// initialised. - pub async fn build( - config: Arc, - events_queue: Arc, - signer: Arc, - ) -> KrillResult { + pub async fn build(config: Arc, mq: Arc, signer: Arc) -> KrillResult { let mut ca_store = AggregateStore::::new(&config.data_dir, CASERVER_DIR)?; if config.always_recover_data { ca_store.recover()?; @@ -131,7 +128,7 @@ impl CaServer { ); ca_store.recover()?; } - ca_store.add_listener(events_queue); + ca_store.add_listener(mq.clone()); let status_store = StatusStore::new(&config.data_dir, STATUS_DIR)?; @@ -143,6 +140,7 @@ impl CaServer { ca_store: Arc::new(ca_store), locks, status_store: Arc::new(Mutex::new(status_store)), + mq, }) } @@ -258,40 +256,26 @@ impl CaServer { /// Refresh all CAs: /// - send pending requests if present, or /// - ask parent for updates and process if present - pub async fn cas_resync_all(&self, actor: &Actor) { - if let Err(e) = self.cas_resync_all_fallible(actor).await { - error!("Failed to refresh CA certificates: {}", e); - } - } - - /// Try to get updates for all embedded CAs, will skip the TA and/or CAs that - /// have no parents. Will try to process all and log possible errors, i.e. do - /// not bail out because of issues with one CA. - /// - /// This method may return with an Error in case there are issues with the KeyStore. - async fn cas_resync_all_fallible(&self, actor: &Actor) -> KrillResult<()> { - for handle in self.ca_store.list()? { - if let Ok(ca) = self.get_ca(&handle).await { - for parent in ca.parents() { - if ca.has_pending_requests(parent) { - if let Err(e) = self.send_requests(&handle, parent, actor).await { - error!( - "Failed to send pending requests for CA '{}' to parent: '{}', error: {}", - &handle, parent, e - ); - } else { - info!("Sent pending requests for CA '{}' to parent '{}'", &handle, parent); - } - } else if let Err(e) = self.get_updates_from_parent(&handle, &parent, actor).await { - error!("Failed to refresh CA certificates for {}, error: {}", &handle, e); - } else { - info!("Synchronised CA '{}' with parent '{}'", &handle, &parent); + pub async fn cas_refresh_all(&self) { + if let Ok(cas) = self.ca_store.list() { + for ca_handle in cas { + if let Ok(ca) = self.get_ca(&ca_handle).await { + for parent in ca.parents() { + self.mq.push_sync_parent(ca_handle.clone(), parent.clone()) } } } } + } - Ok(()) + pub async fn ca_sync_parent(&self, handle: &Handle, parent: &ParentHandle, actor: &Actor) -> KrillResult<()> { + let ca = self.get_ca(handle).await?; + + if ca.has_pending_requests(parent) { + self.send_requests(&handle, parent, actor).await + } else { + self.get_updates_from_parent(&handle, &parent, actor).await + } } /// Adds a child under an embedded CA diff --git a/src/daemon/http/server.rs b/src/daemon/http/server.rs index 70bb8b8d..c4dc69fd 100644 --- a/src/daemon/http/server.rs +++ b/src/daemon/http/server.rs @@ -152,13 +152,17 @@ struct ApiCallLogger { impl ApiCallLogger { fn new(req: &Request) -> Self { if log_enabled!(log::Level::Trace) { - trace!("Request: method={} path={} headers={:?}", - &req.method(), &req.path(), &req.headers()); + trace!( + "Request: method={} path={} headers={:?}", + &req.method(), + &req.path(), + &req.headers() + ); } - + ApiCallLogger { req_method: req.method().clone(), - req_path: req.path.full().to_string() + req_path: req.path.full().to_string(), } } @@ -1556,8 +1560,7 @@ async fn api_resync_all(req: Request) -> RoutingResult { async fn api_refresh_all(req: Request) -> RoutingResult { match *req.method() { Method::POST => aa!(req, CA_UPDATE, { - let actor = req.actor(); - render_empty_res(req.state().read().await.refresh_all(&actor).await) + render_empty_res(req.state().read().await.cas_refresh_all().await) }), _ => aa!(req, LOGIN, render_unknown_method()), } diff --git a/src/daemon/krillserver.rs b/src/daemon/krillserver.rs index ce61967a..b6a4956d 100644 --- a/src/daemon/krillserver.rs +++ b/src/daemon/krillserver.rs @@ -38,7 +38,7 @@ use crate::daemon::ca::{ }; use crate::daemon::config::{AuthType, Config}; use crate::daemon::http::HttpResponse; -use crate::daemon::mq::EventQueueListener; +use crate::daemon::mq::MessageQueue; use crate::daemon::scheduler::Scheduler; use crate::pubd::{PubServer, RepoStats}; use crate::publish::CaPublisher; @@ -203,7 +203,7 @@ impl KrillServer { let pubserver: Option> = pubserver.map(Arc::new); // Used to have a shared queue for the caserver and the background job scheduler. - let event_queue = Arc::new(EventQueueListener::default()); + let event_queue = Arc::new(MessageQueue::default()); let caserver = if mode.cas_enabled() { let caserver = Arc::new(ca::CaServer::build(config.clone(), event_queue.clone(), signer).await?); @@ -661,8 +661,8 @@ impl KrillServer { } /// Refresh all CAs: ask for updates and shrink as needed. - pub async fn refresh_all(&self, actor: &Actor) -> KrillEmptyResult { - self.get_caserver()?.cas_resync_all(actor).await; + pub async fn cas_refresh_all(&self) -> KrillEmptyResult { + self.get_caserver()?.cas_refresh_all().await; Ok(()) } diff --git a/src/daemon/mq.rs b/src/daemon/mq.rs index 434b40cd..fd6c7606 100644 --- a/src/daemon/mq.rs +++ b/src/daemon/mq.rs @@ -20,66 +20,65 @@ use crate::daemon::ca::{CertAuth, Evt, EvtDet}; #[derive(Clone, Debug, Eq, PartialEq)] #[allow(clippy::large_enum_variant)] pub enum QueueEvent { - Delta(Handle, u64), - ParentAdded(Handle, u64, ParentHandle), - RepositoryConfigured(Handle, u64), - RequestsPending(Handle, u64), - ResourceClassRemoved( - Handle, - u64, - ParentHandle, - HashMap>, - ), - UnexpectedKey(Handle, u64, ResourceClassName, RevocationRequest), - CleanOldRepo(Handle, u64), - ReschedulePublish(Handle, Time), ServerStarted, + + SyncRepo(Handle), + RescheduleSyncRepo(Handle, Time), + + SyncParent(Handle, ParentHandle), + RescheduleSyncParent(Handle, ParentHandle, Time), + + ResourceClassRemoved(Handle, ParentHandle, HashMap>), + UnexpectedKey(Handle, ResourceClassName, RevocationRequest), + CleanOldRepo(Handle), } impl fmt::Display for QueueEvent { fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result { match self { - QueueEvent::Delta(ca, version) => write!(f, "delta for '{}' version '{}'", ca, version), - QueueEvent::ParentAdded(ca, version, _parent) => { - write!(f, "parent added to '{}' version '{}'", ca, version) - } - QueueEvent::RepositoryConfigured(ca, version) => { - write!(f, "configured repository for '{}' version '{}'", ca, version) - } - QueueEvent::RequestsPending(ca, version) => { - write!(f, "requests pending for '{}' version '{}'", ca, version) - } - QueueEvent::ResourceClassRemoved(ca, version, _, _) => { - write!(f, "resource class removed for '{}' version '{}'", ca, version) - } - QueueEvent::UnexpectedKey(ca, version, rcn, _) => write!( - f, - "unexpected key found for '{}' version '{}' resource class: '{}'", - ca, version, rcn - ), - QueueEvent::CleanOldRepo(ca, version) => { - write!(f, "clean up old repo *if it exists* for '{}' version '{}'", ca, version) - } - QueueEvent::ReschedulePublish(ca, _time) => write!(f, "reschedule failed publication for '{}'", ca), QueueEvent::ServerStarted => write!(f, "Server just started"), + QueueEvent::SyncRepo(ca) => write!(f, "synchronize repo for '{}'", ca), + QueueEvent::RescheduleSyncRepo(ca, time) => write!( + f, + "reschedule failed synchronize repo for '{}' at: {}", + ca, + time.to_rfc3339() + ), + QueueEvent::SyncParent(ca, parent) => write!(f, "synchronize CA '{}' with parent '{}'", ca, parent), + QueueEvent::RescheduleSyncParent(ca, parent, time) => write!( + f, + "reschedule failed synchronize CA '{}' with parent '{}' for {}", + ca, + parent, + time.to_rfc3339() + ), + QueueEvent::ResourceClassRemoved(ca, _, _) => { + write!(f, "resource class removed for '{}' ", ca) + } + QueueEvent::UnexpectedKey(ca, rcn, _) => { + write!(f, "unexpected key found for '{}' resource class: '{}'", ca, rcn) + } + QueueEvent::CleanOldRepo(ca) => { + write!(f, "clean up old repo *if it exists* for '{}'", ca) + } } } } #[derive(Debug)] -pub struct EventQueueListener { +pub struct MessageQueue { q: RwLock>, } -impl Default for EventQueueListener { +impl Default for MessageQueue { fn default() -> Self { let mut vec = VecDeque::new(); vec.push_back(QueueEvent::ServerStarted); - EventQueueListener { q: RwLock::new(vec) } + MessageQueue { q: RwLock::new(vec) } } } -impl EventQueueListener { +impl MessageQueue { pub fn pop_all(&self) -> Vec { let mut res = vec![]; let mut q = self.q.write().unwrap(); @@ -89,40 +88,102 @@ impl EventQueueListener { res } + /// Add a queue event to the back of the queue UNLESS there is + /// already an equivalent event scheduled. pub fn push_back(&self, evt: QueueEvent) { - self.q.write().unwrap().push_back(evt) + let mut q = self.q.write().unwrap(); + + match &evt { + QueueEvent::SyncRepo(ca) | QueueEvent::RescheduleSyncRepo(ca, _) => { + for existing in q.iter() { + match existing { + QueueEvent::SyncRepo(existing_ca) | QueueEvent::RescheduleSyncRepo(existing_ca, _) => { + if existing_ca == ca { + debug!( + "Not (re-)scheduling publication for '{}', because event exists on queue", + ca + ); + return; + } + } + _ => {} + } + } + } + QueueEvent::SyncParent(ca, parent) | QueueEvent::RescheduleSyncParent(ca, parent, _) => { + for existing in q.iter() { + match existing { + QueueEvent::SyncParent(existing_ca, existing_parent) + | QueueEvent::RescheduleSyncParent(existing_ca, existing_parent, _) => { + if existing_ca == ca && existing_parent == parent { + debug!( + "Not (re-)scheduling sync for '{}' with parent '{}', because event exists on queue", + ca, parent + ); + return; + } + } + _ => {} + } + } + } + _ => {} + } + + q.push_back(evt); + } + + pub fn push_sync_repo(&self, ca: Handle) { + self.push_back(QueueEvent::SyncRepo(ca)) + } + + pub fn push_sync_parent(&self, ca: Handle, parent: ParentHandle) { + self.push_back(QueueEvent::SyncParent(ca, parent)) + } + + pub fn drop_sync_parent(&self, ca: &Handle, parent: &ParentHandle) { + let mut q = self.q.write().unwrap(); + + q.retain(|existing| match existing { + QueueEvent::SyncParent(ex_ca, ex_parent) | QueueEvent::RescheduleSyncParent(ex_ca, ex_parent, _) => { + ca != ex_ca || parent != ex_parent + } + _ => true, + }); } } -unsafe impl Send for EventQueueListener {} -unsafe impl Sync for EventQueueListener {} +unsafe impl Send for MessageQueue {} +unsafe impl Sync for MessageQueue {} /// Implement listening for CertAuth Published events. -impl eventsourcing::EventListener for EventQueueListener { - fn listen(&self, _ca: &CertAuth, event: &Evt) { +impl eventsourcing::EventListener for MessageQueue { + fn listen(&self, ca: &CertAuth, event: &Evt) { trace!("Seen CertAuth event '{}'", event); let handle = event.handle(); - let version = event.version(); match event.details() { EvtDet::ObjectSetUpdated(_, _) - | EvtDet::ParentRemoved(_, _) | EvtDet::KeyPendingToNew(_, _, _) | EvtDet::KeyPendingToActive(_, _, _) | EvtDet::KeyRollFinished(_, _) => { - let evt = QueueEvent::Delta(handle.clone(), version); - self.push_back(evt); + self.push_sync_repo(handle.clone()); } + + EvtDet::ParentRemoved(parent, _) => { + self.drop_sync_parent(&handle, parent); + self.push_sync_repo(handle.clone()); + } + EvtDet::ResourceClassRemoved(class_name, _delta, parent, revocations) => { - self.push_back(QueueEvent::Delta(handle.clone(), version)); + self.push_sync_repo(handle.clone()); let mut revocations_map = HashMap::new(); revocations_map.insert(class_name.clone(), revocations.clone()); self.push_back(QueueEvent::ResourceClassRemoved( handle.clone(), - version, parent.clone(), revocations_map, )) @@ -130,30 +191,29 @@ impl eventsourcing::EventListener for EventQueueListener { EvtDet::UnexpectedKeyFound(rcn, revocation) => self.push_back(QueueEvent::UnexpectedKey( handle.clone(), - version, rcn.clone(), revocation.clone(), )), EvtDet::ParentAdded(parent, _contact) => { - let evt = QueueEvent::ParentAdded(handle.clone(), version, parent.clone()); - self.push_back(evt); + self.push_sync_parent(handle.clone(), parent.clone()); } EvtDet::RepoUpdated(_) => { - let evt = QueueEvent::RepositoryConfigured(handle.clone(), version); - self.push_back(evt); + for parent in ca.parents() { + self.push_sync_parent(handle.clone(), parent.clone()); + } } - EvtDet::CertificateRequested(_, _, _) => { - let evt = QueueEvent::RequestsPending(handle.clone(), version); - self.push_back(evt); - } - EvtDet::KeyRollActivated(_, _) => { - let evt = QueueEvent::RequestsPending(handle.clone(), version); - self.push_back(evt); + EvtDet::CertificateRequested(rcn, _, _) | EvtDet::KeyRollActivated(rcn, _) => { + if let Ok(parent) = ca.parent_for_rc(rcn) { + self.push_sync_parent(handle.clone(), parent.clone()); + } } + EvtDet::CertificateReceived(_, _, _) => { - let evt = QueueEvent::CleanOldRepo(handle.clone(), version); - self.push_back(evt); + if ca.old_repository_contact().is_some() { + let evt = QueueEvent::CleanOldRepo(handle.clone()); + self.push_back(evt); + } } _ => {} } diff --git a/src/daemon/scheduler.rs b/src/daemon/scheduler.rs index cb104973..733e74ce 100644 --- a/src/daemon/scheduler.rs +++ b/src/daemon/scheduler.rs @@ -9,14 +9,15 @@ use tokio::runtime::Runtime; use rpki::x509::Time; +use crate::commons::actor::Actor; +use crate::commons::api::{Handle, ParentHandle}; use crate::commons::bgp::BgpAnalyser; -use crate::commons::{actor::Actor, api::Handle}; -use crate::constants::test_mode_enabled; +use crate::constants::{test_mode_enabled, REQUEUE_DELAY_SECONDS}; #[cfg(feature = "multi-user")] use crate::daemon::auth::common::session::LoginSessionCache; use crate::daemon::ca::CaServer; use crate::daemon::config::Config; -use crate::daemon::mq::{EventQueueListener, QueueEvent}; +use crate::daemon::mq::{MessageQueue, QueueEvent}; use crate::pubd::PubServer; use crate::publish::CaPublisher; @@ -52,7 +53,7 @@ pub struct Scheduler { impl Scheduler { pub fn build( - event_queue: Arc, + event_queue: Arc, caserver: Option>, pubserver: Option>, bgp_analyser: Arc, @@ -73,7 +74,7 @@ impl Scheduler { )); cas_republish = Some(make_cas_republish(caserver.clone(), actor.clone())); - cas_refresh = Some(make_cas_refresh(caserver.clone(), config.ca_refresh, actor.clone())); + cas_refresh = Some(make_cas_refresh(caserver.clone(), config.ca_refresh)); } let announcements_refresh = make_announcements_refresh(bgp_analyser); @@ -96,7 +97,7 @@ impl Scheduler { #[allow(clippy::cognitive_complexity)] fn make_cas_event_triggers( - event_queue: Arc, + event_queue: Arc, caserver: Arc, pubserver: Option>, actor: Actor, @@ -109,7 +110,7 @@ fn make_cas_event_triggers( match evt { QueueEvent::ServerStarted => { info!("Will re-sync all CAs with their parents and repository after startup"); - caserver.cas_resync_all(&actor).await; + caserver.cas_refresh_all().await; let publisher = CaPublisher::new(caserver.clone(), pubserver.clone()); match caserver.ca_list(&actor) { Err(e) => error!("Unable to obtain CA list: {}", e), @@ -125,17 +126,29 @@ fn make_cas_event_triggers( } } - QueueEvent::Delta(handle, _version) => { + QueueEvent::SyncRepo(handle) => { try_publish(&event_queue, caserver.clone(), pubserver.clone(), handle, &actor).await } - QueueEvent::ReschedulePublish(handle, last_try) => { - if Time::five_minutes_ago().timestamp() > last_try.timestamp() { + QueueEvent::RescheduleSyncRepo(handle, time) => { + if time > Time::now() { try_publish(&event_queue, caserver.clone(), pubserver.clone(), handle, &actor).await } else { - event_queue.push_back(QueueEvent::ReschedulePublish(handle, last_try)); + event_queue.push_back(QueueEvent::RescheduleSyncRepo(handle, time)); } } - QueueEvent::ResourceClassRemoved(handle, _, parent, revocations) => { + QueueEvent::SyncParent(ca, parent) => { + try_sync_parent(&event_queue, &caserver, ca, parent, &actor).await + } + QueueEvent::RescheduleSyncParent(ca, parent, time) => { + if time > Time::now() { + try_sync_parent(&event_queue, &caserver, ca, parent, &actor).await + } else { + event_queue.push_back(QueueEvent::RescheduleSyncParent(ca, parent, time)) + } + } + + + QueueEvent::ResourceClassRemoved(handle, parent, revocations) => { info!("Trigger send revoke requests for removed RC for '{}' under '{}'",handle,parent); if caserver.send_revoke_requests(&handle, &parent, revocations, &actor).await.is_err() { @@ -144,7 +157,7 @@ fn make_cas_event_triggers( just before removing the resource class entitlements."); } } - QueueEvent::UnexpectedKey(handle, _, rcn, revocation) => { + QueueEvent::UnexpectedKey(handle, rcn, revocation) => { info!( "Trigger sending revocation requests for unexpected key with id '{}' in RC '{}'", revocation.key(), @@ -155,39 +168,7 @@ fn make_cas_event_triggers( error!("Could not revoke unexpected surplus key at parent: {}", e); } } - QueueEvent::ParentAdded(handle, _, parent) => { - info!( - "Get updates for '{}' from added parent '{}'.", - handle, - parent - ); - if let Err(e) = caserver.get_updates_from_parent(&handle, &parent, &actor).await { - error!( - "Error getting updates for '{}', from parent '{}', error: '{}'", - &handle, &parent, e - ) - } - } - QueueEvent::RepositoryConfigured(ca, _) => { - info!("Repository configured for '{}'", ca); - if let Err(e) = caserver.get_delayed_updates(&ca, &actor).await { - error!( - "Error getting updates after configuring repository for '{}', error: '{}'", - &ca, e - ) - } - } - - QueueEvent::RequestsPending(handle, _) => { - info!("Get updates for pending requests for '{}'.", handle); - if let Err(e) = caserver.send_all_requests(&handle, &actor).await { - error!( - "Failed to send pending requests for '{}', error '{}'", - &handle, e - ); - } - } - QueueEvent::CleanOldRepo(handle, _) => { + QueueEvent::CleanOldRepo(handle) => { let publisher = CaPublisher::new(caserver.clone(), pubserver.clone()); if let Err(e) = publisher.clean_up(&handle, &actor).await { info!( @@ -208,8 +189,12 @@ fn make_cas_event_triggers( }) } +fn requeue_time() -> Time { + Time::now() + chrono::Duration::seconds(REQUEUE_DELAY_SECONDS) +} + async fn try_publish( - event_queue: &Arc, + event_queue: &Arc, caserver: Arc, pubserver: Option>, ca: Handle, @@ -223,11 +208,29 @@ async fn try_publish( error!("Failed to publish for '{}', error: {}", ca, e); } else { error!("Failed to publish for '{}' will reschedule, error: {}", ca, e); - event_queue.push_back(QueueEvent::ReschedulePublish(ca, Time::now())); + event_queue.push_back(QueueEvent::RescheduleSyncRepo(ca, requeue_time())); } } } +/// Try to synchronize a CA with its parents, reschedule if this fails +async fn try_sync_parent( + event_queue: &Arc, + caserver: &CaServer, + ca: Handle, + parent: ParentHandle, + actor: &Actor, +) { + info!("Synchronize CA '{}' with its parent '{}'", ca, parent); + if let Err(e) = caserver.ca_sync_parent(&ca, &parent, actor).await { + error!( + "Failed to synchronize CA '{}' with its parent '{}', error: {}", + ca, parent, e + ); + event_queue.push_back(QueueEvent::RescheduleSyncParent(ca, parent, requeue_time())); + } +} + fn make_cas_republish(caserver: Arc, actor: Actor) -> ScheduleHandle { SkippingScheduler::run(120, "CA certificate republish", move || { let mut rt = Runtime::new().unwrap(); @@ -240,12 +243,12 @@ fn make_cas_republish(caserver: Arc, actor: Actor) -> ScheduleHandle { }) } -fn make_cas_refresh(caserver: Arc, refresh_rate: u32, actor: Actor) -> ScheduleHandle { +fn make_cas_refresh(caserver: Arc, refresh_rate: u32) -> ScheduleHandle { SkippingScheduler::run(refresh_rate, "CA certificate refresh", move || { let mut rt = Runtime::new().unwrap(); rt.block_on(async { info!("Triggering background refresh for all CAs"); - caserver.cas_resync_all(&actor).await; + caserver.cas_refresh_all().await; }); }) }