1use std::collections::BTreeMap;
37
38use bynk_syntax::ast::{EventDecl, TypeRef};
39use bynk_syntax::error::CompileError;
40use bynk_syntax::span::Span;
41
42use bynk_project::schema_registry::{EventEntry, FieldShape, SchemaRegistry};
43
44use crate::symbols::UnitTable;
45
46pub fn parse_or_diagnose(
52 existing: Option<&str>,
53 project_root: &std::path::Path,
54) -> Result<SchemaRegistry, CompileError> {
55 bynk_project::schema_registry::parse(existing, project_root).map_err(|msg| {
56 CompileError::new("bynk.project.schema_registry_corrupt", Span::default(), msg)
57 })
58}
59
60fn snapshot(event: &EventDecl) -> Vec<FieldShape> {
67 let mut fields: Vec<FieldShape> = event
68 .body
69 .fields
70 .iter()
71 .map(|f| FieldShape {
72 name: f.name.name.clone(),
73 ty: canon_type(&f.type_ref),
74 default: f.init.is_some(),
75 })
76 .collect();
77 fields.sort_by(|a, b| a.name.cmp(&b.name));
78 fields
79}
80
81fn canon_type(t: &TypeRef) -> String {
90 match t {
91 TypeRef::Base(b, _) => b.name().to_string(),
92 TypeRef::Named(id) => id.name.clone(),
93 TypeRef::Result(a, b, _) => format!("Result[{}, {}]", canon_type(a), canon_type(b)),
94 TypeRef::Option(t, _) => format!("Option[{}]", canon_type(t)),
95 TypeRef::Effect(t, _) => format!("Effect[{}]", canon_type(t)),
96 TypeRef::HttpResult(t, _) => format!("HttpResult[{}]", canon_type(t)),
97 TypeRef::QueueResult(_) => "QueueResult".to_string(),
98 TypeRef::List(t, _) => format!("List[{}]", canon_type(t)),
99 TypeRef::Map(k, v, _) => format!("Map[{}, {}]", canon_type(k), canon_type(v)),
100 TypeRef::Query(t, _) => format!("Query[{}]", canon_type(t)),
101 TypeRef::Stream(t, _) => format!("Stream[{}]", canon_type(t)),
102 TypeRef::Connection(t, _) => format!("Connection[{}]", canon_type(t)),
103 TypeRef::History(t, _) => format!("History[{}]", canon_type(t)),
104 TypeRef::ValidationError(_) => "ValidationError".to_string(),
105 TypeRef::JsonError(_) => "JsonError".to_string(),
106 TypeRef::Unit(_) => "()".to_string(),
107 TypeRef::Fn(params, ret, _) => format!(
108 "({}) -> {}",
109 params.iter().map(canon_type).collect::<Vec<_>>().join(", "),
110 canon_type(ret)
111 ),
112 TypeRef::App { name, args, .. } => format!(
113 "{}[{}]",
114 name.name,
115 args.iter().map(canon_type).collect::<Vec<_>>().join(", ")
116 ),
117 }
118}
119
120enum Reconciled {
123 Baseline { fields: Vec<FieldShape> },
126 Unchanged {
128 fields: Vec<FieldShape>,
129 stored: i64,
130 },
131 Additive {
133 fields: Vec<FieldShape>,
134 bumped: i64,
135 },
136 NonAdditive {
139 removed: Vec<String>,
140 retyped: Vec<String>,
141 added_without_default: Vec<String>,
142 lost_default: Vec<String>,
143 },
144}
145
146fn reconcile_one(current: &[FieldShape], stored: Option<&EventEntry>) -> Reconciled {
147 let Some(stored) = stored else {
148 return Reconciled::Baseline {
149 fields: current.to_vec(),
150 };
151 };
152 if current == stored.fields.as_slice() {
153 return Reconciled::Unchanged {
154 fields: current.to_vec(),
155 stored: stored.schema,
156 };
157 }
158
159 fn by_name(fields: &[FieldShape]) -> BTreeMap<&str, &FieldShape> {
160 fields.iter().map(|f| (f.name.as_str(), f)).collect()
161 }
162 let old = by_name(&stored.fields);
163 let new = by_name(current);
164
165 let removed: Vec<String> = old
166 .keys()
167 .filter(|n| !new.contains_key(*n))
168 .map(|n| n.to_string())
169 .collect();
170 let retyped: Vec<String> = old
171 .iter()
172 .filter_map(|(n, old_field)| {
173 new.get(n)
174 .filter(|new_field| new_field.ty != old_field.ty)
175 .map(|_| n.to_string())
176 })
177 .collect();
178 let added_without_default: Vec<String> = new
179 .iter()
180 .filter(|(n, f)| !old.contains_key(*n) && !f.default)
181 .map(|(n, _)| n.to_string())
182 .collect();
183 let lost_default: Vec<String> = old
188 .iter()
189 .filter_map(|(n, old_field)| {
190 new.get(n).filter(|new_field| {
191 new_field.ty == old_field.ty && old_field.default && !new_field.default
192 })
193 })
194 .map(|f| f.name.clone())
195 .collect();
196
197 if removed.is_empty()
198 && retyped.is_empty()
199 && added_without_default.is_empty()
200 && lost_default.is_empty()
201 {
202 Reconciled::Additive {
203 fields: current.to_vec(),
204 bumped: stored.schema + 1,
205 }
206 } else {
207 Reconciled::NonAdditive {
208 removed,
209 retyped,
210 lost_default,
211 added_without_default,
212 }
213 }
214}
215
216pub fn reconcile(
222 existing: &SchemaRegistry,
223 unit_tables: &std::collections::HashMap<String, UnitTable>,
224 errors: &mut Vec<CompileError>,
225) -> (SchemaRegistry, std::collections::HashMap<String, i64>) {
226 let mut updated = SchemaRegistry::new();
227 let mut effective = std::collections::HashMap::new();
228
229 let mut units: Vec<_> = unit_tables.iter().collect();
230 units.sort_by_key(|(name, _)| *name);
231
232 for (unit_name, table) in units {
233 let mut events: Vec<_> = table.events.iter().collect();
234 events.sort_by_key(|(name, _)| *name);
235 for (event_name, event) in events {
236 let key = format!("{unit_name}.{event_name}");
237 let fields = snapshot(event);
238 let declared = event.schema_version();
239 let annotation_span = event
240 .annotations
241 .iter()
242 .find(|a| a.name.name == "schema")
243 .map(|a| a.span);
244 let has_annotation = annotation_span.is_some();
245
246 let (effective_version, entry) = match reconcile_one(&fields, existing.get(&key)) {
247 Reconciled::Baseline { fields } => (
248 declared,
249 EventEntry {
250 schema: declared,
251 fields,
252 },
253 ),
254 Reconciled::Unchanged { fields, stored } => {
255 if has_annotation && declared != stored {
256 errors.push(mismatch_error(
257 event_name,
258 annotation_span.unwrap(),
259 declared,
260 stored,
261 ));
262 }
263 (
264 stored,
265 EventEntry {
266 schema: stored,
267 fields,
268 },
269 )
270 }
271 Reconciled::Additive { fields, bumped } => {
272 if has_annotation && declared != bumped {
273 errors.push(mismatch_error(
274 event_name,
275 annotation_span.unwrap(),
276 declared,
277 bumped,
278 ));
279 }
280 (
281 bumped,
282 EventEntry {
283 schema: bumped,
284 fields,
285 },
286 )
287 }
288 Reconciled::NonAdditive {
289 removed,
290 retyped,
291 added_without_default,
292 lost_default,
293 } => {
294 errors.push(non_additive_error(
295 event,
296 event_name,
297 &removed,
298 &retyped,
299 &added_without_default,
300 &lost_default,
301 ));
302 let old = existing
307 .get(&key)
308 .expect("a NonAdditive verdict only fires against a stored entry")
309 .clone();
310 (old.schema, old)
311 }
312 };
313
314 effective.insert(key.clone(), effective_version);
315 updated.insert(key, entry);
316 }
317 }
318
319 (updated, effective)
320}
321
322fn mismatch_error(
323 event_name: &str,
324 span: bynk_syntax::span::Span,
325 declared: i64,
326 computed: i64,
327) -> CompileError {
328 CompileError::new(
329 "bynk.event.schema_version_mismatch",
330 span,
331 format!(
332 "`{event_name}`'s `@schema({declared})` disagrees with the schema \
333 registry, which computes version {computed} from the event's \
334 build history"
335 ),
336 )
337 .with_note(format!(
338 "update the annotation to `@schema({computed})`, or remove it to let \
339 the compiler track the version automatically"
340 ))
341}
342
343fn non_additive_error(
344 event: &EventDecl,
345 event_name: &str,
346 removed: &[String],
347 retyped: &[String],
348 added_without_default: &[String],
349 lost_default: &[String],
350) -> CompileError {
351 let mut parts = Vec::new();
352 if !removed.is_empty() {
353 parts.push(format!("field(s) removed: {}", removed.join(", ")));
354 }
355 if !retyped.is_empty() {
356 parts.push(format!("field(s) retyped: {}", retyped.join(", ")));
357 }
358 if !added_without_default.is_empty() {
359 parts.push(format!(
360 "field(s) added without a default: {}",
361 added_without_default.join(", ")
362 ));
363 }
364 if !lost_default.is_empty() {
365 parts.push(format!(
366 "field(s) lost their default: {}",
367 lost_default.join(", ")
368 ));
369 }
370 CompileError::new(
371 "bynk.event.non_additive_schema_change",
372 event.span,
373 format!(
374 "`{event_name}` changed in a way the schema registry cannot \
375 evolve additively — {}",
376 parts.join("; ")
377 ),
378 )
379 .with_note(
380 "an additive change adds only fields that carry a default; give a \
381 breaking change a new event type name instead",
382 )
383}
384
385#[cfg(test)]
386mod tests {
387 use super::*;
388 use bynk_syntax::ast::{
389 Annotation, AnnotationArg, BaseType, Expr, ExprId, ExprKind, Ident, RecordBody, Trivia,
390 };
391 use bynk_syntax::span::Span;
392 use std::collections::HashMap as StdHashMap;
393
394 fn ident(name: &str) -> Ident {
395 Ident {
396 name: name.to_string(),
397 span: Span::default(),
398 }
399 }
400
401 fn int_lit(value: i64) -> Expr {
402 Expr {
403 id: ExprId::SYNTHETIC,
404 kind: ExprKind::IntLit {
405 value,
406 lexeme: value.to_string(),
407 },
408 span: Span::default(),
409 }
410 }
411
412 fn field(name: &str, ty: TypeRef, has_default: bool) -> bynk_syntax::ast::RecordField {
413 bynk_syntax::ast::RecordField {
414 trivia: Default::default(),
415 name: ident(name),
416 type_ref: ty,
417 refinement: None,
418 init: has_default.then(|| int_lit(0)),
419 span: Span::default(),
420 }
421 }
422
423 fn schema_annotation(n: i64) -> Annotation {
424 Annotation {
425 name: ident("schema"),
426 args: vec![AnnotationArg {
427 label: None,
428 value: int_lit(n),
429 span: Span::default(),
430 }],
431 span: Span::default(),
432 }
433 }
434
435 fn event(
436 name: &str,
437 annotations: Vec<Annotation>,
438 fields: Vec<bynk_syntax::ast::RecordField>,
439 ) -> EventDecl {
440 EventDecl {
441 name: ident(name),
442 annotations,
443 body: RecordBody {
444 trailing_comments: Default::default(),
445 fields,
446 span: Span::default(),
447 },
448 documentation: None,
449 span: Span::default(),
450 trivia: Trivia::default(),
451 }
452 }
453
454 fn int_ty() -> TypeRef {
455 TypeRef::Base(BaseType::Int, Span::default())
456 }
457 fn string_ty() -> TypeRef {
458 TypeRef::Base(BaseType::String, Span::default())
459 }
460
461 fn shape(name: &str, ty: &str, default: bool) -> FieldShape {
462 FieldShape {
463 name: name.to_string(),
464 ty: ty.to_string(),
465 default,
466 }
467 }
468
469 #[test]
472 fn canon_type_renders_base_and_generic_shapes() {
473 assert_eq!(canon_type(&int_ty()), "Int");
474 assert_eq!(
475 canon_type(&TypeRef::Option(Box::new(string_ty()), Span::default())),
476 "Option[String]"
477 );
478 assert_eq!(
479 canon_type(&TypeRef::List(Box::new(int_ty()), Span::default())),
480 "List[Int]"
481 );
482 }
483
484 #[test]
487 fn no_entry_baselines_silently() {
488 let current = vec![shape("orderId", "String", false)];
489 match reconcile_one(¤t, None) {
490 Reconciled::Baseline { fields } => assert_eq!(fields, current),
491 _ => panic!("expected Baseline"),
492 }
493 }
494
495 #[test]
496 fn unchanged_shape_keeps_stored_version() {
497 let current = vec![shape("orderId", "String", false)];
498 let stored = EventEntry {
499 schema: 2,
500 fields: current.clone(),
501 };
502 match reconcile_one(¤t, Some(&stored)) {
503 Reconciled::Unchanged { stored: v, .. } => assert_eq!(v, 2),
504 _ => panic!("expected Unchanged"),
505 }
506 }
507
508 #[test]
509 fn additive_field_with_default_bumps_version() {
510 let old = vec![shape("orderId", "String", false)];
511 let new = vec![
512 shape("orderId", "String", false),
513 shape("region", "Region", true),
514 ];
515 let stored = EventEntry {
516 schema: 1,
517 fields: old,
518 };
519 match reconcile_one(&new, Some(&stored)) {
520 Reconciled::Additive { bumped, .. } => assert_eq!(bumped, 2),
521 _ => panic!("expected Additive"),
522 }
523 }
524
525 #[test]
526 fn field_removed_is_non_additive() {
527 let old = vec![
528 shape("orderId", "String", false),
529 shape("region", "Region", false),
530 ];
531 let new = vec![shape("orderId", "String", false)];
532 let stored = EventEntry {
533 schema: 1,
534 fields: old,
535 };
536 match reconcile_one(&new, Some(&stored)) {
537 Reconciled::NonAdditive { removed, .. } => {
538 assert_eq!(removed, vec!["region".to_string()])
539 }
540 _ => panic!("expected NonAdditive"),
541 }
542 }
543
544 #[test]
545 fn field_retyped_is_non_additive() {
546 let old = vec![shape("orderId", "String", false)];
547 let new = vec![shape("orderId", "Int", false)];
548 let stored = EventEntry {
549 schema: 1,
550 fields: old,
551 };
552 match reconcile_one(&new, Some(&stored)) {
553 Reconciled::NonAdditive { retyped, .. } => {
554 assert_eq!(retyped, vec!["orderId".to_string()])
555 }
556 _ => panic!("expected NonAdditive"),
557 }
558 }
559
560 #[test]
561 fn field_added_without_default_is_non_additive() {
562 let old = vec![shape("orderId", "String", false)];
563 let new = vec![
564 shape("orderId", "String", false),
565 shape("region", "Region", false),
566 ];
567 let stored = EventEntry {
568 schema: 1,
569 fields: old,
570 };
571 match reconcile_one(&new, Some(&stored)) {
572 Reconciled::NonAdditive {
573 added_without_default,
574 ..
575 } => assert_eq!(added_without_default, vec!["region".to_string()]),
576 _ => panic!("expected NonAdditive"),
577 }
578 }
579
580 #[test]
581 fn a_field_losing_its_default_is_non_additive() {
582 let old = vec![
587 shape("orderId", "String", false),
588 shape("region", "String", true),
589 ];
590 let new = vec![
591 shape("orderId", "String", false),
592 shape("region", "String", false),
593 ];
594 let stored = EventEntry {
595 schema: 1,
596 fields: old,
597 };
598 match reconcile_one(&new, Some(&stored)) {
599 Reconciled::NonAdditive { lost_default, .. } => {
600 assert_eq!(lost_default, vec!["region".to_string()])
601 }
602 _ => panic!("expected NonAdditive"),
603 }
604 }
605
606 fn table_with(events: Vec<(&str, EventDecl)>) -> UnitTable {
609 UnitTable {
610 kind: None,
611 types: StdHashMap::new(),
612 fns: StdHashMap::new(),
613 methods: StdHashMap::new(),
614 capabilities: StdHashMap::new(),
615 providers: StdHashMap::new(),
616 services: StdHashMap::new(),
617 agents: StdHashMap::new(),
618 actors: StdHashMap::new(),
619 exported_capabilities: Default::default(),
620 events: events
621 .into_iter()
622 .map(|(n, e)| (n.to_string(), e))
623 .collect(),
624 flattened_caps: StdHashMap::new(),
625 }
626 }
627
628 #[test]
629 fn reconcile_baselines_a_brand_new_event_at_its_declared_annotation() {
630 let e = event(
631 "PaymentConfirmed",
632 vec![schema_annotation(3)],
633 vec![field("orderId", string_ty(), false)],
634 );
635 let mut units = StdHashMap::new();
636 units.insert(
637 "commerce.order".to_string(),
638 table_with(vec![("PaymentConfirmed", e)]),
639 );
640 let existing = SchemaRegistry::new();
641 let mut errors = Vec::new();
642 let (updated, effective) = reconcile(&existing, &units, &mut errors);
643 assert!(errors.is_empty(), "a first-ever compile must not error");
644 assert_eq!(effective.get("commerce.order.PaymentConfirmed"), Some(&3));
645 assert_eq!(
646 updated
647 .get("commerce.order.PaymentConfirmed")
648 .unwrap()
649 .schema,
650 3
651 );
652 }
653
654 #[test]
655 fn reconcile_rejects_a_mismatched_schema_annotation_on_an_unchanged_event() {
656 let e = event(
657 "PaymentConfirmed",
658 vec![schema_annotation(5)],
659 vec![field("orderId", string_ty(), false)],
660 );
661 let mut units = StdHashMap::new();
662 units.insert(
663 "commerce.order".to_string(),
664 table_with(vec![("PaymentConfirmed", e)]),
665 );
666 let mut existing = SchemaRegistry::new();
667 existing.insert(
668 "commerce.order.PaymentConfirmed".to_string(),
669 EventEntry {
670 schema: 2,
671 fields: vec![shape("orderId", "String", false)],
672 },
673 );
674 let mut errors = Vec::new();
675 let (_, effective) = reconcile(&existing, &units, &mut errors);
676 assert_eq!(errors.len(), 1);
677 assert_eq!(errors[0].category, "bynk.event.schema_version_mismatch");
678 assert_eq!(effective.get("commerce.order.PaymentConfirmed"), Some(&2));
681 }
682
683 #[test]
684 fn reconcile_auto_bumps_an_unannotated_additive_change() {
685 let e = event(
686 "OrderCancelled",
687 vec![],
688 vec![
689 field("orderId", string_ty(), false),
690 field("reason", string_ty(), true),
691 ],
692 );
693 let mut units = StdHashMap::new();
694 units.insert(
695 "commerce.order".to_string(),
696 table_with(vec![("OrderCancelled", e)]),
697 );
698 let mut existing = SchemaRegistry::new();
699 existing.insert(
700 "commerce.order.OrderCancelled".to_string(),
701 EventEntry {
702 schema: 1,
703 fields: vec![shape("orderId", "String", false)],
704 },
705 );
706 let mut errors = Vec::new();
707 let (updated, effective) = reconcile(&existing, &units, &mut errors);
708 assert!(errors.is_empty());
709 assert_eq!(effective.get("commerce.order.OrderCancelled"), Some(&2));
710 assert_eq!(
711 updated.get("commerce.order.OrderCancelled").unwrap().schema,
712 2
713 );
714 }
715
716 #[test]
717 fn reconcile_rejects_a_non_additive_change_and_keeps_the_old_entry() {
718 let e = event(
719 "OrderCancelled",
720 vec![],
721 vec![field("orderId", string_ty(), false)],
722 );
723 let mut units = StdHashMap::new();
724 units.insert(
725 "commerce.order".to_string(),
726 table_with(vec![("OrderCancelled", e)]),
727 );
728 let mut existing = SchemaRegistry::new();
729 existing.insert(
730 "commerce.order.OrderCancelled".to_string(),
731 EventEntry {
732 schema: 4,
733 fields: vec![
734 shape("orderId", "String", false),
735 shape("reason", "String", false),
736 ],
737 },
738 );
739 let mut errors = Vec::new();
740 let (updated, effective) = reconcile(&existing, &units, &mut errors);
741 assert_eq!(errors.len(), 1);
742 assert_eq!(errors[0].category, "bynk.event.non_additive_schema_change");
743 assert_eq!(effective.get("commerce.order.OrderCancelled"), Some(&4));
744 assert_eq!(
745 updated.get("commerce.order.OrderCancelled").unwrap().schema,
746 4
747 );
748 }
749
750 #[test]
751 fn a_stale_key_for_a_renamed_event_is_dropped_silently() {
752 let e = event(
757 "PaymentConfirmedV2",
758 vec![],
759 vec![field("orderId", string_ty(), false)],
760 );
761 let mut units = StdHashMap::new();
762 units.insert(
763 "commerce.order".to_string(),
764 table_with(vec![("PaymentConfirmedV2", e)]),
765 );
766 let mut existing = SchemaRegistry::new();
767 existing.insert(
768 "commerce.order.PaymentConfirmed".to_string(),
769 EventEntry {
770 schema: 3,
771 fields: vec![shape("orderId", "String", false)],
772 },
773 );
774 let mut errors = Vec::new();
775 let (updated, _) = reconcile(&existing, &units, &mut errors);
776 assert!(errors.is_empty());
777 assert!(updated.get("commerce.order.PaymentConfirmed").is_none());
778 assert!(updated.get("commerce.order.PaymentConfirmedV2").is_some());
779 }
780
781 #[test]
788 fn parse_or_diagnose_passes_through_a_valid_registry() {
789 let mut reg = SchemaRegistry::new();
790 reg.insert(
791 "commerce.order.PaymentConfirmed".to_string(),
792 EventEntry {
793 schema: 1,
794 fields: vec![shape("orderId", "String", false)],
795 },
796 );
797 let text = bynk_project::schema_registry::serialize(®);
798 let parsed = parse_or_diagnose(Some(&text), std::path::Path::new("/tmp"))
799 .expect("a freshly serialized registry must parse");
800 assert_eq!(
801 parsed
802 .get("commerce.order.PaymentConfirmed")
803 .map(|e| e.schema),
804 Some(1)
805 );
806 }
807
808 #[test]
809 fn parse_or_diagnose_reports_a_corrupt_registry_under_its_own_code() {
810 let err = parse_or_diagnose(Some("not valid toml {{{"), std::path::Path::new("/tmp"))
811 .expect_err("garbage content must not parse");
812 assert_eq!(err.category, "bynk.project.schema_registry_corrupt");
813 }
814
815 #[test]
816 fn parse_or_diagnose_with_no_existing_content_baselines_empty() {
817 let parsed = parse_or_diagnose(None, std::path::Path::new("/tmp"))
818 .expect("no lock file yet is not corruption");
819 assert!(parsed.get("anything.at_all").is_none());
820 }
821}