1use std::fmt;
18
19use crate::common_state::CommonFamily;
20use crate::encoding::WireEncoding;
21use crate::grammar::{BlobTier, Class};
22use crate::origin::ServiceOrigin;
23use crate::qos::QosProfile;
24use crate::registry_doc::{RawTable, RawValue, SliceFormat, negotiate, parse_raw, write_kdl};
25
26pub trait SliceToken: Sized {
50 fn from_token(token: &str) -> Option<Self>;
53
54 fn token(&self) -> &str;
57}
58
59#[derive(Debug, Clone, PartialEq, Eq, Hash)]
67pub enum Declared<T> {
68 Known(T),
70 Other(String),
72}
73
74impl<T: SliceToken> Declared<T> {
75 pub fn parse(token: &str) -> Self {
77 match T::from_token(token) {
78 Some(known) => Declared::Known(known),
79 None => Declared::Other(token.to_string()),
80 }
81 }
82
83 pub fn token(&self) -> &str {
86 match self {
87 Declared::Known(k) => k.token(),
88 Declared::Other(s) => s,
89 }
90 }
91
92 pub fn known(&self) -> Option<&T> {
94 match self {
95 Declared::Known(k) => Some(k),
96 Declared::Other(_) => None,
97 }
98 }
99
100 pub fn is(&self, other: &T) -> bool
104 where
105 T: PartialEq,
106 {
107 matches!(self, Declared::Known(k) if k == other)
108 }
109}
110
111impl<T: SliceToken> From<T> for Declared<T> {
112 fn from(value: T) -> Self {
113 Declared::Known(value)
114 }
115}
116
117impl<T: SliceToken> fmt::Display for Declared<T> {
118 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
119 f.write_str(self.token())
120 }
121}
122
123#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
126pub enum ProcedureKind {
127 Read,
129 Write,
132}
133
134impl SliceToken for ProcedureKind {
135 fn from_token(token: &str) -> Option<Self> {
136 match token {
137 "read" => Some(ProcedureKind::Read),
138 "write" => Some(ProcedureKind::Write),
139 _ => None,
140 }
141 }
142
143 fn token(&self) -> &str {
144 match self {
145 ProcedureKind::Read => "read",
146 ProcedureKind::Write => "write",
147 }
148 }
149}
150
151#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
154pub enum Fanout {
155 Allowed,
157 Forbidden,
159}
160
161impl SliceToken for Fanout {
162 fn from_token(token: &str) -> Option<Self> {
163 match token {
164 "allowed" => Some(Fanout::Allowed),
165 "forbidden" => Some(Fanout::Forbidden),
166 _ => None,
167 }
168 }
169
170 fn token(&self) -> &str {
171 match self {
172 Fanout::Allowed => "allowed",
173 Fanout::Forbidden => "forbidden",
174 }
175 }
176}
177
178#[derive(Debug, Clone, PartialEq, Eq, Hash)]
187pub enum RateClass {
188 Rare,
190 Low,
192 Burst(u64),
194 Other(String),
196}
197
198impl RateClass {
199 pub fn parse(token: &str) -> RateClass {
202 match token {
203 "rare" => RateClass::Rare,
204 "low" => RateClass::Low,
205 other => match other
206 .strip_prefix("burst(")
207 .and_then(|r| r.strip_suffix("/h)"))
208 .and_then(|n| n.parse().ok())
209 {
210 Some(n) => RateClass::Burst(n),
211 None => RateClass::Other(other.to_string()),
212 },
213 }
214 }
215
216 pub fn cap_per_hour(&self) -> Option<u64> {
221 match self {
222 RateClass::Rare => Some(1),
223 RateClass::Low => Some(60),
224 RateClass::Burst(n) => Some(*n),
225 RateClass::Other(_) => None,
226 }
227 }
228
229 pub fn token(&self) -> String {
231 match self {
232 RateClass::Rare => "rare".to_string(),
233 RateClass::Low => "low".to_string(),
234 RateClass::Burst(n) => format!("burst({n}/h)"),
235 RateClass::Other(s) => s.clone(),
236 }
237 }
238}
239
240impl fmt::Display for RateClass {
241 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
242 f.write_str(&self.token())
243 }
244}
245
246impl SliceToken for Class {
247 fn from_token(token: &str) -> Option<Self> {
248 Class::from_chunk(token)
249 }
250
251 fn token(&self) -> &str {
252 self.chunk()
253 }
254}
255
256impl SliceToken for QosProfile {
257 fn from_token(token: &str) -> Option<Self> {
258 QosProfile::from_name(token)
259 }
260
261 fn token(&self) -> &str {
262 self.name()
263 }
264}
265
266impl SliceToken for BlobTier {
267 fn from_token(token: &str) -> Option<Self> {
268 BlobTier::from_chunk(token)
269 }
270
271 fn token(&self) -> &str {
272 self.chunk()
273 }
274}
275
276impl SliceToken for CommonFamily {
277 fn from_token(token: &str) -> Option<Self> {
278 CommonFamily::ALL.into_iter().find(|f| f.token() == token)
279 }
280
281 fn token(&self) -> &str {
282 CommonFamily::token(*self)
283 }
284}
285
286impl SliceToken for ServiceOrigin {
287 fn from_token(token: &str) -> Option<Self> {
288 ServiceOrigin::new(token).ok()
289 }
290
291 fn token(&self) -> &str {
292 self.as_str()
293 }
294}
295
296#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
302pub enum SubjectKind {
303 Counter,
307 Gauge,
309 Text,
311 Bool,
314 Histogram,
318}
319
320impl SubjectKind {
321 pub const ALL: [SubjectKind; 5] = [
324 SubjectKind::Counter,
325 SubjectKind::Gauge,
326 SubjectKind::Text,
327 SubjectKind::Bool,
328 SubjectKind::Histogram,
329 ];
330
331 #[must_use]
335 pub fn payload_tag(self) -> &'static str {
336 match self {
337 SubjectKind::Bool => "boolean",
338 other => other.token_str(),
339 }
340 }
341
342 #[must_use]
346 pub fn from_payload_tag(tag: &str) -> Option<Self> {
347 SubjectKind::ALL
348 .into_iter()
349 .find(|k| k.payload_tag() == tag)
350 }
351
352 const fn token_str(self) -> &'static str {
353 match self {
354 SubjectKind::Counter => "counter",
355 SubjectKind::Gauge => "gauge",
356 SubjectKind::Text => "text",
357 SubjectKind::Bool => "bool",
358 SubjectKind::Histogram => "histogram",
359 }
360 }
361}
362
363#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
368pub enum Semantic {
369 Temperature,
370 Power,
371 Bytes,
372 Duration,
373 Ratio,
375 Count,
377 Identity,
379 State,
381}
382
383impl Semantic {
384 pub const ALL: [Semantic; 8] = [
386 Semantic::Temperature,
387 Semantic::Power,
388 Semantic::Bytes,
389 Semantic::Duration,
390 Semantic::Ratio,
391 Semantic::Count,
392 Semantic::Identity,
393 Semantic::State,
394 ];
395
396 const fn token_str(self) -> &'static str {
397 match self {
398 Semantic::Temperature => "temperature",
399 Semantic::Power => "power",
400 Semantic::Bytes => "bytes",
401 Semantic::Duration => "duration",
402 Semantic::Ratio => "ratio",
403 Semantic::Count => "count",
404 Semantic::Identity => "identity",
405 Semantic::State => "state",
406 }
407 }
408}
409
410impl SliceToken for Semantic {
411 fn from_token(token: &str) -> Option<Self> {
412 Semantic::ALL.into_iter().find(|k| k.token_str() == token)
413 }
414
415 fn token(&self) -> &str {
416 self.token_str()
417 }
418}
419
420#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
428pub enum Exposure {
429 Host,
432 Link,
435 Fleet,
438}
439
440impl Exposure {
441 pub const ALL: [Exposure; 3] = [Exposure::Host, Exposure::Link, Exposure::Fleet];
443
444 const fn token_str(self) -> &'static str {
445 match self {
446 Exposure::Host => "host",
447 Exposure::Link => "link",
448 Exposure::Fleet => "fleet",
449 }
450 }
451}
452
453impl SliceToken for Exposure {
454 fn from_token(token: &str) -> Option<Self> {
455 Exposure::ALL.into_iter().find(|e| e.token_str() == token)
456 }
457
458 fn token(&self) -> &str {
459 self.token_str()
460 }
461}
462
463#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
471pub enum PredicateKind {
472 Feature,
475 Config,
478 Capability,
481}
482
483impl PredicateKind {
484 pub const ALL: [PredicateKind; 3] = [
486 PredicateKind::Feature,
487 PredicateKind::Config,
488 PredicateKind::Capability,
489 ];
490
491 const fn token_str(self) -> &'static str {
492 match self {
493 PredicateKind::Feature => "feature",
494 PredicateKind::Config => "config",
495 PredicateKind::Capability => "capability",
496 }
497 }
498
499 #[must_use]
503 pub const fn gated_error(self) -> &'static str {
504 match self {
505 PredicateKind::Feature => "error/unsupported",
506 PredicateKind::Config | PredicateKind::Capability => "error/gated",
507 }
508 }
509}
510
511impl SliceToken for PredicateKind {
512 fn from_token(token: &str) -> Option<Self> {
513 PredicateKind::ALL
514 .into_iter()
515 .find(|k| k.token_str() == token)
516 }
517
518 fn token(&self) -> &str {
519 self.token_str()
520 }
521}
522
523#[derive(Debug, Clone, PartialEq, Eq, Hash)]
531pub struct Predicate {
532 pub kind: Declared<PredicateKind>,
533 pub name: String,
536}
537
538impl Predicate {
539 #[must_use]
541 pub fn new(kind: PredicateKind, name: impl Into<String>) -> Self {
542 Predicate {
543 kind: Declared::Known(kind),
544 name: name.into(),
545 }
546 }
547
548 #[must_use]
550 pub fn parse(token: &str) -> Self {
551 match token.split_once(':') {
552 Some((kind, name)) => Predicate {
553 kind: Declared::parse(kind),
554 name: name.to_string(),
555 },
556 None => Predicate {
557 kind: Declared::Other(token.to_string()),
558 name: String::new(),
559 },
560 }
561 }
562
563 #[must_use]
565 pub fn token(&self) -> String {
566 match (&self.kind, self.name.is_empty()) {
567 (Declared::Other(raw), true) => raw.clone(),
568 (kind, _) => format!("{}:{}", kind.token(), self.name),
569 }
570 }
571}
572
573impl fmt::Display for Predicate {
574 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
575 f.write_str(&self.token())
576 }
577}
578
579fn when_toml(out: &mut String, when: Option<&[Predicate]>) {
581 if let Some(preds) = when {
582 let items: Vec<String> = preds.iter().map(|p| toml_quote(&p.token())).collect();
583 out.push_str(&format!("when = [{}]\n", items.join(", ")));
584 }
585}
586
587#[derive(Debug, Clone)]
594pub struct Buckets(Vec<f64>);
595
596impl Buckets {
597 #[must_use]
598 pub fn new(bounds: Vec<f64>) -> Self {
599 Buckets(bounds)
600 }
601
602 #[must_use]
604 pub fn as_slice(&self) -> &[f64] {
605 &self.0
606 }
607
608 #[must_use]
611 pub fn is_well_formed(&self) -> bool {
612 !self.0.is_empty()
613 && self.0.iter().all(|b| b.is_finite())
614 && self.0.windows(2).all(|w| w[0] < w[1])
615 }
616}
617
618impl PartialEq for Buckets {
619 fn eq(&self, other: &Self) -> bool {
620 self.0.len() == other.0.len()
621 && self
622 .0
623 .iter()
624 .zip(&other.0)
625 .all(|(a, b)| a.to_bits() == b.to_bits())
626 }
627}
628
629impl Eq for Buckets {}
630
631impl SliceToken for SubjectKind {
632 fn from_token(token: &str) -> Option<Self> {
633 SubjectKind::ALL
634 .into_iter()
635 .find(|k| k.token_str() == token)
636 }
637
638 fn token(&self) -> &str {
639 self.token_str()
640 }
641}
642
643#[derive(Debug, Clone, PartialEq, Eq)]
651#[non_exhaustive]
652pub struct SubjectDecl {
653 pub path: String,
656 pub class: Declared<Class>,
658 pub type_name: String,
660 pub kind: Option<Declared<SubjectKind>>,
664 pub common: Option<Declared<CommonFamily>>,
674 pub since: Option<String>,
676 pub description: Option<String>,
677 pub qos: Option<Declared<QosProfile>>,
679 pub ttl_s: Option<i64>,
681 pub unit: Option<String>,
683 pub rate: Option<RateClass>,
685 pub cardinality: Option<i64>,
687 pub encoding: Option<WireEncoding>,
692 pub buckets: Option<Buckets>,
695 pub semantic: Option<Declared<Semantic>>,
697 pub when: Option<Vec<Predicate>>,
702 pub gate_note: Option<String>,
705 pub exposure: Option<Declared<Exposure>>,
708}
709
710#[derive(Debug, Clone, PartialEq, Eq)]
718#[non_exhaustive]
719pub struct ProcedureDecl {
720 pub path: String,
722 pub kind: Option<Declared<ProcedureKind>>,
724 pub reply: Option<String>,
725 pub request: Option<String>,
727 pub encoding: Option<WireEncoding>,
729 pub fanout: Option<Declared<Fanout>>,
734 pub idempotent: Option<bool>,
736 pub cardinality: Option<i64>,
741 pub when: Option<Vec<Predicate>>,
746 pub gate_note: Option<String>,
748 pub exposure: Option<Declared<Exposure>>,
751 pub sensitive: Option<bool>,
754 pub since: Option<String>,
755 pub description: Option<String>,
756}
757
758#[derive(Debug, Clone, PartialEq, Eq)]
766#[non_exhaustive]
767pub struct ErrorDecl {
768 pub name: String,
770 pub procedures: Vec<String>,
773 pub since: Option<String>,
774 pub description: Option<String>,
775}
776
777impl ErrorDecl {
778 #[must_use]
780 pub fn new(name: impl Into<String>) -> Self {
781 ErrorDecl {
782 name: name.into(),
783 procedures: Vec::new(),
784 since: None,
785 description: None,
786 }
787 }
788
789 #[must_use]
791 pub fn wire_name(&self, producer: &str) -> String {
792 crate::rpc_error::producer_error(producer, &self.name)
793 }
794}
795
796#[derive(Debug, Clone, PartialEq, Eq)]
815#[non_exhaustive]
816pub struct BlobDecl {
817 pub tier: Declared<BlobTier>,
819 pub endpoints: Vec<String>,
822 pub algo: Option<String>,
824 pub reference: Option<String>,
827 pub encoding: Option<WireEncoding>,
829 pub since: Option<String>,
830 pub description: Option<String>,
831}
832
833#[derive(Debug, Clone, PartialEq, Eq)]
848#[non_exhaustive]
849pub struct MediaDecl {
850 pub path: String,
853 pub encoding: WireEncoding,
857 pub attachment: Option<String>,
859 pub cardinality: Option<i64>,
861 pub since: Option<String>,
862 pub description: Option<String>,
863}
864
865#[derive(Debug, Clone, PartialEq, Eq)]
883#[non_exhaustive]
884pub struct BudgetDecl {
885 pub rss_mb: Option<i64>,
887 pub tables: Vec<TableBudget>,
889}
890
891#[derive(Debug, Clone, PartialEq, Eq)]
900#[non_exhaustive]
901pub struct TableBudget {
902 pub name: String,
905 pub max_entries: Option<i64>,
907 pub max_bytes: Option<i64>,
909}
910
911#[derive(Debug, Clone, PartialEq, Eq)]
919#[non_exhaustive]
920pub struct DeprecationDecl {
921 pub path: String,
922 pub kind: DeprecatedKind,
929 pub since: Option<String>,
931 pub replaced_by: Option<String>,
933}
934
935#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
942pub enum DeprecatedKind {
943 #[default]
944 Subject,
945 Procedure,
946 Error,
949}
950
951impl DeprecatedKind {
952 pub fn as_str(self) -> &'static str {
954 match self {
955 DeprecatedKind::Subject => "subject",
956 DeprecatedKind::Procedure => "procedure",
957 DeprecatedKind::Error => "error",
958 }
959 }
960}
961
962impl std::str::FromStr for DeprecatedKind {
963 type Err = ();
964
965 fn from_str(s: &str) -> Result<DeprecatedKind, ()> {
966 match s {
967 "subject" => Ok(DeprecatedKind::Subject),
968 "procedure" => Ok(DeprecatedKind::Procedure),
969 "error" => Ok(DeprecatedKind::Error),
970 _ => Err(()),
971 }
972 }
973}
974
975#[derive(Debug, Clone, PartialEq, Eq)]
983#[non_exhaustive]
984pub struct RegistrySlice {
985 pub version: String,
987 pub app: String,
988 pub convention: i64,
990 pub name: String,
992 pub service_origin: Option<Declared<ServiceOrigin>>,
995 pub description: Option<String>,
996 pub subjects: Vec<SubjectDecl>,
997 pub procedures: Vec<ProcedureDecl>,
998 pub blob: Vec<BlobDecl>,
1003 pub media: Vec<MediaDecl>,
1007 pub errors: Vec<ErrorDecl>,
1011 pub deprecated: Vec<DeprecationDecl>,
1012 pub budget: Option<BudgetDecl>,
1016}
1017
1018impl SubjectDecl {
1019 #[must_use]
1023 pub fn new(path: impl Into<String>, class: impl Into<Declared<Class>>) -> Self {
1024 SubjectDecl {
1025 path: path.into(),
1026 class: class.into(),
1027 type_name: String::new(),
1028 kind: None,
1029 common: None,
1030 since: None,
1031 description: None,
1032 qos: None,
1033 ttl_s: None,
1034 unit: None,
1035 rate: None,
1036 cardinality: None,
1037 encoding: None,
1038 buckets: None,
1039 semantic: None,
1040 when: None,
1041 gate_note: None,
1042 exposure: None,
1043 }
1044 }
1045}
1046
1047impl ProcedureDecl {
1048 #[must_use]
1050 pub fn new(path: impl Into<String>) -> Self {
1051 ProcedureDecl {
1052 path: path.into(),
1053 kind: None,
1054 reply: None,
1055 request: None,
1056 encoding: None,
1057 fanout: None,
1058 idempotent: None,
1059 cardinality: None,
1060 when: None,
1061 gate_note: None,
1062 exposure: None,
1063 sensitive: None,
1064 since: None,
1065 description: None,
1066 }
1067 }
1068}
1069
1070impl BlobDecl {
1071 #[must_use]
1073 pub fn new(tier: impl Into<Declared<BlobTier>>) -> Self {
1074 BlobDecl {
1075 tier: tier.into(),
1076 endpoints: Vec::new(),
1077 algo: None,
1078 reference: None,
1079 encoding: None,
1080 since: None,
1081 description: None,
1082 }
1083 }
1084}
1085
1086impl MediaDecl {
1087 #[must_use]
1089 pub fn new(path: impl Into<String>, encoding: WireEncoding) -> Self {
1090 MediaDecl {
1091 path: path.into(),
1092 encoding,
1093 attachment: None,
1094 cardinality: None,
1095 since: None,
1096 description: None,
1097 }
1098 }
1099}
1100
1101impl BudgetDecl {
1102 #[must_use]
1104 pub fn new() -> Self {
1105 BudgetDecl {
1106 rss_mb: None,
1107 tables: Vec::new(),
1108 }
1109 }
1110}
1111
1112impl Default for BudgetDecl {
1113 fn default() -> Self {
1114 BudgetDecl::new()
1115 }
1116}
1117
1118impl TableBudget {
1119 #[must_use]
1122 pub fn new(name: impl Into<String>) -> Self {
1123 TableBudget {
1124 name: name.into(),
1125 max_entries: None,
1126 max_bytes: None,
1127 }
1128 }
1129}
1130
1131impl DeprecationDecl {
1132 #[must_use]
1135 pub fn new(path: impl Into<String>) -> Self {
1136 DeprecationDecl::of(DeprecatedKind::Subject, path)
1137 }
1138
1139 #[must_use]
1141 pub fn of(kind: DeprecatedKind, path: impl Into<String>) -> Self {
1142 DeprecationDecl {
1143 path: path.into(),
1144 kind,
1145 since: None,
1146 replaced_by: None,
1147 }
1148 }
1149}
1150
1151#[derive(Debug, Clone, PartialEq)]
1157pub struct Bound<'s> {
1158 pub decl: &'s SubjectDecl,
1159 pub vars: Vec<(&'s str, String)>,
1160}
1161
1162impl Bound<'_> {
1163 #[must_use]
1165 pub fn var(&self, name: &str) -> Option<&str> {
1166 self.vars
1167 .iter()
1168 .find(|(n, _)| *n == name)
1169 .map(|(_, v)| v.as_str())
1170 }
1171}
1172
1173impl RegistrySlice {
1174 #[must_use]
1191 pub fn bind(&self, class: Class, tail: &[&str]) -> Option<Bound<'_>> {
1192 let (decls, patterns): (Vec<&SubjectDecl>, Vec<crate::pattern::SubjectPattern>) = self
1193 .subjects
1194 .iter()
1195 .filter(|s| s.class == Declared::Known(class))
1196 .filter_map(|s| {
1197 crate::pattern::SubjectPattern::parse(&s.path)
1198 .ok()
1199 .map(|p| (s, p))
1200 })
1201 .unzip();
1202 let (idx, binds) = crate::pattern::best_match(&patterns, tail)?;
1203 let decl = decls[idx];
1204 let vars = binds
1208 .into_iter()
1209 .map(|(name, value)| (var_name_in(&decl.path, name), value))
1210 .collect();
1211 Some(Bound { decl, vars })
1212 }
1213}
1214
1215fn var_name_in<'p>(path: &'p str, name: &str) -> &'p str {
1219 path.split('/')
1220 .filter_map(|c| c.strip_prefix('{'))
1221 .map(|c| c.trim_end_matches('}').trim_end_matches("..."))
1222 .find(|c| *c == name)
1223 .expect("a bound variable is spelled in the path it was parsed from")
1224}
1225
1226impl RegistrySlice {
1227 #[must_use]
1234 pub fn new(
1235 version: impl Into<String>,
1236 app: impl Into<String>,
1237 name: impl Into<String>,
1238 ) -> Self {
1239 RegistrySlice {
1240 version: version.into(),
1241 app: app.into(),
1242 convention: 1,
1243 name: name.into(),
1244 service_origin: None,
1245 description: None,
1246 subjects: Vec::new(),
1247 procedures: Vec::new(),
1248 blob: Vec::new(),
1249 media: Vec::new(),
1250 errors: Vec::new(),
1251 deprecated: Vec::new(),
1252 budget: None,
1253 }
1254 }
1255
1256 pub fn subjects_in(
1264 &self,
1265 class: impl Into<Declared<Class>>,
1266 ) -> impl Iterator<Item = &SubjectDecl> {
1267 let class = class.into();
1268 self.subjects.iter().filter(move |s| s.class == class)
1269 }
1270
1271 pub fn serves_subject(&self, path: &str) -> bool {
1273 self.subjects.iter().any(|s| s.path == path)
1274 }
1275
1276 pub fn serves_procedure(&self, path: &str) -> bool {
1278 self.procedures.iter().any(|p| p.path == path)
1279 }
1280
1281 pub fn serves_blob_tier(&self, tier: impl Into<Declared<BlobTier>>) -> bool {
1287 let tier = tier.into();
1288 self.blob.iter().any(|b| b.tier == tier)
1289 }
1290
1291 pub fn blob_tiers(&self) -> impl Iterator<Item = &Declared<BlobTier>> {
1297 self.blob.iter().map(|b| &b.tier)
1298 }
1299
1300 pub fn serves_media(&self, path: &str) -> bool {
1303 self.media.iter().any(|m| m.path == path)
1304 }
1305}
1306
1307#[derive(Debug, thiserror::Error)]
1317#[non_exhaustive]
1318pub enum SliceError {
1319 #[error("malformed registry slice: {0}")]
1321 Toml(#[from] toml::de::Error),
1322 #[error("malformed registry slice: not KDL 2.0 (RFC 08 §5.1): {}", kdl_detail(.0))]
1327 Kdl(#[from] kdl::KdlError),
1328 #[error("malformed registry slice: {0}")]
1333 Shape(String),
1334 #[error(
1339 "unreadable registry slice: the reply declares encoding {0:?}, which is neither \
1340 application/toml nor application/kdl (RFC 08 §6)"
1341 )]
1342 Encoding(String),
1343}
1344
1345fn kdl_detail(e: &kdl::KdlError) -> String {
1347 let mut parts = Vec::new();
1348 for d in &e.diagnostics {
1349 let offset = d.span.offset().min(e.input.len());
1350 let before = &e.input.as_bytes()[..offset];
1351 let line = before.iter().filter(|&&b| b == b'\n').count() + 1;
1352 let col = offset
1353 - before
1354 .iter()
1355 .rposition(|&b| b == b'\n')
1356 .map_or(0, |p| p + 1)
1357 + 1;
1358 let mut part = format!(
1359 "line {line}:{col}: {}",
1360 d.message.as_deref().unwrap_or("unexpected input")
1361 );
1362 if let Some(label) = &d.label {
1363 part.push_str(&format!(" ({label})"));
1364 }
1365 if let Some(help) = &d.help {
1366 part.push_str(&format!(" — {help}"));
1367 }
1368 parts.push(part);
1369 }
1370 if parts.is_empty() {
1371 "the document does not parse".to_string()
1372 } else {
1373 parts.join("; ")
1374 }
1375}
1376
1377pub fn parse_slice(src: &str) -> Result<RegistrySlice, SliceError> {
1392 parse_slice_as(src, SliceFormat::sniff(src))
1393}
1394
1395pub fn parse_slice_as(src: &str, format: SliceFormat) -> Result<RegistrySlice, SliceError> {
1398 slice_from_raw(&parse_raw(src, format)?)
1399}
1400
1401pub fn parse_served(declared: Option<&str>, src: &str) -> Result<RegistrySlice, SliceError> {
1405 parse_slice_as(src, negotiate(declared, src)?)
1406}
1407
1408pub fn slice_from_raw(doc: &RawTable) -> Result<RegistrySlice, SliceError> {
1412 let err = |m: &str| SliceError::Shape(m.to_string());
1413 let s = |v: Option<&RawValue>| v.and_then(|v| v.as_str()).map(str::to_string);
1414 fn tok<T: SliceToken>(v: Option<&RawValue>) -> Option<Declared<T>> {
1418 v.and_then(|v| v.as_str()).map(Declared::parse)
1419 }
1420 fn enc(v: Option<&RawValue>) -> Option<WireEncoding> {
1421 v.and_then(|v| v.as_str())
1422 .map(WireEncoding::from_encoding_str)
1423 }
1424
1425 let header = doc
1426 .get("registry")
1427 .ok_or_else(|| err("missing [registry]"))?;
1428 let version = s(header.get("version")).ok_or_else(|| err("[registry] missing version"))?;
1429 let app = s(header.get("app")).ok_or_else(|| err("[registry] missing app"))?;
1430 let convention = header
1431 .get("convention")
1432 .and_then(|v| v.as_integer())
1433 .ok_or_else(|| err("[registry] missing convention"))?;
1434
1435 let (name, service_origin, description) = if let Some(svc) = doc.get("service") {
1436 (
1437 s(svc.get("name")).ok_or_else(|| err("[service] missing name"))?,
1438 Some(tok(svc.get("origin")).ok_or_else(|| err("[service] missing origin"))?),
1439 s(svc.get("description")),
1440 )
1441 } else if let Some(prod) = doc.get("producer") {
1442 (
1443 s(prod.get("name")).ok_or_else(|| err("[producer] missing name"))?,
1444 None,
1445 s(prod.get("description")),
1446 )
1447 } else {
1448 return Err(err("missing [producer] or [service]"));
1449 };
1450
1451 let budget = match doc.get("budget") {
1457 None => None,
1458 Some(b) => {
1459 let mut tables = Vec::new();
1460 for e in b
1461 .get("tables")
1462 .and_then(|v| v.as_array())
1463 .into_iter()
1464 .flatten()
1465 {
1466 tables.push(TableBudget {
1467 name: s(e.get("name")).ok_or_else(|| err("[[budget.tables]] missing name"))?,
1468 max_entries: e.get("max_entries").and_then(|v| v.as_integer()),
1469 max_bytes: e.get("max_bytes").and_then(|v| v.as_integer()),
1470 });
1471 }
1472 Some(BudgetDecl {
1473 rss_mb: b.get("rss_mb").and_then(|v| v.as_integer()),
1474 tables,
1475 })
1476 }
1477 };
1478
1479 let array = |key: &str| -> Vec<&RawValue> {
1480 doc.get(key)
1481 .and_then(|v| v.as_array())
1482 .map(|a| a.iter().collect())
1483 .unwrap_or_default()
1484 };
1485
1486 let when_of = |e: &RawValue| -> Option<Vec<Predicate>> {
1490 e.get("when").and_then(|v| v.as_array()).map(|a| {
1491 a.iter()
1492 .filter_map(|v| v.as_str())
1493 .map(Predicate::parse)
1494 .collect()
1495 })
1496 };
1497
1498 let mut subjects = Vec::new();
1499 for e in array("subject") {
1500 subjects.push(SubjectDecl {
1501 path: s(e.get("path")).ok_or_else(|| err("[[subject]] missing path"))?,
1502 class: tok(e.get("class")).ok_or_else(|| err("[[subject]] missing class"))?,
1503 type_name: s(e.get("type")).unwrap_or_default(),
1504 kind: tok(e.get("kind")),
1505 common: tok(e.get("common")),
1506 since: s(e.get("since")),
1507 description: s(e.get("description")),
1508 qos: tok(e.get("qos")),
1509 ttl_s: e.get("ttl_s").and_then(|v| v.as_integer()),
1510 unit: s(e.get("unit")),
1511 rate: e.get("rate").and_then(|v| v.as_str()).map(RateClass::parse),
1512 cardinality: e.get("cardinality").and_then(|v| v.as_integer()),
1513 encoding: enc(e.get("encoding")),
1514 buckets: e.get("buckets").and_then(|v| v.as_array()).and_then(|a| {
1517 a.iter()
1518 .map(|v| v.as_float().or_else(|| v.as_integer().map(|i| i as f64)))
1519 .collect::<Option<Vec<f64>>>()
1520 .map(Buckets::new)
1521 }),
1522 semantic: tok(e.get("semantic")),
1523 when: when_of(e),
1524 gate_note: s(e.get("gate_note")),
1525 exposure: tok(e.get("exposure")),
1526 });
1527 }
1528
1529 let mut procedures = Vec::new();
1530 for e in array("procedure") {
1531 procedures.push(ProcedureDecl {
1532 path: s(e.get("path")).ok_or_else(|| err("[[procedure]] missing path"))?,
1533 kind: tok(e.get("kind")),
1534 reply: s(e.get("reply")),
1535 request: s(e.get("request")),
1536 fanout: tok(e.get("fanout")),
1537 idempotent: e.get("idempotent").and_then(|v| v.as_bool()),
1538 cardinality: e.get("cardinality").and_then(|v| v.as_integer()),
1539 encoding: enc(e.get("encoding")),
1540 when: when_of(e),
1541 gate_note: s(e.get("gate_note")),
1542 exposure: tok(e.get("exposure")),
1543 sensitive: e.get("sensitive").and_then(|v| v.as_bool()),
1544 since: s(e.get("since")),
1545 description: s(e.get("description")),
1546 });
1547 }
1548
1549 let mut blob = Vec::new();
1556 for e in array("blob") {
1557 blob.push(BlobDecl {
1558 tier: tok(e.get("tier")).ok_or_else(|| err("[[blob]] missing tier"))?,
1559 endpoints: e
1560 .get("endpoints")
1561 .and_then(|v| v.as_array())
1562 .map(|a| {
1563 a.iter()
1564 .filter_map(|v| v.as_str())
1565 .map(str::to_string)
1566 .collect()
1567 })
1568 .unwrap_or_default(),
1569 algo: s(e.get("algo")),
1570 reference: s(e.get("reference")),
1571 encoding: enc(e.get("encoding")),
1572 since: s(e.get("since")),
1573 description: s(e.get("description")),
1574 });
1575 }
1576
1577 let mut media = Vec::new();
1582 for e in array("media") {
1583 media.push(MediaDecl {
1584 path: s(e.get("path")).ok_or_else(|| err("[[media]] missing path"))?,
1585 encoding: enc(e.get("encoding")).ok_or_else(|| err("[[media]] missing encoding"))?,
1586 attachment: s(e.get("attachment")),
1587 cardinality: e.get("cardinality").and_then(|v| v.as_integer()),
1588 since: s(e.get("since")),
1589 description: s(e.get("description")),
1590 });
1591 }
1592
1593 let mut errors = Vec::new();
1596 for e in array("error") {
1597 errors.push(ErrorDecl {
1598 name: s(e.get("name")).ok_or_else(|| err("[[error]] missing name"))?,
1599 procedures: e
1600 .get("procedures")
1601 .and_then(|v| v.as_array())
1602 .map(|a| {
1603 a.iter()
1604 .filter_map(|v| v.as_str().map(str::to_string))
1605 .collect()
1606 })
1607 .unwrap_or_default(),
1608 since: s(e.get("since")),
1609 description: s(e.get("description")),
1610 });
1611 }
1612
1613 let mut deprecated = Vec::new();
1614 for e in array("deprecated") {
1615 let path = s(e.get("path")).ok_or_else(|| err("[[deprecated]] missing path"))?;
1616 let kind = match s(e.get("kind")) {
1617 None => DeprecatedKind::Subject,
1618 Some(k) => k.parse().map_err(|()| {
1619 err(&format!(
1620 "[[deprecated]] {path:?} has kind = {k:?} — it is `subject` \
1621 (the default), `procedure` or `error`"
1622 ))
1623 })?,
1624 };
1625 deprecated.push(DeprecationDecl {
1626 path,
1627 kind,
1628 since: s(e.get("since")),
1629 replaced_by: s(e.get("replaced_by")),
1630 });
1631 }
1632
1633 Ok(RegistrySlice {
1634 version,
1635 app,
1636 convention,
1637 name,
1638 service_origin,
1639 description,
1640 subjects,
1641 procedures,
1642 blob,
1643 media,
1644 errors,
1645 deprecated,
1646 budget,
1647 })
1648}
1649
1650pub fn toml_quote(value: &str) -> String {
1668 let mut out = String::with_capacity(value.len() + 2);
1669 out.push('"');
1670 for c in value.chars() {
1671 match c {
1672 '"' => out.push_str("\\\""),
1673 '\\' => out.push_str("\\\\"),
1674 '\n' => out.push_str("\\n"),
1675 '\r' => out.push_str("\\r"),
1676 '\t' => out.push_str("\\t"),
1677 c => out.push(c),
1678 }
1679 }
1680 out.push('"');
1681 out
1682}
1683
1684pub fn to_toml(slice: &RegistrySlice) -> String {
1685 fn s(value: &str) -> String {
1686 toml_quote(value)
1687 }
1688 fn opt(out: &mut String, key: &str, value: Option<&str>) {
1689 if let Some(v) = value {
1690 out.push_str(&format!("{key} = {}\n", s(v)));
1691 }
1692 }
1693 fn opt_tok<T: SliceToken>(out: &mut String, key: &str, value: Option<&Declared<T>>) {
1697 if let Some(v) = value {
1698 out.push_str(&format!("{key} = {}\n", s(v.token())));
1699 }
1700 }
1701 fn opt_enc(out: &mut String, key: &str, value: Option<&WireEncoding>) {
1702 if let Some(v) = value {
1703 out.push_str(&format!("{key} = {}\n", s(v.as_encoding_str())));
1704 }
1705 }
1706 fn opt_int(out: &mut String, key: &str, value: Option<i64>) {
1707 if let Some(v) = value {
1708 out.push_str(&format!("{key} = {v}\n"));
1709 }
1710 }
1711
1712 let mut out = String::new();
1713 out.push_str("[registry]\n");
1714 out.push_str(&format!("version = {}\n", s(&slice.version)));
1715 out.push_str(&format!("app = {}\n", s(&slice.app)));
1716 out.push_str(&format!("convention = {}\n", slice.convention));
1717
1718 match &slice.service_origin {
1719 Some(origin) => {
1720 out.push_str("\n[service]\n");
1721 out.push_str(&format!("name = {}\n", s(&slice.name)));
1722 out.push_str(&format!("origin = {}\n", s(origin.token())));
1723 }
1724 None => {
1725 out.push_str("\n[producer]\n");
1726 out.push_str(&format!("name = {}\n", s(&slice.name)));
1727 }
1728 }
1729 opt(&mut out, "description", slice.description.as_deref());
1730
1731 if let Some(b) = &slice.budget {
1736 out.push_str("\n[budget]\n");
1737 opt_int(&mut out, "rss_mb", b.rss_mb);
1738 for t in &b.tables {
1739 out.push_str("\n[[budget.tables]]\n");
1740 out.push_str(&format!("name = {}\n", s(&t.name)));
1741 opt_int(&mut out, "max_entries", t.max_entries);
1742 opt_int(&mut out, "max_bytes", t.max_bytes);
1743 }
1744 }
1745
1746 for d in &slice.subjects {
1747 out.push_str("\n[[subject]]\n");
1748 out.push_str(&format!("path = {}\n", s(&d.path)));
1749 out.push_str(&format!("class = {}\n", s(d.class.token())));
1750 if !d.type_name.is_empty() {
1751 out.push_str(&format!("type = {}\n", s(&d.type_name)));
1752 }
1753 opt_tok(&mut out, "kind", d.kind.as_ref());
1756 opt_tok(&mut out, "common", d.common.as_ref());
1757 opt_tok(&mut out, "qos", d.qos.as_ref());
1758 opt_int(&mut out, "ttl_s", d.ttl_s);
1759 opt(&mut out, "unit", d.unit.as_deref());
1760 opt(
1761 &mut out,
1762 "rate",
1763 d.rate.as_ref().map(RateClass::token).as_deref(),
1764 );
1765 opt_int(&mut out, "cardinality", d.cardinality);
1766 opt_enc(&mut out, "encoding", d.encoding.as_ref());
1767 if let Some(b) = &d.buckets {
1768 let items: Vec<String> = b.as_slice().iter().map(|v| format!("{v:?}")).collect();
1771 out.push_str(&format!("buckets = [{}]\n", items.join(", ")));
1772 }
1773 opt_tok(&mut out, "semantic", d.semantic.as_ref());
1774 when_toml(&mut out, d.when.as_deref());
1775 opt(&mut out, "gate_note", d.gate_note.as_deref());
1776 opt_tok(&mut out, "exposure", d.exposure.as_ref());
1777 opt(&mut out, "since", d.since.as_deref());
1778 opt(&mut out, "description", d.description.as_deref());
1779 }
1780
1781 for d in &slice.procedures {
1782 out.push_str("\n[[procedure]]\n");
1783 out.push_str(&format!("path = {}\n", s(&d.path)));
1784 opt_tok(&mut out, "kind", d.kind.as_ref());
1785 opt(&mut out, "request", d.request.as_deref());
1786 opt(&mut out, "reply", d.reply.as_deref());
1787 opt_enc(&mut out, "encoding", d.encoding.as_ref());
1788 opt_tok(&mut out, "fanout", d.fanout.as_ref());
1789 if let Some(i) = d.idempotent {
1790 out.push_str(&format!("idempotent = {i}\n"));
1791 }
1792 opt_int(&mut out, "cardinality", d.cardinality);
1793 when_toml(&mut out, d.when.as_deref());
1794 opt(&mut out, "gate_note", d.gate_note.as_deref());
1795 opt_tok(&mut out, "exposure", d.exposure.as_ref());
1796 if let Some(s) = d.sensitive {
1797 out.push_str(&format!("sensitive = {s}\n"));
1798 }
1799 opt(&mut out, "since", d.since.as_deref());
1800 opt(&mut out, "description", d.description.as_deref());
1801 }
1802
1803 for d in &slice.blob {
1804 out.push_str("\n[[blob]]\n");
1805 out.push_str(&format!("tier = {}\n", s(d.tier.token())));
1806 if !d.endpoints.is_empty() {
1807 let items: Vec<String> = d.endpoints.iter().map(|e| s(e)).collect();
1808 out.push_str(&format!("endpoints = [{}]\n", items.join(", ")));
1809 }
1810 opt(&mut out, "algo", d.algo.as_deref());
1811 opt(&mut out, "reference", d.reference.as_deref());
1812 opt_enc(&mut out, "encoding", d.encoding.as_ref());
1813 opt(&mut out, "since", d.since.as_deref());
1814 opt(&mut out, "description", d.description.as_deref());
1815 }
1816
1817 for d in &slice.media {
1818 out.push_str("\n[[media]]\n");
1819 out.push_str(&format!("path = {}\n", s(&d.path)));
1820 out.push_str(&format!("encoding = {}\n", s(d.encoding.as_encoding_str())));
1821 opt(&mut out, "attachment", d.attachment.as_deref());
1822 if let Some(c) = d.cardinality {
1823 out.push_str(&format!("cardinality = {c}\n"));
1824 }
1825 opt(&mut out, "since", d.since.as_deref());
1826 opt(&mut out, "description", d.description.as_deref());
1827 }
1828
1829 for d in &slice.errors {
1830 out.push_str("\n[[error]]\n");
1831 out.push_str(&format!("name = {}\n", s(&d.name)));
1832 if !d.procedures.is_empty() {
1833 let items: Vec<String> = d.procedures.iter().map(|p| s(p)).collect();
1834 out.push_str(&format!("procedures = [{}]\n", items.join(", ")));
1835 }
1836 opt(&mut out, "since", d.since.as_deref());
1837 opt(&mut out, "description", d.description.as_deref());
1838 }
1839
1840 for d in &slice.deprecated {
1841 out.push_str("\n[[deprecated]]\n");
1842 out.push_str(&format!("path = {}\n", s(&d.path)));
1843 if d.kind != DeprecatedKind::Subject {
1846 out.push_str(&format!("kind = {}\n", s(d.kind.as_str())));
1847 }
1848 opt(&mut out, "since", d.since.as_deref());
1849 opt(&mut out, "replaced_by", d.replaced_by.as_deref());
1850 }
1851
1852 out
1853}
1854
1855#[must_use]
1863pub fn to_kdl(slice: &RegistrySlice) -> String {
1864 write_kdl(&slice_to_raw(slice))
1865 .expect("every column a slice carries has a KDL spelling (RFC 08 §5.1)")
1866}
1867
1868#[must_use]
1870pub fn slice_to_raw(slice: &RegistrySlice) -> RawTable {
1871 fn str_(t: &mut RawTable, key: &str, v: Option<&str>) {
1872 if let Some(v) = v {
1873 t.fields
1874 .push((key.to_string(), RawValue::Str(v.to_string())));
1875 }
1876 }
1877 fn tok<T: SliceToken>(t: &mut RawTable, key: &str, v: Option<&Declared<T>>) {
1878 str_(t, key, v.map(Declared::token));
1879 }
1880 fn enc(t: &mut RawTable, key: &str, v: Option<&WireEncoding>) {
1881 str_(t, key, v.map(WireEncoding::as_encoding_str));
1882 }
1883 fn int(t: &mut RawTable, key: &str, v: Option<i64>) {
1884 if let Some(v) = v {
1885 t.fields.push((key.to_string(), RawValue::Int(v)));
1886 }
1887 }
1888 fn boolean(t: &mut RawTable, key: &str, v: Option<bool>) {
1889 if let Some(v) = v {
1890 t.fields.push((key.to_string(), RawValue::Bool(v)));
1891 }
1892 }
1893 fn list(t: &mut RawTable, key: &str, v: Vec<RawValue>) {
1894 t.fields.push((key.to_string(), RawValue::List(v)));
1895 }
1896 fn when(t: &mut RawTable, v: Option<&[Predicate]>) {
1897 if let Some(preds) = v {
1898 list(
1899 t,
1900 "when",
1901 preds.iter().map(|p| RawValue::Str(p.token())).collect(),
1902 );
1903 }
1904 }
1905 fn strs(v: &[String]) -> Vec<RawValue> {
1906 v.iter().map(|s| RawValue::Str(s.clone())).collect()
1907 }
1908
1909 let mut doc = RawTable::new();
1910 let mut header = RawTable::new();
1911 str_(&mut header, "version", Some(&slice.version));
1912 str_(&mut header, "app", Some(&slice.app));
1913 int(&mut header, "convention", Some(slice.convention));
1914 doc.insert("registry", RawValue::Table(header));
1915
1916 let mut owner = RawTable::new();
1917 str_(&mut owner, "name", Some(&slice.name));
1918 let owner_key = match &slice.service_origin {
1919 Some(origin) => {
1920 tok(&mut owner, "origin", Some(origin));
1921 "service"
1922 }
1923 None => "producer",
1924 };
1925 str_(&mut owner, "description", slice.description.as_deref());
1926 doc.insert(owner_key, RawValue::Table(owner));
1927
1928 if let Some(b) = &slice.budget {
1929 let mut budget = RawTable::new();
1930 int(&mut budget, "rss_mb", b.rss_mb);
1931 if !b.tables.is_empty() {
1932 let rows = b
1933 .tables
1934 .iter()
1935 .map(|t| {
1936 let mut row = RawTable::new();
1937 str_(&mut row, "name", Some(&t.name));
1938 int(&mut row, "max_entries", t.max_entries);
1939 int(&mut row, "max_bytes", t.max_bytes);
1940 RawValue::Table(row)
1941 })
1942 .collect();
1943 list(&mut budget, "tables", rows);
1944 }
1945 doc.insert("budget", RawValue::Table(budget));
1946 }
1947
1948 let mut rows = Vec::new();
1949 for d in &slice.subjects {
1950 let mut t = RawTable::new();
1951 str_(&mut t, "path", Some(&d.path));
1952 str_(&mut t, "class", Some(d.class.token()));
1953 if !d.type_name.is_empty() {
1954 str_(&mut t, "type", Some(&d.type_name));
1955 }
1956 tok(&mut t, "kind", d.kind.as_ref());
1957 tok(&mut t, "common", d.common.as_ref());
1958 tok(&mut t, "qos", d.qos.as_ref());
1959 int(&mut t, "ttl_s", d.ttl_s);
1960 str_(&mut t, "unit", d.unit.as_deref());
1961 str_(
1962 &mut t,
1963 "rate",
1964 d.rate.as_ref().map(RateClass::token).as_deref(),
1965 );
1966 int(&mut t, "cardinality", d.cardinality);
1967 enc(&mut t, "encoding", d.encoding.as_ref());
1968 if let Some(b) = &d.buckets {
1969 list(
1970 &mut t,
1971 "buckets",
1972 b.as_slice().iter().map(|v| RawValue::Float(*v)).collect(),
1973 );
1974 }
1975 tok(&mut t, "semantic", d.semantic.as_ref());
1976 when(&mut t, d.when.as_deref());
1977 str_(&mut t, "gate_note", d.gate_note.as_deref());
1978 tok(&mut t, "exposure", d.exposure.as_ref());
1979 str_(&mut t, "since", d.since.as_deref());
1980 str_(&mut t, "description", d.description.as_deref());
1981 rows.push(RawValue::Table(t));
1982 }
1983 if !rows.is_empty() {
1984 doc.insert("subject", RawValue::List(std::mem::take(&mut rows)));
1985 }
1986
1987 for d in &slice.procedures {
1988 let mut t = RawTable::new();
1989 str_(&mut t, "path", Some(&d.path));
1990 tok(&mut t, "kind", d.kind.as_ref());
1991 str_(&mut t, "request", d.request.as_deref());
1992 str_(&mut t, "reply", d.reply.as_deref());
1993 enc(&mut t, "encoding", d.encoding.as_ref());
1994 tok(&mut t, "fanout", d.fanout.as_ref());
1995 boolean(&mut t, "idempotent", d.idempotent);
1996 int(&mut t, "cardinality", d.cardinality);
1997 when(&mut t, d.when.as_deref());
1998 str_(&mut t, "gate_note", d.gate_note.as_deref());
1999 tok(&mut t, "exposure", d.exposure.as_ref());
2000 boolean(&mut t, "sensitive", d.sensitive);
2001 str_(&mut t, "since", d.since.as_deref());
2002 str_(&mut t, "description", d.description.as_deref());
2003 rows.push(RawValue::Table(t));
2004 }
2005 if !rows.is_empty() {
2006 doc.insert("procedure", RawValue::List(std::mem::take(&mut rows)));
2007 }
2008
2009 for d in &slice.blob {
2010 let mut t = RawTable::new();
2011 str_(&mut t, "tier", Some(d.tier.token()));
2012 if !d.endpoints.is_empty() {
2013 list(&mut t, "endpoints", strs(&d.endpoints));
2014 }
2015 str_(&mut t, "algo", d.algo.as_deref());
2016 str_(&mut t, "reference", d.reference.as_deref());
2017 enc(&mut t, "encoding", d.encoding.as_ref());
2018 str_(&mut t, "since", d.since.as_deref());
2019 str_(&mut t, "description", d.description.as_deref());
2020 rows.push(RawValue::Table(t));
2021 }
2022 if !rows.is_empty() {
2023 doc.insert("blob", RawValue::List(std::mem::take(&mut rows)));
2024 }
2025
2026 for d in &slice.media {
2027 let mut t = RawTable::new();
2028 str_(&mut t, "path", Some(&d.path));
2029 str_(&mut t, "encoding", Some(d.encoding.as_encoding_str()));
2030 str_(&mut t, "attachment", d.attachment.as_deref());
2031 int(&mut t, "cardinality", d.cardinality);
2032 str_(&mut t, "since", d.since.as_deref());
2033 str_(&mut t, "description", d.description.as_deref());
2034 rows.push(RawValue::Table(t));
2035 }
2036 if !rows.is_empty() {
2037 doc.insert("media", RawValue::List(std::mem::take(&mut rows)));
2038 }
2039
2040 for d in &slice.errors {
2041 let mut t = RawTable::new();
2042 str_(&mut t, "name", Some(&d.name));
2043 if !d.procedures.is_empty() {
2044 list(&mut t, "procedures", strs(&d.procedures));
2045 }
2046 str_(&mut t, "since", d.since.as_deref());
2047 str_(&mut t, "description", d.description.as_deref());
2048 rows.push(RawValue::Table(t));
2049 }
2050 if !rows.is_empty() {
2051 doc.insert("error", RawValue::List(std::mem::take(&mut rows)));
2052 }
2053
2054 for d in &slice.deprecated {
2055 let mut t = RawTable::new();
2056 str_(&mut t, "path", Some(&d.path));
2057 if d.kind != DeprecatedKind::Subject {
2058 str_(&mut t, "kind", Some(d.kind.as_str()));
2059 }
2060 str_(&mut t, "since", d.since.as_deref());
2061 str_(&mut t, "replaced_by", d.replaced_by.as_deref());
2062 rows.push(RawValue::Table(t));
2063 }
2064 if !rows.is_empty() {
2065 doc.insert("deprecated", RawValue::List(rows));
2066 }
2067 doc
2068}
2069
2070#[derive(Debug, Clone, PartialEq, Eq)]
2075pub enum SliceFinding {
2076 VersionSkew {
2078 served: String,
2079 local: String,
2080 },
2081 UnknownSubject {
2083 path: String,
2084 class: Declared<Class>,
2085 },
2086 MissingSubject {
2088 path: String,
2089 class: Declared<Class>,
2090 },
2091 UnknownProcedure {
2093 path: String,
2094 },
2095 MissingProcedure {
2096 path: String,
2097 },
2098 UnknownBlobTier {
2100 tier: Declared<BlobTier>,
2101 },
2102 MissingBlobTier {
2104 tier: Declared<BlobTier>,
2105 },
2106 UnknownMediaStream {
2110 path: String,
2111 },
2112 MissingMediaStream {
2114 path: String,
2115 },
2116 ServesDeprecated {
2118 path: String,
2119 replaced_by: Option<String>,
2120 },
2121}
2122
2123impl SliceFinding {
2124 pub fn summary(&self) -> String {
2126 match self {
2127 Self::VersionSkew { served, local } => {
2128 format!("registry {served} (we compiled {local})")
2129 }
2130 Self::UnknownSubject { path, class } => format!("serves unknown {class} {path}"),
2131 Self::MissingSubject { path, class } => format!("does not serve {class} {path}"),
2132 Self::UnknownProcedure { path } => format!("serves unknown procedure {path}"),
2133 Self::MissingProcedure { path } => format!("does not serve procedure {path}"),
2134 Self::UnknownBlobTier { tier } => format!("serves unknown @blob tier {tier}"),
2135 Self::MissingBlobTier { tier } => format!("does not serve @blob tier {tier}"),
2136 Self::UnknownMediaStream { path } => format!("serves unknown @media stream {path}"),
2137 Self::MissingMediaStream { path } => format!("does not serve @media stream {path}"),
2138 Self::ServesDeprecated { path, replaced_by } => match replaced_by {
2139 Some(r) => format!("serves deprecated {path} (use {r})"),
2140 None => format!("serves deprecated {path}"),
2141 },
2142 }
2143 }
2144}
2145
2146pub fn diff(served: &RegistrySlice, local: &RegistrySlice) -> Vec<SliceFinding> {
2151 let mut out = Vec::new();
2152 if served.version != local.version {
2153 out.push(SliceFinding::VersionSkew {
2154 served: served.version.clone(),
2155 local: local.version.clone(),
2156 });
2157 }
2158 for s in &served.subjects {
2159 if !local.serves_subject(&s.path) {
2160 out.push(SliceFinding::UnknownSubject {
2161 path: s.path.clone(),
2162 class: s.class.clone(),
2163 });
2164 }
2165 }
2166 for s in &local.subjects {
2167 if !served.serves_subject(&s.path) {
2168 out.push(SliceFinding::MissingSubject {
2169 path: s.path.clone(),
2170 class: s.class.clone(),
2171 });
2172 }
2173 }
2174 for p in &served.procedures {
2175 if !local.serves_procedure(&p.path) {
2176 out.push(SliceFinding::UnknownProcedure {
2177 path: p.path.clone(),
2178 });
2179 }
2180 }
2181 for p in &local.procedures {
2182 if !served.serves_procedure(&p.path) {
2183 out.push(SliceFinding::MissingProcedure {
2184 path: p.path.clone(),
2185 });
2186 }
2187 }
2188 for b in &served.blob {
2189 if !local.serves_blob_tier(b.tier.clone()) {
2190 out.push(SliceFinding::UnknownBlobTier {
2191 tier: b.tier.clone(),
2192 });
2193 }
2194 }
2195 for b in &local.blob {
2196 if !served.serves_blob_tier(b.tier.clone()) {
2197 out.push(SliceFinding::MissingBlobTier {
2198 tier: b.tier.clone(),
2199 });
2200 }
2201 }
2202 for m in &served.media {
2203 if !local.serves_media(&m.path) {
2204 out.push(SliceFinding::UnknownMediaStream {
2205 path: m.path.clone(),
2206 });
2207 }
2208 }
2209 for m in &local.media {
2210 if !served.serves_media(&m.path) {
2211 out.push(SliceFinding::MissingMediaStream {
2212 path: m.path.clone(),
2213 });
2214 }
2215 }
2216 for d in &served.deprecated {
2219 let still_served = match d.kind {
2220 DeprecatedKind::Subject => served.serves_subject(&d.path),
2221 DeprecatedKind::Procedure => served.serves_procedure(&d.path),
2222 DeprecatedKind::Error => served.errors.iter().any(|e| e.name == d.path),
2223 };
2224 if still_served {
2225 out.push(SliceFinding::ServesDeprecated {
2226 path: d.path.clone(),
2227 replaced_by: d.replaced_by.clone(),
2228 });
2229 }
2230 }
2231 out
2232}
2233
2234#[cfg(test)]
2235mod tests {
2236 use super::*;
2237
2238 fn bind_fixture() -> RegistrySlice {
2239 let mut slice = RegistrySlice::new("1.0.0", "test", "bmc");
2240 for (path, class) in [
2241 ("{chassis}/thermal/{sensor}/celsius", Class::Telemetry),
2242 (
2243 "{chassis}/thermal/{sensor}/upper_warning_c",
2244 Class::Telemetry,
2245 ),
2246 ("{target}/total", Class::Telemetry),
2247 ("targets/total", Class::Telemetry),
2248 ("{device}/{metric...}", Class::Telemetry),
2249 ("chassis/{chassis}", Class::State),
2250 ] {
2251 slice.subjects.push(SubjectDecl::new(path, class));
2252 }
2253 slice
2254 }
2255
2256 #[test]
2259 fn histogram_buckets_and_semantic_round_trip() {
2260 let src = r#"
2261[registry]
2262version = "1.0.0"
2263app = "test"
2264convention = 1
2265
2266[producer]
2267name = "probe"
2268
2269[[subject]]
2270path = "{target}/rtt"
2271class = "telemetry"
2272type = "TelemetryPoint"
2273kind = "histogram"
2274buckets = [0.005, 0.01, 1, 2.5, 10]
2275semantic = "duration"
2276description = "round-trip time"
2277
2278[[subject]]
2279path = "{target}/ok"
2280class = "telemetry"
2281type = "TelemetryPoint"
2282kind = "bool"
2283semantic = "state"
2284description = "reachable"
2285"#;
2286 let parsed = parse_slice(src).unwrap();
2287 let rtt = &parsed.subjects[0];
2288 assert_eq!(rtt.kind, Some(Declared::Known(SubjectKind::Histogram)));
2289 assert_eq!(
2290 rtt.buckets.as_ref().map(Buckets::as_slice),
2291 Some(&[0.005, 0.01, 1.0, 2.5, 10.0][..])
2292 );
2293 assert!(rtt.buckets.as_ref().unwrap().is_well_formed());
2294 assert_eq!(rtt.semantic, Some(Declared::Known(Semantic::Duration)));
2295 assert_eq!(parsed.subjects[1].buckets, None);
2296 let back = parse_slice(&to_toml(&parsed)).unwrap();
2297 assert_eq!(back, parsed, "parse → emit → parse is the identity");
2298 assert_eq!(SubjectKind::Histogram.payload_tag(), "histogram");
2300 assert_eq!(
2301 SubjectKind::from_payload_tag("histogram"),
2302 Some(SubjectKind::Histogram)
2303 );
2304 let odd = parse_slice(&src.replace("\"duration\"", "\"loudness\"")).unwrap();
2306 assert_eq!(
2307 odd.subjects[0].semantic,
2308 Some(Declared::Other("loudness".into()))
2309 );
2310 }
2311
2312 #[test]
2313 fn buckets_well_formedness_is_strictly_ascending_finite_and_non_empty() {
2314 assert!(Buckets::new(vec![0.1, 1.0]).is_well_formed());
2315 assert!(!Buckets::new(vec![]).is_well_formed());
2316 assert!(!Buckets::new(vec![1.0, 1.0]).is_well_formed());
2317 assert!(!Buckets::new(vec![2.0, 1.0]).is_well_formed());
2318 assert!(
2319 !Buckets::new(vec![1.0, f64::INFINITY]).is_well_formed(),
2320 "+Inf is implicit"
2321 );
2322 let nan = Buckets::new(vec![f64::NAN]);
2323 assert_eq!(nan, nan.clone(), "bit equality is reflexive, NaN included");
2324 }
2325
2326 #[test]
2328 fn bind_finds_the_declaration_and_its_vars() {
2329 let slice = bind_fixture();
2330 let b = slice
2331 .bind(
2332 Class::Telemetry,
2333 &["rack-a-1", "thermal", "inlet", "celsius"],
2334 )
2335 .expect("declared");
2336 assert_eq!(b.decl.path, "{chassis}/thermal/{sensor}/celsius");
2337 assert_eq!(
2338 b.vars,
2339 vec![
2340 ("chassis", "rack-a-1".to_string()),
2341 ("sensor", "inlet".to_string())
2342 ]
2343 );
2344 assert_eq!(
2345 b.var("sensor"),
2346 Some("inlet"),
2347 "a value with `-` binds whole"
2348 );
2349 let w = slice
2351 .bind(
2352 Class::Telemetry,
2353 &["rack-a-1", "thermal", "inlet", "upper_warning_c"],
2354 )
2355 .unwrap();
2356 assert_eq!(w.decl.path, "{chassis}/thermal/{sensor}/upper_warning_c");
2357 }
2358
2359 #[test]
2361 fn bind_prefers_the_most_literal_declaration() {
2362 let slice = bind_fixture();
2363 let lit = slice.bind(Class::Telemetry, &["targets", "total"]).unwrap();
2364 assert_eq!(lit.decl.path, "targets/total");
2365 assert!(lit.vars.is_empty());
2366 let var = slice.bind(Class::Telemetry, &["web-1", "total"]).unwrap();
2367 assert_eq!(var.decl.path, "{target}/total");
2368 assert_eq!(var.var("target"), Some("web-1"));
2369 }
2370
2371 #[test]
2373 fn bind_joins_a_rest_variable() {
2374 let slice = bind_fixture();
2375 let b = slice
2376 .bind(Class::Telemetry, &["x-sw1", "if", "ge-0", "in_octets"])
2377 .unwrap();
2378 assert_eq!(b.decl.path, "{device}/{metric...}");
2379 assert_eq!(
2380 b.var("device"),
2381 Some("x-sw1"),
2382 "a slugged `x-…` chunk is just a value"
2383 );
2384 assert_eq!(b.var("metric"), Some("if/ge-0/in_octets"));
2385 assert_eq!(b.vars.last().map(|(n, _)| *n), Some("metric"));
2386 assert!(slice.bind(Class::Telemetry, &["lonely"]).is_none());
2387 }
2388
2389 #[test]
2391 fn bind_is_none_for_an_undeclared_tail_or_another_class() {
2392 let slice = bind_fixture();
2393 assert!(
2394 slice
2395 .bind(Class::Events, &["rack-a-1", "thermal", "inlet", "celsius"])
2396 .is_none()
2397 );
2398 assert!(slice.bind(Class::State, &["targets", "total"]).is_none());
2399 let state = slice.bind(Class::State, &["chassis", "rack-a-1"]).unwrap();
2400 assert_eq!(state.var("chassis"), Some("rack-a-1"));
2401 assert!(slice.bind(Class::Telemetry, &[]).is_none());
2402 }
2403
2404 #[test]
2406 fn bind_skips_an_unparsable_declaration() {
2407 let mut slice = RegistrySlice::new("1.0.0", "test", "odd");
2408 slice
2409 .subjects
2410 .push(SubjectDecl::new("{rest...}/after", Class::Telemetry));
2411 slice
2412 .subjects
2413 .push(SubjectDecl::new("ok/{x}", Class::Telemetry));
2414 assert!(slice.bind(Class::Telemetry, &["a", "after"]).is_none());
2415 assert_eq!(
2416 slice.bind(Class::Telemetry, &["ok", "1"]).unwrap().var("x"),
2417 Some("1")
2418 );
2419 }
2420
2421 #[test]
2432 fn error_entries_round_trip() {
2433 let source = "[registry]\nversion = \"2.0\"\napp = \"acme\"\nconvention = 1\n\n[producer]\nname = \"modem\"\n\n[[error]]\nname = \"restart-required\"\nprocedures = [\"config/{device}/{group}/set\"]\nsince = \"2.0\"\ndescription = \"the parameter is part of what the transport was started against\"\n\n[[error]]\nname = \"device-refused\"\n\n[[deprecated]]\npath = \"not-ready\"\nkind = \"error\"\nsince = \"2.0\"\n";
2434 let slice = parse_slice(source).expect("parses");
2435 assert_eq!(slice.errors.len(), 2);
2436 assert_eq!(slice.errors[0].name, "restart-required");
2437 assert_eq!(slice.errors[0].procedures, ["config/{device}/{group}/set"]);
2438 assert_eq!(
2439 slice.errors[0].wire_name("modem"),
2440 "error/modem/restart-required"
2441 );
2442 assert!(slice.errors[1].procedures.is_empty());
2443 assert_eq!(slice.deprecated[0].kind, DeprecatedKind::Error);
2444 let again = parse_slice(&to_toml(&slice)).expect("the export parses");
2445 assert_eq!(again, slice);
2446 let bare = parse_slice("[registry]\nversion = \"1.0\"\napp = \"a\"\nconvention = 1\n[producer]\nname = \"p\"\n").expect("parses");
2447 assert!(bare.errors.is_empty());
2448 assert!(!to_toml(&bare).contains("[[error]]"));
2449 }
2450
2451 #[test]
2455 fn when_predicates_round_trip_and_bind_their_error() {
2456 let source = "[registry]\nversion = \"2.0\"\napp = \"acme\"\nconvention = 1\n\n[producer]\nname = \"modem\"\n\n[[subject]]\npath = \"{device}/radio/rssi\"\nclass = \"telemetry\"\ntype = \"Gauge\"\ncardinality = 4\nwhen = [\"capability:rssi\", \"config:collect.radio\", \"weather:sunny\", \"nocolon\"]\ngate_note = \"only a device that measures dBm\"\n\n[[procedure]]\npath = \"device/{device}/sdu/lanes\"\nkind = \"read\"\nreply = \"LaneTable\"\ncardinality = 4\nwhen = [\"feature:mgmt\"]\n";
2457 let slice = parse_slice(source).expect("parses");
2458 let when = slice.subjects[0].when.as_ref().expect("declared");
2459 assert_eq!(when.len(), 4);
2460 assert_eq!(when[0], Predicate::new(PredicateKind::Capability, "rssi"));
2461 assert_eq!(when[1].name, "collect.radio");
2462 assert_eq!(when[2].kind, Declared::Other("weather".into()));
2463 assert_eq!(when[2].token(), "weather:sunny");
2464 assert_eq!(
2465 when[3].token(),
2466 "nocolon",
2467 "a colon-less token is kept whole"
2468 );
2469 assert_eq!(
2470 slice.subjects[0].gate_note.as_deref(),
2471 Some("only a device that measures dBm")
2472 );
2473 let pw = slice.procedures[0].when.as_ref().expect("declared");
2474 assert_eq!(pw, &[Predicate::new(PredicateKind::Feature, "mgmt")]);
2475 assert_eq!(PredicateKind::Feature.gated_error(), "error/unsupported");
2476 assert_eq!(PredicateKind::Config.gated_error(), "error/gated");
2477 assert_eq!(PredicateKind::Capability.gated_error(), "error/gated");
2478 let again = parse_slice(&to_toml(&slice)).expect("the export parses");
2479 assert_eq!(again, slice);
2480 let bare = parse_slice("[registry]\nversion = \"1.0\"\napp = \"a\"\nconvention = 1\n[producer]\nname = \"p\"\n[[subject]]\npath = \"x\"\nclass = \"state\"\ntype = \"T\"\n").expect("parses");
2481 assert!(bare.subjects[0].when.is_none());
2482 assert!(!to_toml(&bare).contains("when"));
2483 }
2484
2485 #[test]
2488 fn exposure_and_sensitive_round_trip() {
2489 let source = "[registry]\nversion = \"2.0\"\napp = \"acme\"\nconvention = 1\n\n[producer]\nname = \"modem\"\n\n[[subject]]\npath = \"{device}/tx_sdus_total\"\nclass = \"telemetry\"\ntype = \"Counter\"\ncardinality = 4\nexposure = \"host\"\n\n[[subject]]\npath = \"{device}/queue_depth\"\nclass = \"telemetry\"\ntype = \"Gauge\"\ncardinality = 4\nexposure = \"orbit\"\n\n[[procedure]]\npath = \"config/{device}/access/set\"\nkind = \"write\"\nrequest = \"ConfigChange\"\nreply = \"ConfigView\"\ncardinality = 4\nexposure = \"host\"\nsensitive = true\n";
2490 let slice = parse_slice(source).expect("parses");
2491 assert_eq!(
2492 slice.subjects[0].exposure,
2493 Some(Declared::Known(Exposure::Host))
2494 );
2495 assert_eq!(
2496 slice.subjects[1].exposure,
2497 Some(Declared::Other("orbit".into()))
2498 );
2499 assert_eq!(
2500 slice.procedures[0].exposure,
2501 Some(Declared::Known(Exposure::Host))
2502 );
2503 assert_eq!(slice.procedures[0].sensitive, Some(true));
2504 let again = parse_slice(&to_toml(&slice)).expect("the export parses");
2505 assert_eq!(again, slice);
2506 assert_eq!(Exposure::from_token("link"), Some(Exposure::Link));
2507 assert_eq!(Exposure::Fleet.token(), "fleet");
2508 let bare = parse_slice("[registry]\nversion = \"1.0\"\napp = \"a\"\nconvention = 1\n[producer]\nname = \"p\"\n[[subject]]\npath = \"x\"\nclass = \"state\"\ntype = \"T\"\n").expect("parses");
2509 assert!(bare.subjects[0].exposure.is_none());
2510 assert!(!to_toml(&bare).contains("exposure"));
2511 }
2512
2513 #[test]
2514 fn toml_export_round_trips_every_carried_field() {
2515 let source = r#"
2516 [registry]
2517 version = "2.1"
2518 app = "acme"
2519 convention = 1
2520 [producer]
2521 name = "netring"
2522 description = "flow capture"
2523 [budget]
2524 rss_mb = 64
2525 [[budget.tables]]
2526 name = "flows"
2527 max_entries = 65536
2528 max_bytes = 16777216
2529 [[budget.tables]]
2530 name = "names"
2531 max_entries = 16384
2532 [[subject]]
2533 path = "flows/{proto}/count"
2534 class = "telemetry"
2535 type = "TelemetryPoint"
2536 kind = "counter"
2537 qos = "sampled"
2538 unit = "packets"
2539 cardinality = 512
2540 encoding = "application/cbor"
2541 since = "1.0"
2542 description = "per-protocol flow counter"
2543 [[subject]]
2544 path = "health"
2545 class = "state"
2546 type = "Health"
2547 common = "health"
2548 ttl_s = 60
2549 rate = "burst"
2550 [[procedure]]
2551 path = "capture/trigger"
2552 kind = "write"
2553 request = "CaptureSpec"
2554 reply = "Ack"
2555 encoding = "application/json"
2556 fanout = "forbidden"
2557 idempotent = false
2558 since = "1.1"
2559 description = "start a capture"
2560 [[procedure]]
2561 path = "capture/{port}/drain"
2562 kind = "write"
2563 reply = "Ack"
2564 cardinality = 8
2565 since = "1.2"
2566 description = "drain one port's ring"
2567 [[blob]]
2568 tier = "artifact"
2569 endpoints = ["manifest", "slice", "have"]
2570 reference = "ArtifactRef"
2571 encoding = "application/octet-stream"
2572 since = "1.2"
2573 description = "captured pcaps"
2574 [[blob]]
2575 tier = "store"
2576 algo = "blake3"
2577 since = "1.2"
2578 [[media]]
2579 path = "{stream}/preview/jpeg"
2580 encoding = "image/jpeg"
2581 attachment = "FrameMeta"
2582 cardinality = 16
2583 since = "1.3"
2584 description = "preview rung"
2585 [[deprecated]]
2586 path = "flows/legacy"
2587 since = "2.0"
2588 replaced_by = "flows/{proto}/count"
2589 [[deprecated]]
2590 kind = "procedure"
2591 path = "flows/reset"
2592 since = "2.0"
2593 "#;
2594 let parsed = parse_slice(source).unwrap();
2595 let emitted = to_toml(&parsed);
2596 let back = parse_slice(&emitted)
2597 .unwrap_or_else(|e| panic!("exported TOML must re-parse: {e}\n---\n{emitted}"));
2598 assert_eq!(back, parsed, "exported TOML:\n{emitted}");
2599 }
2600
2601 #[test]
2607 fn a_budget_is_carried_and_its_absence_stays_unwritten() {
2608 let source = "[registry]\nversion = \"1.0\"\napp = \"t\"\nconvention = 1\n\n\
2609 [producer]\nname = \"netring\"\n\n[budget]\nrss_mb = 64\n\n\
2610 [[budget.tables]]\nname = \"flows\"\nmax_entries = 65536\n\
2611 max_bytes = 16777216\n\n[[budget.tables]]\nname = \"names\"\n\
2612 max_entries = 16384\n";
2613 let parsed = parse_slice(source).unwrap();
2614 let budget = parsed.budget.as_ref().expect("[budget] carried");
2615 assert_eq!(budget.rss_mb, Some(64));
2616 assert_eq!(budget.tables.len(), 2);
2617 assert_eq!(budget.tables[0].name, "flows");
2618 assert_eq!(budget.tables[0].max_entries, Some(65536));
2619 assert_eq!(budget.tables[0].max_bytes, Some(16_777_216));
2620 assert_eq!(budget.tables[1].name, "names");
2621 assert_eq!(budget.tables[1].max_bytes, None);
2622 assert_eq!(
2623 to_toml(&parsed),
2624 source,
2625 "the emitter is the parser's inverse"
2626 );
2627
2628 let bare = "[registry]\nversion = \"1.0\"\napp = \"t\"\nconvention = 1\n\n\
2629 [producer]\nname = \"netring\"\n";
2630 let parsed = parse_slice(bare).unwrap();
2631 assert_eq!(parsed.budget, None, "no [budget] is not an empty one");
2632 assert_eq!(to_toml(&parsed), bare);
2633 }
2634
2635 #[test]
2639 fn a_budget_table_without_a_name_is_refused() {
2640 let header = "[registry]\nversion = \"1.0\"\napp = \"t\"\nconvention = 1\n\
2641 [producer]\nname = \"netring\"\n";
2642 let e = parse_slice(&format!("{header}[[budget.tables]]\nmax_entries = 4\n"))
2643 .expect_err("a nameless row");
2644 assert!(
2645 e.to_string().contains("[[budget.tables]] missing name"),
2646 "{e}"
2647 );
2648 let ok = parse_slice(&format!("{header}[budget]\n")).unwrap();
2649 assert_eq!(
2650 ok.budget,
2651 Some(BudgetDecl::new()),
2652 "an empty table is carried"
2653 );
2654 }
2655
2656 #[test]
2659 fn a_deprecation_without_a_kind_is_a_subject_and_stays_unwritten() {
2660 let source = r#"
2661 [registry]
2662 version = "1.0"
2663 app = "t"
2664 convention = 1
2665 [producer]
2666 name = "p"
2667 [[deprecated]]
2668 path = "old"
2669 [[deprecated]]
2670 kind = "procedure"
2671 path = "reset"
2672 "#;
2673 let parsed = parse_slice(source).unwrap();
2674 assert_eq!(parsed.deprecated[0].kind, DeprecatedKind::Subject);
2675 assert_eq!(parsed.deprecated[1].kind, DeprecatedKind::Procedure);
2676 let emitted = to_toml(&parsed);
2677 assert!(
2680 emitted.contains("[[deprecated]]\npath = \"old\"\n"),
2681 "{emitted}"
2682 );
2683 assert!(
2684 emitted.contains("path = \"reset\"\nkind = \"procedure\"\n"),
2685 "{emitted}"
2686 );
2687 }
2688
2689 #[test]
2691 fn a_deprecation_with_an_unknown_kind_is_refused() {
2692 let source = r#"
2693 [registry]
2694 version = "1.0"
2695 app = "t"
2696 convention = 1
2697 [producer]
2698 name = "p"
2699 [[deprecated]]
2700 kind = "media"
2701 path = "old"
2702 "#;
2703 let e = parse_slice(source).expect_err("kind is a closed vocabulary");
2704 assert!(e.to_string().contains("`procedure`"), "{e}");
2705 }
2706
2707 #[test]
2710 fn toml_export_keeps_a_service_origin() {
2711 let parsed = parse_slice(
2712 r#"
2713 [registry]
2714 version = "1.0"
2715 app = "acme"
2716 convention = 1
2717 [service]
2718 name = "catalog"
2719 origin = "@catalog"
2720 "#,
2721 )
2722 .unwrap();
2723 let emitted = to_toml(&parsed);
2724 assert!(emitted.contains("[service]"), "{emitted}");
2725 assert_eq!(parse_slice(&emitted).unwrap(), parsed);
2726 }
2727
2728 #[test]
2731 fn toml_export_escapes_free_text() {
2732 let mut parsed = parse_slice(
2733 r#"
2734 [registry]
2735 version = "1.0"
2736 app = "acme"
2737 convention = 1
2738 [producer]
2739 name = "p"
2740 "#,
2741 )
2742 .unwrap();
2743 parsed.description = Some("a \"quoted\" \\ back\nslash".into());
2744 let emitted = to_toml(&parsed);
2745 assert_eq!(parse_slice(&emitted).unwrap(), parsed, "{emitted}");
2746 }
2747
2748 #[test]
2749 fn a_service_slice_carries_its_origin() {
2750 let slice = parse_slice(
2751 r#"
2752 [registry]
2753 version = "1.0"
2754 app = "acme"
2755 convention = 1
2756 [service]
2757 name = "catalog"
2758 origin = "@catalog"
2759 [[subject]]
2760 path = "entity/{entity_id}"
2761 class = "state"
2762 type = "Entity"
2763 [[procedure]]
2764 path = "introspect"
2765 kind = "read"
2766 "#,
2767 )
2768 .unwrap();
2769 assert_eq!(
2770 slice.service_origin.as_ref().map(Declared::token),
2771 Some("@catalog")
2772 );
2773 assert!(slice.serves_procedure("introspect"));
2774 }
2775
2776 #[test]
2780 fn a_malformed_slice_keeps_the_toml_error_as_its_source() {
2781 use std::error::Error as _;
2782
2783 let err = parse_slice("[registry\nthis is not = = toml").unwrap_err();
2785 assert!(matches!(err, SliceError::Toml(_)), "{err:?}");
2786
2787 let source = err.source().expect("a toml error underneath");
2790 assert!(
2791 source.downcast_ref::<toml::de::Error>().is_some(),
2792 "source was {source:?}"
2793 );
2794
2795 let shape = parse_slice("[registry]\nversion = \"1.0\"\n").unwrap_err();
2798 assert!(matches!(shape, SliceError::Shape(_)), "{shape:?}");
2799 assert!(shape.source().is_none());
2800 assert!(shape.to_string().contains("missing app"), "{shape}");
2801 }
2802
2803 #[test]
2808 fn a_declared_column_keeps_what_it_does_not_recognise() {
2809 let src = r#"
2810 [registry]
2811 version = "1.0"
2812 app = "acme"
2813 convention = 1
2814 [producer]
2815 name = "netring"
2816 [[subject]]
2817 path = "a"
2818 class = "telemetry"
2819 qos = "sampled"
2820 [[subject]]
2821 path = "b"
2822 class = "metrics"
2823 qos = "urgent"
2824 "#;
2825 let slice = parse_slice(src).unwrap();
2826
2827 assert_eq!(slice.subjects[0].class, Declared::Known(Class::Telemetry));
2828 assert_eq!(
2829 slice.subjects[0].qos,
2830 Some(Declared::Known(QosProfile::Sampled))
2831 );
2832 assert_eq!(slice.subjects[1].class, Declared::Other("metrics".into()));
2834 assert_eq!(
2835 slice.subjects[1].qos,
2836 Some(Declared::Other("urgent".into()))
2837 );
2838 assert_eq!(slice.subjects[1].class.token(), "metrics");
2839 assert_eq!(slice.subjects[1].class.known(), None);
2840
2841 assert_eq!(slice.subjects_in(Class::Telemetry).count(), 1);
2844 assert_eq!(
2845 slice
2846 .subjects_in(Declared::Other("metrics".to_string()))
2847 .count(),
2848 1
2849 );
2850
2851 assert_eq!(parse_slice(&to_toml(&slice)).unwrap(), slice);
2853 }
2854
2855 #[test]
2860 fn a_rate_class_keeps_its_burst_budget() {
2861 assert_eq!(RateClass::parse("rare").cap_per_hour(), Some(1));
2862 assert_eq!(RateClass::parse("low").cap_per_hour(), Some(60));
2863 assert_eq!(RateClass::parse("burst(240/h)"), RateClass::Burst(240));
2864 assert_eq!(RateClass::parse("burst(240/h)").cap_per_hour(), Some(240));
2865 assert_eq!(RateClass::parse("burst(240)").cap_per_hour(), None);
2869
2870 let odd = RateClass::parse("whenever");
2873 assert_eq!(odd, RateClass::Other("whenever".into()));
2874 assert_eq!(odd.cap_per_hour(), None);
2875
2876 for token in ["rare", "low", "burst(240/h)", "whenever"] {
2878 assert_eq!(RateClass::parse(token).token(), token);
2879 }
2880 }
2881
2882 #[test]
2888 fn blob_entries_parse_lax_with_only_tier_required() {
2889 let header = r#"
2890 [registry]
2891 version = "1.8"
2892 app = "acme"
2893 convention = 1
2894 [producer]
2895 name = "netring"
2896 "#;
2897 let slice = parse_slice(&format!(
2898 r#"{header}
2899 [[blob]]
2900 tier = "artifact"
2901 endpoints = ["manifest", "have"]
2902 reference = "Delivery"
2903 [[blob]]
2904 tier = "flux"
2905 "#
2906 ))
2907 .unwrap();
2908 assert!(slice.serves_blob_tier(BlobTier::Artifact));
2909 let decl = slice
2910 .blob
2911 .iter()
2912 .find(|b| b.tier.is(&BlobTier::Artifact))
2913 .unwrap();
2914 assert_eq!(decl.endpoints, ["manifest", "have"]);
2915 assert_eq!(decl.reference.as_deref(), Some("Delivery"));
2916 assert_eq!(decl.algo, None);
2917 assert!(slice.serves_blob_tier(Declared::Other("flux".into())));
2919
2920 assert!(parse_slice(&format!("{header}\n[[blob]]\nalgo = \"blake3\"\n")).is_err());
2922
2923 let old = parse_slice(header).unwrap();
2925 assert!(old.blob.is_empty());
2926 assert!(!old.serves_blob_tier(BlobTier::Artifact));
2927 }
2928
2929 #[test]
2933 fn blob_tier_drift_is_a_finding() {
2934 let with = |tiers: &[&str]| {
2935 let mut src = String::from(
2936 "[registry]\nversion = \"1.8\"\napp = \"acme\"\nconvention = 1\n\
2937 [producer]\nname = \"netring\"\n",
2938 );
2939 for t in tiers {
2940 src.push_str(&format!("[[blob]]\ntier = {t:?}\n"));
2941 }
2942 parse_slice(&src).unwrap()
2943 };
2944 let served = with(&["artifact", "tree"]);
2945 let local = with(&["tree", "store"]);
2946 let findings = diff(&served, &local);
2947 assert!(findings.iter().any(
2948 |f| matches!(f, SliceFinding::UnknownBlobTier { tier } if tier.is(&BlobTier::Artifact))
2949 ));
2950 assert!(findings.iter().any(
2951 |f| matches!(f, SliceFinding::MissingBlobTier { tier } if tier.is(&BlobTier::Store))
2952 ));
2953 assert!(diff(&served, &served).is_empty());
2954 }
2955
2956 #[test]
2961 fn media_stream_drift_is_a_finding() {
2962 let with = |paths: &[&str]| {
2963 let mut src = String::from(
2964 "[registry]\nversion = \"1.16\"\napp = \"acme\"\nconvention = 1\n\
2965 [producer]\nname = \"netring\"\n",
2966 );
2967 for p in paths {
2968 src.push_str(&format!(
2969 "[[media]]\npath = {p:?}\nencoding = \"image/jpeg\"\n"
2970 ));
2971 }
2972 parse_slice(&src).unwrap()
2973 };
2974 let served = with(&["{stream}/preview/jpeg", "{stream}/live/h264"]);
2975 let local = with(&["{stream}/live/h264", "{stream}/still/png"]);
2976 let findings = diff(&served, &local);
2977 assert!(findings.iter().any(|f| matches!(
2978 f,
2979 SliceFinding::UnknownMediaStream { path } if path == "{stream}/preview/jpeg"
2980 )));
2981 assert!(findings.iter().any(|f| matches!(
2982 f,
2983 SliceFinding::MissingMediaStream { path } if path == "{stream}/still/png"
2984 )));
2985 assert!(!findings.iter().any(|f| matches!(
2987 f,
2988 SliceFinding::UnknownMediaStream { path }
2989 | SliceFinding::MissingMediaStream { path } if path == "{stream}/live/h264"
2990 )));
2991 assert!(diff(&served, &served).is_empty());
2992
2993 let old = with(&[]);
2996 assert!(diff(&old, &old).is_empty());
2997 assert!(
2998 diff(&old, &local)
2999 .iter()
3000 .all(|f| matches!(f, SliceFinding::MissingMediaStream { .. }))
3001 );
3002 }
3003
3004 #[test]
3005 fn a_slice_identical_to_ours_is_no_finding() {
3006 let slice = parse_slice(
3007 r#"
3008 [registry]
3009 version = "1.0"
3010 app = "acme"
3011 convention = 1
3012 [producer]
3013 name = "sysinfo"
3014 [[subject]]
3015 path = "cpu/usage"
3016 class = "telemetry"
3017 type = "TelemetryPoint"
3018 [[procedure]]
3019 path = "introspect"
3020 kind = "read"
3021 "#,
3022 )
3023 .unwrap();
3024 assert!(diff(&slice, &slice).is_empty());
3025 }
3026
3027 #[test]
3028 fn skew_and_drift_are_findings() {
3029 let local = parse_slice(
3030 r#"
3031 [registry]
3032 version = "1.1"
3033 app = "zensight"
3034 convention = 1
3035 [producer]
3036 name = "sysinfo"
3037 [[subject]]
3038 path = "cpu/usage"
3039 class = "telemetry"
3040 type = "TelemetryPoint"
3041 [[procedure]]
3042 path = "introspect"
3043 kind = "read"
3044 "#,
3045 )
3046 .unwrap();
3047 let served = parse_slice(
3048 r#"
3049 [registry]
3050 version = "1.2"
3051 app = "zensight"
3052 convention = 1
3053 [producer]
3054 name = "sysinfo"
3055 [[subject]]
3056 path = "cpu/temperature"
3057 class = "telemetry"
3058 type = "TelemetryPoint"
3059 [[procedure]]
3060 path = "introspect"
3061 kind = "read"
3062 "#,
3063 )
3064 .unwrap();
3065
3066 let findings = diff(&served, &local);
3067 assert!(findings.iter().any(|f| matches!(
3068 f,
3069 SliceFinding::VersionSkew { served, local } if served == "1.2" && local == "1.1"
3070 )));
3071 assert!(findings.iter().any(
3072 |f| matches!(f, SliceFinding::UnknownSubject { path, .. } if path == "cpu/temperature")
3073 ));
3074 assert!(findings.iter().any(
3075 |f| matches!(f, SliceFinding::MissingSubject { path, .. } if path == "cpu/usage")
3076 ));
3077 }
3078
3079 #[test]
3083 fn unknown_fields_do_not_break_the_parse() {
3084 let slice = parse_slice(
3085 r#"
3086 [registry]
3087 version = "9.9"
3088 app = "zensight"
3089 convention = 1
3090 future_knob = true
3091 [producer]
3092 name = "sysinfo"
3093 [[subject]]
3094 path = "cpu/usage"
3095 class = "telemetry"
3096 type = "TelemetryPoint"
3097 unheard_of = "whatever"
3098 "#,
3099 )
3100 .unwrap();
3101 assert_eq!(slice.version, "9.9");
3102 assert!(slice.serves_subject("cpu/usage"));
3103 }
3104
3105 #[test]
3106 fn a_slice_without_a_version_cannot_be_diffed_and_is_rejected() {
3107 let e = parse_slice(
3108 r#"
3109 [registry]
3110 app = "zensight"
3111 convention = 1
3112 [producer]
3113 name = "sysinfo"
3114 "#,
3115 )
3116 .unwrap_err();
3117 assert!(e.to_string().contains("version"));
3118 }
3119}