diff --git a/src/bin/rpki_query_service.rs b/src/bin/rpki_query_service.rs index 5e9947d..c90b21d 100644 --- a/src/bin/rpki_query_service.rs +++ b/src/bin/rpki_query_service.rs @@ -8,6 +8,7 @@ use std::sync::{Arc, Mutex}; use rpki::blob_store::ExternalRepoBytesDb; use rpki::query::report_stream::{ObjectFilter, ObjectScope}; +use rpki::query::vrp::VrpLookup; use rpki::query_db::{ ChainEdgeRecord, ExportJobRecord, ObjectInstanceRecord, ObjectUriIndexRecord, QueryDb, QueryDbError, ValidationExplainRecord, @@ -598,6 +599,22 @@ fn route_request( Some(run_id), ) } + ["runs", raw_run_id, "vrps", "lookup"] => { + let run_id = resolve_run(db, raw_run_id)?; + let lookup = vrp_lookup_from_query(&query)?; + let asn = query + .get("asn") + .map(|value| { + value.parse::().map_err(|_| { + ApiError::new(400, format!("invalid ASN query parameter: {value}")) + }) + }) + .transpose()?; + page_response( + db.lookup_vrps(&run_id, &lookup, asn, limit(&query), cursor(&query))?, + Some(run_id), + ) + } ["runs", raw_run_id, "repos"] => { let run_id = resolve_run(db, raw_run_id)?; page_response( @@ -2136,6 +2153,27 @@ fn object_filter_from_query(query: &BTreeMap) -> Result) -> Result { + let ip = query.get("ip"); + let prefix = query.get("prefix"); + match (ip, prefix) { + (Some(_), Some(_)) => Err(ApiError::new( + 400, + "ip and prefix query parameters are mutually exclusive", + )), + (None, None) => Err(ApiError::new( + 400, + "one of ip or prefix query parameters is required", + )), + (Some(value), None) => { + VrpLookup::parse_ip(value).map_err(|error| ApiError::new(400, error)) + } + (None, Some(value)) => { + VrpLookup::parse_prefix(value).map_err(|error| ApiError::new(400, error)) + } + } +} + fn cursor(query: &BTreeMap) -> Option<&str> { query.get("cursor").map(String::as_str) } @@ -2481,6 +2519,28 @@ mod tests { .as_u64(), Some(1) ); + let vrp_ip = test_route_request( + &db, + Some(&repo_bytes), + "/api/v1/runs/latest/vrps/lookup?ip=2001:da8::1", + ) + .expect("VRP IP lookup"); + assert_eq!(vrp_ip["data"][0]["asn"].as_u64(), Some(4538)); + assert_eq!(vrp_ip["meta"]["runId"].as_str(), Some("run_0001")); + let vrp_prefix = test_route_request( + &db, + Some(&repo_bytes), + "/api/v1/runs/latest/vrps/lookup?prefix=2001%3Ada8%3A%3A%2F32&asn=4538", + ) + .expect("VRP prefix lookup"); + assert_eq!(vrp_prefix["data"].as_array().unwrap().len(), 1); + let invalid = test_route_request( + &db, + Some(&repo_bytes), + "/api/v1/runs/latest/vrps/lookup?ip=192.0.2.1&prefix=192.0.2.0%2F24", + ) + .expect_err("conflicting VRP lookup parameters"); + assert_eq!(invalid.status, 400); for target in [ "/api/v1/runs/latest/repos", diff --git a/src/query/mod.rs b/src/query/mod.rs index b28f6e1..086aa9b 100644 --- a/src/query/mod.rs +++ b/src/query/mod.rs @@ -1,3 +1,4 @@ pub mod artifact_manifest; pub mod object_resolver; pub mod report_stream; +pub mod vrp; diff --git a/src/query/report_stream.rs b/src/query/report_stream.rs index deef72c..466c143 100644 --- a/src/query/report_stream.rs +++ b/src/query/report_stream.rs @@ -9,6 +9,7 @@ use serde_json::value::RawValue; use serde_json::{Value, json}; use sha2::{Digest, Sha256}; +use crate::query::vrp::VrpRecord; use crate::query_db::{ ObjectInstanceRecord, PublicationPointRecord, QUERY_DB_SCHEMA_VERSION, QueryDbError, QueryDbResult, QueryPage, RepositoryRecord, StatsRecord, @@ -23,6 +24,7 @@ pub struct ReportSummary { pub query_audit: Option, pub download_stats: Option, pub vrps_count: u64, + pub vrps: Vec, pub aspas_count: u64, pub warnings_count: u64, pub objects_count: u64, @@ -287,7 +289,13 @@ impl<'de> Visitor<'de> for ReportSummaryVisitor<'_> { out.repos = pps.repos; out.stats = pps.stats; } - "vrps" => out.vrps_count = map.next_value::()?.0, + "vrps" => { + let vrps = map.next_value_seed(VrpSequenceSeed { + run_id: self.run_id, + })?; + out.vrps_count = vrps.count; + out.vrps = vrps.records; + } "aspas" => out.aspas_count = map.next_value::()?.0, "download_stats" => out.download_stats = Some(map.next_value()?), "queryAudit" => out.query_audit = Some(map.next_value()?), @@ -325,6 +333,62 @@ struct PublicationPointsSummary { objects_count: u64, } +#[derive(Deserialize)] +struct RawVrpRecord { + asn: u32, + prefix: String, + max_length: u8, +} + +struct VrpSequenceSeed<'a> { + run_id: &'a str, +} + +struct VrpSequence { + count: u64, + records: Vec, +} + +impl<'de> DeserializeSeed<'de> for VrpSequenceSeed<'_> { + type Value = VrpSequence; + + fn deserialize(self, deserializer: D) -> Result + where + D: Deserializer<'de>, + { + deserializer.deserialize_seq(VrpSequenceVisitor { + run_id: self.run_id, + }) + } +} + +struct VrpSequenceVisitor<'a> { + run_id: &'a str, +} + +impl<'de> Visitor<'de> for VrpSequenceVisitor<'_> { + type Value = VrpSequence; + + fn expecting(&self, formatter: &mut std::fmt::Formatter) -> std::fmt::Result { + formatter.write_str("VRP array") + } + + fn visit_seq(self, mut seq: A) -> Result + where + A: SeqAccess<'de>, + { + let mut count = 0u64; + let mut records = Vec::new(); + while let Some(raw) = seq.next_element::()? { + let record = VrpRecord::from_report(self.run_id, raw.asn, &raw.prefix, raw.max_length) + .map_err(de::Error::custom)?; + count += 1; + records.push(record); + } + Ok(VrpSequence { count, records }) + } +} + struct PublicationPointsSummaryVisitor<'a> { run_id: &'a str, } @@ -1247,7 +1311,10 @@ mod tests { {"rsync_uri":"rsync://repo.example/rpki/a.roa","sha256_hex":"22","kind":"roa","result":"error","detail":"bad roa"} ] }], - "vrps":[{},{}], + "vrps":[ + {"asn":64496,"prefix":"192.0.2.0/24","max_length":24}, + {"asn":64497,"prefix":"2001:db8::/32","max_length":48} + ], "aspas":[{}], "download_stats":{"eventsTotal":0}, "queryAudit":{"status":"complete","eventsPath":"validation-events.jsonl"} diff --git a/src/query/vrp.rs b/src/query/vrp.rs new file mode 100644 index 0000000..5b66456 --- /dev/null +++ b/src/query/vrp.rs @@ -0,0 +1,275 @@ +use std::net::{IpAddr, Ipv4Addr, Ipv6Addr}; + +use serde::{Deserialize, Serialize}; + +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct VrpRecord { + pub run_id: String, + pub asn: u32, + pub prefix: String, + pub max_length: u8, +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct IpPrefix { + addr: IpAddr, + prefix_length: u8, +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub enum VrpLookup { + Ip(IpAddr), + Prefix(IpPrefix), +} + +impl VrpRecord { + pub fn from_report( + run_id: impl Into, + asn: u32, + prefix: &str, + max_length: u8, + ) -> Result { + let network = IpPrefix::parse(prefix)?; + if max_length < network.prefix_length { + return Err(format!( + "VRP max_length {max_length} is shorter than prefix length {}", + network.prefix_length + )); + } + if max_length > network.max_prefix_length() { + return Err(format!( + "VRP max_length {max_length} exceeds address family limit {}", + network.max_prefix_length() + )); + } + Ok(Self { + run_id: run_id.into(), + asn, + prefix: network.to_cidr(), + max_length, + }) + } + + pub fn prefix_data(&self) -> Result { + IpPrefix::parse(&self.prefix) + } + + pub fn index_key(&self) -> Result { + let prefix = self.prefix_data()?; + Ok(format!( + "vrp/{}/{}/{:03}/{}/{:010}/{:03}", + self.run_id, + prefix.family(), + prefix.prefix_length, + prefix.network_hex(), + self.asn, + self.max_length + )) + } +} + +impl IpPrefix { + pub fn parse(value: &str) -> Result { + let (addr, raw_length) = value + .rsplit_once('/') + .ok_or_else(|| format!("invalid CIDR prefix: {value}"))?; + let ip = addr + .parse::() + .map_err(|_| format!("invalid IP address in prefix: {value}"))?; + let prefix_length = raw_length + .parse::() + .map_err(|_| format!("invalid prefix length in: {value}"))?; + Self::new(ip, prefix_length) + } + + pub fn new(addr: IpAddr, prefix_length: u8) -> Result { + if prefix_length > max_prefix_length(addr) { + return Err(format!( + "prefix length {prefix_length} exceeds address family limit {}", + max_prefix_length(addr) + )); + } + Ok(Self { + addr: network_address(addr, prefix_length), + prefix_length, + }) + } + + pub fn prefix_length(&self) -> u8 { + self.prefix_length + } + + pub fn max_prefix_length(&self) -> u8 { + max_prefix_length(self.addr) + } + + pub fn family(&self) -> &'static str { + match self.addr { + IpAddr::V4(_) => "v4", + IpAddr::V6(_) => "v6", + } + } + + pub fn network_hex(&self) -> String { + match self.addr { + IpAddr::V4(addr) => hex::encode(addr.octets()), + IpAddr::V6(addr) => hex::encode(addr.octets()), + } + } + + pub fn to_cidr(&self) -> String { + format!("{}/{}", self.addr, self.prefix_length) + } + + pub fn contains_ip(&self, ip: IpAddr) -> bool { + if std::mem::discriminant(&self.addr) != std::mem::discriminant(&ip) { + return false; + } + network_address(ip, self.prefix_length) == self.addr + } + + pub fn contains_prefix(&self, other: &Self) -> bool { + self.family() == other.family() + && self.prefix_length <= other.prefix_length + && self.contains_ip(other.addr) + } + + pub fn ancestor_prefixes(&self) -> Vec { + (0..=self.prefix_length) + .filter_map(|length| Self::new(self.addr, length).ok()) + .collect() + } +} + +impl VrpLookup { + pub fn parse_ip(value: &str) -> Result { + value + .parse::() + .map(VrpLookup::Ip) + .map_err(|_| format!("invalid IP address: {value}")) + } + + pub fn parse_prefix(value: &str) -> Result { + IpPrefix::parse(value).map(VrpLookup::Prefix) + } + + pub fn candidate_prefixes(&self) -> Vec { + match self { + VrpLookup::Ip(ip) => (0..=max_prefix_length(*ip)) + .filter_map(|length| IpPrefix::new(*ip, length).ok()) + .collect(), + VrpLookup::Prefix(prefix) => prefix.ancestor_prefixes(), + } + } + + pub fn matches(&self, record: &VrpRecord) -> Result { + let record_prefix = record.prefix_data()?; + Ok(match self { + VrpLookup::Ip(ip) => record_prefix.contains_ip(*ip), + VrpLookup::Prefix(prefix) => { + record_prefix.contains_prefix(prefix) && record.max_length >= prefix.prefix_length + } + }) + } +} + +pub fn vrp_key_prefix(run_id: &str, prefix: &IpPrefix) -> String { + format!( + "vrp/{}/{}/{:03}/{}/", + run_id, + prefix.family(), + prefix.prefix_length, + prefix.network_hex() + ) +} + +pub fn parse_cursor(cursor: Option<&str>) -> Result { + let Some(cursor) = cursor else { + return Ok(0); + }; + let offset = cursor + .strip_prefix("vrp:") + .ok_or_else(|| format!("invalid VRP cursor: {cursor}"))? + .parse::() + .map_err(|_| format!("invalid VRP cursor: {cursor}"))?; + Ok(offset) +} + +pub fn next_cursor(offset: usize, limit: usize, total: usize) -> Option { + let next = offset.saturating_add(limit); + (next < total).then(|| format!("vrp:{next}")) +} + +fn max_prefix_length(addr: IpAddr) -> u8 { + match addr { + IpAddr::V4(_) => 32, + IpAddr::V6(_) => 128, + } +} + +fn network_address(addr: IpAddr, prefix_length: u8) -> IpAddr { + match addr { + IpAddr::V4(addr) => { + let value = u32::from(addr); + let mask = if prefix_length == 0 { + 0 + } else { + u32::MAX << (32 - prefix_length) + }; + IpAddr::V4(Ipv4Addr::from(value & mask)) + } + IpAddr::V6(addr) => { + let mut octets = addr.octets(); + let full_bytes = usize::from(prefix_length / 8); + let remaining_bits = prefix_length % 8; + if full_bytes < octets.len() { + if remaining_bits > 0 { + octets[full_bytes] &= 0xff << (8 - remaining_bits); + } + let first_zero = full_bytes + usize::from(remaining_bits > 0); + for byte in &mut octets[first_zero..] { + *byte = 0; + } + } + IpAddr::V6(Ipv6Addr::from(octets)) + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn normalizes_ipv4_and_ipv6_networks() { + assert_eq!( + IpPrefix::parse("192.0.2.129/24").unwrap().to_cidr(), + "192.0.2.0/24" + ); + assert_eq!( + IpPrefix::parse("2001:db8:1::1234/48").unwrap().to_cidr(), + "2001:db8:1::/48" + ); + } + + #[test] + fn prefix_lookup_uses_covering_and_max_length_semantics() { + let lookup = VrpLookup::parse_prefix("192.0.2.0/24").unwrap(); + let covering = VrpRecord::from_report("run_1", 64496, "192.0.0.0/16", 24).unwrap(); + let too_short = VrpRecord::from_report("run_1", 64497, "192.0.0.0/16", 23).unwrap(); + let intersecting = VrpRecord::from_report("run_1", 64498, "192.0.2.128/25", 25).unwrap(); + assert!(lookup.matches(&covering).unwrap()); + assert!(!lookup.matches(&too_short).unwrap()); + assert!(!lookup.matches(&intersecting).unwrap()); + } + + #[test] + fn cursor_is_scoped_and_bounded() { + assert_eq!(parse_cursor(None).unwrap(), 0); + assert_eq!(parse_cursor(Some("vrp:12")).unwrap(), 12); + assert!(parse_cursor(Some("12")).is_err()); + assert_eq!(next_cursor(0, 10, 12).as_deref(), Some("vrp:10")); + assert_eq!(next_cursor(10, 10, 12), None); + } +} diff --git a/src/query_db.rs b/src/query_db.rs index a8a6bdb..eaaf3b8 100644 --- a/src/query_db.rs +++ b/src/query_db.rs @@ -9,10 +9,11 @@ use sha2::{Digest, Sha256}; use crate::query::artifact_manifest::build_artifact_manifest; use crate::query::report_stream::{self, ObjectLookup, ObjectScope, ReportSummary}; +use crate::query::vrp::{self, VrpLookup, VrpRecord}; use crate::blob_store::ExternalRepoBytesDb; -pub const QUERY_DB_SCHEMA_VERSION: u32 = 1; +pub const QUERY_DB_SCHEMA_VERSION: u32 = 2; pub const CF_META: &str = "meta"; pub const CF_RUNS: &str = "runs"; @@ -26,6 +27,7 @@ pub const CF_VALIDATION_EXPLAIN_CACHE: &str = "validation_explain_cache"; pub const CF_EXPORT_JOBS: &str = "export_jobs"; pub const CF_STATS: &str = "stats"; pub const CF_REASON_INDEX: &str = "reason_index"; +pub const CF_VRPS: &str = "vrps"; pub const QUERY_DB_COLUMN_FAMILIES: &[&str] = &[ CF_META, @@ -40,6 +42,7 @@ pub const QUERY_DB_COLUMN_FAMILIES: &[&str] = &[ CF_EXPORT_JOBS, CF_STATS, CF_REASON_INDEX, + CF_VRPS, ]; const KEY_SCHEMA_VERSION: &[u8] = b"schema_version"; @@ -92,6 +95,7 @@ pub struct QueryIndexSummary { pub object_instances_indexed: u64, pub object_projections_indexed: u64, pub stats_indexed: u64, + pub vrps_indexed: u64, pub latest_ready_run: Option, pub errors: Vec, } @@ -369,6 +373,52 @@ impl QueryDb { self.list_json_by_prefix(CF_REPOS, &format!("repo/{run_id}/"), limit, cursor) } + pub fn lookup_vrps( + &self, + run_id: &str, + lookup: &VrpLookup, + asn: Option, + limit: usize, + cursor: Option<&str>, + ) -> QueryDbResult> { + let cf = self.cf(CF_VRPS)?; + let mut records = BTreeMap::::new(); + for candidate in lookup.candidate_prefixes() { + let key_prefix = vrp::vrp_key_prefix(run_id, &candidate); + let mode = IteratorMode::From(key_prefix.as_bytes(), rocksdb::Direction::Forward); + for item in self.db.iterator_cf(cf, mode) { + let (key, value) = item?; + if !key.starts_with(key_prefix.as_bytes()) { + break; + } + let key_string = String::from_utf8_lossy(&key).to_string(); + let record: VrpRecord = serde_json::from_slice(&value)?; + if asn.is_some_and(|expected| expected != record.asn) + || !lookup + .matches(&record) + .map_err(QueryDbError::InvalidArtifact)? + { + continue; + } + records.insert(key_string, record); + } + } + + let offset = vrp::parse_cursor(cursor).map_err(QueryDbError::InvalidArtifact)?; + let bounded_limit = limit.clamp(1, 1000); + let total = records.len(); + let data = records + .into_values() + .skip(offset) + .take(bounded_limit) + .collect::>(); + Ok(QueryPage { + data, + next_cursor: vrp::next_cursor(offset, bounded_limit, total), + limit: bounded_limit, + }) + } + pub fn get_repo(&self, run_id: &str, repo_id: &str) -> QueryDbResult> { self.get_json_cf(CF_REPOS, repo_key(run_id, repo_id).as_bytes()) } @@ -899,12 +949,20 @@ impl QueryDb { (CF_EXPORT_JOBS, format!("export/{}/", run.run_id)), (CF_STATS, format!("stats/{}/", run.run_id)), (CF_REASON_INDEX, format!("reason/{}/", run.run_id)), + (CF_VRPS, format!("vrp/{}/", run.run_id)), ] { self.delete_prefix_range(&mut batch, cf_name, &prefix)?; } self.write_batch(batch) } + fn clear_vrp_index(&self, run_id: &str) -> QueryDbResult<()> { + let mut batch = WriteBatch::default(); + let prefix = format!("vrp/{run_id}/"); + self.delete_prefix_range(&mut batch, CF_VRPS, &prefix)?; + self.write_batch(batch) + } + fn cf(&self, name: &'static str) -> QueryDbResult<&rocksdb::ColumnFamily> { self.db .cf_handle(name) @@ -1080,6 +1138,7 @@ pub fn index_artifacts_with_open_db( summary.object_instances_indexed += run_summary.object_instances_indexed; summary.object_projections_indexed += run_summary.object_projections_indexed; summary.stats_indexed += run_summary.stats_indexed; + summary.vrps_indexed += run_summary.vrps_indexed; summary.latest_ready_run = run_summary.latest_ready_run; } Err(err) => summary.errors.push(format!("{}: {err}", run_dir.display())), @@ -1152,6 +1211,7 @@ struct SingleRunIndexSummary { object_instances_indexed: u64, object_projections_indexed: u64, stats_indexed: u64, + vrps_indexed: u64, latest_ready_run: Option, } @@ -1193,6 +1253,7 @@ fn index_run_dir( run_record.index_status = "building".to_string(); run_record.index_error = None; db.put_json_cf(CF_RUNS, run_key(&run_id).as_bytes(), &run_record)?; + db.clear_vrp_index(&run_id)?; let indexed = match write_summary_index_records(db, &run_id, &report_summary, &artifact_manifest) { @@ -1236,6 +1297,7 @@ fn index_run_dir( object_instances_indexed: indexed.object_instances_indexed, object_projections_indexed: indexed.object_projections_indexed, stats_indexed: indexed.stats_indexed, + vrps_indexed: indexed.vrps_indexed, latest_ready_run: if should_update_latest { Some(run_id) } else { @@ -1275,6 +1337,7 @@ struct IndexWriteSummary { object_instances_indexed: u64, object_projections_indexed: u64, stats_indexed: u64, + vrps_indexed: u64, } fn write_summary_index_records( @@ -1287,6 +1350,7 @@ fn write_summary_index_records( let mut batch = WriteBatch::default(); let mut pending = 0usize; let mut summary = IndexWriteSummary::default(); + let mut seen_vrp_keys = BTreeSet::new(); macro_rules! flush_if_needed { () => { @@ -1324,6 +1388,18 @@ fn write_summary_index_records( flush_if_needed!(); } + for vrp_record in &report_summary.vrps { + let key = vrp_record + .index_key() + .map_err(QueryDbError::InvalidArtifact)?; + if seen_vrp_keys.insert(key.clone()) { + put_json_batch(&mut batch, db, CF_VRPS, key.as_bytes(), vrp_record)?; + pending += 1; + summary.vrps_indexed += 1; + flush_if_needed!(); + } + } + let artifacts_value = serde_json::to_value(artifact_manifest)?; let stats = report_stream::stats_records_from_summary(run_id, report_summary, artifacts_value); for record in &stats { @@ -1584,6 +1660,56 @@ mod tests { assert_eq!(db.count_cf(CF_OBJECT_INSTANCES).unwrap(), 0); } + #[test] + fn vrp_index_supports_ip_prefix_asn_and_pagination_queries() { + let temp = tempfile::tempdir().expect("tempdir"); + let run_dir = temp.path().join("runs/run_0001"); + fs::create_dir_all(&run_dir).expect("run dir"); + write_sample_run(&run_dir, "run_0001", 1); + let mut report = read_json_file(&run_dir.join("report.json")).expect("read report"); + let duplicate = report["vrps"][0].clone(); + report["vrps"].as_array_mut().unwrap().push(duplicate); + fs::write( + run_dir.join("report.json"), + serde_json::to_vec(&report).unwrap(), + ) + .expect("write duplicate VRP report"); + + let query_db_path = temp.path().join("query-db"); + let summary = index_artifacts(&ArtifactIndexerConfig { + query_db_path: query_db_path.clone(), + run_root: Some(temp.path().to_path_buf()), + run_dir: None, + repo_bytes_db_path: None, + projection_entry_limit: 50, + min_run_seq: None, + retain_indexed_runs: None, + }) + .expect("index"); + assert_eq!(summary.vrps_indexed, 1); + + let db = QueryDb::open(&query_db_path).expect("query db"); + assert_eq!(db.count_cf(CF_VRPS).unwrap(), 1); + let ip = VrpLookup::parse_ip("192.0.2.1").expect("ip lookup"); + let page = db + .lookup_vrps("run_0001", &ip, None, 1, None) + .expect("ip query"); + assert_eq!(page.data.len(), 1); + assert_eq!(page.data[0].asn, 64496); + assert_eq!(page.data[0].prefix, "192.0.2.0/24"); + assert!(page.next_cursor.is_none()); + + let prefix = VrpLookup::parse_prefix("192.0.2.0/24").expect("prefix lookup"); + let filtered = db + .lookup_vrps("run_0001", &prefix, Some(64497), 100, None) + .expect("asn query"); + assert!(filtered.data.is_empty()); + let latest = db + .lookup_vrps("run_0001", &prefix, Some(64496), 100, None) + .expect("filtered prefix query"); + assert_eq!(latest.data.len(), 1); + } + #[test] fn failed_run_does_not_replace_previous_latest() { let temp = tempfile::tempdir().expect("tempdir");