mirror of
https://github.com/NLnetLabs/krill.git
synced 2026-09-28 12:24:52 +02:00
Show BGP vs ROA analysis in CLI. (#233)
This commit is contained in:
@@ -13,6 +13,7 @@ use crate::commons::api::{
|
||||
AllCertAuthIssues, CaRepoDetails, CertAuthIssues, ChildCaInfo, CurrentRepoState,
|
||||
ParentCaContact, PublisherDetails, PublisherList, Token,
|
||||
};
|
||||
use crate::commons::bgp::{RoaSummary, RoaTable};
|
||||
use crate::commons::remote::rfc8183;
|
||||
use crate::commons::util::httpclient;
|
||||
use crate::constants::KRILL_CLI_API_ENV;
|
||||
@@ -211,6 +212,19 @@ impl KrillClient {
|
||||
Ok(ApiResponse::Empty)
|
||||
}
|
||||
|
||||
CaCommand::RouteAuthorizationsBgpDetails(handle) => {
|
||||
let uri = format!("api/v1/cas/{}/routes/bgp", handle);
|
||||
let roa_table = self.get_json(&uri).await?;
|
||||
Ok(ApiResponse::RouteAuthorizationsBgpDetails(roa_table))
|
||||
}
|
||||
|
||||
CaCommand::RouteAuthorizationsBgpSummary(handle) => {
|
||||
let uri = format!("api/v1/cas/{}/routes/bgp", handle);
|
||||
let roa_table: RoaTable = self.get_json(&uri).await?;
|
||||
let summary: RoaSummary = roa_table.into();
|
||||
Ok(ApiResponse::RouteAuthorizationsBgpSummary(summary))
|
||||
}
|
||||
|
||||
CaCommand::Show(handle) => {
|
||||
let uri = format!("api/v1/cas/{}", handle);
|
||||
let ca_info = self.get_json(&uri).await?;
|
||||
|
||||
@@ -607,11 +607,29 @@ impl Options {
|
||||
app.subcommand(sub)
|
||||
}
|
||||
|
||||
fn make_cas_routes_bgp_sc<'a, 'b>(app: App<'a, 'b>) -> App<'a, 'b> {
|
||||
let mut sub = SubCommand::with_name("bgp")
|
||||
.about("Show current authorizations in relation to known announcements.");
|
||||
|
||||
sub = Self::add_general_args(sub);
|
||||
sub = Self::add_my_ca_arg(sub);
|
||||
|
||||
sub = sub.arg(
|
||||
Arg::with_name("full")
|
||||
.long("full")
|
||||
.help("Show detailed view instead of summary")
|
||||
.required(false),
|
||||
);
|
||||
|
||||
app.subcommand(sub)
|
||||
}
|
||||
|
||||
fn make_cas_routes_sc<'a, 'b>(app: App<'a, 'b>) -> App<'a, 'b> {
|
||||
let mut sub = SubCommand::with_name("roas").about("Manage ROAs for your CA.");
|
||||
|
||||
sub = Self::make_cas_routes_list_sc(sub);
|
||||
sub = Self::make_cas_routes_update_sc(sub);
|
||||
sub = Self::make_cas_routes_bgp_sc(sub);
|
||||
|
||||
app.subcommand(sub)
|
||||
}
|
||||
@@ -1287,11 +1305,26 @@ impl Options {
|
||||
Ok(Options::make(general_args, command))
|
||||
}
|
||||
|
||||
fn parse_matches_cas_routes_bgp(matches: &ArgMatches) -> Result<Options, Error> {
|
||||
let general_args = GeneralArgs::from_matches(matches)?;
|
||||
let my_ca = Self::parse_my_ca(matches)?;
|
||||
|
||||
let command = if matches.is_present("full") {
|
||||
Command::CertAuth(CaCommand::RouteAuthorizationsBgpDetails(my_ca))
|
||||
} else {
|
||||
Command::CertAuth(CaCommand::RouteAuthorizationsBgpSummary(my_ca))
|
||||
};
|
||||
|
||||
Ok(Options::make(general_args, command))
|
||||
}
|
||||
|
||||
fn parse_matches_cas_routes(matches: &ArgMatches) -> Result<Options, Error> {
|
||||
if let Some(m) = matches.subcommand_matches("list") {
|
||||
Self::parse_matches_cas_routes_list(m)
|
||||
} else if let Some(m) = matches.subcommand_matches("update") {
|
||||
Self::parse_matches_cas_routes_update(m)
|
||||
} else if let Some(m) = matches.subcommand_matches("bgp") {
|
||||
Self::parse_matches_cas_routes_bgp(m)
|
||||
} else {
|
||||
Err(Error::UnrecognisedSubCommand)
|
||||
}
|
||||
@@ -1618,12 +1651,19 @@ pub enum CaCommand {
|
||||
#[display(fmt = "activate key roll for ca: '{}'", _0)]
|
||||
KeyRollActivate(Handle),
|
||||
|
||||
// Authorizations
|
||||
#[display(fmt = "list ROAS for ca: '{}'", _0)]
|
||||
RouteAuthorizationsList(Handle),
|
||||
|
||||
#[display(fmt = "Update ROAS for ca: '{}' -> {}", _0, _1)]
|
||||
RouteAuthorizationsUpdate(Handle, RoaDefinitionUpdates),
|
||||
|
||||
#[display(fmt = "Show detailed ROA vs BGP analysis for ca: '{}'", _0)]
|
||||
RouteAuthorizationsBgpDetails(Handle),
|
||||
|
||||
#[display(fmt = "Show summary of ROA vs BGP analysis for ca: '{}'", _0)]
|
||||
RouteAuthorizationsBgpSummary(Handle),
|
||||
|
||||
// Show details for this CA
|
||||
#[display(fmt = "Show details for ca: '{}'", _0)]
|
||||
Show(Handle),
|
||||
|
||||
@@ -11,6 +11,7 @@ use crate::commons::api::{
|
||||
ParentCaContact, PublisherDetails, PublisherList, RepositoryContact, RoaDefinition, ServerInfo,
|
||||
StoredEffect,
|
||||
};
|
||||
use crate::commons::bgp::{RoaSummary, RoaTable};
|
||||
use crate::commons::eventsourcing::WithStorableDetails;
|
||||
use crate::commons::remote::api::ClientInfo;
|
||||
use crate::commons::remote::rfc8183;
|
||||
@@ -30,6 +31,8 @@ pub enum ApiResponse {
|
||||
CertAuthAction(CaCommandDetails),
|
||||
CertAuths(CertAuthList),
|
||||
RouteAuthorizations(Vec<RoaDefinition>),
|
||||
RouteAuthorizationsBgpDetails(RoaTable),
|
||||
RouteAuthorizationsBgpSummary(RoaSummary),
|
||||
|
||||
ParentCaContact(ParentCaContact),
|
||||
|
||||
@@ -69,6 +72,10 @@ impl ApiResponse {
|
||||
ApiResponse::CertAuthIssues(issues) => Ok(Some(issues.report(fmt)?)),
|
||||
ApiResponse::AllCertAuthIssues(issues) => Ok(Some(issues.report(fmt)?)),
|
||||
ApiResponse::RouteAuthorizations(auths) => Ok(Some(auths.report(fmt)?)),
|
||||
ApiResponse::RouteAuthorizationsBgpDetails(table) => Ok(Some(table.report(fmt)?)),
|
||||
ApiResponse::RouteAuthorizationsBgpSummary(summary) => {
|
||||
Ok(Some(summary.report(fmt)?))
|
||||
}
|
||||
ApiResponse::ParentCaContact(contact) => Ok(Some(contact.report(fmt)?)),
|
||||
ApiResponse::ChildInfo(info) => Ok(Some(info.report(fmt)?)),
|
||||
ApiResponse::PublisherList(list) => Ok(Some(list.report(fmt)?)),
|
||||
@@ -407,6 +414,18 @@ impl Report for Vec<RoaDefinition> {
|
||||
}
|
||||
}
|
||||
|
||||
impl Report for RoaTable {
|
||||
fn text(&self) -> Result<String, ReportError> {
|
||||
Ok(self.to_string())
|
||||
}
|
||||
}
|
||||
|
||||
impl Report for RoaSummary {
|
||||
fn text(&self) -> Result<String, ReportError> {
|
||||
Ok(self.to_string())
|
||||
}
|
||||
}
|
||||
|
||||
impl Report for CaRepoDetails {
|
||||
fn text(&self) -> Result<String, ReportError> {
|
||||
let mut res = String::new();
|
||||
|
||||
@@ -31,18 +31,30 @@ impl BgpAnalyser {
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn update(&self) -> Result<(), BgpAnalyserError> {
|
||||
pub async fn update(&self) -> Result<bool, BgpAnalyserError> {
|
||||
if let Some(loader) = &self.dumploader {
|
||||
let mut seen = self.seen.write().unwrap();
|
||||
if let Some(last_time) = seen.last_updated() {
|
||||
if (last_time + Duration::minutes(BGP_RIS_REFRESH_MINUTES)) > Time::now() {
|
||||
return Ok(()); // no need to update yet
|
||||
debug!("Will not check BGP Ris Dumps until the refresh interval has passed");
|
||||
return Ok(false); // no need to update yet
|
||||
}
|
||||
}
|
||||
let announcements = loader.download_updates().await?;
|
||||
seen.update(announcements);
|
||||
if seen.equivalent(&announcements) {
|
||||
info!("BGP Ris Dumps unchanged");
|
||||
Ok(false)
|
||||
} else {
|
||||
info!(
|
||||
"Found {} announcements based on BGP Ris Dumps",
|
||||
announcements.len()
|
||||
);
|
||||
seen.update(announcements);
|
||||
Ok(true)
|
||||
}
|
||||
} else {
|
||||
Ok(false)
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub fn analyse(&self, roas: &[RoaDefinition], scope: &ResourceSet) -> RoaTable {
|
||||
@@ -83,7 +95,10 @@ impl BgpAnalyser {
|
||||
} else {
|
||||
let allows: Vec<Announcement> = covered
|
||||
.iter()
|
||||
.filter(|va| va.validity() == AnnouncementValidity::Valid)
|
||||
.filter(|va| {
|
||||
va.validity() == AnnouncementValidity::Valid
|
||||
&& va.announcement().asn() == &roa.asn() // and covered by *this* ROA
|
||||
})
|
||||
.map(|va| va.announcement())
|
||||
.collect();
|
||||
|
||||
|
||||
@@ -1,4 +1,6 @@
|
||||
use std::collections::HashSet;
|
||||
use std::fmt;
|
||||
use std::iter::FromIterator;
|
||||
use std::str::FromStr;
|
||||
|
||||
use rpki::x509::Time;
|
||||
@@ -167,6 +169,12 @@ impl Announcements {
|
||||
self.last_updated = Some(Time::now());
|
||||
}
|
||||
|
||||
pub fn equivalent(&self, announcements: &Vec<Announcement>) -> bool {
|
||||
let current_set: HashSet<&Announcement> = HashSet::from_iter(self.seen.all().into_iter());
|
||||
let new_set: HashSet<&Announcement> = HashSet::from_iter(announcements.iter());
|
||||
current_set == new_set
|
||||
}
|
||||
|
||||
pub fn all(&self) -> Vec<&Announcement> {
|
||||
self.seen.all()
|
||||
}
|
||||
|
||||
@@ -1,3 +1,6 @@
|
||||
use std::collections::HashMap;
|
||||
use std::fmt;
|
||||
|
||||
use crate::commons::api::RoaDefinition;
|
||||
use crate::commons::bgp::Announcement;
|
||||
|
||||
@@ -5,8 +8,8 @@ use crate::commons::bgp::Announcement;
|
||||
#[serde(rename_all = "snake_case")]
|
||||
pub enum RoaTableEntryState {
|
||||
RoaAuthorizing,
|
||||
RoaStale,
|
||||
RoaDisallowing,
|
||||
RoaStale,
|
||||
AnnouncementInvalidLength,
|
||||
AnnouncementInvalidAsn,
|
||||
AnnouncementNotFound,
|
||||
@@ -127,8 +130,243 @@ impl RoaTable {
|
||||
pub fn entries(&self) -> &Vec<RoaTableEntry> {
|
||||
&self.0
|
||||
}
|
||||
|
||||
fn matching_defs(&self, state: RoaTableEntryState) -> Vec<&RoaDefinition> {
|
||||
self.matching_entries(state)
|
||||
.into_iter()
|
||||
.map(|e| &e.definition)
|
||||
.collect()
|
||||
}
|
||||
|
||||
fn matching_entries(&self, state: RoaTableEntryState) -> Vec<&RoaTableEntry> {
|
||||
self.0.iter().filter(|e| e.state == state).collect()
|
||||
}
|
||||
}
|
||||
|
||||
//------------ Tests -------------------------------------------------------
|
||||
impl fmt::Display for RoaTable {
|
||||
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
|
||||
let entries = self.entries();
|
||||
|
||||
// tested as part of tests in analyser.rs
|
||||
let mut entry_map: HashMap<RoaTableEntryState, Vec<&RoaTableEntry>> = HashMap::new();
|
||||
for entry in entries.into_iter() {
|
||||
let state = entry.state();
|
||||
if !entry_map.contains_key(&state) {
|
||||
entry_map.insert(state, vec![]);
|
||||
}
|
||||
entry_map.get_mut(&state).unwrap().push(entry);
|
||||
}
|
||||
|
||||
if entry_map.contains_key(&RoaTableEntryState::RoaNoAnnouncementInfo) {
|
||||
write!(f, "no BGP announcements known")
|
||||
} else {
|
||||
if let Some(valids) = entry_map.get(&RoaTableEntryState::RoaAuthorizing) {
|
||||
writeln!(f, "Authorizations causing VALID announcements:")?;
|
||||
for roa in valids {
|
||||
writeln!(f)?;
|
||||
writeln!(f, "\tDefinition: {}", roa.definition)?;
|
||||
writeln!(f)?;
|
||||
writeln!(f, "\t\tAuthorises:")?;
|
||||
for ann in roa.authorizes.iter() {
|
||||
writeln!(f, "\t\t{}", ann)?;
|
||||
}
|
||||
|
||||
if !roa.disallows.is_empty() {
|
||||
writeln!(f)?;
|
||||
writeln!(f, "\t\tDisallows:")?;
|
||||
for ann in roa.disallows.iter() {
|
||||
writeln!(f, "\t\t{}", ann)?;
|
||||
}
|
||||
}
|
||||
}
|
||||
writeln!(f)?;
|
||||
}
|
||||
|
||||
if let Some(invalids) = entry_map.get(&RoaTableEntryState::RoaDisallowing) {
|
||||
writeln!(f, "Authorizations causing INVALID announcements only:")?;
|
||||
for roa in invalids {
|
||||
writeln!(f)?;
|
||||
writeln!(f, "\tDefinition: {}", roa.definition)?;
|
||||
writeln!(f)?;
|
||||
writeln!(f, "\t\tDisallows:")?;
|
||||
for ann in roa.disallows.iter() {
|
||||
writeln!(f, "\t\t{}", ann)?;
|
||||
}
|
||||
}
|
||||
writeln!(f)?;
|
||||
}
|
||||
|
||||
if let Some(stales) = entry_map.get(&RoaTableEntryState::RoaStale) {
|
||||
writeln!(
|
||||
f,
|
||||
"Authorizations for which no announcements are found (possibly stale):"
|
||||
)?;
|
||||
writeln!(f)?;
|
||||
for roa in stales {
|
||||
writeln!(f, "\tDefinition: {}", roa.definition)?;
|
||||
}
|
||||
writeln!(f)?;
|
||||
}
|
||||
|
||||
if let Some(invalid_asn) = entry_map.get(&RoaTableEntryState::AnnouncementInvalidAsn) {
|
||||
writeln!(f, "Announcements from an unauthorized ASN:")?;
|
||||
for ann in invalid_asn {
|
||||
writeln!(f)?;
|
||||
writeln!(f, "\tAnnouncement: {}", ann.definition)?;
|
||||
writeln!(f)?;
|
||||
writeln!(f, "\t\tDisallowed by authorization(s):")?;
|
||||
for roa in ann.disallowed_by.iter() {
|
||||
writeln!(f, "\t\t{}", roa)?;
|
||||
}
|
||||
}
|
||||
writeln!(f)?;
|
||||
}
|
||||
|
||||
if let Some(invalid_length) =
|
||||
entry_map.get(&RoaTableEntryState::AnnouncementInvalidLength)
|
||||
{
|
||||
writeln!(f, "Announcements from an authorized ASN, which are too specific (not allowed by max length):")?;
|
||||
for ann in invalid_length {
|
||||
writeln!(f)?;
|
||||
writeln!(f, "\tAnnouncement: {}", ann.definition)?;
|
||||
writeln!(f)?;
|
||||
writeln!(f, "\t\tDisallowed by authorization(s):")?;
|
||||
for roa in ann.disallowed_by.iter() {
|
||||
writeln!(f, "\t\t{}", roa)?;
|
||||
}
|
||||
}
|
||||
writeln!(f)?;
|
||||
}
|
||||
|
||||
if let Some(not_found) = entry_map.get(&RoaTableEntryState::AnnouncementNotFound) {
|
||||
writeln!(f, "Announcements which are 'not found' (not covered by any of your authorizations):")?;
|
||||
for ann in not_found {
|
||||
writeln!(f, "\tAnnouncement: {}", ann.definition)?;
|
||||
}
|
||||
writeln!(f)?;
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
//------------ RoaSummary --------------------------------------------------
|
||||
|
||||
#[derive(Clone, Debug, Deserialize, Eq, Hash, PartialEq, Serialize)]
|
||||
pub struct RoaSummary(Vec<RoaSummmaryEntry>);
|
||||
|
||||
#[derive(Clone, Debug, Deserialize, Eq, Hash, PartialEq, Serialize)]
|
||||
pub struct RoaSummmaryEntry {
|
||||
definition: RoaDefinition,
|
||||
state: RoaSummaryState,
|
||||
}
|
||||
|
||||
impl fmt::Display for RoaSummmaryEntry {
|
||||
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
|
||||
let state_str = match self.state {
|
||||
RoaSummaryState::Valid => "announcement 'valid'",
|
||||
RoaSummaryState::InvalidAsn => "announcement 'invalid': unauthorised asn",
|
||||
RoaSummaryState::InvalidLength => "announcement 'invalid': more specific than allowed",
|
||||
RoaSummaryState::NotFound => "announcement 'not found': not covered by your ROAs",
|
||||
RoaSummaryState::Stale => {
|
||||
"ROA does not cover any known announcement (stale or backup?)"
|
||||
}
|
||||
RoaSummaryState::NoInfo => "ROA exists, but no bgp info currently available",
|
||||
};
|
||||
write!(f, "{}\t{}", self.definition, state_str)
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, Deserialize, Eq, Hash, PartialEq, Serialize)]
|
||||
#[serde(rename_all = "snake_case")]
|
||||
pub enum RoaSummaryState {
|
||||
Valid,
|
||||
InvalidAsn,
|
||||
InvalidLength,
|
||||
NotFound,
|
||||
Stale,
|
||||
NoInfo,
|
||||
}
|
||||
|
||||
impl From<RoaTable> for RoaSummary {
|
||||
fn from(table: RoaTable) -> Self {
|
||||
let mut entries: Vec<RoaSummmaryEntry> = vec![];
|
||||
for valid in table
|
||||
.matching_entries(RoaTableEntryState::RoaAuthorizing)
|
||||
.into_iter()
|
||||
.flat_map(|e| &e.authorizes)
|
||||
{
|
||||
entries.push(RoaSummmaryEntry {
|
||||
definition: valid.clone().into(),
|
||||
state: RoaSummaryState::Valid,
|
||||
})
|
||||
}
|
||||
|
||||
for def in table.matching_defs(RoaTableEntryState::AnnouncementInvalidAsn) {
|
||||
entries.push(RoaSummmaryEntry {
|
||||
definition: def.clone(),
|
||||
state: RoaSummaryState::InvalidAsn,
|
||||
})
|
||||
}
|
||||
|
||||
for def in table.matching_defs(RoaTableEntryState::AnnouncementInvalidLength) {
|
||||
entries.push(RoaSummmaryEntry {
|
||||
definition: def.clone(),
|
||||
state: RoaSummaryState::InvalidLength,
|
||||
})
|
||||
}
|
||||
|
||||
for def in table.matching_defs(RoaTableEntryState::AnnouncementNotFound) {
|
||||
entries.push(RoaSummmaryEntry {
|
||||
definition: def.clone(),
|
||||
state: RoaSummaryState::NotFound,
|
||||
})
|
||||
}
|
||||
for def in table.matching_defs(RoaTableEntryState::RoaStale) {
|
||||
entries.push(RoaSummmaryEntry {
|
||||
definition: def.clone(),
|
||||
state: RoaSummaryState::Stale,
|
||||
})
|
||||
}
|
||||
for def in table.matching_defs(RoaTableEntryState::RoaNoAnnouncementInfo) {
|
||||
entries.push(RoaSummmaryEntry {
|
||||
definition: def.clone(),
|
||||
state: RoaSummaryState::NoInfo,
|
||||
})
|
||||
}
|
||||
RoaSummary(entries)
|
||||
}
|
||||
}
|
||||
|
||||
impl fmt::Display for RoaSummary {
|
||||
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
|
||||
for e in self.0.iter() {
|
||||
writeln!(f, "{}", e)?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
//------------ Tests --------------------------------------------------------
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use crate::commons::bgp::{RoaSummary, RoaTable};
|
||||
|
||||
// #[test]
|
||||
// fn print_roa_table() {
|
||||
// let json = include_str!("../../../test-resources/bgp/expected_roa_table.json");
|
||||
// let table: RoaTable = serde_json::from_str(json).unwrap();
|
||||
//
|
||||
// println!("{}", table)
|
||||
// }
|
||||
//
|
||||
// #[test]
|
||||
// fn print_roa_table_summary() {
|
||||
// let json = include_str!("../../../test-resources/bgp/expected_roa_table.json");
|
||||
// let table: RoaTable = serde_json::from_str(json).unwrap();
|
||||
// let summary: RoaSummary = table.into();
|
||||
//
|
||||
// println!("{}", summary)
|
||||
// }
|
||||
}
|
||||
|
||||
@@ -95,7 +95,7 @@ impl ConfigDefaults {
|
||||
}
|
||||
|
||||
fn bgp_risdumps_enabled() -> bool {
|
||||
false
|
||||
true
|
||||
}
|
||||
|
||||
fn bgp_risdumps_v4_uri() -> String {
|
||||
|
||||
@@ -490,6 +490,10 @@ async fn api_ca_routes(req: Request, path: &mut RequestPath, ca: Handle) -> Rout
|
||||
Method::POST => ca_routes_update(req, ca).await,
|
||||
_ => render_unknown_method(),
|
||||
},
|
||||
Some("bgp") => match *req.method() {
|
||||
Method::GET => ca_routes_bgp_analysis(req, ca).await,
|
||||
_ => render_unknown_method(),
|
||||
},
|
||||
_ => render_unknown_method(),
|
||||
}
|
||||
}
|
||||
@@ -975,8 +979,6 @@ async fn ca_kr_activate(req: Request, handle: Handle) -> RoutingResult {
|
||||
render_empty_res(req.state().read().await.ca_keyroll_activate(handle))
|
||||
}
|
||||
|
||||
//------------ Admin: Force republish ----------------------------------------
|
||||
|
||||
/// Update the route authorizations for this CA
|
||||
async fn ca_routes_update(req: Request, handle: Handle) -> RoutingResult {
|
||||
let state = req.state().clone();
|
||||
@@ -995,6 +997,11 @@ async fn ca_routes_show(req: Request, handle: Handle) -> RoutingResult {
|
||||
}
|
||||
}
|
||||
|
||||
/// Show the state of ROAs vs BGP for this CA
|
||||
async fn ca_routes_bgp_analysis(req: Request, handle: Handle) -> RoutingResult {
|
||||
render_json_res(req.state().read().await.ca_routes_bgp_analysis(&handle))
|
||||
}
|
||||
|
||||
//------------ Admin: Force republish ----------------------------------------
|
||||
|
||||
async fn republish_all(req: Request) -> RoutingResult {
|
||||
|
||||
@@ -18,7 +18,7 @@ use crate::commons::api::{
|
||||
RepositoryContact, RepositoryUpdate, RoaDefinition, RoaDefinitionUpdates, ServerInfo,
|
||||
TaCertDetails, UpdateChildRequest,
|
||||
};
|
||||
use crate::commons::bgp::BgpAnalyser;
|
||||
use crate::commons::bgp::{BgpAnalyser, RoaTable};
|
||||
use crate::commons::error::Error;
|
||||
use crate::commons::eventsourcing::CommandKey;
|
||||
use crate::commons::remote::rfc8183;
|
||||
@@ -53,6 +53,9 @@ pub struct KrillServer {
|
||||
// Handles the internal TA and/or CAs
|
||||
caserver: Arc<ca::CaServer<OpenSslSigner>>,
|
||||
|
||||
// Handles the internal TA and/or CAs
|
||||
bgp_analyser: Arc<BgpAnalyser>,
|
||||
|
||||
// Responsible for background tasks, e.g. re-publishing
|
||||
#[allow(dead_code)] // just need to keep this in scope
|
||||
scheduler: Scheduler,
|
||||
@@ -174,7 +177,7 @@ impl KrillServer {
|
||||
}
|
||||
}
|
||||
|
||||
let _bgpanalyser = Arc::new(BgpAnalyser::new(
|
||||
let bgp_analyser = Arc::new(BgpAnalyser::new(
|
||||
config.bgp_risdumps_enabled,
|
||||
&config.bgp_risdumps_v4_uri,
|
||||
&config.bgp_risdumps_v6_uri,
|
||||
@@ -184,6 +187,7 @@ impl KrillServer {
|
||||
event_queue,
|
||||
caserver.clone(),
|
||||
pubserver.clone(),
|
||||
bgp_analyser.clone(),
|
||||
ca_refresh_rate,
|
||||
);
|
||||
|
||||
@@ -199,6 +203,7 @@ impl KrillServer {
|
||||
authorizer,
|
||||
pubserver,
|
||||
caserver,
|
||||
bgp_analyser,
|
||||
scheduler,
|
||||
started: Time::now(),
|
||||
post_limits,
|
||||
@@ -674,6 +679,15 @@ impl KrillServer {
|
||||
let ca = self.caserver.get_ca(handle)?;
|
||||
Ok(ca.roa_definitions())
|
||||
}
|
||||
|
||||
pub fn ca_routes_bgp_analysis(&self, handle: &Handle) -> KrillResult<RoaTable> {
|
||||
let ca = self.caserver.get_ca(handle)?;
|
||||
let definitions = ca.roa_definitions();
|
||||
let resources = ca.all_resources();
|
||||
Ok(self
|
||||
.bgp_analyser
|
||||
.analyse(definitions.as_slice(), &resources))
|
||||
}
|
||||
}
|
||||
|
||||
/// # Handle publication requests
|
||||
|
||||
@@ -10,6 +10,7 @@ use tokio::runtime::Runtime;
|
||||
use rpki::x509::Time;
|
||||
|
||||
use crate::commons::api::Handle;
|
||||
use crate::commons::bgp::BgpAnalyser;
|
||||
use crate::commons::util::softsigner::OpenSslSigner;
|
||||
use crate::daemon::ca::CaServer;
|
||||
use crate::daemon::mq::{EventQueueListener, QueueEvent};
|
||||
@@ -31,6 +32,10 @@ pub struct Scheduler {
|
||||
/// they are not renewed within the configured grace period.
|
||||
#[allow(dead_code)] // just need to keep this in scope
|
||||
ca_refresh_sh: ScheduleHandle,
|
||||
|
||||
/// Responsible for refreshing announcement information
|
||||
#[allow(dead_code)] // just need to keep this in scope
|
||||
announcements_refresh_sh: ScheduleHandle,
|
||||
}
|
||||
|
||||
impl Scheduler {
|
||||
@@ -38,16 +43,19 @@ impl Scheduler {
|
||||
event_queue: Arc<EventQueueListener>,
|
||||
caserver: Arc<CaServer<OpenSslSigner>>,
|
||||
pubserver: Option<Arc<PubServer>>,
|
||||
bgp_analyser: Arc<BgpAnalyser>,
|
||||
ca_refresh_rate: u32,
|
||||
) -> Self {
|
||||
let event_sh = make_event_sh(event_queue, caserver.clone(), pubserver);
|
||||
let republish_sh = make_republish_sh(caserver.clone());
|
||||
let ca_refresh_sh = make_ca_refresh_sh(caserver, ca_refresh_rate);
|
||||
let announcements_refresh_sh = make_announcements_refresh_sh(bgp_analyser);
|
||||
|
||||
Scheduler {
|
||||
event_sh,
|
||||
republish_sh,
|
||||
ca_refresh_sh,
|
||||
announcements_refresh_sh,
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -207,3 +215,16 @@ fn make_ca_refresh_sh(caserver: Arc<CaServer<OpenSslSigner>>, refresh_rate: u32)
|
||||
});
|
||||
scheduler.watch_thread(Duration::from_millis(100))
|
||||
}
|
||||
|
||||
fn make_announcements_refresh_sh(bgp_analyser: Arc<BgpAnalyser>) -> ScheduleHandle {
|
||||
let mut scheduler = clokwerk::Scheduler::new();
|
||||
scheduler.every(1.seconds()).run(move || {
|
||||
let mut rt = Runtime::new().unwrap();
|
||||
rt.block_on(async {
|
||||
if let Err(e) = bgp_analyser.update().await {
|
||||
error!("Failed to update BGP announcements: {}", e)
|
||||
}
|
||||
})
|
||||
});
|
||||
scheduler.watch_thread(Duration::from_millis(100))
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user