Skip to main content

reddb_server/runtime/
analytics_schema_registry.rs

1//! Issue #577 — Analytics slice 2: `AnalyticsSchemaRegistry`.
2//! Issue #581 — Analytics slice 3: additive schema evolution +
3//! breaking-change rejection.
4//!
5//! Owns `(event_name, version) → schema_json` mappings persisted in
6//! `red_config`, validates payloads at insert time, and exposes the
7//! registered set for the `red.schema_registry` virtual table.
8//!
9//! Re-registering an existing `event_name` is allowed iff the change
10//! is *additive*: new optional fields only (with or without default),
11//! widening string `maxLength`. Anything else — rename, retype, drop,
12//! optional→required, brand-new required field — is rejected with a
13//! typed `SchemaError::BreakingChange { offenders }` whose `offenders`
14//! list names every offending field together with the kind of break,
15//! so the caller can pick a new `event_name` rather than smuggle the
16//! incompatible change through the same one.
17//!
18//! Persistence shape: a single JSON document stored under
19//! `red.analytics.schema_registry.entries_json` as a Text value. It
20//! contains an array of entries `{event_name, version, schema_json,
21//! registered_at_ms}`. `red_config` is append-only, so we read by
22//! scanning that collection and keeping the row with the largest
23//! engine-assigned `EntityId` (most recent write) — same trick
24//! `signed_writes_kind` uses.
25//!
26//! The schema language is a minimal JSON Schema subset:
27//! ```json
28//! { "type": "object",
29//!   "properties": { "url": { "type": "string" } },
30//!   "required": ["url"] }
31//! ```
32//! Validation rules (v1):
33//! * payload must parse to a JSON object,
34//! * every key in `required` must be present,
35//! * every key in the payload must appear in `properties` (unknown
36//!   field rejected — strict mode),
37//! * for keys present in both, the type tag must match.
38
39use 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/// One registered schema row.
48#[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    /// Schema text did not parse as JSON.
59    InvalidSchemaJson(String),
60    /// Schema parsed but did not match the expected
61    /// `{type:"object", properties:{}, required:[]}` shape.
62    InvalidSchemaShape(String),
63    /// A re-registration would break wire compatibility with the
64    /// previously registered version. `offenders` carries every
65    /// breaking change found in the diff so the caller can fix the
66    /// schema or pick a different `event_name` in one shot.
67    BreakingChange {
68        event_name: String,
69        previous_version: u32,
70        offenders: Vec<BreakingChange>,
71    },
72}
73
74/// One reason why a candidate schema is not an additive successor of
75/// the previous version. Used inside [`SchemaError::BreakingChange`].
76#[derive(Debug, Clone, PartialEq)]
77pub enum BreakingChange {
78    /// A field present in the previous version disappeared, and a new
79    /// field of the same declared type appeared in the candidate.
80    /// Treated as a rename rather than two separate changes because
81    /// the caller almost certainly meant to rename — the error
82    /// message tells them which pair we paired up.
83    Rename { from: String, to: String },
84    /// A field changed declared `type`.
85    Retype {
86        field: String,
87        from: String,
88        to: String,
89    },
90    /// A previously declared field is gone in the candidate.
91    Drop { field: String },
92    /// A field that was previously optional became required, or a new
93    /// field appeared in the candidate's `required` list (existing
94    /// rows wouldn't carry it).
95    RequiredAdd { field: String },
96}
97
98impl BreakingChange {
99    /// Short, machine-parseable description used in error bodies.
100    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    /// No schema is registered for this `event_name`. Callers may
117    /// silently treat this as "no validation" (insert path) — the
118    /// variant is here so library consumers can branch on it.
119    UnknownEventName(String),
120    InvalidPayloadJson(String),
121    /// Payload parsed but is not a JSON object.
122    PayloadNotObject,
123    /// Payload omitted a field listed in `required`.
124    MissingRequiredField {
125        event_name: String,
126        version: u32,
127        field: String,
128    },
129    /// Payload included a field that the registered schema does not
130    /// declare in `properties`. Strict mode — slice 2 has no
131    /// `additionalProperties: true` escape hatch.
132    UnknownField {
133        event_name: String,
134        version: u32,
135        field: String,
136    },
137    /// Field's JSON type does not match the property's declared type.
138    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
154/// Read the *latest* Text payload for the registry key out of
155/// `red_config`. `red_config` is append-only — `UnifiedStore::get_config`
156/// returns the first matching row, but the most-recent write wins
157/// for us. We sort by `EntityId` descending and keep the first
158/// matching row — identical to `signed_writes_kind::read_latest_config`.
159fn 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    // Store the array as one Text value, not as a flattened tree,
244    // so `set_config_tree` writes a single row whose `value` column
245    // round-trips back into the same JSON bytes.
246    let wrapped = crate::serde_json::Value::String(arr.to_string());
247    store.set_config_tree(REGISTRY_KEY, &wrapped);
248}
249
250/// Parse + minimal shape check on the schema string. Returns the
251/// canonical re-serialised form, so the registry stores a normalised
252/// representation regardless of caller whitespace / key ordering.
253fn 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
299/// Register a schema for `event_name`.
300///
301/// * First registration → returns version `1`.
302/// * Additive successor → returns `previous_version + 1`.
303/// * Anything else → `SchemaError::BreakingChange { offenders }` with
304///   every break the diff turned up so the caller can fix them all in
305///   one round-trip.
306pub 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
351/// Extract `(field, type, is_required)` triples from a parsed schema
352/// object. `type` is the declared JSON-Schema `type` string for the
353/// property, or `""` when none was declared.
354fn 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
384/// Diff a previously registered schema against a candidate and return
385/// every breaking change. Empty result == additive (or identical).
386///
387/// The diff intentionally pairs unmatched drops + adds of the same
388/// declared type as a [`BreakingChange::Rename`] — the caller is told
389/// which pair we associated so they can disambiguate if our guess is
390/// wrong.
391fn 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    // (name, type, required) for fields present in next but not prev.
398    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    // Pair drops with same-typed additions first → rename. A paired
428    // addition is *not* also reported as RequiredAdd even if the new
429    // version flagged it required: the user's intent was a rename,
430    // and surfacing both would just be noise for the same root cause.
431    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    // Unpaired added fields: required-add is breaking, optional-add
448    // is additive (the happy path).
449    for (name, _, required) in added {
450        if required {
451            breaks.push(BreakingChange::RequiredAdd { field: name });
452        }
453    }
454
455    breaks
456}
457
458/// Return `(version, schema_json)` for the latest registered schema
459/// of `event_name`, or `None` if nothing is registered. Since slice
460/// 2 only allows version 1 per event, "latest" == "the one row that
461/// exists". Once evolution lands, the resolver will keep the
462/// max-version row per event_name.
463pub 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
472/// Snapshot every registered schema. Used by the
473/// `red.schema_registry` virtual table.
474pub 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
508/// Validate `payload` (a JSON string) against the latest schema
509/// registered for `event_name`. Returns `Ok(())` if the payload
510/// matches; `Err(ValidationError)` with a typed reason otherwise.
511///
512/// `UnknownEventName` is returned when no schema is registered —
513/// the insert path treats that as "no validation, accept" for
514/// back-compat with `timeseries` rows that don't carry an
515/// `event_name` registered yet.
516pub 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    // Required-field check first so callers see the missing-field
545    // error before the unknown-field error when both could fire.
546    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    // Strict mode: every payload key must appear in properties.
556    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
585/// Map a [`ValidationError`] onto a [`RedDBError`] with a marker
586/// prefix the transport layer can pattern-match for status codes.
587/// The exact HTTP mapping is wired up alongside the broader analytics
588/// transport work; here we keep the body shape stable so callers can
589/// already parse it.
590pub 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        // Slice 3 (#581): re-registering an identical schema is the
651        // degenerate additive case — no fields changed, so it must
652        // be accepted as v2.
653        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    // --- slice 3 (#581): additive evolution + breaking-change rejection ---
660
661    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        // Retype `a` (number → string), drop `c`, and add brand-new
833        // required field `d`. Three independent breaks in one diff.
834        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        // After an additive evolution, validate() must use v2's
861        // strict-properties set — a payload using only v1 fields
862        // still passes; a payload using v2's new optional field
863        // also passes; an unknown field still rejects.
864        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        // red.schema_registry virtual table is fed by list(); slice 3
885        // contract is "every version, not just the latest".
886        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        // Slice-2 "engine restart" is simulated by handing the same
1010        // store handle to a second `latest()` call after the
1011        // original `register` returns. The real engine restart wires
1012        // through the same `UnifiedStore` API — we exercise the
1013        // serialise/deserialise path here, which is what survives
1014        // process restart on a durable backend.
1015        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        // Round-trip through a fresh load that reuses only the
1021        // public read path:
1022        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}