20260718 新增查询服务 VRP 检索接口
This commit is contained in:
parent
e80fba44c3
commit
6d9eedd680
@ -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::<u32>().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<String, String>) -> Result<ObjectFi
|
||||
Ok(filter)
|
||||
}
|
||||
|
||||
fn vrp_lookup_from_query(query: &BTreeMap<String, String>) -> Result<VrpLookup, ApiError> {
|
||||
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<String, String>) -> 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",
|
||||
|
||||
@ -1,3 +1,4 @@
|
||||
pub mod artifact_manifest;
|
||||
pub mod object_resolver;
|
||||
pub mod report_stream;
|
||||
pub mod vrp;
|
||||
|
||||
@ -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<Value>,
|
||||
pub download_stats: Option<Value>,
|
||||
pub vrps_count: u64,
|
||||
pub vrps: Vec<VrpRecord>,
|
||||
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::<CountSeq>()?.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::<CountSeq>()?.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<VrpRecord>,
|
||||
}
|
||||
|
||||
impl<'de> DeserializeSeed<'de> for VrpSequenceSeed<'_> {
|
||||
type Value = VrpSequence;
|
||||
|
||||
fn deserialize<D>(self, deserializer: D) -> Result<Self::Value, D::Error>
|
||||
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<A>(self, mut seq: A) -> Result<Self::Value, A::Error>
|
||||
where
|
||||
A: SeqAccess<'de>,
|
||||
{
|
||||
let mut count = 0u64;
|
||||
let mut records = Vec::new();
|
||||
while let Some(raw) = seq.next_element::<RawVrpRecord>()? {
|
||||
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"}
|
||||
|
||||
275
src/query/vrp.rs
Normal file
275
src/query/vrp.rs
Normal file
@ -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<String>,
|
||||
asn: u32,
|
||||
prefix: &str,
|
||||
max_length: u8,
|
||||
) -> Result<Self, String> {
|
||||
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, String> {
|
||||
IpPrefix::parse(&self.prefix)
|
||||
}
|
||||
|
||||
pub fn index_key(&self) -> Result<String, String> {
|
||||
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<Self, String> {
|
||||
let (addr, raw_length) = value
|
||||
.rsplit_once('/')
|
||||
.ok_or_else(|| format!("invalid CIDR prefix: {value}"))?;
|
||||
let ip = addr
|
||||
.parse::<IpAddr>()
|
||||
.map_err(|_| format!("invalid IP address in prefix: {value}"))?;
|
||||
let prefix_length = raw_length
|
||||
.parse::<u8>()
|
||||
.map_err(|_| format!("invalid prefix length in: {value}"))?;
|
||||
Self::new(ip, prefix_length)
|
||||
}
|
||||
|
||||
pub fn new(addr: IpAddr, prefix_length: u8) -> Result<Self, String> {
|
||||
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<Self> {
|
||||
(0..=self.prefix_length)
|
||||
.filter_map(|length| Self::new(self.addr, length).ok())
|
||||
.collect()
|
||||
}
|
||||
}
|
||||
|
||||
impl VrpLookup {
|
||||
pub fn parse_ip(value: &str) -> Result<Self, String> {
|
||||
value
|
||||
.parse::<IpAddr>()
|
||||
.map(VrpLookup::Ip)
|
||||
.map_err(|_| format!("invalid IP address: {value}"))
|
||||
}
|
||||
|
||||
pub fn parse_prefix(value: &str) -> Result<Self, String> {
|
||||
IpPrefix::parse(value).map(VrpLookup::Prefix)
|
||||
}
|
||||
|
||||
pub fn candidate_prefixes(&self) -> Vec<IpPrefix> {
|
||||
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<bool, String> {
|
||||
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<usize, String> {
|
||||
let Some(cursor) = cursor else {
|
||||
return Ok(0);
|
||||
};
|
||||
let offset = cursor
|
||||
.strip_prefix("vrp:")
|
||||
.ok_or_else(|| format!("invalid VRP cursor: {cursor}"))?
|
||||
.parse::<usize>()
|
||||
.map_err(|_| format!("invalid VRP cursor: {cursor}"))?;
|
||||
Ok(offset)
|
||||
}
|
||||
|
||||
pub fn next_cursor(offset: usize, limit: usize, total: usize) -> Option<String> {
|
||||
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);
|
||||
}
|
||||
}
|
||||
128
src/query_db.rs
128
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<String>,
|
||||
pub errors: Vec<String>,
|
||||
}
|
||||
@ -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<u32>,
|
||||
limit: usize,
|
||||
cursor: Option<&str>,
|
||||
) -> QueryDbResult<QueryPage<VrpRecord>> {
|
||||
let cf = self.cf(CF_VRPS)?;
|
||||
let mut records = BTreeMap::<String, VrpRecord>::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::<Vec<_>>();
|
||||
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<Option<RepositoryRecord>> {
|
||||
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<String>,
|
||||
}
|
||||
|
||||
@ -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");
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user