use std::fmt;
use crate::key::Key;
pub const VERSION_CHUNK: &str = "v1";
pub const CLASS_TELEMETRY: &str = "telemetry";
pub const CLASS_STATE: &str = "state";
pub const CLASS_EVENTS: &str = "events";
pub const PLANE_RPC: &str = "@rpc";
pub const PLANE_MEDIA: &str = "@media";
pub const PLANE_BLOB: &str = "@blob";
pub const SERVICE_CATALOG: &str = "@catalog";
pub const BLOB_TIER_ARTIFACT: &str = "artifact";
pub const BLOB_TIER_TREE: &str = "tree";
pub const BLOB_TIER_STORE: &str = "store";
pub const SUBJECT_ALIVE: &str = "alive";
#[derive(Debug, thiserror::Error, PartialEq, Eq)]
pub enum KeyError {
#[error("invalid plain chunk {0:?}: must match [a-z0-9]([a-z0-9._-]*[a-z0-9])? (RFC 03 §2)")]
InvalidPlainChunk(String),
#[error("invalid verbatim chunk {0:?}: must match @[a-z0-9][a-z0-9_-]* (RFC 03 §2)")]
InvalidVerbatimChunk(String),
#[error("invalid host origin {0:?}: must match h-[0-9a-f]{{12}} (RFC 03 §1.3)")]
InvalidHostOrigin(String),
#[error("invalid producer {0:?}: {1} (RFC 03 §1.5)")]
InvalidProducer(String, &'static str),
#[error("empty subject: keys need >= 1 subject chunk (RFC 03 §1.6)")]
EmptySubject,
#[error("blob tier token expected (artifact|tree|store), got {0:?} (RFC 03 §1.5)")]
InvalidBlobTier(String),
#[error("reserved token {0:?} may not be used as a {1} (RFC 03 §3)")]
ReservedToken(String, &'static str),
#[error("not a v1 key: {0}")]
Parse(String),
}
pub fn is_valid_plain_chunk(chunk: &str) -> bool {
let bytes = chunk.as_bytes();
let alnum = |b: u8| b.is_ascii_lowercase() || b.is_ascii_digit();
match bytes {
[] => false,
[one] => alnum(*one),
[first, mid @ .., last] => {
alnum(*first)
&& alnum(*last)
&& mid
.iter()
.all(|&b| alnum(b) || b == b'.' || b == b'_' || b == b'-')
}
}
}
pub fn is_valid_verbatim_chunk(chunk: &str) -> bool {
let Some(rest) = chunk.strip_prefix('@') else {
return false;
};
let bytes = rest.as_bytes();
let alnum = |b: u8| b.is_ascii_lowercase() || b.is_ascii_digit();
match bytes {
[] => false,
[first, rest @ ..] => {
alnum(*first) && rest.iter().all(|&b| alnum(b) || b == b'_' || b == b'-')
}
}
}
pub fn is_valid_host_origin(chunk: &str) -> bool {
let Some(hex) = chunk.strip_prefix("h-") else {
return false;
};
hex.len() == 12
&& hex
.bytes()
.all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub enum Origin {
Host(crate::origin::HostId),
Service(String),
}
impl Origin {
pub fn catalog() -> Self {
Origin::Service(SERVICE_CATALOG.to_string())
}
pub fn service(name: &str) -> Result<Self, KeyError> {
if !is_valid_verbatim_chunk(name) {
return Err(KeyError::InvalidVerbatimChunk(name.to_string()));
}
Ok(Origin::Service(name.to_string()))
}
pub fn chunk(&self) -> &str {
match self {
Origin::Host(id) => id.as_str(),
Origin::Service(s) => s,
}
}
pub fn has_producer_chunk(&self) -> bool {
matches!(self, Origin::Host(_))
}
}
impl fmt::Display for Origin {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(self.chunk())
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct Producer {
name: String,
instance: Option<u32>,
}
impl Producer {
pub fn new(name: &str) -> Result<Self, KeyError> {
Self::validate_name(name)?;
Ok(Producer {
name: name.to_string(),
instance: None,
})
}
pub fn with_instance(name: &str, instance: u32) -> Result<Self, KeyError> {
Self::validate_name(name)?;
if instance == 0 {
return Err(KeyError::InvalidProducer(
name.to_string(),
"instance numbers start at 1 (the first instance uses the bare name)",
));
}
Ok(Producer {
name: name.to_string(),
instance: Some(instance),
})
}
fn validate_name(name: &str) -> Result<(), KeyError> {
if !is_valid_plain_chunk(name) {
return Err(KeyError::InvalidProducer(
name.to_string(),
"not a valid plain chunk",
));
}
if Self::split_trailing_int(name).is_some() {
return Err(KeyError::InvalidProducer(
name.to_string(),
"base names must not end in -<int> (reserved for instance suffixes)",
));
}
if name == BLOB_TIER_ARTIFACT || name == BLOB_TIER_TREE || name == BLOB_TIER_STORE {
return Err(KeyError::ReservedToken(name.to_string(), "producer name"));
}
Ok(())
}
fn split_trailing_int(chunk: &str) -> Option<(&str, u32)> {
let (base, tail) = chunk.rsplit_once('-')?;
if base.is_empty() || tail.is_empty() || !tail.bytes().all(|b| b.is_ascii_digit()) {
return None;
}
tail.parse().ok().map(|n| (base, n))
}
pub fn parse_chunk(chunk: &str) -> Result<Self, KeyError> {
if !is_valid_plain_chunk(chunk) {
return Err(KeyError::InvalidProducer(
chunk.to_string(),
"not a valid plain chunk",
));
}
match Self::split_trailing_int(chunk) {
Some((base, n)) if n >= 1 => Ok(Producer {
name: base.to_string(),
instance: Some(n),
}),
_ => Ok(Producer {
name: chunk.to_string(),
instance: None,
}),
}
}
pub fn name(&self) -> &str {
&self.name
}
pub fn instance(&self) -> Option<u32> {
self.instance
}
pub(crate) fn push_chunk(&self, out: &mut String) {
out.push_str(&self.name);
if let Some(i) = self.instance {
use std::fmt::Write as _;
let _ = write!(out, "-{i}");
}
}
pub fn chunk(&self) -> String {
match self.instance {
None => self.name.clone(),
Some(n) => format!("{}-{n}", self.name),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum Class {
Telemetry,
State,
Events,
}
impl Class {
pub fn chunk(self) -> &'static str {
match self {
Class::Telemetry => CLASS_TELEMETRY,
Class::State => CLASS_STATE,
Class::Events => CLASS_EVENTS,
}
}
pub fn from_chunk(chunk: &str) -> Option<Self> {
match chunk {
CLASS_TELEMETRY => Some(Class::Telemetry),
CLASS_STATE => Some(Class::State),
CLASS_EVENTS => Some(Class::Events),
_ => None,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum Plane {
Rpc,
Media,
Blob,
}
impl Plane {
pub fn chunk(self) -> &'static str {
match self {
Plane::Rpc => PLANE_RPC,
Plane::Media => PLANE_MEDIA,
Plane::Blob => PLANE_BLOB,
}
}
pub fn from_chunk(chunk: &str) -> Option<Self> {
match chunk {
PLANE_RPC => Some(Plane::Rpc),
PLANE_MEDIA => Some(Plane::Media),
PLANE_BLOB => Some(Plane::Blob),
_ => None,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum ClassOrPlane {
Class(Class),
Plane(Plane),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum BlobTier {
Artifact,
Tree,
Store,
}
impl BlobTier {
pub fn chunk(self) -> &'static str {
match self {
BlobTier::Artifact => BLOB_TIER_ARTIFACT,
BlobTier::Tree => BLOB_TIER_TREE,
BlobTier::Store => BLOB_TIER_STORE,
}
}
pub fn from_chunk(chunk: &str) -> Option<Self> {
match chunk {
BLOB_TIER_ARTIFACT => Some(BlobTier::Artifact),
BLOB_TIER_TREE => Some(BlobTier::Tree),
BLOB_TIER_STORE => Some(BlobTier::Store),
_ => None,
}
}
}
fn validate_subject(subject: &[&str]) -> Result<(), KeyError> {
if subject.is_empty() {
return Err(KeyError::EmptySubject);
}
for chunk in subject {
if !is_valid_plain_chunk(chunk) {
return Err(KeyError::InvalidPlainChunk((*chunk).to_string()));
}
}
Ok(())
}
fn push_key(parts: &mut String, chunk: &str) {
push_key_sep(parts);
parts.push_str(chunk);
}
fn push_key_sep(parts: &mut String) {
if !parts.is_empty() {
parts.push('/');
}
}
pub fn data_key(
origin: &Origin,
class: Class,
producer: Option<&Producer>,
subject: &[&str],
) -> Result<Key, KeyError> {
validate_subject(subject)?;
if origin.has_producer_chunk() != producer.is_some() {
return Err(KeyError::Parse(
"host origins require a producer chunk; service origins forbid one (RFC 03 §1.5)"
.to_string(),
));
}
if class == Class::State && subject.contains(&SUBJECT_ALIVE) {
return Err(KeyError::ReservedToken(
SUBJECT_ALIVE.to_string(),
"data subject chunk",
));
}
let mut key = String::new();
push_key(&mut key, VERSION_CHUNK);
push_key(&mut key, origin.chunk());
push_key(&mut key, class.chunk());
if let Some(p) = producer {
push_key_sep(&mut key);
p.push_chunk(&mut key);
}
for chunk in subject {
push_key(&mut key, chunk);
}
Ok(Key::from_canonical(key))
}
pub fn rpc_key(
origin: &Origin,
producer: Option<&Producer>,
procedure: &[&str],
) -> Result<Key, KeyError> {
validate_subject(procedure)?;
if origin.has_producer_chunk() != producer.is_some() {
return Err(KeyError::Parse(
"host origins require a producer chunk; service origins forbid one (RFC 03 §1.5)"
.to_string(),
));
}
let mut key = String::new();
push_key(&mut key, VERSION_CHUNK);
push_key(&mut key, origin.chunk());
push_key(&mut key, PLANE_RPC);
if let Some(p) = producer {
push_key_sep(&mut key);
p.push_chunk(&mut key);
}
for chunk in procedure {
push_key(&mut key, chunk);
}
Ok(Key::from_canonical(key))
}
pub fn media_key(origin: &Origin, producer: &Producer, stream: &[&str]) -> Result<Key, KeyError> {
validate_subject(stream)?;
let mut key = String::new();
push_key(&mut key, VERSION_CHUNK);
push_key(&mut key, origin.chunk());
push_key(&mut key, PLANE_MEDIA);
push_key_sep(&mut key);
producer.push_chunk(&mut key);
for chunk in stream {
push_key(&mut key, chunk);
}
Ok(Key::from_canonical(key))
}
pub fn blob_key(origin: &Origin, tier: BlobTier, rest: &[&str]) -> Result<Key, KeyError> {
validate_subject(rest)?;
let mut key = String::new();
push_key(&mut key, VERSION_CHUNK);
push_key(&mut key, origin.chunk());
push_key(&mut key, PLANE_BLOB);
push_key(&mut key, tier.chunk());
for chunk in rest {
push_key(&mut key, chunk);
}
Ok(Key::from_canonical(key))
}
pub fn alive_key(origin: &Origin, producer: Option<&Producer>) -> Result<Key, KeyError> {
if origin.has_producer_chunk() != producer.is_some() {
return Err(KeyError::Parse(
"host origins require a producer chunk; service origins forbid one (RFC 03 §1.5)"
.to_string(),
));
}
let mut key = String::new();
push_key(&mut key, VERSION_CHUNK);
push_key(&mut key, origin.chunk());
push_key(&mut key, CLASS_STATE);
if let Some(p) = producer {
push_key_sep(&mut key);
p.push_chunk(&mut key);
}
push_key(&mut key, SUBJECT_ALIVE);
Ok(Key::from_canonical(key))
}
pub fn device_alive_key(
origin: &Origin,
producer: &Producer,
device: &str,
) -> Result<Key, KeyError> {
if !is_valid_plain_chunk(device) {
return Err(KeyError::InvalidPlainChunk(device.to_string()));
}
let mut key = String::new();
push_key(&mut key, VERSION_CHUNK);
push_key(&mut key, origin.chunk());
push_key(&mut key, CLASS_STATE);
push_key_sep(&mut key);
producer.push_chunk(&mut key);
push_key(&mut key, "device");
push_key(&mut key, device);
push_key(&mut key, SUBJECT_ALIVE);
Ok(Key::from_canonical(key))
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct StructuralKey<'k> {
pub origin: Origin,
pub class: ClassOrPlane,
pub producer: Option<Producer>,
pub blob_tier: Option<BlobTier>,
pub subject: Vec<&'k str>,
}
impl StructuralKey<'_> {
pub fn remote_origin(&self) -> Option<crate::origin::RemoteOrigin> {
match &self.origin {
Origin::Host(id) => Some(crate::origin::RemoteOrigin::from_host(id.clone())),
Origin::Service(_) => None,
}
}
}
pub fn parse(key: &str) -> Result<StructuralKey<'_>, KeyError> {
let mut chunks = key.split('/');
let version = chunks
.next()
.ok_or_else(|| KeyError::Parse("empty key".into()))?;
if version != VERSION_CHUNK {
return Err(KeyError::Parse(format!(
"expected {VERSION_CHUNK} first, got {version:?}"
)));
}
let origin_chunk = chunks
.next()
.ok_or_else(|| KeyError::Parse("missing origin chunk".into()))?;
let origin = if is_valid_host_origin(origin_chunk) {
Origin::Host(crate::origin::HostId::parse(origin_chunk).expect("validated"))
} else if is_valid_verbatim_chunk(origin_chunk) {
Origin::Service(origin_chunk.to_string())
} else {
return Err(KeyError::InvalidHostOrigin(origin_chunk.to_string()));
};
let class_chunk = chunks
.next()
.ok_or_else(|| KeyError::Parse("missing class chunk".into()))?;
let class = if let Some(c) = Class::from_chunk(class_chunk) {
ClassOrPlane::Class(c)
} else if let Some(p) = Plane::from_chunk(class_chunk) {
ClassOrPlane::Plane(p)
} else {
return Err(KeyError::Parse(format!(
"unknown class/plane chunk {class_chunk:?}"
)));
};
let mut producer = None;
let mut blob_tier = None;
match (&origin, &class) {
(_, ClassOrPlane::Plane(Plane::Blob)) => {
let tier = chunks
.next()
.ok_or_else(|| KeyError::Parse("missing blob tier".into()))?;
blob_tier = Some(
BlobTier::from_chunk(tier)
.ok_or_else(|| KeyError::InvalidBlobTier(tier.to_string()))?,
);
}
(Origin::Host(_), _) => {
let chunk = chunks
.next()
.ok_or_else(|| KeyError::Parse("missing producer chunk".into()))?;
producer = Some(Producer::parse_chunk(chunk)?);
}
(Origin::Service(_), _) => {}
}
let subject: Vec<&str> = chunks.collect();
if subject.is_empty() {
return Err(KeyError::EmptySubject);
}
Ok(StructuralKey {
origin,
class,
producer,
blob_tier,
subject,
})
}
pub fn with_base(base: &str, key_or_selector: impl AsRef<str>) -> String {
format!("{base}/{}", key_or_selector.as_ref())
}
pub fn strip_base<'k>(base: &str, key: &'k str) -> Option<&'k str> {
key.strip_prefix(base)?.strip_prefix('/')
}
pub fn parse_full<'k>(base: &str, key: &'k str) -> Option<StructuralKey<'k>> {
parse(strip_base(base, key)?).ok()
}
#[cfg(test)]
mod tests {
use super::*;
use crate::origin::HostId;
fn host() -> Origin {
Origin::Host(HostId::parse("h-3fa9c2d41b7e").unwrap())
}
#[test]
fn plain_chunk_rules() {
for ok in [
"a",
"cpu",
"sys_uptime",
"10-0-0-7",
"sshd.service",
"p95_ms",
"h-3fa9c2d41b7e",
] {
assert!(is_valid_plain_chunk(ok), "{ok}");
}
for bad in ["", "-a", "a-", ".a", "A", "Cpu", "a/b", "a*", "@v1", "é"] {
assert!(!is_valid_plain_chunk(bad), "{bad}");
}
}
#[test]
fn verbatim_chunk_rules() {
for ok in ["@v1", "@rpc", "@catalog", "@adv"] {
assert!(is_valid_verbatim_chunk(ok), "{ok}");
}
for bad in ["@", "@-x", "v1", "@V1", "@a/b"] {
assert!(!is_valid_verbatim_chunk(bad), "{bad}");
}
}
#[test]
fn producer_instance_split_is_unambiguous() {
assert!(Producer::new("snmp").is_ok());
assert!(Producer::new("net-ring").is_ok());
assert!(Producer::new("ipv6-2").is_err());
assert!(Producer::with_instance("snmp", 0).is_err());
let p = Producer::with_instance("snmp", 2).unwrap();
assert_eq!(p.chunk(), "snmp-2");
let back = Producer::parse_chunk("snmp-2").unwrap();
assert_eq!(back.name(), "snmp");
assert_eq!(back.instance(), Some(2));
let bare = Producer::parse_chunk("net-ring").unwrap();
assert_eq!(bare.name(), "net-ring");
assert_eq!(bare.instance(), None);
}
#[test]
fn blob_tiers_are_not_producers() {
assert!(Producer::new("store").is_err());
assert!(Producer::new("tree").is_err());
assert!(Producer::new("artifact").is_err());
}
#[test]
fn normative_examples_build_and_roundtrip() {
let p = |n| Producer::new(n).unwrap();
let cases = [
data_key(
&host(),
Class::Telemetry,
Some(&p("sysinfo")),
&["cpu", "usage"],
)
.unwrap(),
data_key(
&host(),
Class::Telemetry,
Some(&p("snmp")),
&["router01", "system", "sys_uptime"],
)
.unwrap(),
data_key(&host(), Class::State, Some(&p("netring")), &["health"]).unwrap(),
data_key(
&host(),
Class::State,
Some(&p("netlink")),
&["alert", "9f2c81ab04d7e3f1"],
)
.unwrap(),
data_key(
&host(),
Class::State,
Some(&p("netring")),
&["evidence", "names", "10-0-0-7"],
)
.unwrap(),
data_key(
&host(),
Class::Events,
Some(&p("netring")),
&["capture", "01jgxqz4yqk8v6txw3m9f2a7cd"],
)
.unwrap(),
rpc_key(&host(), Some(&p("netlink")), &["sockets"]).unwrap(),
media_key(&host(), &p("parallax"), &["cam0", "video", "h264", "high"]).unwrap(),
blob_key(&host(), BlobTier::Store, &["sha256", "ab12cd34ef56"]).unwrap(),
data_key(
&Origin::catalog(),
Class::State,
None,
&["entity", "h-3fa9c2d41b7e"],
)
.unwrap(),
data_key(
&Origin::catalog(),
Class::State,
None,
&["pdns", "93-184-216-34"],
)
.unwrap(),
];
let expected = [
"v1/h-3fa9c2d41b7e/telemetry/sysinfo/cpu/usage",
"v1/h-3fa9c2d41b7e/telemetry/snmp/router01/system/sys_uptime",
"v1/h-3fa9c2d41b7e/state/netring/health",
"v1/h-3fa9c2d41b7e/state/netlink/alert/9f2c81ab04d7e3f1",
"v1/h-3fa9c2d41b7e/state/netring/evidence/names/10-0-0-7",
"v1/h-3fa9c2d41b7e/events/netring/capture/01jgxqz4yqk8v6txw3m9f2a7cd",
"v1/h-3fa9c2d41b7e/@rpc/netlink/sockets",
"v1/h-3fa9c2d41b7e/@media/parallax/cam0/video/h264/high",
"v1/h-3fa9c2d41b7e/@blob/store/sha256/ab12cd34ef56",
"v1/@catalog/state/entity/h-3fa9c2d41b7e",
"v1/@catalog/state/pdns/93-184-216-34",
];
for (built, want) in cases.iter().zip(expected) {
assert_eq!(built, want);
let parsed = parse(built).unwrap();
let subject = &parsed.subject;
let rebuilt = match parsed.class {
ClassOrPlane::Class(c) => {
data_key(&parsed.origin, c, parsed.producer.as_ref(), subject).unwrap()
}
ClassOrPlane::Plane(Plane::Rpc) => {
rpc_key(&parsed.origin, parsed.producer.as_ref(), subject).unwrap()
}
ClassOrPlane::Plane(Plane::Media) => {
media_key(&parsed.origin, parsed.producer.as_ref().unwrap(), subject).unwrap()
}
ClassOrPlane::Plane(Plane::Blob) => {
blob_key(&parsed.origin, parsed.blob_tier.unwrap(), subject).unwrap()
}
};
assert_eq!(&rebuilt, want);
}
}
#[test]
fn alive_is_liveliness_only() {
assert!(
data_key(
&host(),
Class::State,
Some(&Producer::new("netlink").unwrap()),
&["alive"]
)
.is_err()
);
assert_eq!(
alive_key(&host(), Some(&Producer::new("netlink").unwrap())).unwrap(),
"v1/h-3fa9c2d41b7e/state/netlink/alive"
);
assert_eq!(
alive_key(&Origin::catalog(), None).unwrap(),
"v1/@catalog/state/alive"
);
assert_eq!(
device_alive_key(&host(), &Producer::new("snmp").unwrap(), "router01").unwrap(),
"v1/h-3fa9c2d41b7e/state/snmp/device/router01/alive"
);
}
#[test]
fn service_origin_omits_producer() {
assert!(
data_key(
&Origin::catalog(),
Class::State,
Some(&Producer::new("x").unwrap()),
&["entity", "a"]
)
.is_err()
);
assert!(data_key(&host(), Class::State, None, &["health"]).is_err());
}
#[test]
fn parse_rejects_foreign_keys() {
assert!(parse("zensight/netlink/host/@/health").is_err());
assert!(parse("@v2/h-3fa9c2d41b7e/state/x/health").is_err());
assert!(parse("v1/h-3fa9c2d41b7e/bogus/x/health").is_err());
assert!(parse("v1/h-3fa9c2d41b7e/@blob/bogus/x").is_err());
}
}