Reschedule failed sync to parent requests, use single queue for parent syncs. (#363)

This commit is contained in:
Tim Bruijnzeels
2020-12-21 15:31:39 +01:00
parent 5bbeb074f2
commit 982e1d3bff
5 changed files with 210 additions and 160 deletions
+20 -36
View File
@@ -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<AggregateStore<CertAuth>>,
locks: Arc<CaLocks>,
status_store: Arc<Mutex<StatusStore>>,
mq: Arc<MessageQueue>,
}
impl CaServer {
/// Builds a new CaServer. Will return an error if the TA store cannot be
/// initialised.
pub async fn build(
config: Arc<Config>,
events_queue: Arc<EventQueueListener>,
signer: Arc<KrillSigner>,
) -> KrillResult<Self> {
pub async fn build(config: Arc<Config>, mq: Arc<MessageQueue>, signer: Arc<KrillSigner>) -> KrillResult<Self> {
let mut ca_store = AggregateStore::<CertAuth>::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
+9 -6
View File
@@ -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()),
}
+4 -4
View File
@@ -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<Arc<PubServer>> = 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(())
}
+124 -64
View File
@@ -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<ResourceClassName, Vec<RevocationRequest>>,
),
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<ResourceClassName, Vec<RevocationRequest>>),
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<VecDeque<QueueEvent>>,
}
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<QueueEvent> {
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<CertAuth> for EventQueueListener {
fn listen(&self, _ca: &CertAuth, event: &Evt) {
impl eventsourcing::EventListener<CertAuth> 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<CertAuth> 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);
}
}
_ => {}
}
+53 -50
View File
@@ -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<EventQueueListener>,
event_queue: Arc<MessageQueue>,
caserver: Option<Arc<CaServer>>,
pubserver: Option<Arc<PubServer>>,
bgp_analyser: Arc<BgpAnalyser>,
@@ -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<EventQueueListener>,
event_queue: Arc<MessageQueue>,
caserver: Arc<CaServer>,
pubserver: Option<Arc<PubServer>>,
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<EventQueueListener>,
event_queue: &Arc<MessageQueue>,
caserver: Arc<CaServer>,
pubserver: Option<Arc<PubServer>>,
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<MessageQueue>,
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<CaServer>, 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<CaServer>, actor: Actor) -> ScheduleHandle {
})
}
fn make_cas_refresh(caserver: Arc<CaServer>, refresh_rate: u32, actor: Actor) -> ScheduleHandle {
fn make_cas_refresh(caserver: Arc<CaServer>, 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;
});
})
}