1use crate::storage::schema::Value;
40use crate::storage::unified::{EntityData, UnifiedStore};
41use crate::utils::json::{parse_json, JsonValue};
42
43use std::time::{SystemTime, UNIX_EPOCH};
44
45const REGISTRY_KEY: &str = "red.analytics.schema_registry.entries_json";
46
47#[derive(Debug, Clone, PartialEq)]
49pub struct SchemaEntry {
50 pub event_name: String,
51 pub version: u32,
52 pub schema_json: String,
53 pub registered_at_ms: u128,
54}
55
56#[derive(Debug, Clone, PartialEq)]
57pub enum SchemaError {
58 InvalidSchemaJson(String),
60 InvalidSchemaShape(String),
63 BreakingChange {
68 event_name: String,
69 previous_version: u32,
70 offenders: Vec<BreakingChange>,
71 },
72}
73
74#[derive(Debug, Clone, PartialEq)]
77pub enum BreakingChange {
78 Rename { from: String, to: String },
84 Retype {
86 field: String,
87 from: String,
88 to: String,
89 },
90 Drop { field: String },
92 RequiredAdd { field: String },
96}
97
98impl BreakingChange {
99 pub fn describe(&self) -> String {
101 match self {
102 BreakingChange::Rename { from, to } => format!("renamed field '{from}' to '{to}'"),
103 BreakingChange::Retype { field, from, to } => {
104 format!("retyped field '{field}' from {from} to {to}")
105 }
106 BreakingChange::Drop { field } => format!("dropped field '{field}'"),
107 BreakingChange::RequiredAdd { field } => {
108 format!("required-add for field '{field}'")
109 }
110 }
111 }
112}
113
114#[derive(Debug, Clone, PartialEq)]
115pub enum ValidationError {
116 UnknownEventName(String),
120 InvalidPayloadJson(String),
121 PayloadNotObject,
123 MissingRequiredField {
125 event_name: String,
126 version: u32,
127 field: String,
128 },
129 UnknownField {
133 event_name: String,
134 version: u32,
135 field: String,
136 },
137 TypeMismatch {
139 event_name: String,
140 version: u32,
141 field: String,
142 expected: String,
143 got: String,
144 },
145}
146
147fn now_ms() -> u128 {
148 SystemTime::now()
149 .duration_since(UNIX_EPOCH)
150 .map(|d| d.as_millis())
151 .unwrap_or(0)
152}
153
154fn read_latest_registry_json(store: &UnifiedStore) -> Option<String> {
160 let manager = store.get_collection("red_config")?;
161 let mut all = manager.query_all(|_| true);
162 all.sort_by_key(|b| std::cmp::Reverse(b.id.raw()));
163 for entity in all {
164 let EntityData::Row(row) = &entity.data else {
165 continue;
166 };
167 let Some(named) = &row.named else { continue };
168 let matches = matches!(
169 named.get("key"),
170 Some(Value::Text(s)) if s.as_ref() == REGISTRY_KEY
171 );
172 if matches {
173 if let Some(Value::Text(s)) = named.get("value") {
174 return Some(s.to_string());
175 }
176 }
177 }
178 None
179}
180
181fn load(store: &UnifiedStore) -> Vec<SchemaEntry> {
182 let raw = match read_latest_registry_json(store) {
183 Some(s) => s,
184 None => return Vec::new(),
185 };
186 let Ok(parsed) = parse_json(&raw) else {
187 return Vec::new();
188 };
189 let Some(arr) = parsed.as_array() else {
190 return Vec::new();
191 };
192 let mut out = Vec::with_capacity(arr.len());
193 for item in arr {
194 let Some(obj) = item.as_object() else {
195 continue;
196 };
197 let lookup = |k: &str| obj.iter().find(|(key, _)| key == k).map(|(_, v)| v);
198 let Some(event_name) = lookup("event_name").and_then(JsonValue::as_str) else {
199 continue;
200 };
201 let Some(version) = lookup("version").and_then(JsonValue::as_f64) else {
202 continue;
203 };
204 let Some(schema_json) = lookup("schema_json").and_then(JsonValue::as_str) else {
205 continue;
206 };
207 let Some(registered_at_ms) = lookup("registered_at_ms").and_then(JsonValue::as_f64) else {
208 continue;
209 };
210 out.push(SchemaEntry {
211 event_name: event_name.to_string(),
212 version: version as u32,
213 schema_json: schema_json.to_string(),
214 registered_at_ms: registered_at_ms as u128,
215 });
216 }
217 out
218}
219
220fn entry_to_json(e: &SchemaEntry) -> crate::serde_json::Value {
221 let mut obj = crate::serde_json::Map::new();
222 obj.insert(
223 "event_name".to_string(),
224 crate::serde_json::Value::String(e.event_name.clone()),
225 );
226 obj.insert(
227 "version".to_string(),
228 crate::serde_json::Value::Number(e.version as f64),
229 );
230 obj.insert(
231 "schema_json".to_string(),
232 crate::serde_json::Value::String(e.schema_json.clone()),
233 );
234 obj.insert(
235 "registered_at_ms".to_string(),
236 crate::serde_json::Value::Number(e.registered_at_ms as f64),
237 );
238 crate::serde_json::Value::Object(obj)
239}
240
241fn save(store: &UnifiedStore, entries: &[SchemaEntry]) {
242 let arr = crate::serde_json::Value::Array(entries.iter().map(entry_to_json).collect());
243 let wrapped = crate::serde_json::Value::String(arr.to_string());
247 store.set_config_tree(REGISTRY_KEY, &wrapped);
248}
249
250fn validate_schema_shape(schema_json: &str) -> Result<JsonValue, SchemaError> {
254 let parsed =
255 parse_json(schema_json).map_err(|err| SchemaError::InvalidSchemaJson(err.to_string()))?;
256 let Some(obj) = parsed.as_object() else {
257 return Err(SchemaError::InvalidSchemaShape(
258 "schema must be a JSON object".to_string(),
259 ));
260 };
261 let lookup = |k: &str| obj.iter().find(|(key, _)| key == k).map(|(_, v)| v);
262 match lookup("type").and_then(JsonValue::as_str) {
263 Some("object") => {}
264 Some(other) => {
265 return Err(SchemaError::InvalidSchemaShape(format!(
266 "schema `type` must be \"object\", got \"{other}\""
267 )));
268 }
269 None => {
270 return Err(SchemaError::InvalidSchemaShape(
271 "schema must declare `type`".to_string(),
272 ));
273 }
274 }
275 if let Some(props) = lookup("properties") {
276 if props.as_object().is_none() {
277 return Err(SchemaError::InvalidSchemaShape(
278 "schema `properties` must be an object".to_string(),
279 ));
280 }
281 }
282 if let Some(req) = lookup("required") {
283 let Some(arr) = req.as_array() else {
284 return Err(SchemaError::InvalidSchemaShape(
285 "schema `required` must be an array of strings".to_string(),
286 ));
287 };
288 for item in arr {
289 if item.as_str().is_none() {
290 return Err(SchemaError::InvalidSchemaShape(
291 "schema `required` must be an array of strings".to_string(),
292 ));
293 }
294 }
295 }
296 Ok(parsed)
297}
298
299pub fn register(
307 store: &UnifiedStore,
308 event_name: &str,
309 schema_json: &str,
310) -> Result<u32, SchemaError> {
311 let candidate = validate_schema_shape(schema_json)?;
312 let mut entries = load(store);
313
314 let previous = entries
315 .iter()
316 .filter(|e| e.event_name == event_name)
317 .max_by_key(|e| e.version)
318 .cloned();
319
320 let next_version = match previous {
321 None => 1,
322 Some(prev) => {
323 let prev_schema = parse_json(&prev.schema_json).map_err(|e| {
324 SchemaError::InvalidSchemaShape(format!(
325 "previously registered schema for {event_name} v{} is corrupt: {e}",
326 prev.version
327 ))
328 })?;
329 let offenders = diff_for_breaking_changes(&prev_schema, &candidate);
330 if !offenders.is_empty() {
331 return Err(SchemaError::BreakingChange {
332 event_name: event_name.to_string(),
333 previous_version: prev.version,
334 offenders,
335 });
336 }
337 prev.version + 1
338 }
339 };
340
341 entries.push(SchemaEntry {
342 event_name: event_name.to_string(),
343 version: next_version,
344 schema_json: schema_json.to_string(),
345 registered_at_ms: now_ms(),
346 });
347 save(store, &entries);
348 Ok(next_version)
349}
350
351fn schema_fields(schema: &JsonValue) -> Vec<(String, String, bool)> {
355 let Some(obj) = schema.as_object() else {
356 return Vec::new();
357 };
358 let properties: &[(String, JsonValue)] = obj
359 .iter()
360 .find(|(k, _)| k == "properties")
361 .and_then(|(_, v)| v.as_object())
362 .unwrap_or(&[]);
363 let required: Vec<&str> = obj
364 .iter()
365 .find(|(k, _)| k == "required")
366 .and_then(|(_, v)| v.as_array())
367 .map(|arr| arr.iter().filter_map(JsonValue::as_str).collect())
368 .unwrap_or_default();
369 properties
370 .iter()
371 .map(|(name, prop)| {
372 let ty = prop
373 .as_object()
374 .and_then(|entries| entries.iter().find(|(k, _)| k == "type"))
375 .and_then(|(_, v)| v.as_str())
376 .unwrap_or("")
377 .to_string();
378 let req = required.contains(&name.as_str());
379 (name.clone(), ty, req)
380 })
381 .collect()
382}
383
384fn diff_for_breaking_changes(prev: &JsonValue, next: &JsonValue) -> Vec<BreakingChange> {
392 let prev_fields = schema_fields(prev);
393 let next_fields = schema_fields(next);
394
395 let mut breaks = Vec::new();
396 let mut dropped: Vec<(String, String)> = Vec::new();
397 let mut added: Vec<(String, String, bool)> = Vec::new();
399
400 for (name, prev_type, prev_required) in &prev_fields {
401 match next_fields.iter().find(|(n, _, _)| n == name) {
402 Some((_, next_type, next_required)) => {
403 if prev_type != next_type && !prev_type.is_empty() && !next_type.is_empty() {
404 breaks.push(BreakingChange::Retype {
405 field: name.clone(),
406 from: prev_type.clone(),
407 to: next_type.clone(),
408 });
409 }
410 if !prev_required && *next_required {
411 breaks.push(BreakingChange::RequiredAdd {
412 field: name.clone(),
413 });
414 }
415 }
416 None => dropped.push((name.clone(), prev_type.clone())),
417 }
418 }
419
420 for (name, next_type, next_required) in &next_fields {
421 if prev_fields.iter().any(|(n, _, _)| n == name) {
422 continue;
423 }
424 added.push((name.clone(), next_type.clone(), *next_required));
425 }
426
427 for (drop_name, drop_type) in dropped {
432 let paired = added
433 .iter()
434 .position(|(_, ty, _)| ty == &drop_type && !drop_type.is_empty());
435 match paired {
436 Some(idx) => {
437 let (add_name, _, _) = added.remove(idx);
438 breaks.push(BreakingChange::Rename {
439 from: drop_name,
440 to: add_name,
441 });
442 }
443 None => breaks.push(BreakingChange::Drop { field: drop_name }),
444 }
445 }
446
447 for (name, _, required) in added {
450 if required {
451 breaks.push(BreakingChange::RequiredAdd { field: name });
452 }
453 }
454
455 breaks
456}
457
458pub fn latest(store: &UnifiedStore, event_name: &str) -> Option<(u32, String)> {
464 let entries = load(store);
465 entries
466 .into_iter()
467 .filter(|e| e.event_name == event_name)
468 .max_by_key(|e| e.version)
469 .map(|e| (e.version, e.schema_json))
470}
471
472pub fn list(store: &UnifiedStore) -> Vec<SchemaEntry> {
475 load(store)
476}
477
478fn json_type_name(v: &JsonValue) -> &'static str {
479 match v {
480 JsonValue::Null => "null",
481 JsonValue::Bool(_) => "boolean",
482 JsonValue::Integer(_) => "number",
483 JsonValue::Number(_) => "number",
484 JsonValue::Decimal(_) => "number",
485 JsonValue::String(_) => "string",
486 JsonValue::Array(_) => "array",
487 JsonValue::Object(_) => "object",
488 }
489}
490
491fn type_matches(expected: &str, got: &JsonValue) -> bool {
492 match expected {
493 "string" => matches!(got, JsonValue::String(_)),
494 "boolean" => matches!(got, JsonValue::Bool(_)),
495 "array" => matches!(got, JsonValue::Array(_)),
496 "object" => matches!(got, JsonValue::Object(_)),
497 "null" => matches!(got, JsonValue::Null),
498 "number" => matches!(got, JsonValue::Integer(_) | JsonValue::Number(_)),
499 "integer" => match got {
500 JsonValue::Integer(_) => true,
501 JsonValue::Number(n) => *n == n.trunc(),
502 _ => false,
503 },
504 _ => false,
505 }
506}
507
508pub fn validate(
517 store: &UnifiedStore,
518 event_name: &str,
519 payload_json: &str,
520) -> Result<(), ValidationError> {
521 let Some((version, schema_json)) = latest(store, event_name) else {
522 return Err(ValidationError::UnknownEventName(event_name.to_string()));
523 };
524 let schema = parse_json(&schema_json)
525 .map_err(|e| ValidationError::InvalidPayloadJson(format!("schema corrupt: {e}")))?;
526 let payload =
527 parse_json(payload_json).map_err(|e| ValidationError::InvalidPayloadJson(e.to_string()))?;
528 let Some(payload_obj) = payload.as_object() else {
529 return Err(ValidationError::PayloadNotObject);
530 };
531 let schema_obj = schema.as_object().unwrap_or(&[]);
532 let properties: &[(String, JsonValue)] = schema_obj
533 .iter()
534 .find(|(k, _)| k == "properties")
535 .and_then(|(_, v)| v.as_object())
536 .unwrap_or(&[]);
537 let required: Vec<&str> = schema_obj
538 .iter()
539 .find(|(k, _)| k == "required")
540 .and_then(|(_, v)| v.as_array())
541 .map(|arr| arr.iter().filter_map(JsonValue::as_str).collect())
542 .unwrap_or_default();
543
544 for req in &required {
547 if !payload_obj.iter().any(|(k, _)| k == *req) {
548 return Err(ValidationError::MissingRequiredField {
549 event_name: event_name.to_string(),
550 version,
551 field: (*req).to_string(),
552 });
553 }
554 }
555 for (key, value) in payload_obj {
557 let Some((_, prop)) = properties.iter().find(|(k, _)| k == key) else {
558 return Err(ValidationError::UnknownField {
559 event_name: event_name.to_string(),
560 version,
561 field: key.clone(),
562 });
563 };
564 let expected_type = prop
565 .as_object()
566 .and_then(|entries| entries.iter().find(|(k, _)| k == "type"))
567 .and_then(|(_, v)| v.as_str())
568 .unwrap_or("");
569 if expected_type.is_empty() {
570 continue;
571 }
572 if !type_matches(expected_type, value) {
573 return Err(ValidationError::TypeMismatch {
574 event_name: event_name.to_string(),
575 version,
576 field: key.clone(),
577 expected: expected_type.to_string(),
578 got: json_type_name(value).to_string(),
579 });
580 }
581 }
582 Ok(())
583}
584
585pub fn validation_error_to_reddb(err: ValidationError) -> crate::api::RedDBError {
591 let body = match &err {
592 ValidationError::UnknownEventName(name) => {
593 format!("AnalyticsSchemaError:UnknownEventName:{name}")
594 }
595 ValidationError::InvalidPayloadJson(reason) => {
596 format!("AnalyticsSchemaError:InvalidPayloadJson:{reason}")
597 }
598 ValidationError::PayloadNotObject => "AnalyticsSchemaError:PayloadNotObject".to_string(),
599 ValidationError::MissingRequiredField {
600 event_name,
601 version,
602 field,
603 } => format!("AnalyticsSchemaError:MissingRequiredField:{event_name}:v{version}:{field}"),
604 ValidationError::UnknownField {
605 event_name,
606 version,
607 field,
608 } => format!("AnalyticsSchemaError:UnknownField:{event_name}:v{version}:{field}"),
609 ValidationError::TypeMismatch {
610 event_name,
611 version,
612 field,
613 expected,
614 got,
615 } => format!(
616 "AnalyticsSchemaError:TypeMismatch:{event_name}:v{version}:{field}:{expected}:{got}"
617 ),
618 };
619 crate::api::RedDBError::InvalidOperation(body)
620}
621
622#[cfg(test)]
623mod tests {
624 use super::*;
625
626 fn store() -> UnifiedStore {
627 UnifiedStore::new()
628 }
629
630 const PAGE_VIEW_SCHEMA: &str = r#"{
631 "type": "object",
632 "properties": {
633 "url": {"type": "string"},
634 "user_id": {"type": "integer"}
635 },
636 "required": ["url"]
637 }"#;
638
639 #[test]
640 fn first_registration_is_version_1() {
641 let s = store();
642 let v = register(&s, "page_view", PAGE_VIEW_SCHEMA).expect("register ok");
643 assert_eq!(v, 1);
644 let (latest_v, _) = latest(&s, "page_view").expect("latest present");
645 assert_eq!(latest_v, 1);
646 }
647
648 #[test]
649 fn re_registering_identical_schema_bumps_to_next_version() {
650 let s = store();
654 register(&s, "page_view", PAGE_VIEW_SCHEMA).unwrap();
655 let v = register(&s, "page_view", PAGE_VIEW_SCHEMA).expect("identical is additive");
656 assert_eq!(v, 2);
657 }
658
659 const PURCHASE_V1: &str =
662 r#"{"type":"object","properties":{"amount":{"type":"number"}},"required":["amount"]}"#;
663
664 #[test]
665 fn additive_optional_field_is_accepted_as_v2() {
666 let s = store();
667 register(&s, "purchase", PURCHASE_V1).unwrap();
668 let v2 = register(
669 &s,
670 "purchase",
671 r#"{"type":"object",
672 "properties":{"amount":{"type":"number"},
673 "discount_code":{"type":"string"}},
674 "required":["amount"]}"#,
675 )
676 .expect("optional add is additive");
677 assert_eq!(v2, 2);
678 let (latest_v, _) = latest(&s, "purchase").unwrap();
679 assert_eq!(latest_v, 2);
680 }
681
682 #[test]
683 fn additive_optional_field_with_default_is_accepted() {
684 let s = store();
685 register(&s, "purchase", PURCHASE_V1).unwrap();
686 let v2 = register(
687 &s,
688 "purchase",
689 r#"{"type":"object",
690 "properties":{"amount":{"type":"number"},
691 "currency":{"type":"string","default":"USD"}},
692 "required":["amount"]}"#,
693 )
694 .expect("optional add with default is additive");
695 assert_eq!(v2, 2);
696 }
697
698 #[test]
699 fn widening_string_max_length_is_accepted() {
700 let s = store();
701 register(
702 &s,
703 "ev",
704 r#"{"type":"object","properties":{"name":{"type":"string","maxLength":32}},"required":["name"]}"#,
705 )
706 .unwrap();
707 let v2 = register(
708 &s,
709 "ev",
710 r#"{"type":"object","properties":{"name":{"type":"string","maxLength":128}},"required":["name"]}"#,
711 )
712 .expect("widening maxLength is additive");
713 assert_eq!(v2, 2);
714 }
715
716 #[test]
717 fn breaking_rename_is_rejected() {
718 let s = store();
719 register(&s, "purchase", PURCHASE_V1).unwrap();
720 let err = register(
721 &s,
722 "purchase",
723 r#"{"type":"object","properties":{"total":{"type":"number"}},"required":["total"]}"#,
724 )
725 .unwrap_err();
726 match err {
727 SchemaError::BreakingChange {
728 event_name,
729 previous_version,
730 offenders,
731 } => {
732 assert_eq!(event_name, "purchase");
733 assert_eq!(previous_version, 1);
734 assert!(
735 offenders.iter().any(|b| matches!(
736 b,
737 BreakingChange::Rename { from, to }
738 if from == "amount" && to == "total"
739 )),
740 "expected Rename(amount->total), got {offenders:?}"
741 );
742 }
743 other => panic!("expected BreakingChange, got {other:?}"),
744 }
745 }
746
747 #[test]
748 fn breaking_retype_is_rejected() {
749 let s = store();
750 register(&s, "purchase", PURCHASE_V1).unwrap();
751 let err = register(
752 &s,
753 "purchase",
754 r#"{"type":"object","properties":{"amount":{"type":"string"}},"required":["amount"]}"#,
755 )
756 .unwrap_err();
757 let SchemaError::BreakingChange { offenders, .. } = err else {
758 panic!("expected BreakingChange");
759 };
760 assert!(offenders.iter().any(|b| matches!(
761 b,
762 BreakingChange::Retype { field, from, to }
763 if field == "amount" && from == "number" && to == "string"
764 )));
765 }
766
767 #[test]
768 fn breaking_drop_is_rejected() {
769 let s = store();
770 register(
771 &s,
772 "ev",
773 r#"{"type":"object",
774 "properties":{"a":{"type":"number"},"b":{"type":"boolean"}},
775 "required":["a"]}"#,
776 )
777 .unwrap();
778 let err = register(
779 &s,
780 "ev",
781 r#"{"type":"object","properties":{"a":{"type":"number"}},"required":["a"]}"#,
782 )
783 .unwrap_err();
784 let SchemaError::BreakingChange { offenders, .. } = err else {
785 panic!("expected BreakingChange");
786 };
787 assert!(offenders
788 .iter()
789 .any(|b| matches!(b, BreakingChange::Drop { field } if field == "b")));
790 }
791
792 #[test]
793 fn breaking_optional_to_required_is_rejected() {
794 let s = store();
795 register(
796 &s,
797 "ev",
798 r#"{"type":"object",
799 "properties":{"a":{"type":"number"},"b":{"type":"string"}},
800 "required":["a"]}"#,
801 )
802 .unwrap();
803 let err = register(
804 &s,
805 "ev",
806 r#"{"type":"object",
807 "properties":{"a":{"type":"number"},"b":{"type":"string"}},
808 "required":["a","b"]}"#,
809 )
810 .unwrap_err();
811 let SchemaError::BreakingChange { offenders, .. } = err else {
812 panic!("expected BreakingChange");
813 };
814 assert!(offenders
815 .iter()
816 .any(|b| matches!(b, BreakingChange::RequiredAdd { field } if field == "b")));
817 }
818
819 #[test]
820 fn multi_field_break_reports_every_offender() {
821 let s = store();
822 register(
823 &s,
824 "ev",
825 r#"{"type":"object",
826 "properties":{"a":{"type":"number"},
827 "b":{"type":"string"},
828 "c":{"type":"boolean"}},
829 "required":["a"]}"#,
830 )
831 .unwrap();
832 let err = register(
835 &s,
836 "ev",
837 r#"{"type":"object",
838 "properties":{"a":{"type":"string"},
839 "b":{"type":"string"},
840 "d":{"type":"integer"}},
841 "required":["a","d"]}"#,
842 )
843 .unwrap_err();
844 let SchemaError::BreakingChange { offenders, .. } = err else {
845 panic!("expected BreakingChange");
846 };
847 assert!(offenders
848 .iter()
849 .any(|b| matches!(b, BreakingChange::Retype { field, .. } if field == "a")));
850 assert!(offenders
851 .iter()
852 .any(|b| matches!(b, BreakingChange::Drop { field } if field == "c")));
853 assert!(offenders
854 .iter()
855 .any(|b| matches!(b, BreakingChange::RequiredAdd { field } if field == "d")));
856 }
857
858 #[test]
859 fn validate_resolves_to_latest_version_after_evolution() {
860 let s = store();
865 register(&s, "purchase", PURCHASE_V1).unwrap();
866 register(
867 &s,
868 "purchase",
869 r#"{"type":"object",
870 "properties":{"amount":{"type":"number"},
871 "discount_code":{"type":"string"}},
872 "required":["amount"]}"#,
873 )
874 .unwrap();
875 validate(&s, "purchase", r#"{"amount":1.0}"#).expect("v1-shape still valid");
876 validate(&s, "purchase", r#"{"amount":1.0,"discount_code":"X"}"#)
877 .expect("v2-only field accepted");
878 let err = validate(&s, "purchase", r#"{"amount":1.0,"mystery":1}"#).unwrap_err();
879 assert!(matches!(err, ValidationError::UnknownField { version, .. } if version == 2));
880 }
881
882 #[test]
883 fn list_returns_every_version_not_just_latest() {
884 let s = store();
887 register(&s, "purchase", PURCHASE_V1).unwrap();
888 register(
889 &s,
890 "purchase",
891 r#"{"type":"object",
892 "properties":{"amount":{"type":"number"},
893 "discount_code":{"type":"string"}},
894 "required":["amount"]}"#,
895 )
896 .unwrap();
897 let purchase_versions: Vec<u32> = list(&s)
898 .into_iter()
899 .filter(|e| e.event_name == "purchase")
900 .map(|e| e.version)
901 .collect();
902 let mut sorted = purchase_versions.clone();
903 sorted.sort();
904 assert_eq!(
905 sorted,
906 vec![1, 2],
907 "expected both versions, got {purchase_versions:?}"
908 );
909 }
910
911 #[test]
912 fn invalid_schema_json_rejected_at_register() {
913 let s = store();
914 let err = register(&s, "x", "{not json").unwrap_err();
915 assert!(matches!(err, SchemaError::InvalidSchemaJson(_)));
916 }
917
918 #[test]
919 fn schema_must_be_type_object() {
920 let s = store();
921 let err = register(&s, "x", r#"{"type":"string"}"#).unwrap_err();
922 assert!(matches!(err, SchemaError::InvalidSchemaShape(_)));
923 }
924
925 #[test]
926 fn validate_happy_path_accepts_known_fields() {
927 let s = store();
928 register(&s, "page_view", PAGE_VIEW_SCHEMA).unwrap();
929 validate(&s, "page_view", r#"{"url":"/x","user_id":42}"#).expect("ok");
930 validate(&s, "page_view", r#"{"url":"/y"}"#).expect("ok without optional");
931 }
932
933 #[test]
934 fn validate_rejects_unknown_field() {
935 let s = store();
936 register(&s, "page_view", PAGE_VIEW_SCHEMA).unwrap();
937 let err = validate(&s, "page_view", r#"{"url":"/x","mystery":1}"#).unwrap_err();
938 match err {
939 ValidationError::UnknownField { field, .. } => assert_eq!(field, "mystery"),
940 other => panic!("expected UnknownField, got {other:?}"),
941 }
942 }
943
944 #[test]
945 fn validate_rejects_missing_required_field() {
946 let s = store();
947 register(&s, "page_view", PAGE_VIEW_SCHEMA).unwrap();
948 let err = validate(&s, "page_view", r#"{}"#).unwrap_err();
949 match err {
950 ValidationError::MissingRequiredField { field, .. } => assert_eq!(field, "url"),
951 other => panic!("expected MissingRequiredField, got {other:?}"),
952 }
953 }
954
955 #[test]
956 fn validate_rejects_type_mismatch() {
957 let s = store();
958 register(&s, "page_view", PAGE_VIEW_SCHEMA).unwrap();
959 let err = validate(&s, "page_view", r#"{"url":123}"#).unwrap_err();
960 match err {
961 ValidationError::TypeMismatch {
962 field,
963 expected,
964 got,
965 ..
966 } => {
967 assert_eq!(field, "url");
968 assert_eq!(expected, "string");
969 assert_eq!(got, "number");
970 }
971 other => panic!("expected TypeMismatch, got {other:?}"),
972 }
973 }
974
975 #[test]
976 fn validate_unknown_event_name() {
977 let s = store();
978 let err = validate(&s, "nope", r#"{}"#).unwrap_err();
979 assert!(matches!(err, ValidationError::UnknownEventName(name) if name == "nope"));
980 }
981
982 #[test]
983 fn validate_payload_must_be_object() {
984 let s = store();
985 register(&s, "page_view", PAGE_VIEW_SCHEMA).unwrap();
986 let err = validate(&s, "page_view", r#""hello""#).unwrap_err();
987 assert!(matches!(err, ValidationError::PayloadNotObject));
988 }
989
990 #[test]
991 fn list_returns_every_registered_event() {
992 let s = store();
993 register(&s, "page_view", PAGE_VIEW_SCHEMA).unwrap();
994 register(
995 &s,
996 "signup",
997 r#"{"type":"object","properties":{"email":{"type":"string"}},"required":["email"]}"#,
998 )
999 .unwrap();
1000 let mut names: Vec<String> = list(&s).into_iter().map(|e| e.event_name).collect();
1001 names.sort();
1002 assert_eq!(names, vec!["page_view".to_string(), "signup".to_string()]);
1003 assert!(list(&s).iter().all(|e| e.version == 1));
1004 assert!(list(&s).iter().all(|e| e.registered_at_ms > 0));
1005 }
1006
1007 #[test]
1008 fn persistence_smoke_latest_survives_restart() {
1009 let s = store();
1016 register(&s, "page_view", PAGE_VIEW_SCHEMA).unwrap();
1017 let raw =
1018 read_latest_registry_json(&s).expect("registry json must be persisted on register");
1019 assert!(raw.contains("page_view"));
1020 let (v, schema) = latest(&s, "page_view").expect("latest after persist");
1023 assert_eq!(v, 1);
1024 assert!(schema.contains("\"url\""));
1025 }
1026
1027 #[test]
1028 fn validation_error_maps_to_invalid_operation_with_typed_marker() {
1029 let err = validation_error_to_reddb(ValidationError::MissingRequiredField {
1030 event_name: "page_view".to_string(),
1031 version: 1,
1032 field: "url".to_string(),
1033 });
1034 match err {
1035 crate::api::RedDBError::InvalidOperation(body) => {
1036 assert!(
1037 body.starts_with("AnalyticsSchemaError:MissingRequiredField:page_view:v1:url"),
1038 "unexpected body: {body}"
1039 );
1040 }
1041 other => panic!("expected InvalidOperation, got {other:?}"),
1042 }
1043 }
1044}