use crate::common_state::CommonFamily;
use crate::grammar::{
CLASS_EVENTS, CLASS_STATE, CLASS_TELEMETRY, Class, PLANE_RPC, SUBJECT_ALIVE, VERSION_CHUNK,
is_valid_plain_chunk,
};
use crate::key::{Key, Selector};
use crate::origin::{ConcreteOrigin, Fleet, HostOrigin, ServiceOrigin};
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Scope(ScopeInner);
#[derive(Debug, Clone, PartialEq, Eq)]
enum ScopeInner {
Fleet,
Origin(String),
}
impl Scope {
pub fn fleet() -> Scope {
Scope(ScopeInner::Fleet)
}
pub fn origin(o: &impl ConcreteOrigin) -> Scope {
Scope(ScopeInner::Origin(o.chunk().to_string()))
}
pub fn chunk(&self) -> &str {
match &self.0 {
ScopeInner::Fleet => "*",
ScopeInner::Origin(c) => c,
}
}
}
impl From<Fleet> for Scope {
fn from(_: Fleet) -> Scope {
Scope::fleet()
}
}
fn legal_chunk<'c>(chunk: &'c str, what: &str) -> &'c str {
assert!(
is_valid_plain_chunk(chunk),
"{what} {chunk:?} violates RFC 03 §2"
);
chunk
}
fn class_selector(scope: Scope, class: &str) -> Selector {
Selector::from_canonical(format!("{VERSION_CHUNK}/{}/{class}/**", scope.chunk()))
}
#[must_use]
pub fn all_state(scope: Scope) -> Selector {
class_selector(scope, CLASS_STATE)
}
#[must_use]
pub fn all_telemetry(scope: Scope) -> Selector {
class_selector(scope, CLASS_TELEMETRY)
}
#[must_use]
pub fn all_events(scope: Scope) -> Selector {
class_selector(scope, CLASS_EVENTS)
}
#[must_use]
pub fn all_of_class(scope: Scope, class: Class) -> Selector {
class_selector(scope, class.chunk())
}
#[must_use]
pub fn all_liveliness(scope: Scope) -> Selector {
Selector::from_canonical(format!(
"{VERSION_CHUNK}/{}/{CLASS_STATE}/*/{SUBJECT_ALIVE}",
scope.chunk()
))
}
#[must_use]
pub fn producer_state(scope: Scope, producer: &str, prefix: &[&str]) -> Selector {
let mut out = format!(
"{VERSION_CHUNK}/{}/{CLASS_STATE}/{}",
scope.chunk(),
legal_chunk(producer, "producer")
);
for chunk in prefix {
out.push('/');
out.push_str(legal_chunk(chunk, "subject chunk"));
}
out.push_str("/**");
Selector::from_canonical(out)
}
#[must_use]
pub fn common_family(scope: Scope, family: CommonFamily) -> Selector {
let mut out = format!("{VERSION_CHUNK}/{}/{CLASS_STATE}/*", scope.chunk());
for chunk in family.prefix() {
out.push('/');
out.push_str(chunk);
}
if family.var().is_some() {
out.push_str("/*");
}
Selector::from_canonical(out)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Producers(ProducersInner);
#[derive(Debug, Clone, PartialEq, Eq)]
enum ProducersInner {
All,
Named(String),
}
impl Producers {
pub fn all() -> Producers {
Producers(ProducersInner::All)
}
pub fn named(producer: &str) -> Producers {
Producers(ProducersInner::Named(
legal_chunk(producer, "producer").to_string(),
))
}
pub fn chunk(&self) -> &str {
match &self.0 {
ProducersInner::All => "*",
ProducersInner::Named(p) => p,
}
}
}
#[must_use]
pub fn rpc(scope: Scope, producer: Producers, procedure: &[&str]) -> Selector {
let p = producer.chunk();
let mut out = format!("{VERSION_CHUNK}/{}/{PLANE_RPC}/{p}", scope.chunk());
for chunk in procedure {
out.push('/');
out.push_str(legal_chunk(chunk, "procedure chunk"));
}
Selector::from_canonical(out)
}
#[must_use]
pub fn fleet_rpc(producer: &str, procedure: &[&str]) -> Selector {
rpc(Scope::fleet(), Producers::named(producer), procedure)
}
#[must_use]
pub fn rpc_at(origin: &impl HostOrigin, producer: &str, procedure: &[&str]) -> Key {
let mut out = format!(
"{VERSION_CHUNK}/{}/{PLANE_RPC}/{}",
origin.chunk(),
legal_chunk(producer, "producer")
);
for chunk in procedure {
out.push('/');
out.push_str(legal_chunk(chunk, "procedure chunk"));
}
Key::from_canonical(out)
}
#[must_use]
pub fn service_rpc(origin: &ServiceOrigin, procedure: &[&str]) -> Key {
let mut out = format!("{VERSION_CHUNK}/{}/{PLANE_RPC}", origin.as_str());
for chunk in procedure {
out.push('/');
out.push_str(legal_chunk(chunk, "procedure chunk"));
}
Key::from_canonical(out)
}
#[must_use]
pub fn service_alive(origin: &ServiceOrigin) -> Key {
Key::from_canonical(format!(
"{VERSION_CHUNK}/{}/{CLASS_STATE}/{SUBJECT_ALIVE}",
origin.as_str()
))
}
#[cfg(test)]
mod tests {
use super::*;
use crate::origin::RemoteOrigin;
#[test]
fn framework_selector_shapes() {
assert_eq!(all_state(Scope::fleet()), "v1/*/state/**");
assert_eq!(all_telemetry(Scope::fleet()), "v1/*/telemetry/**");
assert_eq!(all_events(Scope::fleet()), "v1/*/events/**");
assert_eq!(all_liveliness(Scope::fleet()), "v1/*/state/*/alive");
let o = RemoteOrigin::parse("h-3fa9c2d41b7e").unwrap();
assert_eq!(all_state(Scope::origin(&o)), "v1/h-3fa9c2d41b7e/state/**");
assert_eq!(
producer_state(Scope::origin(&o), "tc", &["config"]),
"v1/h-3fa9c2d41b7e/state/tc/config/**"
);
}
#[test]
fn common_family_shapes() {
assert_eq!(
common_family(Scope::fleet(), CommonFamily::Health),
"v1/*/state/*/health"
);
assert_eq!(
common_family(Scope::fleet(), CommonFamily::Alert),
"v1/*/state/*/alert/*"
);
assert_eq!(
common_family(Scope::fleet(), CommonFamily::EvidenceSelf),
"v1/*/state/*/evidence/self"
);
assert_eq!(
common_family(Scope::fleet(), CommonFamily::EvidenceNames),
"v1/*/state/*/evidence/names/*"
);
let o = RemoteOrigin::parse("h-3fa9c2d41b7e").unwrap();
assert_eq!(
common_family(Scope::origin(&o), CommonFamily::Errors),
"v1/h-3fa9c2d41b7e/state/*/errors"
);
for f in CommonFamily::ALL {
let sel = common_family(Scope::fleet(), f);
assert!(all_state(Scope::fleet()).includes(&sel), "{sel}");
}
}
#[test]
fn common_family_never_sees_a_service() {
let entity = Key::from_canonical("v1/@catalog/state/entity/abc".to_string());
for f in CommonFamily::ALL {
assert!(!common_family(Scope::fleet(), f).intersects(&entity));
}
}
#[test]
fn rpc_shapes() {
assert_eq!(
rpc(Scope::fleet(), Producers::all(), &["introspect"]),
"v1/*/@rpc/*/introspect"
);
assert_eq!(fleet_rpc("netring", &["flows"]), "v1/*/@rpc/netring/flows");
let o = RemoteOrigin::parse("h-3fa9c2d41b7e").unwrap();
assert_eq!(
rpc(Scope::origin(&o), Producers::all(), &["introspect"]),
"v1/h-3fa9c2d41b7e/@rpc/*/introspect"
);
assert!(
rpc(Scope::fleet(), Producers::all(), &["introspect"]).includes(&rpc(
Scope::origin(&o),
Producers::all(),
&["introspect"]
))
);
assert_eq!(
rpc_at(&o, "netring", &["capture_disk", "set"]),
"v1/h-3fa9c2d41b7e/@rpc/netring/capture_disk/set"
);
let cat = ServiceOrigin::catalog();
assert_eq!(
service_rpc(&cat, &["introspect"]),
"v1/@catalog/@rpc/introspect"
);
assert_eq!(service_alive(&cat), "v1/@catalog/state/alive");
}
#[test]
fn fleet_liveliness_excludes_services() {
let fleet = all_liveliness(Scope::fleet());
let svc = service_alive(&ServiceOrigin::catalog());
assert!(!fleet.intersects(&svc));
}
}