1use crate::error::{PyQLError, PyQLFragmentError};
21use crate::schema::{
22 DeleteAction, DeleteSide, FunctionDescriptor, OnDeletePolicy, SchemaDescriptor, TypeConstraint, TypeDescriptor,
23 resolved_pg_type,
24};
25use std::collections::{BTreeSet, HashMap, HashSet};
26
27pub mod python_snippet;
28
29pub fn export_schema(schema: &SchemaDescriptor) -> Result<String, PyQLError> {
47 let mut out = String::new();
48
49 let type_map: HashMap<String, (&str, &str)> = schema
51 .types
52 .iter()
53 .map(|t| {
54 (
55 format!("{}::{}", t.module, t.name),
56 (t.module.as_str(), t.table.as_str()),
57 )
58 })
59 .collect();
60
61 emit_schemas(schema, &mut out);
62 emit_enums(schema, &mut out);
63 emit_scalars(schema, &mut out);
64 emit_scalar_functions(schema, &mut out)?;
65 emit_tables(schema, &mut out);
66 emit_fk_constraints(schema, &type_map, &mut out);
67 emit_link_source_triggers(schema, &type_map, &mut out);
68 emit_junction_tables(schema, &mut out);
69 emit_junction_fk_constraints(schema, &type_map, &mut out);
70 emit_multilink_deletion_triggers(schema, &type_map, &mut out);
71 emit_interface_link_triggers(schema, &mut out);
72 emit_signal_triggers(schema, &mut out);
73 emit_unique_indexes(schema, &mut out);
74 emit_check_constraints(schema, &mut out)?;
75 emit_plain_indexes(schema, &mut out);
76 emit_triggers(schema, &mut out)?;
77 emit_interface_views(schema, &mut out);
78 emit_interface_junction_views(schema, &mut out);
79 emit_interface_exclusive_triggers(schema, &mut out);
80 emit_object_functions(schema, &mut out)?;
81 emit_vector_columns(schema, &mut out);
82 emit_vector_indexes(schema, &mut out);
83 emit_search_columns(schema, &mut out);
84 emit_search_indexes(schema, &mut out);
85
86 Ok(out)
87}
88
89fn qi(s: &str) -> String {
93 format!("\"{}\"", s.replace('"', "\"\""))
94}
95
96fn pg_schema(module: &str) -> String {
97 if module == "default" {
98 "\"public\"".into()
99 } else {
100 qi(module)
101 }
102}
103
104fn qn(module: &str, name: &str) -> String {
105 format!("{}.{}", pg_schema(module), qi(name))
106}
107
108pub(crate) fn fnv(parts: &[&str]) -> String {
111 let mut h: u64 = 0xcbf29ce484222325;
112 for p in parts {
113 for b in p.bytes() {
114 h ^= b as u64;
115 h = h.wrapping_mul(0x100000001b3);
116 }
117 h ^= b'|' as u64;
118 h = h.wrapping_mul(0x100000001b3);
119 }
120 format!("{:016x}", h)
121}
122
123fn trigger_names(module: &str, table: &str, pointer: &str, suffix: &str, body: &str) -> (String, String) {
138 let hash = fnv(&[table, pointer, suffix, body]);
139 let fname = format!("{}_{}_{}", table, pointer, &hash[..8]);
140 let fn_qname = qn(module, &fname);
141 (fname, fn_qname)
142}
143
144fn trigger_events(on: u8) -> String {
147 let mut events = Vec::new();
148 if on & 1 != 0 {
149 events.push("INSERT");
150 }
151 if on & 2 != 0 {
152 events.push("UPDATE");
153 }
154 if on & 4 != 0 {
155 events.push("DELETE");
156 }
157 events.join(" OR ")
158}
159
160fn trigger_timing(timing: &str) -> &str {
161 match timing {
162 "Before" => "BEFORE",
163 "After" => "AFTER",
164 "InsteadOf" => "INSTEAD OF",
165 t => t,
166 }
167}
168
169fn emit_schemas(schema: &SchemaDescriptor, out: &mut String) {
172 let mut modules: BTreeSet<&str> = BTreeSet::new();
179 for t in &schema.types {
180 modules.insert(&t.module);
181 }
182 for s in &schema.scalars {
183 modules.insert(&s.module);
184 }
185 for e in &schema.enums {
186 modules.insert(&e.module);
187 }
188 for f in &schema.functions {
189 modules.insert(&f.module);
190 }
191 for g in &schema.globals {
192 modules.insert(&g.module);
193 }
194 for a in &schema.aliases {
195 modules.insert(&a.module);
196 }
197 let non_default: Vec<&str> = modules.into_iter().filter(|m| *m != "default").collect();
198 for module in &non_default {
199 out.push_str(&format!("CREATE SCHEMA IF NOT EXISTS {};\n", pg_schema(module)));
200 }
201 if !non_default.is_empty() {
202 out.push('\n');
203 }
204}
205
206fn emit_enums(schema: &SchemaDescriptor, out: &mut String) {
209 for e in &schema.enums {
210 let members: Vec<String> = e
211 .members
212 .iter()
213 .map(|m| format!("'{}'", m.replace('\'', "''")))
214 .collect();
215 out.push_str(&format!(
216 "DO $$ BEGIN CREATE TYPE {}.{} AS ENUM ({}); EXCEPTION WHEN duplicate_object THEN NULL; END $$;\n",
217 pg_schema(&e.module),
218 qi(&e.name),
219 members.join(", "),
220 ));
221 }
222 if !schema.enums.is_empty() {
223 out.push('\n');
224 }
225}
226
227fn emit_scalars(schema: &SchemaDescriptor, out: &mut String) {
230 for s in &schema.scalars {
231 if s.is_sequence {
232 out.push_str(&format!(
233 "CREATE SEQUENCE {}.{};\n",
234 pg_schema(&s.module),
235 qi(&format!("{}_seq", s.name)),
236 ));
237 }
238 let check_clause = scalar_check_clauses(schema, &s.module, &s.name);
239 out.push_str(&format!(
240 "CREATE DOMAIN {}.{} AS {}{};\n",
241 pg_schema(&s.module),
242 qi(&s.name),
243 s.pg_type,
244 check_clause,
245 ));
246 }
247 if !schema.scalars.is_empty() {
248 out.push('\n');
249 }
250}
251
252fn emit_tables(schema: &SchemaDescriptor, out: &mut String) {
255 for t in &schema.types {
256 if t.abstract_ || t.junction {
257 continue;
258 }
259 emit_one_table(t, Some(schema), out);
260 }
261}
262
263fn column_default(p: &crate::schema::PropertyDescriptor, schema: Option<&SchemaDescriptor>) -> Option<String> {
270 if let Some(sql) = &p.default_sql {
271 return Some(sql.clone());
272 }
273 let pyql = p.default_pyql.as_deref()?;
274 crate::ir::compile_scalar_default(pyql, schema?).ok()
275}
276
277fn emit_one_table(t: &TypeDescriptor, schema: Option<&SchemaDescriptor>, out: &mut String) {
278 out.push_str(&format!("CREATE TABLE {} (\n", qn(&t.module, &t.table)));
279
280 let mut lines: Vec<String> = Vec::new();
281
282 for p in &t.properties {
284 let not_null = if p.nullable { "" } else { " NOT NULL" };
285 let default = column_default(p, schema)
286 .map(|d| format!(" DEFAULT {}", d))
287 .unwrap_or_default();
288 let col_type = p.column_type.as_deref().unwrap_or_else(|| resolved_pg_type(&p.pg_type));
289 lines.push(format!(" {} {}{}{}", qi(&p.name), col_type, not_null, default));
290 }
291
292 for l in &t.links {
296 if l.is_junction_backed() {
297 continue;
298 }
299 let not_null = if l.nullable { "" } else { " NOT NULL" };
300 lines.push(format!(" {} uuid{}", qi(&format!("{}_id", l.name)), not_null));
301 }
302
303 let mut pk_cols: Vec<String> = t.properties.iter().filter(|p| p.is_pk).map(|p| qi(&p.name)).collect();
308 if let Some(part) = &t.partition {
309 let key = qi(&part.pointer);
310 if !pk_cols.contains(&key) {
311 pk_cols.push(key);
312 }
313 }
314 if !pk_cols.is_empty() {
315 lines.push(format!(" PRIMARY KEY ({})", pk_cols.join(", ")));
316 }
317
318 out.push_str(&lines.join(",\n"));
319 out.push_str("\n)");
320 if let Some(part) = &t.partition {
321 out.push_str(&format!(" PARTITION BY RANGE ({})", qi(&part.pointer)));
322 }
323 out.push_str(";\n\n");
324 if let Some(part) = &t.partition {
325 out.push_str(&partman_setup_sql(&t.module, &t.table, part));
326 out.push_str("\n\n");
327 }
328 out.push_str(&cache_invalidate_trigger_sql(&qn(&t.module, &t.table)));
329 out.push_str("\n\n");
330}
331
332fn partman_setup_sql(module: &str, table: &str, part: &crate::schema::PartitionDescriptor) -> String {
345 let pg_schema = if module == "default" { "public" } else { module };
346 let parent = format!("{pg_schema}.{table}").replace('\'', "''");
349 let mut out = String::new();
350
351 out.push_str("DO $$ BEGIN\n");
352 out.push_str(&format!(
353 " IF NOT EXISTS (SELECT 1 FROM partman.part_config WHERE parent_table = '{parent}') THEN\n"
354 ));
355 out.push_str(&format!(
356 " PERFORM partman.create_parent(\n\
357 \x20 p_parent_table := '{parent}',\n\
358 \x20 p_control := '{control}',\n\
359 \x20 p_interval := '{interval}',\n\
360 \x20 p_premake := {premake}\n\
361 \x20 );\n",
362 control = part.pointer.replace('\'', "''"),
363 interval = part.interval.as_pg_interval(),
364 premake = part.premake,
365 ));
366 out.push_str(" END IF;\nEND $$;\n");
367
368 match part.retention_interval() {
369 Some(retention) => {
370 out.push_str(&format!(
371 "UPDATE partman.part_config\n\
372 \x20 SET retention = '{retention}', retention_keep_table = false\n\
373 \x20 WHERE parent_table = '{parent}';",
374 ));
375 }
376 None => {
377 out.push_str(&format!(
381 "UPDATE partman.part_config\n SET retention = NULL\n WHERE parent_table = '{parent}';",
382 ));
383 }
384 }
385 out
386}
387
388fn cache_invalidate_trigger_sql(qualified_table: &str) -> String {
394 format!(
395 "CREATE OR REPLACE TRIGGER pylon_cache_invalidate\n AFTER INSERT OR UPDATE OR DELETE ON {}\n FOR EACH STATEMENT EXECUTE FUNCTION _pylon.notify_cache_invalidate();",
396 qualified_table
397 )
398}
399
400pub(crate) fn polymorphic_types(schema: &SchemaDescriptor) -> HashSet<String> {
411 let extended: HashSet<&str> = schema
412 .types
413 .iter()
414 .flat_map(|t| t.bases.iter().map(String::as_str))
415 .collect();
416 schema
417 .types
418 .iter()
419 .map(|t| format!("{}::{}", t.module, t.name))
420 .zip(&schema.types)
421 .filter(|(qname, t)| (t.abstract_ && t.materialized) || extended.contains(qname.as_str()))
422 .map(|(qname, _)| qname)
423 .collect()
424}
425
426pub(crate) fn interface_implementors(schema: &SchemaDescriptor) -> HashMap<String, Vec<&TypeDescriptor>> {
430 let mut implementors: HashMap<String, Vec<&TypeDescriptor>> = HashMap::new();
431 for t in &schema.types {
432 if !t.abstract_ {
433 for iface in &t.interfaces {
434 implementors.entry(iface.clone()).or_default().push(t);
435 }
436 }
437 }
438 for t in &schema.types {
440 for base in &t.bases {
441 if let Some(base_td) = schema
442 .types
443 .iter()
444 .find(|b| format!("{}::{}", b.module, b.name) == *base)
445 {
446 let entry = implementors.entry(base.clone()).or_default();
447 if entry.is_empty() {
448 entry.push(base_td);
449 }
450 entry.push(t);
451 }
452 }
453 }
454 implementors
455}
456
457fn policy_for<'a>(policies: &'a [OnDeletePolicy], side: &DeleteSide) -> Option<&'a DeleteAction> {
460 policies.iter().find(|p| &p.side == side).map(|p| &p.action)
461}
462
463pub(crate) fn needs_deferred_target_fk(policies: &[OnDeletePolicy]) -> bool {
475 policies.iter().any(|p| {
476 p.side == DeleteSide::Source
477 && matches!(
478 p.action,
479 DeleteAction::DeleteTarget | DeleteAction::DeleteTargetIfOrphan
480 )
481 })
482}
483
484fn target_fk_suffix(policies: &[OnDeletePolicy]) -> String {
487 match policy_for(policies, &DeleteSide::Target).unwrap_or(&DeleteAction::Restrict) {
488 DeleteAction::Restrict if needs_deferred_target_fk(policies) => " DEFERRABLE INITIALLY DEFERRED".into(),
489 DeleteAction::Restrict => " ON DELETE RESTRICT".into(),
490 DeleteAction::DeferredRestrict => " DEFERRABLE INITIALLY DEFERRED".into(),
491 DeleteAction::DeleteSource => " ON DELETE CASCADE".into(),
492 DeleteAction::Allow => " ON DELETE SET NULL".into(),
493 _ => " ON DELETE RESTRICT".into(),
494 }
495}
496
497fn source_jt_fk_suffix(_policies: &[OnDeletePolicy]) -> &'static str {
502 " ON DELETE CASCADE"
503}
504
505fn target_jt_fk_suffix(policies: &[OnDeletePolicy]) -> String {
507 match policy_for(policies, &DeleteSide::Target).unwrap_or(&DeleteAction::Restrict) {
508 DeleteAction::Restrict if needs_deferred_target_fk(policies) => " DEFERRABLE INITIALLY DEFERRED".into(),
509 DeleteAction::Restrict => " ON DELETE RESTRICT".into(),
510 DeleteAction::DeferredRestrict => " DEFERRABLE INITIALLY DEFERRED".into(),
511 DeleteAction::Allow => " ON DELETE CASCADE".into(),
512 DeleteAction::DeleteSource => " ON DELETE CASCADE".into(),
514 _ => " ON DELETE RESTRICT".into(),
515 }
516}
517
518fn emit_fk_constraints(schema: &SchemaDescriptor, type_map: &HashMap<String, (&str, &str)>, out: &mut String) {
521 let polymorphic = polymorphic_types(schema);
522 let mut emitted = false;
523 for t in &schema.types {
524 if t.abstract_ || t.junction {
525 continue;
526 }
527 for l in &t.links {
528 if l.is_junction_backed() {
529 continue;
530 }
531 if polymorphic.contains(&l.target) {
532 continue;
533 }
534 let Some((tgt_module, tgt_table)) = type_map.get(&l.target) else {
535 continue;
536 };
537 let cname = qi(&format!("{}_{}_fkey", t.table, l.name));
538 let suffix = target_fk_suffix(&l.on_delete);
539 out.push_str(&format!(
540 "ALTER TABLE {} ADD CONSTRAINT {} FOREIGN KEY ({}) REFERENCES {}(id){};\n",
541 qn(&t.module, &t.table),
542 cname,
543 qi(&format!("{}_id", l.name)),
544 qn(tgt_module, tgt_table),
545 suffix,
546 ));
547 emitted = true;
548 }
549 }
550 if emitted {
551 out.push('\n');
552 }
553}
554
555pub fn junction_fk_constraints(
567 schema: &SchemaDescriptor,
568 type_map: &HashMap<String, (&str, &str)>,
569) -> Vec<(String, String, String, String)> {
570 let polymorphic = polymorphic_types(schema);
571 let mut result = Vec::new();
572 for t in &schema.types {
573 if t.abstract_ || t.junction {
574 continue;
575 }
576 let pointers = t
577 .multilinks
578 .iter()
579 .map(|ml| (ml.name.as_str(), &ml.target, &ml.on_delete))
580 .chain(
581 t.links
582 .iter()
583 .filter(|l| l.is_junction_backed())
584 .map(|l| (l.name.as_str(), &l.target, &l.on_delete)),
585 );
586 for (name, target, on_delete) in pointers {
587 if polymorphic.contains(target) {
588 continue;
589 }
590 let Some((tgt_module, tgt_table)) = type_map.get(target) else {
591 continue;
592 };
593 let jt_name = format!("{}.{}", t.table, name);
594 let cname = format!("{}_{}_target_fkey", t.table, name);
598 let ddl = format!(
599 "ALTER TABLE {} ADD CONSTRAINT {} FOREIGN KEY (target) REFERENCES {}(id){};",
600 qn(&t.module, &jt_name),
601 qi(&cname),
602 qn(tgt_module, tgt_table),
603 target_jt_fk_suffix(on_delete),
604 );
605 result.push((t.module.clone(), jt_name, cname, ddl));
606 }
607 }
608 result
609}
610
611fn emit_junction_fk_constraints(schema: &SchemaDescriptor, type_map: &HashMap<String, (&str, &str)>, out: &mut String) {
612 let constraints = junction_fk_constraints(schema, type_map);
613 for (_, _, _, ddl) in &constraints {
614 out.push_str(ddl);
615 out.push('\n');
616 }
617 if !constraints.is_empty() {
618 out.push('\n');
619 }
620}
621
622fn emit_before_delete_trigger(fn_qname: &str, trigger_name: &str, table_qname: &str, body: &str, out: &mut String) {
625 out.push_str(&format!(
626 "CREATE OR REPLACE FUNCTION {fn_qname}()\n\
627 RETURNS trigger LANGUAGE plpgsql AS $$\n\
628 BEGIN\n"
629 ));
630 out.push_str(body);
631 out.push_str(&format!(
632 "\n RETURN OLD;\nEND;\n$$;\n\n\
633 CREATE TRIGGER {trigger_name}\n\
634 BEFORE DELETE ON {table_qname}\n\
635 FOR EACH ROW EXECUTE FUNCTION {fn_qname}();\n\n"
636 ));
637}
638
639fn emit_after_delete_trigger(fn_qname: &str, trigger_name: &str, table_qname: &str, body: &str, out: &mut String) {
651 out.push_str(&format!(
652 "CREATE OR REPLACE FUNCTION {fn_qname}()\n\
653 RETURNS trigger LANGUAGE plpgsql AS $$\n\
654 BEGIN\n"
655 ));
656 out.push_str(body);
657 out.push_str(&format!(
658 "\n RETURN NULL;\nEND;\n$$;\n\n\
659 CREATE TRIGGER {trigger_name}\n\
660 AFTER DELETE ON {table_qname}\n\
661 FOR EACH ROW EXECUTE FUNCTION {fn_qname}();\n\n"
662 ));
663}
664
665fn emit_after_mutation_trigger(
673 fn_qname: &str,
674 trigger_name: &str,
675 table_qname: &str,
676 events: &str,
677 body: &str,
678 out: &mut String,
679) {
680 out.push_str(&format!(
681 "CREATE OR REPLACE FUNCTION {fn_qname}()\n\
682 RETURNS trigger LANGUAGE plpgsql AS $$\n\
683 BEGIN\n"
684 ));
685 out.push_str(body);
686 out.push_str(&format!(
687 "\n RETURN NULL;\nEND;\n$$;\n\n\
688 CREATE OR REPLACE TRIGGER {trigger_name}\n\
689 AFTER {events} ON {table_qname}\n\
690 FOR EACH ROW EXECUTE FUNCTION {fn_qname}();\n\n"
691 ));
692}
693
694fn target_table_qnames(
700 target: &str,
701 type_map: &HashMap<String, (&str, &str)>,
702 polymorphic: &HashSet<String>,
703 implementors: &HashMap<String, Vec<&TypeDescriptor>>,
704) -> Option<Vec<String>> {
705 if polymorphic.contains(target) {
706 let impls = implementors.get(target)?;
707 if impls.is_empty() {
708 return None;
709 }
710 return Some(impls.iter().map(|i| qn(&i.module, &i.table)).collect());
711 }
712 let (tgt_module, tgt_table) = type_map.get(target)?;
713 Some(vec![qn(tgt_module, tgt_table)])
714}
715
716pub struct DeletionTriggerInfo {
730 pub table_module: String,
731 pub table_name: String,
735 pub trigger_name: String,
736 pub ddl: String,
737}
738
739fn link_source_trigger_infos(
740 schema: &SchemaDescriptor,
741 type_map: &HashMap<String, (&str, &str)>,
742) -> Vec<DeletionTriggerInfo> {
743 let polymorphic = polymorphic_types(schema);
744 let implementors = interface_implementors(schema);
745 let mut result = Vec::new();
746 for t in &schema.types {
747 if t.abstract_ || t.junction {
748 continue;
749 }
750 for l in &t.links {
751 if l.is_junction_backed() {
756 continue;
757 }
758 let src_action = policy_for(&l.on_delete, &DeleteSide::Source).unwrap_or(&DeleteAction::Allow);
759 match src_action {
760 DeleteAction::Allow => continue,
761 DeleteAction::DeleteTarget | DeleteAction::DeleteTargetIfOrphan => {}
762 _ => continue,
763 }
764 let suffix = if matches!(src_action, DeleteAction::DeleteTargetIfOrphan) {
765 "del_orphan"
766 } else {
767 "del_target"
768 };
769 let tbl_qname = qn(&t.module, &t.table);
770 let col = qi(&format!("{}_id", l.name));
771 let Some(tgt_qnames) = target_table_qnames(&l.target, type_map, &polymorphic, &implementors) else {
772 continue;
773 };
774 let deletes = tgt_qnames
775 .iter()
776 .map(|tgt_qname| format!("DELETE FROM {tgt_qname} WHERE id = OLD.{col};"))
777 .collect::<Vec<_>>();
778
779 let body = if matches!(src_action, DeleteAction::DeleteTargetIfOrphan) {
780 let indented = deletes
781 .iter()
782 .map(|d| format!(" {d}"))
783 .collect::<Vec<_>>()
784 .join("\n");
785 format!(
786 " IF NOT EXISTS (\n SELECT 1 FROM {tbl_qname} WHERE {col} = OLD.{col} AND id != OLD.id\n ) THEN\n{indented}\n END IF;"
787 )
788 } else {
789 deletes
790 .iter()
791 .map(|d| format!(" {d}"))
792 .collect::<Vec<_>>()
793 .join("\n")
794 };
795
796 let (fname, fn_qname) = trigger_names(&t.module, &t.table, &l.name, suffix, &body);
797
798 let mut ddl = String::new();
799 emit_before_delete_trigger(&fn_qname, &qi(&fname), &tbl_qname, &body, &mut ddl);
800 result.push(DeletionTriggerInfo {
801 table_module: t.module.clone(),
802 table_name: t.table.clone(),
803 trigger_name: fname,
804 ddl,
805 });
806 }
807 }
808 result
809}
810
811fn emit_link_source_triggers(schema: &SchemaDescriptor, type_map: &HashMap<String, (&str, &str)>, out: &mut String) {
812 for info in link_source_trigger_infos(schema, type_map) {
813 out.push_str(&info.ddl);
814 }
815}
816
817fn emit_junction_tables(schema: &SchemaDescriptor, out: &mut String) {
820 for t in &schema.types {
821 if t.abstract_ || t.junction {
822 continue;
823 }
824 for ml in &t.multilinks {
825 emit_one_junction_table(
826 schema,
827 t,
828 &ml.name,
829 &ml.on_delete,
830 ml.through.as_deref(),
831 false,
832 ml.is_exclusive,
833 out,
834 );
835 }
836 for l in &t.links {
842 if !l.is_junction_backed() {
843 continue;
844 }
845 emit_one_junction_table(
846 schema,
847 t,
848 &l.name,
849 &l.on_delete,
850 l.through.as_deref(),
851 true,
852 l.is_exclusive,
853 out,
854 );
855 }
856 }
857}
858
859#[allow(clippy::too_many_arguments)]
862fn emit_one_junction_table(
863 schema: &SchemaDescriptor,
864 t: &TypeDescriptor,
865 name: &str,
866 on_delete: &[OnDeletePolicy],
867 through: Option<&str>,
868 single: bool,
869 exclusive: bool,
870 out: &mut String,
871) {
872 let jt_name = format!("{}.{}", t.table, name);
873 let src_suffix = source_jt_fk_suffix(on_delete);
874 out.push_str(&format!(
875 "CREATE TABLE {} (\n source uuid NOT NULL REFERENCES {}(id){},\n",
876 qn(&t.module, &jt_name),
877 qn(&t.module, &t.table),
878 src_suffix,
879 ));
880 out.push_str(" target uuid NOT NULL,\n");
884
885 if let Some(through_qname) = through {
887 let through_td = schema
888 .types
889 .iter()
890 .find(|td| format!("{}::{}", td.module, td.name) == *through_qname);
891 if let Some(td) = through_td
892 && td.junction
893 {
894 for p in &td.properties {
895 if p.name == "id" {
896 continue;
897 }
898 let not_null = if p.nullable { "" } else { " NOT NULL" };
899 out.push_str(&format!(" {} {}{},\n", qi(&p.name), p.pg_type, not_null));
900 }
901 }
902 }
903
904 if single {
905 out.push_str(" PRIMARY KEY (source)");
906 } else {
907 out.push_str(" PRIMARY KEY (source, target)");
908 }
909 if exclusive {
910 out.push_str(",\n UNIQUE (target)\n);\n\n");
911 } else {
912 out.push_str("\n);\n\n");
913 }
914 out.push_str(&cache_invalidate_trigger_sql(&qn(&t.module, &jt_name)));
915 out.push_str("\n\n");
916}
917
918fn multilink_deletion_trigger_infos(
921 schema: &SchemaDescriptor,
922 type_map: &HashMap<String, (&str, &str)>,
923) -> Vec<DeletionTriggerInfo> {
924 let polymorphic = polymorphic_types(schema);
925 let implementors = interface_implementors(schema);
926 let mut result = Vec::new();
927 for t in &schema.types {
928 if t.abstract_ || t.junction {
929 continue;
930 }
931 for ml in &t.multilinks {
932 push_junction_deletion_triggers(
933 t,
934 &ml.name,
935 &ml.target,
936 &ml.on_delete,
937 type_map,
938 &polymorphic,
939 &implementors,
940 &mut result,
941 );
942 }
943 for l in &t.links {
947 if !l.is_junction_backed() {
948 continue;
949 }
950 push_junction_deletion_triggers(
951 t,
952 &l.name,
953 &l.target,
954 &l.on_delete,
955 type_map,
956 &polymorphic,
957 &implementors,
958 &mut result,
959 );
960 }
961 }
962 result
963}
964
965#[allow(clippy::too_many_arguments)]
966fn push_junction_deletion_triggers(
967 t: &TypeDescriptor,
968 name: &str,
969 target: &str,
970 on_delete: &[OnDeletePolicy],
971 type_map: &HashMap<String, (&str, &str)>,
972 polymorphic: &HashSet<String>,
973 implementors: &HashMap<String, Vec<&TypeDescriptor>>,
974 result: &mut Vec<DeletionTriggerInfo>,
975) {
976 let jt_name = format!("{}.{}", t.table, name);
977 let jt_qname = qn(&t.module, &jt_name);
978
979 let src_action = policy_for(on_delete, &DeleteSide::Source).unwrap_or(&DeleteAction::Allow);
981 if matches!(
982 src_action,
983 DeleteAction::DeleteTarget | DeleteAction::DeleteTargetIfOrphan
984 ) && let Some(tgt_qnames) = target_table_qnames(target, type_map, polymorphic, implementors)
985 {
986 let suffix = if matches!(src_action, DeleteAction::DeleteTargetIfOrphan) {
987 "del_orphan"
988 } else {
989 "del_target"
990 };
991 let owner_qname = qn(&t.module, &t.table);
998 let orphan = matches!(src_action, DeleteAction::DeleteTargetIfOrphan);
999 let body = tgt_qnames
1000 .iter()
1001 .map(|tgt_qname| {
1002 let mut delete = format!(
1003 " DELETE FROM {tgt_qname} AS _tgt\n WHERE _tgt.id IN (SELECT target FROM {jt_qname} WHERE source = OLD.id)"
1004 );
1005 if orphan {
1006 delete.push_str(&format!(
1007 "\n AND NOT EXISTS (SELECT 1 FROM {jt_qname} WHERE target = _tgt.id AND source != OLD.id)"
1008 ));
1009 }
1010 delete.push(';');
1011 delete
1012 })
1013 .collect::<Vec<_>>()
1014 .join("\n");
1015
1016 let hash = fnv(&[&t.table, name, suffix, &body]);
1019 let fname = format!("{}_{}_{}_{}", t.table, name, suffix, &hash[..8]);
1020 let fn_qname = qn(&t.module, &fname);
1021
1022 let mut ddl = String::new();
1023 emit_before_delete_trigger(&fn_qname, &qi(&fname), &owner_qname, &body, &mut ddl);
1024 result.push(DeletionTriggerInfo {
1025 table_module: t.module.clone(),
1026 table_name: t.table.clone(),
1027 trigger_name: fname,
1028 ddl,
1029 });
1030 }
1031
1032 let tgt_action = policy_for(on_delete, &DeleteSide::Target).unwrap_or(&DeleteAction::Restrict);
1035 if matches!(tgt_action, DeleteAction::DeleteSource) {
1036 let src_qname = qn(&t.module, &t.table);
1037
1038 let body = format!(" DELETE FROM {src_qname} WHERE id = OLD.source;");
1039 let (fname, fn_qname) = trigger_names(&t.module, &t.table, name, "del_source", &body);
1040
1041 let mut ddl = String::new();
1042 emit_after_delete_trigger(&fn_qname, &qi(&fname), &jt_qname, &body, &mut ddl);
1043 result.push(DeletionTriggerInfo {
1044 table_module: t.module.clone(),
1045 table_name: jt_name.clone(),
1046 trigger_name: fname,
1047 ddl,
1048 });
1049 }
1050}
1051
1052#[allow(clippy::too_many_arguments)]
1063fn push_interface_link_triggers(
1064 module: &str,
1065 referencing: &str,
1066 column: &str,
1067 pointer: &str,
1068 src_table: &str,
1069 via_junction: bool,
1070 on_delete: &[OnDeletePolicy],
1071 impls: &[&TypeDescriptor],
1072 result: &mut Vec<DeletionTriggerInfo>,
1073) {
1074 let action = policy_for(on_delete, &DeleteSide::Target).unwrap_or(&DeleteAction::Restrict);
1075 let col = qi(column);
1076
1077 let deferred = matches!(action, DeleteAction::DeferredRestrict) || needs_deferred_target_fk(on_delete);
1082
1083 let body = match action {
1084 DeleteAction::Allow if via_junction => {
1085 format!(" DELETE FROM {referencing} WHERE {col} = OLD.id;")
1088 }
1089 DeleteAction::Allow => format!(" UPDATE {referencing} SET {col} = NULL WHERE {col} = OLD.id;"),
1090 DeleteAction::DeleteSource => format!(" DELETE FROM {referencing} WHERE {col} = OLD.id;"),
1091 _ => format!(
1092 " IF EXISTS (SELECT 1 FROM {referencing} WHERE {col} = OLD.id) THEN\n\
1093 \x20 RAISE foreign_key_violation\n\
1094 \x20 USING MESSAGE = 'update or delete on table \"' || TG_TABLE_NAME || '\" violates foreign key constraint on table {src_table}',\n\
1095 \x20 DETAIL = format('Key (id)=(%s) is still referenced from table \"{src_table}\".', OLD.id);\n\
1096 \x20 END IF;"
1097 ),
1098 };
1099
1100 let suffix = match action {
1101 DeleteAction::Allow => "ifl_allow",
1102 DeleteAction::DeleteSource => "ifl_del_source",
1103 _ => "ifl_restrict",
1104 };
1105
1106 for impl_t in impls {
1107 let hash = fnv(&[&impl_t.table, src_table, pointer, suffix, &body]);
1108 let fname = format!("_ifl_{}_{}_{}", src_table, pointer, &hash[..8]);
1109 let fn_qname = qn(module, &fname);
1110 let impl_qname = qn(&impl_t.module, &impl_t.table);
1111
1112 let mut ddl = String::new();
1113 if deferred {
1114 ddl.push_str(&format!(
1118 "CREATE OR REPLACE FUNCTION {fn_qname}()\n\
1119 RETURNS trigger LANGUAGE plpgsql AS $$\n\
1120 BEGIN\n{body}\n RETURN NULL;\nEND;\n$$;\n\n\
1121 CREATE CONSTRAINT TRIGGER {}\n\
1122 AFTER DELETE ON {impl_qname}\n\
1123 DEFERRABLE INITIALLY DEFERRED\n\
1124 FOR EACH ROW EXECUTE FUNCTION {fn_qname}();\n\n",
1125 qi(&fname),
1126 ));
1127 } else {
1128 emit_before_delete_trigger(&fn_qname, &qi(&fname), &impl_qname, &body, &mut ddl);
1129 }
1130
1131 result.push(DeletionTriggerInfo {
1132 table_module: impl_t.module.clone(),
1133 table_name: impl_t.table.clone(),
1134 trigger_name: fname,
1135 ddl,
1136 });
1137 }
1138}
1139
1140fn interface_link_trigger_infos(schema: &SchemaDescriptor) -> Vec<DeletionTriggerInfo> {
1144 let polymorphic = polymorphic_types(schema);
1145 let implementors = interface_implementors(schema);
1146 let mut result = Vec::new();
1147
1148 for t in &schema.types {
1149 if t.abstract_ || t.junction {
1150 continue;
1151 }
1152 for l in &t.links {
1153 if !polymorphic.contains(&l.target) {
1154 continue;
1155 }
1156 let Some(impls) = implementors.get(&l.target) else {
1157 continue;
1158 };
1159 let via_junction = l.is_junction_backed();
1160 let (referencing, column) = if via_junction {
1161 (qn(&t.module, &format!("{}.{}", t.table, l.name)), "target".to_string())
1162 } else {
1163 (qn(&t.module, &t.table), format!("{}_id", l.name))
1164 };
1165 push_interface_link_triggers(
1166 &t.module,
1167 &referencing,
1168 &column,
1169 &l.name,
1170 &t.table,
1171 via_junction,
1172 &l.on_delete,
1173 impls,
1174 &mut result,
1175 );
1176 }
1177 for ml in &t.multilinks {
1178 if !polymorphic.contains(&ml.target) {
1179 continue;
1180 }
1181 let Some(impls) = implementors.get(&ml.target) else {
1182 continue;
1183 };
1184 let referencing = qn(&t.module, &format!("{}.{}", t.table, ml.name));
1185 push_interface_link_triggers(
1186 &t.module,
1187 &referencing,
1188 "target",
1189 &ml.name,
1190 &t.table,
1191 true,
1192 &ml.on_delete,
1193 impls,
1194 &mut result,
1195 );
1196 }
1197 }
1198 result
1199}
1200
1201fn emit_interface_link_triggers(schema: &SchemaDescriptor, out: &mut String) {
1202 for info in interface_link_trigger_infos(schema) {
1203 out.push_str(&info.ddl);
1204 }
1205}
1206
1207fn emit_multilink_deletion_triggers(
1208 schema: &SchemaDescriptor,
1209 type_map: &HashMap<String, (&str, &str)>,
1210 out: &mut String,
1211) {
1212 for info in multilink_deletion_trigger_infos(schema, type_map) {
1213 out.push_str(&info.ddl);
1214 }
1215}
1216
1217pub fn deletion_policy_trigger_infos(
1220 schema: &SchemaDescriptor,
1221 type_map: &HashMap<String, (&str, &str)>,
1222) -> Vec<DeletionTriggerInfo> {
1223 let mut result = link_source_trigger_infos(schema, type_map);
1224 result.extend(multilink_deletion_trigger_infos(schema, type_map));
1225 result.extend(interface_link_trigger_infos(schema));
1226 result
1227}
1228
1229pub fn signal_trigger_infos(schema: &SchemaDescriptor) -> Vec<DeletionTriggerInfo> {
1238 use crate::schema::SearchBackend;
1239
1240 let mut result = Vec::new();
1241 for t in &schema.types {
1242 if t.abstract_ || t.junction || t.signals.is_empty() {
1243 continue;
1244 }
1245
1246 let qname = format!("{}::{}", t.module, t.name);
1247 let qname_literal = format!("'{}'", qname.replace('\'', "''"));
1248 let tbl_qname = qn(&t.module, &t.table);
1249
1250 let combined_on = t.signals.iter().fold(0u8, |acc, s| acc | s.on);
1256 let mut events = Vec::new();
1257 if combined_on & 1 != 0 {
1258 events.push("INSERT");
1259 }
1260 if combined_on & 2 != 0 {
1261 events.push("UPDATE");
1262 }
1263 if combined_on & 4 != 0 {
1264 events.push("DELETE");
1265 }
1266 let events_str = events.join(" OR ");
1267
1268 let index_cols: Vec<String> = t
1281 .vector_indexes
1282 .iter()
1283 .map(|vi| vi.column_name())
1284 .chain(
1285 t.search_indexes
1286 .iter()
1287 .filter(|si| si.backend == SearchBackend::Postgres)
1288 .map(|si| si.column_name()),
1289 )
1290 .collect();
1291
1292 let update_guard = if combined_on & 2 != 0 && !index_cols.is_empty() {
1293 let strip: String = index_cols
1294 .iter()
1295 .map(|c| format!(" - '{}'", c.replace('\'', "''")))
1296 .collect();
1297 format!(
1298 " IF TG_OP = 'UPDATE' AND (to_jsonb(OLD){strip}) = (to_jsonb(NEW){strip}) THEN\n \
1299 RETURN NULL;\n \
1300 END IF;\n"
1301 )
1302 } else {
1303 String::new()
1304 };
1305
1306 let body = format!(
1307 "{update_guard} INSERT INTO _pylon.\"SignalOutbox\" (type_name, operation, old_row, new_row)\n \
1308 VALUES (\n \
1309 {qname_literal},\n \
1310 TG_OP,\n \
1311 CASE WHEN TG_OP IN ('UPDATE', 'DELETE') THEN to_jsonb(OLD) ELSE NULL END,\n \
1312 CASE WHEN TG_OP IN ('INSERT', 'UPDATE') THEN to_jsonb(NEW) ELSE NULL END\n \
1313 );"
1314 );
1315
1316 let hash = fnv(&[&t.table, "signal", &events_str, &body]);
1320 let fname = format!("{}_signal_{}", t.table, &hash[..8]);
1321 let fn_qname = qn(&t.module, &fname);
1322
1323 let mut ddl = String::new();
1324 emit_after_mutation_trigger(&fn_qname, &qi(&fname), &tbl_qname, &events_str, &body, &mut ddl);
1325 result.push(DeletionTriggerInfo {
1326 table_module: t.module.clone(),
1327 table_name: t.table.clone(),
1328 trigger_name: fname,
1329 ddl,
1330 });
1331 }
1332 result
1333}
1334
1335fn emit_signal_triggers(schema: &SchemaDescriptor, out: &mut String) {
1336 for info in signal_trigger_infos(schema) {
1337 out.push_str(&info.ddl);
1338 }
1339}
1340
1341fn constraint_column(t: &TypeDescriptor, pointer: &str) -> String {
1349 if t.links.iter().any(|l| l.name == pointer && !l.is_junction_backed()) {
1350 format!("{}_id", pointer)
1351 } else {
1352 pointer.to_string()
1353 }
1354}
1355
1356fn emit_unique_indexes(schema: &SchemaDescriptor, out: &mut String) {
1359 let mut emitted = false;
1360 for t in &schema.types {
1361 if t.abstract_ || t.junction {
1362 continue;
1363 }
1364 let qname = qn(&t.module, &t.table);
1365
1366 for p in &t.properties {
1367 if p.is_exclusive && !p.is_pk {
1368 out.push_str(&format!("CREATE UNIQUE INDEX ON {} ({});\n", qname, qi(&p.name),));
1369 emitted = true;
1370 }
1371 }
1372 for l in &t.links {
1373 if l.is_exclusive && !l.is_junction_backed() {
1377 out.push_str(&format!(
1378 "CREATE UNIQUE INDEX ON {} ({});\n",
1379 qname,
1380 qi(&format!("{}_id", l.name)),
1381 ));
1382 emitted = true;
1383 }
1384 }
1385 for c in &t.constraints {
1386 if let TypeConstraint::Exclusive {
1387 pointers: fields,
1388 unless,
1389 } = c
1390 {
1391 let cols: Vec<String> = fields.iter().map(|f| qi(&constraint_column(t, f))).collect();
1392 let where_clause = match unless.as_deref() {
1396 Some(u) => {
1397 let qualified = format!("{}::{}", t.module, t.name);
1398 match crate::ir::compile_constraint_expr(u, &qualified, schema) {
1399 Ok(compiled) => format!(" WHERE NOT ({compiled})"),
1400 Err(_) => continue,
1405 }
1406 }
1407 None => String::new(),
1408 };
1409 out.push_str(&format!(
1410 "CREATE UNIQUE INDEX ON {} ({}){};\n",
1411 qname,
1412 cols.join(", "),
1413 where_clause,
1414 ));
1415 emitted = true;
1416 }
1417 }
1418 }
1419 if emitted {
1420 out.push('\n');
1421 }
1422}
1423
1424pub fn scalar_check_constraints(schema: &SchemaDescriptor) -> Vec<(String, String, String, String)> {
1442 let mut result = Vec::new();
1443 for s in &schema.scalars {
1444 for expr in &s.check_constraints {
1445 let hash = fnv(&[&s.name, expr.as_str()]);
1446 result.push((
1447 s.module.clone(),
1448 s.name.clone(),
1449 format!("{}_{}_check", s.name, &hash[..8]),
1450 expr.clone(),
1451 ));
1452 }
1453 }
1454 result
1455}
1456
1457pub fn scalar_check_clauses(schema: &SchemaDescriptor, module: &str, name: &str) -> String {
1459 let clauses: Vec<String> = scalar_check_constraints(schema)
1460 .into_iter()
1461 .filter(|(m, n, _, _)| m == module && n == name)
1462 .map(|(_, _, cname, expr)| format!(" CONSTRAINT {} CHECK ({})", qi(&cname), expr))
1463 .collect();
1464 if clauses.is_empty() {
1465 String::new()
1466 } else {
1467 format!("\n{}", clauses.join("\n"))
1468 }
1469}
1470
1471pub fn check_constraints(schema: &SchemaDescriptor) -> Result<Vec<(String, String, String, String)>, PyQLError> {
1472 let mut result = Vec::new();
1473 for t in &schema.types {
1474 if t.abstract_ || t.junction {
1475 continue;
1476 }
1477 let qname = qn(&t.module, &t.table);
1478 for p in &t.properties {
1479 for expr in &p.check_constraints {
1480 let hash = fnv(&[&t.table, &p.name, expr.as_str()]);
1481 let cname = format!("{}_{}_{}_check", t.table, p.name, &hash[..8]);
1482 result.push((
1483 t.module.clone(),
1484 t.table.clone(),
1485 cname.clone(),
1486 format!("ALTER TABLE {} ADD CONSTRAINT {} CHECK ({});", qname, qi(&cname), expr),
1487 ));
1488 }
1489 }
1490 for c in &t.constraints {
1491 if let TypeConstraint::Expression { expr } = c {
1492 let qualified = format!("{}::{}", t.module, t.name);
1493 let sql = crate::ir::compile_constraint_expr(expr, &qualified, schema).map_err(|e| {
1494 PyQLError::Fragment(PyQLFragmentError {
1495 message: format!("error in constraint on '{}': {}", qualified, e),
1496 context: qualified.clone(),
1497 position: crate::error::Position { line: 0, col: 0 },
1498 })
1499 })?;
1500 let hash = fnv(&[&t.table, expr.as_str()]);
1501 let cname = format!("{}_{}_check", t.table, &hash[..8]);
1502 result.push((
1503 t.module.clone(),
1504 t.table.clone(),
1505 cname.clone(),
1506 format!("ALTER TABLE {} ADD CONSTRAINT {} CHECK ({});", qname, qi(&cname), sql),
1507 ));
1508 }
1509 }
1510 }
1511 Ok(result)
1512}
1513
1514fn emit_check_constraints(schema: &SchemaDescriptor, out: &mut String) -> Result<(), PyQLError> {
1515 let constraints = check_constraints(schema)?;
1516 for (_, _, _, ddl) in &constraints {
1517 out.push_str(ddl);
1518 out.push('\n');
1519 }
1520 if !constraints.is_empty() {
1521 out.push('\n');
1522 }
1523 Ok(())
1524}
1525
1526pub(crate) fn index_body_and_predicate(
1536 t: &TypeDescriptor,
1537 pointers: &[String],
1538 expression: Option<&str>,
1539 unless: Option<&str>,
1540 schema: &SchemaDescriptor,
1541) -> Result<(String, String), PyQLError> {
1542 let qualified = format!("{}::{}", t.module, t.name);
1543 let body = match expression {
1544 Some(expr) => format!("({})", crate::ir::compile_constraint_expr(expr, &qualified, schema)?),
1545 None => {
1546 let mut cols = Vec::with_capacity(pointers.len());
1547 for pointer in pointers {
1548 cols.push(index_pointer_sql(t, pointer, &qualified, schema)?);
1549 }
1550 format!("({})", cols.join(", "))
1551 }
1552 };
1553 let predicate = match unless {
1554 Some(u) => format!(
1555 " WHERE NOT ({})",
1556 crate::ir::compile_constraint_expr(u, &qualified, schema)?
1557 ),
1558 None => String::new(),
1559 };
1560 Ok((body, predicate))
1561}
1562
1563pub(crate) fn index_pointer_sql(
1565 t: &TypeDescriptor,
1566 pointer: &str,
1567 qualified: &str,
1568 schema: &SchemaDescriptor,
1569) -> Result<String, PyQLError> {
1570 if t.properties.iter().any(|p| p.name == pointer) {
1571 return Ok(qi(pointer));
1572 }
1573 if t.links.iter().any(|l| l.name == pointer && !l.is_junction_backed()) {
1574 return Ok(qi(&format!("{pointer}_id")));
1575 }
1576 Ok(format!(
1578 "({})",
1579 crate::ir::compile_constraint_expr(&format!(".{pointer}"), qualified, schema)?
1580 ))
1581}
1582
1583fn emit_plain_indexes(schema: &SchemaDescriptor, out: &mut String) {
1584 let mut emitted = false;
1585 for t in &schema.types {
1586 if t.abstract_ || t.junction {
1587 continue;
1588 }
1589 let qname = qn(&t.module, &t.table);
1590
1591 for idx in &t.indexes {
1592 let unique = if idx.unique { "UNIQUE " } else { "" };
1593 let Ok((body, where_clause)) = index_body_and_predicate(
1594 t,
1595 &idx.pointers,
1596 idx.expression.as_deref(),
1597 idx.unless.as_deref(),
1598 schema,
1599 ) else {
1600 continue;
1603 };
1604 out.push_str(&format!(
1605 "CREATE {}INDEX ON {} {}{};\n",
1606 unique, qname, body, where_clause,
1607 ));
1608 emitted = true;
1609 }
1610 }
1611 if emitted {
1612 out.push('\n');
1613 }
1614}
1615
1616fn emit_triggers(schema: &SchemaDescriptor, out: &mut String) -> Result<(), PyQLError> {
1619 for info in user_trigger_infos(schema)? {
1620 out.push_str(&info.ddl);
1621 }
1622 Ok(())
1623}
1624
1625fn trigger_return_statement(timing: &str, on: u8) -> &'static str {
1635 if timing == "After" {
1636 return "RETURN NULL;";
1637 }
1638 let has_delete = on & 4 != 0;
1639 let has_insert_or_update = on & (1 | 2) != 0;
1640 match (has_delete, has_insert_or_update) {
1641 (true, true) => "IF TG_OP = 'DELETE' THEN RETURN OLD; ELSE RETURN NEW; END IF;",
1642 (true, false) => "RETURN OLD;",
1643 _ => "RETURN NEW;",
1644 }
1645}
1646
1647fn trigger_ddl_name(table: &str, trig: &crate::schema::TriggerDescriptor, body_sql: &str) -> String {
1652 let hash = fnv(&[
1655 table,
1656 &trig.on.to_string(),
1657 trig.timing.as_str(),
1658 trig.handler.as_str(),
1659 body_sql,
1660 ]);
1661 format!("{table}_{}", &hash[..12])
1662}
1663
1664fn trigger_name_body(trig: &crate::schema::TriggerDescriptor, type_name: &str, schema: &SchemaDescriptor) -> String {
1668 crate::query::compile_trigger_handler(&trig.handler, type_name, trig.on, schema)
1669 .unwrap_or_else(|_| trig.handler.clone())
1670}
1671
1672pub fn user_trigger_names(schema: &SchemaDescriptor) -> Vec<(String, String, String)> {
1678 let mut result = Vec::new();
1679 for t in &schema.types {
1680 if t.abstract_ || t.junction {
1681 continue;
1682 }
1683 let type_name = format!("{}::{}", t.module, t.name);
1684 for trig in &t.triggers {
1685 let body = trigger_name_body(trig, &type_name, schema);
1686 result.push((
1687 t.module.clone(),
1688 t.table.clone(),
1689 trigger_ddl_name(&t.table, trig, &body),
1690 ));
1691 }
1692 for (on, _) in REWRITE_EVENTS {
1693 for deferred in [false, true] {
1694 if let Some(name) = rewrite_trigger_name(t, on, schema, deferred) {
1695 result.push((t.module.clone(), t.table.clone(), name));
1696 }
1697 }
1698 }
1699 }
1700 result
1701}
1702
1703const REWRITE_EVENTS: [(u8, &str); 2] = [(1, "INSERT"), (2, "UPDATE")];
1705
1706fn rewrite_trigger_name(t: &TypeDescriptor, on: u8, schema: &SchemaDescriptor, deferred: bool) -> Option<String> {
1710 let handlers: Vec<String> = t
1711 .properties
1712 .iter()
1713 .map(|p| (&p.name, &p.rewrites))
1714 .chain(t.links.iter().map(|l| (&l.name, &l.rewrites)))
1715 .flat_map(|(name, rewrites)| {
1716 rewrites
1717 .iter()
1718 .filter(move |rw| rw.on & on != 0)
1719 .map(move |rw| format!("{name}:{}", rw.handler))
1720 })
1721 .collect();
1722 if handlers.is_empty() {
1723 return None;
1724 }
1725 let event = if on == 1 { "ins" } else { "upd" };
1726 let assignments =
1728 crate::ir::compile_rewrite_assignments(&format!("{}::{}", t.module, t.name), on, schema).unwrap_or_default();
1729 let assignments: Vec<&crate::ir::RewriteAssignment> = assignments
1730 .iter()
1731 .filter(|a| a.reads_a_multi_link == deferred)
1732 .collect();
1733 if assignments.is_empty() {
1734 return None;
1735 }
1736 let compiled = assignments.iter().map(|a| a.sql.clone()).collect::<Vec<_>>().join("\n");
1737 let timing = if deferred { "rwa" } else { "rw" };
1738 let hash = fnv(&[&t.table, event, &handlers.join("\n"), &compiled]);
1739 Some(format!("{}_{timing}_{event}_{}", t.table, &hash[..12]))
1740}
1741
1742fn rewrite_trigger_infos(schema: &SchemaDescriptor) -> Result<Vec<DeletionTriggerInfo>, PyQLError> {
1755 let mut result = Vec::new();
1756 for t in &schema.types {
1757 if t.abstract_ || t.junction {
1758 continue;
1759 }
1760 let type_name = format!("{}::{}", t.module, t.name);
1761 for (on, event) in REWRITE_EVENTS {
1762 let compiled = crate::ir::compile_rewrite_assignments(&type_name, on, schema).map_err(|e| {
1763 PyQLError::Fragment(PyQLFragmentError {
1764 message: format!("error in a rewrite of '{type_name}': {e}"),
1765 context: type_name.clone(),
1766 position: crate::error::Position { line: 0, col: 0 },
1767 })
1768 })?;
1769 for deferred in [false, true] {
1770 let Some(fname) = rewrite_trigger_name(t, on, schema, deferred) else {
1771 continue;
1772 };
1773 let assignments: Vec<&crate::ir::RewriteAssignment> = compiled
1774 .iter()
1775 .filter(|assignment| assignment.reads_a_multi_link == deferred)
1776 .collect();
1777 if assignments.is_empty() {
1778 continue;
1779 }
1780 let values = assignments
1781 .iter()
1782 .enumerate()
1783 .map(|(i, assignment)| format!("{} AS \"v{i}\"", assignment.sql))
1784 .collect::<Vec<_>>()
1785 .join(",\n\t\t");
1786 let fn_qname = qn(&t.module, &fname);
1787 let table_qname = qn(&t.module, &t.table);
1788 let globals_arg = qi(crate::ir::GLOBALS_ARG);
1789 let body = if deferred {
1790 let sets = assignments
1791 .iter()
1792 .enumerate()
1793 .map(|(i, assignment)| format!("\t\t{} = _pylon_rewrites.\"v{i}\"", qi(&assignment.column)))
1794 .collect::<Vec<_>>()
1795 .join(",\n");
1796 let stored = assignments
1797 .iter()
1798 .map(|assignment| qi(&assignment.column))
1799 .collect::<Vec<_>>()
1800 .join(", ");
1801 let computed = (0..assignments.len())
1802 .map(|i| format!("_pylon_rewrites.\"v{i}\""))
1803 .collect::<Vec<_>>()
1804 .join(", ");
1805 format!(
1806 "\tSELECT {values} INTO _pylon_rewrites;\n\
1807 \tUPDATE {table_qname} SET\n\
1808 {sets}\n\
1809 \tWHERE \"id\" = NEW.\"id\"\n\
1810 \t AND ({stored}) IS DISTINCT FROM ({computed});\n\
1811 \tRETURN NULL;"
1812 )
1813 } else {
1814 let sets = assignments
1815 .iter()
1816 .enumerate()
1817 .map(|(i, assignment)| format!("\tNEW.{} := _pylon_rewrites.\"v{i}\";", qi(&assignment.column)))
1818 .collect::<Vec<_>>()
1819 .join("\n");
1820 format!("\tSELECT {values} INTO _pylon_rewrites;\n{sets}\n\tRETURN NEW;")
1821 };
1822 let timing = if deferred { "AFTER" } else { "BEFORE" };
1823 let ddl = format!(
1824 "CREATE OR REPLACE FUNCTION {fn_qname}()\n\
1825 RETURNS trigger LANGUAGE plpgsql AS $$\n\
1826 DECLARE\n\
1827 \t_pylon_rewrites record;\n\
1828 \t{globals_arg} jsonb := nullif(current_setting('pylon.globals', true), '')::jsonb;\n\
1829 BEGIN\n\
1830 {body}\n\
1831 END;\n\
1832 $$;\n\n\
1833 CREATE OR REPLACE TRIGGER {}\n\
1834 {timing} {event} ON {table_qname}\n\
1835 FOR EACH ROW EXECUTE FUNCTION {fn_qname}();\n\n",
1836 qi(&fname),
1837 );
1838 result.push(DeletionTriggerInfo {
1839 table_module: t.module.clone(),
1840 table_name: t.table.clone(),
1841 trigger_name: fname,
1842 ddl,
1843 });
1844 }
1845 }
1846 }
1847 Ok(result)
1848}
1849
1850pub fn user_trigger_infos(schema: &SchemaDescriptor) -> Result<Vec<DeletionTriggerInfo>, PyQLError> {
1865 let mut result = rewrite_trigger_infos(schema)?;
1866 for t in &schema.types {
1867 if t.abstract_ || t.junction {
1868 continue;
1869 }
1870 let table_qname = qn(&t.module, &t.table);
1871 let type_name = format!("{}::{}", t.module, t.name);
1872
1873 for trig in &t.triggers {
1874 let fname = trigger_ddl_name(&t.table, trig, &trigger_name_body(trig, &type_name, schema));
1875 let fn_qname = qn(&t.module, &fname);
1876 let events = trigger_events(trig.on);
1877 let timing = trigger_timing(&trig.timing);
1878 let return_stmt = trigger_return_statement(&trig.timing, trig.on);
1879
1880 let body_sql =
1881 crate::query::compile_trigger_handler(&trig.handler, &type_name, trig.on, schema).map_err(|e| {
1882 let msg = format!("error in trigger handler for '{type_name}': {e}");
1883 PyQLError::Fragment(PyQLFragmentError {
1884 message: msg,
1885 context: type_name.clone(),
1886 position: crate::error::Position { line: 0, col: 0 },
1887 })
1888 })?;
1889 let body_sql = body_sql.trim_end_matches(';');
1890 let globals_arg = qi(crate::ir::GLOBALS_ARG);
1891
1892 let ddl = format!(
1893 "CREATE OR REPLACE FUNCTION {fn_qname}()\n\
1894 RETURNS trigger LANGUAGE plpgsql AS $$\n\
1895 DECLARE\n\
1896 \t_pylon_trigger_result record;\n\
1897 \t{globals_arg} jsonb := nullif(current_setting('pylon.globals', true), '')::jsonb;\n\
1898 BEGIN\n\
1899 {body_sql} INTO _pylon_trigger_result;\n\
1900 {return_stmt}\n\
1901 END;\n\
1902 $$;\n\n\
1903 CREATE OR REPLACE TRIGGER {}\n\
1904 {timing} {events} ON {table_qname}\n\
1905 FOR EACH ROW EXECUTE FUNCTION {fn_qname}();\n\n",
1906 qi(&fname),
1907 );
1908
1909 result.push(DeletionTriggerInfo {
1910 table_module: t.module.clone(),
1911 table_name: t.table.clone(),
1912 trigger_name: fname,
1913 ddl,
1914 });
1915 }
1916 }
1917 Ok(result)
1918}
1919
1920pub fn interface_view_ddl(schema: &SchemaDescriptor) -> Vec<String> {
1922 interface_view_ddl_with_names(schema)
1923 .into_iter()
1924 .map(|(_, _, ddl)| ddl)
1925 .collect()
1926}
1927
1928pub fn interface_view_ddl_with_names(schema: &SchemaDescriptor) -> Vec<(String, String, String)> {
1930 let mut implementors: HashMap<String, Vec<&TypeDescriptor>> = HashMap::new();
1931 for t in &schema.types {
1932 if !t.abstract_ {
1933 for iface in &t.interfaces {
1934 implementors.entry(iface.clone()).or_default().push(t);
1935 }
1936 }
1937 }
1938 let mut result = Vec::new();
1939 for t in &schema.types {
1940 if !(t.abstract_ && t.materialized) {
1941 continue;
1942 }
1943 let key = format!("{}::{}", t.module, t.name);
1944 let Some(impls) = implementors.get(&key) else { continue };
1945 if impls.is_empty() {
1946 continue;
1947 }
1948 let mut ddl = String::new();
1949 emit_one_interface_view(t, impls, &mut ddl);
1950 let ddl = ddl.trim().to_string();
1951 if !ddl.is_empty() {
1952 result.push((t.module.clone(), t.name.clone(), ddl));
1953 }
1954 }
1955 result
1956}
1957
1958fn interface_junction_view_name(iface_table: &str, link_name: &str) -> String {
1964 format!("{}.{}", iface_table, link_name)
1965}
1966
1967pub fn interface_junction_view_ddl_with_names(schema: &SchemaDescriptor) -> Vec<(String, String, String)> {
1983 let implementors = interface_implementors(schema);
1984 let mut result = Vec::new();
1985 for t in &schema.types {
1986 if !(t.abstract_ && t.materialized) {
1987 continue;
1988 }
1989 let key = format!("{}::{}", t.module, t.name);
1990 let Some(impls) = implementors.get(&key) else { continue };
1991 if impls.is_empty() {
1992 continue;
1993 }
1994
1995 let pointers: Vec<(&str, Option<&str>)> = t
1998 .multilinks
1999 .iter()
2000 .map(|ml| (ml.name.as_str(), ml.through.as_deref()))
2001 .chain(
2002 t.links
2003 .iter()
2004 .filter(|l| l.is_junction_backed())
2005 .map(|l| (l.name.as_str(), l.through.as_deref())),
2006 )
2007 .collect();
2008
2009 for (link_name, through) in pointers {
2010 let mut columns = vec!["source".to_string(), "target".to_string()];
2014 if let Some(through_qname) = through
2015 && let Some(td) = schema
2016 .types
2017 .iter()
2018 .find(|td| format!("{}::{}", td.module, td.name) == through_qname && td.junction)
2019 {
2020 columns.extend(td.properties.iter().filter(|p| p.name != "id").map(|p| qi(&p.name)));
2021 }
2022 let column_list = columns.join(", ");
2023
2024 let view_name = interface_junction_view_name(&t.table, link_name);
2025 let selects: Vec<String> = impls
2026 .iter()
2027 .map(|impl_t| {
2028 let jt_name = format!("{}.{}", impl_t.table, link_name);
2029 format!(" SELECT {} FROM {}", column_list, qn(&impl_t.module, &jt_name))
2030 })
2031 .collect();
2032 let ddl = format!(
2033 "CREATE VIEW {} AS\n{};",
2034 qn(&t.module, &view_name),
2035 selects.join("\n UNION ALL\n"),
2036 );
2037 result.push((t.module.clone(), view_name, ddl));
2038 }
2039 }
2040 result
2041}
2042
2043fn emit_interface_junction_views(schema: &SchemaDescriptor, out: &mut String) {
2044 for (_, _, ddl) in interface_junction_view_ddl_with_names(schema) {
2045 out.push_str(&ddl);
2046 out.push_str("\n\n");
2047 }
2048}
2049
2050pub fn function_ddl(schema: &SchemaDescriptor) -> Result<Vec<String>, crate::error::PyQLError> {
2052 function_ddl_with_names(schema).map(|v| v.into_iter().map(|(_, _, ddl)| ddl).collect())
2053}
2054
2055pub fn function_ddl_with_names(
2057 schema: &SchemaDescriptor,
2058) -> Result<Vec<(String, String, String)>, crate::error::PyQLError> {
2059 schema
2060 .functions
2061 .iter()
2062 .map(|fd| emit_one_function(fd, schema).map(|ddl| (fd.module.clone(), fd.name.clone(), ddl)))
2063 .collect()
2064}
2065
2066pub fn scalar_function_ddl_with_names(
2068 schema: &SchemaDescriptor,
2069) -> Result<Vec<(String, String, String)>, crate::error::PyQLError> {
2070 schema
2071 .functions
2072 .iter()
2073 .filter(|fd| !fd.return_is_object)
2074 .map(|fd| emit_one_function(fd, schema).map(|ddl| (fd.module.clone(), fd.name.clone(), ddl)))
2075 .collect()
2076}
2077
2078pub fn object_function_ddl_with_names(
2080 schema: &SchemaDescriptor,
2081) -> Result<Vec<(String, String, String)>, crate::error::PyQLError> {
2082 schema
2083 .functions
2084 .iter()
2085 .filter(|fd| fd.return_is_object)
2086 .map(|fd| emit_one_function(fd, schema).map(|ddl| (fd.module.clone(), fd.name.clone(), ddl)))
2087 .collect()
2088}
2089
2090fn emit_one_interface_view(t: &TypeDescriptor, impls: &[&TypeDescriptor], out: &mut String) {
2093 let cols: Vec<String> = t
2094 .properties
2095 .iter()
2096 .map(|p| qi(&p.name))
2097 .chain(
2098 t.links
2099 .iter()
2100 .filter(|l| !l.is_junction_backed())
2101 .map(|l| qi(&format!("{}_id", l.name))),
2102 )
2103 .collect();
2104 let col_list = cols.join(", ");
2105 let selects: Vec<String> = impls
2106 .iter()
2107 .map(|impl_t| format!(" SELECT {} FROM {}", col_list, qn(&impl_t.module, &impl_t.table)))
2108 .collect();
2109 out.push_str(&format!("CREATE VIEW {} AS\n", qn(&t.module, &t.table)));
2110 out.push_str(&selects.join("\n UNION ALL\n"));
2111 out.push_str(";\n\n");
2112}
2113
2114fn emit_interface_views(schema: &SchemaDescriptor, out: &mut String) {
2115 let implementors = interface_implementors(schema);
2116 for t in &schema.types {
2117 if !(t.abstract_ && t.materialized) {
2118 continue;
2119 }
2120 let key = format!("{}::{}", t.module, t.name);
2121 let Some(impls) = implementors.get(&key) else { continue };
2122 if impls.is_empty() {
2123 continue;
2124 }
2125 emit_one_interface_view(t, impls, out);
2126 }
2127}
2128
2129pub struct ExclTriggerInfo {
2134 pub fn_module: String,
2135 pub fn_name: String,
2136 pub fn_ddl: String,
2137 pub impl_module: String,
2138 pub impl_table: String,
2139 pub ins_trigger_name: String,
2140 pub ins_ddl: String,
2141 pub upd_trigger_name: String,
2142 pub upd_ddl: String,
2143}
2144
2145fn excl_fn_name(iface_table: &str, fields: &[String]) -> String {
2146 format!("_excl_{}_{}", iface_table, fields.join("_"))
2147}
2148
2149fn make_excl_info(
2150 iface: &TypeDescriptor,
2151 fields: &[String],
2152 columns: &[String],
2153 impl_t: &TypeDescriptor,
2154 impls: &[&TypeDescriptor],
2155) -> ExclTriggerInfo {
2156 let fn_name = excl_fn_name(&iface.table, fields);
2157 let fn_qname = format!("{}.{}", pg_schema(&iface.module), qi(&fn_name));
2158 let view_qname = if iface.abstract_ {
2160 qn(&iface.module, &iface.table)
2161 } else {
2162 let spanned_columns = std::iter::once("id".to_string())
2163 .chain(columns.iter().cloned())
2164 .map(|c| qi(&c))
2165 .collect::<Vec<_>>()
2166 .join(", ");
2167 let branches = impls
2168 .iter()
2169 .map(|t| format!("SELECT {spanned_columns} FROM {}", qn(&t.module, &t.table)))
2170 .collect::<Vec<_>>()
2171 .join(" UNION ALL ");
2172 format!("({branches}) AS \"_spanned\"")
2173 };
2174 let tbl_qname = qn(&impl_t.module, &impl_t.table);
2175
2176 let field_conds: Vec<String> = columns.iter().map(|c| format!("{} = NEW.{}", qi(c), qi(c))).collect();
2177 let where_clause = format!("{} AND \"id\" <> NEW.\"id\"", field_conds.join(" AND "));
2178
2179 let detail_keys = columns.join(", ");
2180 let detail_vals = columns
2181 .iter()
2182 .map(|c| format!("NEW.{}::text", qi(c)))
2183 .collect::<Vec<_>>()
2184 .join(" || ', ' || ");
2185
2186 let fn_ddl = format!(
2187 "CREATE OR REPLACE FUNCTION {}()\n\
2188 RETURNS trigger LANGUAGE plpgsql AS $$\n\
2189 BEGIN\n\
2190 IF EXISTS (\n\
2191 SELECT 1 FROM {}\n\
2192 WHERE {}\n\
2193 ) THEN\n\
2194 RAISE unique_violation\n\
2195 USING CONSTRAINT = '{}',\n\
2196 DETAIL = format('Key ({})=(%s) already exists.', {});\n\
2197 END IF;\n\
2198 RETURN NEW;\n\
2199 END;\n\
2200 $$;",
2201 fn_qname, view_qname, where_clause, fn_name, detail_keys, detail_vals,
2202 );
2203
2204 let ins_trigger_name = format!("{}_ins", fn_name);
2205 let upd_trigger_name = format!("{}_upd", fn_name);
2206 let of_cols = columns.iter().map(|c| qi(c)).collect::<Vec<_>>().join(", ");
2207 let when_clause = columns
2208 .iter()
2209 .map(|c| format!("OLD.{} IS DISTINCT FROM NEW.{}", qi(c), qi(c)))
2210 .collect::<Vec<_>>()
2211 .join(" OR ");
2212
2213 let ins_ddl = format!(
2214 "CREATE CONSTRAINT TRIGGER {}\n\
2215 AFTER INSERT ON {}\n\
2216 DEFERRABLE INITIALLY DEFERRED\n\
2217 FOR EACH ROW EXECUTE FUNCTION {}();",
2218 qi(&ins_trigger_name),
2219 tbl_qname,
2220 fn_qname,
2221 );
2222 let upd_ddl = format!(
2223 "CREATE CONSTRAINT TRIGGER {}\n\
2224 AFTER UPDATE OF {} ON {}\n\
2225 DEFERRABLE INITIALLY DEFERRED\n\
2226 FOR EACH ROW WHEN ({})\n\
2227 EXECUTE FUNCTION {}();",
2228 qi(&upd_trigger_name),
2229 of_cols,
2230 tbl_qname,
2231 when_clause,
2232 fn_qname,
2233 );
2234
2235 ExclTriggerInfo {
2236 fn_module: iface.module.clone(),
2237 fn_name,
2238 fn_ddl,
2239 impl_module: impl_t.module.clone(),
2240 impl_table: impl_t.table.clone(),
2241 ins_trigger_name,
2242 ins_ddl,
2243 upd_trigger_name,
2244 upd_ddl,
2245 }
2246}
2247
2248fn make_excl_junction_info(
2256 iface: &TypeDescriptor,
2257 link_name: &str,
2258 impl_t: &TypeDescriptor,
2259 impls: &[&TypeDescriptor],
2260) -> ExclTriggerInfo {
2261 let fn_name = excl_fn_name(&iface.table, std::slice::from_ref(&link_name.to_string()));
2262 let fn_qname = format!("{}.{}", pg_schema(&iface.module), qi(&fn_name));
2263 let view_qname = if iface.abstract_ {
2266 qn(&iface.module, &interface_junction_view_name(&iface.table, link_name))
2267 } else {
2268 let branches = impls
2269 .iter()
2270 .map(|t| {
2271 format!(
2272 "SELECT \"source\", \"target\" FROM {}",
2273 qn(&t.module, &format!("{}.{}", t.table, link_name))
2274 )
2275 })
2276 .collect::<Vec<_>>()
2277 .join(" UNION ALL ");
2278 format!("({branches}) AS \"_spanned\"")
2279 };
2280 let jt_name = format!("{}.{}", impl_t.table, link_name);
2281 let jt_qname = qn(&impl_t.module, &jt_name);
2282
2283 let where_clause = "\"target\" = NEW.\"target\" AND \"source\" <> NEW.\"source\"";
2284
2285 let fn_ddl = format!(
2286 "CREATE OR REPLACE FUNCTION {}()\n\
2287 RETURNS trigger LANGUAGE plpgsql AS $$\n\
2288 BEGIN\n\
2289 IF EXISTS (\n\
2290 SELECT 1 FROM {}\n\
2291 WHERE {}\n\
2292 ) THEN\n\
2293 RAISE unique_violation\n\
2294 USING CONSTRAINT = '{}',\n\
2295 DETAIL = format('Key (target)=(%s) already exists.', NEW.\"target\"::text);\n\
2296 END IF;\n\
2297 RETURN NEW;\n\
2298 END;\n\
2299 $$;",
2300 fn_qname, view_qname, where_clause, fn_name,
2301 );
2302
2303 let ins_trigger_name = format!("{}_ins", fn_name);
2304 let upd_trigger_name = format!("{}_upd", fn_name);
2305
2306 let ins_ddl = format!(
2307 "CREATE CONSTRAINT TRIGGER {}\n\
2308 AFTER INSERT ON {}\n\
2309 DEFERRABLE INITIALLY DEFERRED\n\
2310 FOR EACH ROW EXECUTE FUNCTION {}();",
2311 qi(&ins_trigger_name),
2312 jt_qname,
2313 fn_qname,
2314 );
2315 let upd_ddl = format!(
2316 "CREATE CONSTRAINT TRIGGER {}\n\
2317 AFTER UPDATE OF \"target\" ON {}\n\
2318 DEFERRABLE INITIALLY DEFERRED\n\
2319 FOR EACH ROW WHEN (OLD.\"target\" IS DISTINCT FROM NEW.\"target\")\n\
2320 EXECUTE FUNCTION {}();",
2321 qi(&upd_trigger_name),
2322 jt_qname,
2323 fn_qname,
2324 );
2325
2326 ExclTriggerInfo {
2327 fn_module: iface.module.clone(),
2328 fn_name,
2329 fn_ddl,
2330 impl_module: impl_t.module.clone(),
2331 impl_table: jt_name,
2332 ins_trigger_name,
2333 ins_ddl,
2334 upd_trigger_name,
2335 upd_ddl,
2336 }
2337}
2338
2339pub fn interface_exclusive_trigger_infos(schema: &SchemaDescriptor) -> Vec<ExclTriggerInfo> {
2345 let implementors = interface_implementors(schema);
2346
2347 let mut result = Vec::new();
2348 for t in &schema.types {
2349 let key = format!("{}::{}", t.module, t.name);
2350 let spans_subtypes = !t.abstract_ && schema.types.iter().any(|sub| sub.bases.contains(&key));
2351 if !(spans_subtypes || t.abstract_ && t.materialized) {
2352 continue;
2353 }
2354 let Some(impls) = implementors.get(&key) else { continue };
2355 if impls.is_empty() {
2356 continue;
2357 }
2358 let interfaces: Vec<&TypeDescriptor> = if t.abstract_ {
2361 vec![]
2362 } else {
2363 schema
2364 .types
2365 .iter()
2366 .filter(|i| {
2367 i.abstract_ && i.materialized && t.interfaces.contains(&format!("{}::{}", i.module, i.name))
2368 })
2369 .collect()
2370 };
2371 let declared_by_interface = |name: &str| {
2372 interfaces.iter().any(|i| {
2373 i.properties.iter().any(|p| p.name == name && p.is_exclusive)
2374 || i.links.iter().any(|l| l.name == name && l.is_exclusive)
2375 || i.multilinks.iter().any(|ml| ml.name == name && ml.is_exclusive)
2376 })
2377 };
2378
2379 for p in &t.properties {
2380 if !p.is_exclusive || p.is_pk || declared_by_interface(&p.name) {
2381 continue;
2382 }
2383 let fields = vec![p.name.clone()];
2384 for impl_t in impls {
2385 result.push(make_excl_info(t, &fields, &fields, impl_t, impls));
2386 }
2387 }
2388 for l in &t.links {
2389 if !l.is_exclusive || declared_by_interface(&l.name) {
2390 continue;
2391 }
2392 if l.is_junction_backed() {
2401 for impl_t in impls {
2402 result.push(make_excl_junction_info(t, &l.name, impl_t, impls));
2403 }
2404 continue;
2405 }
2406 let fields = vec![format!("{}_id", l.name)];
2407 for impl_t in impls {
2408 result.push(make_excl_info(t, &fields, &fields, impl_t, impls));
2409 }
2410 }
2411 for ml in &t.multilinks {
2412 if !ml.is_exclusive || declared_by_interface(&ml.name) {
2413 continue;
2414 }
2415 for impl_t in impls {
2416 result.push(make_excl_junction_info(t, &ml.name, impl_t, impls));
2417 }
2418 }
2419 for c in &t.constraints {
2420 if let TypeConstraint::Exclusive { pointers: fields, .. } = c {
2421 let from_interface = interfaces.iter().any(|i| {
2422 i.constraints
2423 .iter()
2424 .any(|ic| matches!(ic, TypeConstraint::Exclusive { pointers, .. } if pointers == fields))
2425 });
2426 if from_interface {
2427 continue;
2428 }
2429 let columns: Vec<String> = fields.iter().map(|f| constraint_column(t, f)).collect();
2430 for impl_t in impls {
2431 result.push(make_excl_info(t, fields, &columns, impl_t, impls));
2432 }
2433 }
2434 }
2435 }
2436 result
2437}
2438
2439fn emit_interface_exclusive_triggers(schema: &SchemaDescriptor, out: &mut String) {
2440 use std::collections::HashSet;
2441 let mut fn_emitted: HashSet<String> = HashSet::new();
2442 for info in interface_exclusive_trigger_infos(schema) {
2443 if fn_emitted.insert(info.fn_name.clone()) {
2444 out.push_str(&info.fn_ddl);
2445 out.push_str("\n\n");
2446 }
2447 out.push_str(&info.ins_ddl);
2448 out.push('\n');
2449 out.push_str(&info.upd_ddl);
2450 out.push_str("\n\n");
2451 }
2452}
2453
2454fn emit_scalar_functions(schema: &SchemaDescriptor, out: &mut String) -> Result<(), PyQLError> {
2457 for fd in schema.functions.iter().filter(|fd| !fd.return_is_object) {
2458 let ddl = emit_one_function(fd, schema)?;
2459 out.push_str(&ddl);
2460 out.push('\n');
2461 }
2462 Ok(())
2463}
2464
2465fn emit_object_functions(schema: &SchemaDescriptor, out: &mut String) -> Result<(), PyQLError> {
2466 for fd in schema.functions.iter().filter(|fd| fd.return_is_object) {
2467 let ddl = emit_one_function(fd, schema)?;
2468 out.push_str(&ddl);
2469 out.push('\n');
2470 }
2471 Ok(())
2472}
2473
2474fn emit_one_function(fd: &FunctionDescriptor, schema: &SchemaDescriptor) -> Result<String, PyQLError> {
2475 use crate::ir::compile_fn_body;
2476 use crate::sql::emit_fn_body;
2477
2478 let ir_output = compile_fn_body(fd, schema).map_err(|e| {
2480 let msg = format!("error in function '{}::{}' body: {}", fd.module, fd.name, e);
2481 PyQLError::Fragment(PyQLFragmentError {
2482 message: msg,
2483 context: format!("{}::{}", fd.module, fd.name),
2484 position: crate::error::Position { line: 0, col: 0 },
2485 })
2486 })?;
2487
2488 let mut body_sql = emit_fn_body(&ir_output);
2490 if fd.return_is_object
2496 && !matches!(
2497 ir_output.stmt,
2498 crate::ir::IrStmt::Insert(_) | crate::ir::IrStmt::Update(_) | crate::ir::IrStmt::Delete(_)
2499 )
2500 && let Some(columns) = fn_return_column_names(fd, schema)
2501 {
2502 body_sql = format!(
2503 "SELECT {} FROM (\n {}\n ) AS \"_returned\"",
2504 columns.join(", "),
2505 body_sql
2506 );
2507 }
2508
2509 let mut param_parts: Vec<String> = Vec::with_capacity(fd.params.len() + 1);
2518 if ir_output.uses_globals_arg {
2519 param_parts.push(format!("{} jsonb", qi(crate::ir::GLOBALS_ARG)));
2520 }
2521 param_parts.extend(fd.params.iter().map(|p| format!("{} {}", qi(&p.name), p.pg_type)));
2522 let params_sql = param_parts.join(", ");
2523
2524 let returns_sql = if fd.return_is_object {
2526 let type_columns = emit_fn_return_table(fd, schema);
2528 if fd.return_is_set {
2529 format!("TABLE({})", type_columns)
2530 } else {
2531 format!("TABLE({})", type_columns)
2534 }
2535 } else if fd.return_is_set {
2536 format!("SETOF {}", fd.return_pg_type)
2537 } else {
2538 fd.return_pg_type.clone()
2539 };
2540
2541 let volatility_kw = match fd.volatility.as_str() {
2542 "immutable" => "IMMUTABLE",
2543 "stable" => "STABLE",
2544 _ => "VOLATILE",
2545 };
2546
2547 Ok(format!(
2548 "CREATE OR REPLACE FUNCTION {fn_name}({params})\nRETURNS {returns}\nLANGUAGE SQL {vol}\nAS $$\n {body}\n$$;\n",
2549 fn_name = qn(&fd.module, &fd.name),
2550 params = params_sql,
2551 returns = returns_sql,
2552 vol = volatility_kw,
2553 body = body_sql,
2554 ))
2555}
2556
2557fn fn_return_column_names(fd: &FunctionDescriptor, schema: &SchemaDescriptor) -> Option<Vec<String>> {
2559 let td = schema
2560 .types
2561 .iter()
2562 .find(|t| format!("{}::{}", t.module, t.name) == fd.return_pg_type)?;
2563 let mut names: Vec<String> = Vec::new();
2564 if fd.return_is_polymorphic {
2565 names.push(qi("__type__"));
2566 }
2567 names.extend(td.properties.iter().map(|p| qi(&p.name)));
2568 names.extend(
2569 td.links
2570 .iter()
2571 .filter(|l| !l.is_junction_backed())
2572 .map(|l| qi(&format!("{}_id", l.name))),
2573 );
2574 Some(names)
2575}
2576
2577fn emit_fn_return_table(fd: &FunctionDescriptor, schema: &SchemaDescriptor) -> String {
2578 let type_name = &fd.return_pg_type; let td = schema
2580 .types
2581 .iter()
2582 .find(|t| format!("{}::{}", t.module, t.name) == *type_name);
2583 let Some(td) = td else {
2584 return "__type__ text, id uuid".to_string();
2585 };
2586
2587 let mut cols: Vec<String> = Vec::new();
2588 if fd.return_is_polymorphic {
2589 cols.push("__type__ text".to_string());
2590 }
2591 for p in &td.properties {
2592 cols.push(format!("{} {}", qi(&p.name), resolved_pg_type(&p.pg_type)));
2593 }
2594 for l in &td.links {
2595 if l.is_junction_backed() {
2596 continue;
2597 }
2598 cols.push(format!("{} uuid", qi(&format!("{}_id", l.name))));
2599 }
2600 cols.join(", ")
2601}
2602
2603fn emit_vector_columns(schema: &SchemaDescriptor, out: &mut String) {
2606 for td in &schema.types {
2607 if td.abstract_ || td.vector_indexes.is_empty() {
2608 continue;
2609 }
2610 for vi in &td.vector_indexes {
2611 out.push_str(&format!(
2612 "ALTER TABLE {} ADD COLUMN IF NOT EXISTS {} vector({});\n",
2613 qn(&td.module, &td.table),
2614 qi(&vi.column_name()),
2615 vi.dimensions,
2616 ));
2617 }
2618 }
2619 if schema
2620 .types
2621 .iter()
2622 .any(|t| !t.abstract_ && !t.vector_indexes.is_empty())
2623 {
2624 out.push('\n');
2625 }
2626}
2627
2628fn emit_vector_indexes(schema: &SchemaDescriptor, out: &mut String) {
2631 for td in &schema.types {
2632 if td.abstract_ || td.vector_indexes.is_empty() {
2633 continue;
2634 }
2635 for vi in &td.vector_indexes {
2636 let index_name = match &vi.index_name {
2637 None => format!("{}__vector__", td.table),
2638 Some(name) => format!("{}__vector_{}__", td.table, name),
2639 };
2640 out.push_str(&format!(
2641 "CREATE INDEX IF NOT EXISTS {} ON {} USING hnsw ({} {});\n",
2642 qi(&index_name),
2643 qn(&td.module, &td.table),
2644 qi(&vi.column_name()),
2645 vi.ops_class(),
2646 ));
2647 }
2648 }
2649}
2650
2651fn emit_search_columns(schema: &SchemaDescriptor, out: &mut String) {
2654 use crate::schema::SearchBackend;
2655
2656 let mut emitted = false;
2657 for td in &schema.types {
2658 if td.abstract_ {
2659 continue;
2660 }
2661 for si in &td.search_indexes {
2662 if si.backend != SearchBackend::Postgres {
2663 continue;
2664 }
2665
2666 let parts: Vec<String> = si
2668 .pointers
2669 .iter()
2670 .map(|sf| {
2671 let col = qi(&sf.name);
2672 let w = sf.weight.as_str();
2673 format!("setweight(to_tsvector('english', coalesce({col}, '')), '{w}')")
2674 })
2675 .collect();
2676
2677 let expr = if parts.len() == 1 {
2678 parts.into_iter().next().unwrap()
2679 } else {
2680 parts.join(" || ")
2681 };
2682
2683 out.push_str(&format!(
2684 "ALTER TABLE {} ADD COLUMN IF NOT EXISTS {} tsvector GENERATED ALWAYS AS ({}) STORED;\n",
2685 qn(&td.module, &td.table),
2686 qi(&si.column_name()),
2687 expr,
2688 ));
2689 emitted = true;
2690 }
2691 }
2692 if emitted {
2693 out.push('\n');
2694 }
2695}
2696
2697fn emit_search_indexes(schema: &SchemaDescriptor, out: &mut String) {
2700 use crate::schema::SearchBackend;
2701
2702 for td in &schema.types {
2703 if td.abstract_ {
2704 continue;
2705 }
2706 for si in &td.search_indexes {
2707 if si.backend != SearchBackend::Postgres {
2708 continue;
2709 }
2710
2711 let col = si.column_name();
2712 let index_name = match &si.index_name {
2713 None => format!("{}__search__", td.table),
2714 Some(name) => format!("{}__search_{}__", td.table, name),
2715 };
2716 out.push_str(&format!(
2717 "CREATE INDEX IF NOT EXISTS {} ON {} USING gin ({});\n",
2718 qi(&index_name),
2719 qn(&td.module, &td.table),
2720 qi(&col),
2721 ));
2722 }
2723 }
2724}
2725
2726pub fn compile_index_fetch(
2740 type_name: &str,
2741 index_name: Option<&str>,
2742 schema: &SchemaDescriptor,
2743) -> Result<String, PyQLError> {
2744 let td = schema
2745 .types
2746 .iter()
2747 .find(|t| format!("{}::{}", t.module, t.name) == type_name)
2748 .ok_or_else(|| {
2749 PyQLError::Fragment(PyQLFragmentError {
2750 message: format!("compile_index_fetch: unknown type '{}'", type_name),
2751 context: type_name.to_string(),
2752 position: crate::error::Position { line: 0, col: 0 },
2753 })
2754 })?;
2755
2756 let vi = td
2757 .vector_indexes
2758 .iter()
2759 .find(|vi| vi.index_name.as_deref() == index_name)
2760 .ok_or_else(|| {
2761 let key = index_name.unwrap_or("<default>");
2762 PyQLError::Fragment(PyQLFragmentError {
2763 message: format!("compile_index_fetch: no vector index '{}' on type '{}'", key, type_name),
2764 context: type_name.to_string(),
2765 position: crate::error::Position { line: 0, col: 0 },
2766 })
2767 })?;
2768
2769 let field_exprs = vi
2770 .pointers
2771 .iter()
2772 .map(|f| {
2773 let pg_type = td
2775 .properties
2776 .iter()
2777 .find(|p| p.name == *f)
2778 .map(|p| p.pg_type.as_str())
2779 .unwrap_or("text");
2780 let col = qi(f);
2781 if pg_type == "text" {
2782 col
2783 } else {
2784 format!("{}::text", col)
2785 }
2786 })
2787 .collect::<Vec<_>>();
2788
2789 let concat = if field_exprs.len() == 1 {
2790 field_exprs.into_iter().next().unwrap()
2791 } else {
2792 format!("concat_ws(E'\\n', {})", field_exprs.join(", "))
2793 };
2794
2795 Ok(format!(
2796 "SELECT \"id\", {} AS source_text\nFROM {}\nWHERE \"id\" = ANY($1::uuid[])",
2797 concat,
2798 qn(&td.module, &td.table),
2799 ))
2800}
2801
2802pub fn compile_search_index_fetch(
2805 type_name: &str,
2806 index_name: Option<&str>,
2807 schema: &SchemaDescriptor,
2808) -> Result<String, PyQLError> {
2809 let td = schema
2810 .types
2811 .iter()
2812 .find(|t| format!("{}::{}", t.module, t.name) == type_name)
2813 .ok_or_else(|| {
2814 PyQLError::Fragment(PyQLFragmentError {
2815 message: format!("compile_search_index_fetch: unknown type '{}'", type_name),
2816 context: type_name.to_string(),
2817 position: crate::error::Position { line: 0, col: 0 },
2818 })
2819 })?;
2820
2821 use crate::schema::SearchBackend;
2822 let si = td
2823 .search_indexes
2824 .iter()
2825 .find(|si| si.index_name.as_deref() == index_name && si.backend != SearchBackend::Postgres)
2826 .ok_or_else(|| {
2827 let key = index_name.unwrap_or("<default>");
2828 PyQLError::Fragment(PyQLFragmentError {
2829 message: format!(
2830 "compile_search_index_fetch: no remote SearchIndex '{}' on type '{}'",
2831 key, type_name
2832 ),
2833 context: type_name.to_string(),
2834 position: crate::error::Position { line: 0, col: 0 },
2835 })
2836 })?;
2837
2838 let field_exprs = si
2839 .pointers
2840 .iter()
2841 .map(|sf| {
2842 let pg_type = td
2843 .properties
2844 .iter()
2845 .find(|p| p.name == sf.name)
2846 .map(|p| p.pg_type.as_str())
2847 .unwrap_or("text");
2848 let col = qi(&sf.name);
2849 if pg_type == "text" {
2850 col
2851 } else {
2852 format!("{}::text", col)
2853 }
2854 })
2855 .collect::<Vec<_>>();
2856
2857 let concat = if field_exprs.len() == 1 {
2858 field_exprs.into_iter().next().unwrap()
2859 } else {
2860 format!("concat_ws(E'\\n', {})", field_exprs.join(", "))
2861 };
2862
2863 Ok(format!(
2864 "SELECT \"id\", {} AS source_text\nFROM {}\nWHERE \"id\" = ANY($1::uuid[])",
2865 concat,
2866 qn(&td.module, &td.table),
2867 ))
2868}
2869
2870#[cfg(test)]
2871mod tests {
2872 use super::*;
2873 use crate::schema::{
2874 DeleteAction, DeleteSide, FunctionDescriptor, FunctionParamDescriptor, LinkDescriptor, MultiLinkDescriptor,
2875 OnDeletePolicy, PropertyDescriptor, SchemaDescriptor, TypeDescriptor,
2876 };
2877
2878 fn person_type() -> TypeDescriptor {
2879 TypeDescriptor {
2880 name: "Person".into(),
2881 module: "default".into(),
2882 table: "Person".into(),
2883 abstract_: false,
2884 materialized: false,
2885 description: None,
2886 parents: vec![],
2887 interfaces: vec![],
2888 bases: vec![],
2889 properties: vec![
2890 PropertyDescriptor {
2891 name: "id".into(),
2892 pg_type: "uuid".into(),
2893 nullable: false,
2894 default_sql: Some("uuidv7()".into()),
2895 default_pyql: None,
2896 description: None,
2897 check_constraints: vec![],
2898 is_exclusive: true,
2899 is_pk: true,
2900 is_readonly: true,
2901 rewrites: vec![],
2902 tuple_members: None,
2903 column_type: None,
2904 },
2905 PropertyDescriptor {
2906 name: "age".into(),
2907 pg_type: "int8".into(),
2908 nullable: true,
2909 default_sql: None,
2910 default_pyql: None,
2911 description: None,
2912 check_constraints: vec![],
2913 is_exclusive: false,
2914 is_pk: false,
2915 is_readonly: false,
2916 rewrites: vec![],
2917 tuple_members: None,
2918 column_type: None,
2919 },
2920 ],
2921 links: vec![],
2922 multilinks: vec![],
2923 computed: vec![],
2924 constraints: vec![],
2925 indexes: vec![],
2926 partition: None,
2927 vector_indexes: vec![],
2928 search_indexes: vec![],
2929 triggers: vec![],
2930 junction: false,
2931 signals: vec![],
2932 }
2933 }
2934
2935 fn minimal_schema(fns: Vec<FunctionDescriptor>) -> SchemaDescriptor {
2936 SchemaDescriptor {
2937 types: vec![person_type()],
2938 scalars: vec![],
2939 enums: vec![],
2940 named_tuples: vec![],
2941 globals: vec![],
2942 functions: fns,
2943 aliases: vec![],
2944 channels: vec![],
2945 ..Default::default()
2946 }
2947 }
2948
2949 fn interface_schema(referencing: Vec<TypeDescriptor>) -> SchemaDescriptor {
2954 fn bare(name: &str, interfaces: Vec<String>, abstract_: bool, materialized: bool) -> TypeDescriptor {
2955 TypeDescriptor {
2956 name: name.into(),
2957 module: "default".into(),
2958 table: name.into(),
2959 abstract_,
2960 materialized,
2961 description: None,
2962 parents: vec![],
2963 interfaces,
2964 bases: vec![],
2965 properties: vec![],
2966 links: vec![],
2967 multilinks: vec![],
2968 computed: vec![],
2969 constraints: vec![],
2970 indexes: vec![],
2971 partition: None,
2972 vector_indexes: vec![],
2973 search_indexes: vec![],
2974 triggers: vec![],
2975 junction: false,
2976 signals: vec![],
2977 }
2978 }
2979 let mut types = vec![
2980 bare("Account", vec![], true, true),
2981 bare("Individual", vec!["default::Account".into()], false, false),
2982 bare("Organization", vec!["default::Account".into()], false, false),
2983 ];
2984 types.extend(referencing);
2985 SchemaDescriptor {
2986 types,
2987 scalars: vec![],
2988 enums: vec![],
2989 named_tuples: vec![],
2990 globals: vec![],
2991 functions: vec![],
2992 aliases: vec![],
2993 channels: vec![],
2994 ..Default::default()
2995 }
2996 }
2997
2998 fn link_to_account(name: &str, on_delete: Vec<OnDeletePolicy>, through: Option<String>) -> LinkDescriptor {
2999 LinkDescriptor {
3000 name: name.into(),
3001 target: "default::Account".into(),
3002 nullable: true,
3003 description: None,
3004 default_pyql: None,
3005 is_exclusive: false,
3006 is_readonly: false,
3007 rewrites: vec![],
3008 on_delete,
3009 through,
3010 }
3011 }
3012
3013 fn referencing_type(
3014 name: &str,
3015 links: Vec<LinkDescriptor>,
3016 multilinks: Vec<MultiLinkDescriptor>,
3017 ) -> TypeDescriptor {
3018 TypeDescriptor {
3019 name: name.into(),
3020 module: "default".into(),
3021 table: name.into(),
3022 abstract_: false,
3023 materialized: false,
3024 description: None,
3025 parents: vec![],
3026 interfaces: vec![],
3027 bases: vec![],
3028 properties: vec![],
3029 links,
3030 multilinks,
3031 computed: vec![],
3032 constraints: vec![],
3033 indexes: vec![],
3034 partition: None,
3035 vector_indexes: vec![],
3036 search_indexes: vec![],
3037 triggers: vec![],
3038 junction: false,
3039 signals: vec![],
3040 }
3041 }
3042
3043 #[test]
3044 fn test_multilink_on_an_interface_gets_a_union_view_over_the_implementors() {
3045 let mut schema = interface_schema(vec![referencing_type("Email", vec![], vec![])]);
3051 for t in schema.types.iter_mut() {
3052 if t.name == "Account" || t.name == "Individual" || t.name == "Organization" {
3053 t.multilinks.push(MultiLinkDescriptor {
3054 name: "emails".into(),
3055 target: "default::Email".into(),
3056 through: None,
3057 nullable: true,
3058 description: None,
3059 default_pyql: None,
3060 on_delete: vec![],
3061 is_exclusive: false,
3062 });
3063 }
3064 }
3065 let ddl = export_schema(&schema).unwrap();
3066 assert!(
3067 ddl.contains("CREATE VIEW \"public\".\"Account.emails\""),
3068 "no union view for the interface's multi-link, got:\n{}",
3069 ddl
3070 );
3071 for implementor in ["Individual", "Organization"] {
3072 assert!(
3073 ddl.contains(&format!("FROM \"public\".\"{}.emails\"", implementor)),
3074 "union view misses {}, got:\n{}",
3075 implementor,
3076 ddl
3077 );
3078 }
3079 }
3080
3081 fn schema_with_a_link_to_a_type_with_subtypes() -> SchemaDescriptor {
3084 let mut link = link_to_account("holder", vec![], None);
3085 link.target = "default::Individual".into();
3086 let mut schema = interface_schema(vec![referencing_type("Session", vec![link], vec![])]);
3087 let mut staff = schema.types[1].clone();
3088 staff.name = "Staff".into();
3089 staff.table = "Staff".into();
3090 staff.bases = vec!["default::Individual".into()];
3091 schema.types.push(staff);
3092 schema
3093 }
3094
3095 #[test]
3096 fn an_exclusive_multilink_on_a_type_with_subtypes_is_unique_across_them() {
3097 let mut schema = schema_with_a_link_to_a_type_with_subtypes();
3098 for t in schema
3099 .types
3100 .iter_mut()
3101 .filter(|t| t.name == "Individual" || t.name == "Staff")
3102 {
3103 t.multilinks.push(MultiLinkDescriptor {
3104 name: "keys".into(),
3105 target: "default::Session".into(),
3106 through: None,
3107 nullable: false,
3108 description: None,
3109 default_pyql: None,
3110 on_delete: vec![],
3111 is_exclusive: true,
3112 });
3113 }
3114 let ddl = export_schema(&schema).unwrap();
3115 for table in ["Individual.keys", "Staff.keys"] {
3116 assert!(
3117 ddl.contains(&format!("CREATE TABLE \"public\".\"{table}\"")),
3118 "got:\n{ddl}"
3119 );
3120 assert!(
3121 ddl.contains(&format!("AFTER INSERT ON \"public\".\"{table}\"")),
3122 "missing exclusive trigger on {table}, got:\n{ddl}"
3123 );
3124 }
3125 assert!(ddl.contains("UNIQUE (target)"), "got:\n{ddl}");
3126 assert!(
3127 ddl.contains(
3128 "SELECT \"source\", \"target\" FROM \"public\".\"Individual.keys\" UNION ALL \
3129 SELECT \"source\", \"target\" FROM \"public\".\"Staff.keys\""
3130 ),
3131 "the trigger must check every subtype's junction, got:\n{ddl}"
3132 );
3133 }
3134
3135 #[test]
3136 fn a_link_to_a_type_with_subtypes_emits_no_foreign_key() {
3137 let ddl = export_schema(&schema_with_a_link_to_a_type_with_subtypes()).unwrap();
3138 assert!(!ddl.contains("Session_holder_fkey"), "got:\n{}", ddl);
3139 for table in ["Individual", "Staff"] {
3140 assert!(
3141 ddl.contains(&format!("BEFORE DELETE ON \"public\".\"{table}\"")),
3142 "missing enforcement trigger on {table}, got:\n{ddl}"
3143 );
3144 }
3145 }
3146
3147 #[test]
3148 fn test_link_to_interface_emits_no_foreign_key() {
3149 let schema = interface_schema(vec![referencing_type(
3152 "Session",
3153 vec![link_to_account("account", vec![], None)],
3154 vec![],
3155 )]);
3156 let ddl = export_schema(&schema).unwrap();
3157 assert!(
3158 !ddl.contains("Session_account_fkey"),
3159 "a link targeting an interface must not get a FK, got:\n{}",
3160 ddl
3161 );
3162 assert!(ddl.contains("CREATE VIEW \"public\".\"Account\""), "got:\n{}", ddl);
3163 }
3164
3165 #[test]
3166 fn test_link_to_interface_enforces_restrict_on_every_implementor() {
3167 let schema = interface_schema(vec![referencing_type(
3168 "Session",
3169 vec![link_to_account("account", vec![], None)],
3170 vec![],
3171 )]);
3172 let ddl = export_schema(&schema).unwrap();
3173 for implementor in ["Individual", "Organization"] {
3174 assert!(
3175 ddl.contains(&format!("BEFORE DELETE ON \"public\".\"{}\"", implementor)),
3176 "missing enforcement trigger on {}, got:\n{}",
3177 implementor,
3178 ddl
3179 );
3180 }
3181 assert!(ddl.contains("RAISE foreign_key_violation"), "got:\n{}", ddl);
3182 assert!(
3183 ddl.contains("SELECT 1 FROM \"public\".\"Session\" WHERE \"account_id\" = OLD.id"),
3184 "got:\n{}",
3185 ddl
3186 );
3187 }
3188
3189 #[test]
3190 fn test_interface_link_allow_nulls_the_column_but_clears_a_junction_row() {
3191 let allow = vec![OnDeletePolicy {
3196 side: DeleteSide::Target,
3197 action: DeleteAction::Allow,
3198 }];
3199 let schema = interface_schema(vec![
3200 referencing_type(
3201 "AuditEntry",
3202 vec![link_to_account("actor", allow.clone(), None)],
3203 vec![],
3204 ),
3205 referencing_type(
3206 "Watchlist",
3207 vec![],
3208 vec![MultiLinkDescriptor {
3209 name: "watched".into(),
3210 target: "default::Account".into(),
3211 through: None,
3212 nullable: true,
3213 description: None,
3214 default_pyql: None,
3215 on_delete: allow,
3216 is_exclusive: false,
3217 }],
3218 ),
3219 ]);
3220 let ddl = export_schema(&schema).unwrap();
3221 assert!(
3222 ddl.contains("UPDATE \"public\".\"AuditEntry\" SET \"actor_id\" = NULL WHERE \"actor_id\" = OLD.id;"),
3223 "single link Allow must null the column, got:\n{}",
3224 ddl
3225 );
3226 assert!(
3227 ddl.contains("DELETE FROM \"public\".\"Watchlist.watched\" WHERE \"target\" = OLD.id;"),
3228 "multi-link Allow must drop the junction row, got:\n{}",
3229 ddl
3230 );
3231 }
3232
3233 #[test]
3234 fn test_multilink_to_interface_junction_target_has_no_reference() {
3235 let schema = interface_schema(vec![referencing_type(
3236 "Watchlist",
3237 vec![],
3238 vec![MultiLinkDescriptor {
3239 name: "watched".into(),
3240 target: "default::Account".into(),
3241 through: None,
3242 nullable: true,
3243 description: None,
3244 default_pyql: None,
3245 on_delete: vec![],
3246 is_exclusive: false,
3247 }],
3248 )]);
3249 let ddl = export_schema(&schema).unwrap();
3250 assert!(
3251 ddl.contains(" target uuid NOT NULL,\n"),
3252 "junction target must be a bare uuid column, got:\n{}",
3253 ddl
3254 );
3255 assert!(
3256 !ddl.contains("target uuid NOT NULL REFERENCES \"public\".\"Account\""),
3257 "got:\n{}",
3258 ddl
3259 );
3260 }
3261
3262 #[test]
3263 fn test_source_side_cascade_deletes_from_implementors_not_the_view() {
3264 let schema = interface_schema(vec![referencing_type(
3267 "Session",
3268 vec![link_to_account(
3269 "account",
3270 vec![OnDeletePolicy {
3271 side: DeleteSide::Source,
3272 action: DeleteAction::DeleteTarget,
3273 }],
3274 None,
3275 )],
3276 vec![],
3277 )]);
3278 let ddl = export_schema(&schema).unwrap();
3279 assert!(
3280 ddl.contains("DELETE FROM \"public\".\"Individual\" WHERE id = OLD.\"account_id\";"),
3281 "got:\n{}",
3282 ddl
3283 );
3284 assert!(
3285 ddl.contains("DELETE FROM \"public\".\"Organization\" WHERE id = OLD.\"account_id\";"),
3286 "got:\n{}",
3287 ddl
3288 );
3289 assert!(
3290 !ddl.contains("DELETE FROM \"public\".\"Account\" WHERE id ="),
3291 "must never delete through the interface view, got:\n{}",
3292 ddl
3293 );
3294 }
3295
3296 #[test]
3297 fn test_link_to_concrete_type_still_gets_its_foreign_key() {
3298 let mut schema = interface_schema(vec![referencing_type(
3300 "Session",
3301 vec![LinkDescriptor {
3302 name: "owner".into(),
3303 target: "default::Individual".into(),
3304 nullable: true,
3305 description: None,
3306 default_pyql: None,
3307 is_exclusive: false,
3308 is_readonly: false,
3309 rewrites: vec![],
3310 on_delete: vec![],
3311 through: None,
3312 }],
3313 vec![],
3314 )]);
3315 schema.types.retain(|t| t.name != "Organization");
3316 let ddl = export_schema(&schema).unwrap();
3317 assert!(
3318 ddl.contains("FOREIGN KEY (\"owner_id\") REFERENCES \"public\".\"Individual\"(id)"),
3319 "got:\n{}",
3320 ddl
3321 );
3322 }
3323
3324 #[test]
3325 fn test_emit_one_table_includes_cache_invalidate_trigger() {
3326 let mut out = String::new();
3327 emit_one_table(&person_type(), None, &mut out);
3328 assert!(
3329 out.contains("CREATE OR REPLACE TRIGGER pylon_cache_invalidate\n AFTER INSERT OR UPDATE OR DELETE ON \"public\".\"Person\""),
3330 "got:\n{out}"
3331 );
3332 }
3333
3334 #[test]
3342 fn test_emit_table_compiles_a_pyql_default() {
3343 let mut td = person_type();
3344 td.properties[1].default_sql = None;
3345 td.properties[1].default_pyql = Some("std::uuid_generate_v7()".into());
3346 td.properties[1].pg_type = "uuid".into();
3347
3348 let schema = SchemaDescriptor {
3349 types: vec![td.clone()],
3350 ..minimal_schema(vec![])
3351 };
3352
3353 let mut out = String::new();
3354 emit_one_table(&td, Some(&schema), &mut out);
3355 assert!(
3356 out.contains("\"age\" uuid NULL DEFAULT uuidv7()") || out.contains("DEFAULT uuidv7()"),
3357 "got:\n{out}"
3358 );
3359 }
3360
3361 #[test]
3362 fn test_export_and_migration_agree_on_a_pyql_default() {
3363 let mut td = person_type();
3364 td.properties[1].default_sql = None;
3365 td.properties[1].default_pyql = Some("std::uuid_generate_v7()".into());
3366 td.properties[1].pg_type = "uuid".into();
3367 let schema = SchemaDescriptor {
3368 types: vec![td.clone()],
3369 ..minimal_schema(vec![])
3370 };
3371
3372 let mut out = String::new();
3373 emit_one_table(&td, Some(&schema), &mut out);
3374 let from_migration = crate::diff::resolve_default_for_test(&td.properties[1], &schema);
3375
3376 assert_eq!(from_migration.as_deref(), Some("uuidv7()"));
3377 assert!(
3378 out.contains(&format!("DEFAULT {}", from_migration.unwrap())),
3379 "export DDL disagrees with the migration path:\n{out}"
3380 );
3381 }
3382
3383 fn trig(on: u8, timing: &str, handler: &str) -> crate::schema::TriggerDescriptor {
3384 crate::schema::TriggerDescriptor {
3385 on,
3386 timing: timing.into(),
3387 handler: handler.into(),
3388 }
3389 }
3390
3391 fn schema_with_trigger(trigger: crate::schema::TriggerDescriptor) -> SchemaDescriptor {
3392 let mut t = person_type();
3393 t.triggers = vec![trigger];
3394 SchemaDescriptor {
3395 types: vec![t],
3396 scalars: vec![],
3397 enums: vec![],
3398 named_tuples: vec![],
3399 globals: vec![],
3400 functions: vec![],
3401 aliases: vec![],
3402 channels: vec![],
3403 ..Default::default()
3404 }
3405 }
3406
3407 #[test]
3408 fn test_trigger_new_anchor_resolves_to_new_alias() {
3409 let schema = schema_with_trigger(trig(1, "After", "update Person set { age := __new__.age }"));
3414 let ddl = export_schema(&schema).unwrap();
3415 assert!(ddl.contains("NEW.\"age\""), "got:\n{ddl}");
3416 }
3417
3418 #[test]
3419 fn test_trigger_that_would_refire_itself_is_rejected() {
3420 let schema = schema_with_trigger(trig(1, "After", "insert Person { age := __new__.age }"));
3423 let err = export_schema(&schema).unwrap_err();
3424 assert!(err.to_string().contains("is recursive"), "got: {err}");
3425 }
3426
3427 #[test]
3428 fn test_recursive_insert_wrapped_in_a_select_shape_is_still_caught() {
3429 let schema = schema_with_trigger(trig(
3432 1,
3433 "After",
3434 "select (insert Person { age := __new__.age }) { age }",
3435 ));
3436 let err = export_schema(&schema).unwrap_err();
3437 assert!(err.to_string().contains("is recursive"), "got: {err}");
3438 }
3439
3440 #[test]
3441 fn test_recursive_check_is_scoped_to_the_triggers_own_events() {
3442 let schema = schema_with_trigger(trig(4, "After", "insert Person { age := __old__.age }"));
3446 assert!(export_schema(&schema).is_ok());
3447 }
3448
3449 #[test]
3450 fn test_trigger_old_anchor_resolves_to_old_alias() {
3451 let schema = schema_with_trigger(trig(4, "After", "update Person set { age := __old__.age }"));
3453 let ddl = export_schema(&schema).unwrap();
3454 assert!(ddl.contains("OLD.\"age\""), "got:\n{ddl}");
3455 }
3456
3457 #[test]
3458 fn test_trigger_update_can_reference_both_new_and_old() {
3459 let schema = schema_with_trigger(trig(2, "After", "select Person filter (__new__.age = __old__.age)"));
3464 let ddl = export_schema(&schema).unwrap();
3465 assert!(ddl.contains("NEW.\"age\""), "got:\n{ddl}");
3466 assert!(ddl.contains("OLD.\"age\""), "got:\n{ddl}");
3467 }
3468
3469 #[test]
3470 fn test_trigger_insert_only_cannot_reference_old() {
3471 let schema = schema_with_trigger(trig(1, "After", "update Person set { age := __old__.age }"));
3472 let err = export_schema(&schema).unwrap_err();
3473 assert!(err.to_string().contains("__old__ cannot be used"), "got: {err}");
3474 }
3475
3476 #[test]
3477 fn test_trigger_delete_only_cannot_reference_new() {
3478 let schema = schema_with_trigger(trig(4, "After", "update Person set { age := __new__.age }"));
3479 let err = export_schema(&schema).unwrap_err();
3480 assert!(err.to_string().contains("__new__ cannot be used"), "got: {err}");
3481 }
3482
3483 #[test]
3484 fn test_trigger_combined_insert_update_cannot_reference_old() {
3485 let ok = schema_with_trigger(trig(3, "After", "select Person filter (__new__.age > 0)"));
3487 assert!(export_schema(&ok).is_ok());
3488
3489 let bad = schema_with_trigger(trig(3, "After", "select Person filter (__old__.age > 0)"));
3490 let err = export_schema(&bad).unwrap_err();
3491 assert!(err.to_string().contains("__old__ cannot be used"), "got: {err}");
3492 }
3493
3494 #[test]
3495 fn test_trigger_after_timing_returns_null() {
3496 let schema = schema_with_trigger(trig(1, "After", "update Person set { age := __new__.age }"));
3497 let ddl = export_schema(&schema).unwrap();
3498 assert!(ddl.contains("RETURN NULL;"), "got:\n{ddl}");
3499 }
3500
3501 #[test]
3502 fn test_trigger_before_timing_returns_new_or_old_appropriately() {
3503 let insert_only = schema_with_trigger(trig(1, "Before", "update Person set { age := __new__.age }"));
3504 let ddl = export_schema(&insert_only).unwrap();
3505 assert!(
3506 ddl.contains("RETURN NEW;"),
3507 "insert-only Before should return NEW, got:\n{ddl}"
3508 );
3509
3510 let delete_only = schema_with_trigger(trig(4, "Before", "update Person set { age := __old__.age }"));
3511 let ddl = export_schema(&delete_only).unwrap();
3512 assert!(
3513 ddl.contains("RETURN OLD;"),
3514 "delete-only Before should return OLD, got:\n{ddl}"
3515 );
3516
3517 let combined = schema_with_trigger(trig(5, "Before", "update Person set { age := 1 }"));
3522 let ddl = export_schema(&combined).unwrap();
3523 assert!(
3524 ddl.contains("IF TG_OP = 'DELETE' THEN RETURN OLD; ELSE RETURN NEW; END IF;"),
3525 "got:\n{ddl}"
3526 );
3527 }
3528
3529 #[test]
3530 fn test_multiple_triggers_get_separate_functions() {
3531 let mut t = person_type();
3532 t.triggers = vec![
3533 trig(1, "After", "update Person set { age := __new__.age }"),
3534 trig(4, "Before", "update Person set { age := __old__.age }"),
3535 ];
3536 let schema = SchemaDescriptor {
3537 types: vec![t],
3538 scalars: vec![],
3539 enums: vec![],
3540 named_tuples: vec![],
3541 globals: vec![],
3542 functions: vec![],
3543 aliases: vec![],
3544 channels: vec![],
3545 ..Default::default()
3546 };
3547 let ddl = export_schema(&schema).unwrap();
3548 let fn_count = ddl.matches("CREATE OR REPLACE FUNCTION \"public\".\"Person_").count();
3549 assert_eq!(fn_count, 2, "expected one function per Trigger(...), got:\n{ddl}");
3550 }
3551
3552 #[test]
3553 fn test_emit_scalar_function_ddl() {
3554 let fd = FunctionDescriptor {
3555 name: "mysum".into(),
3556 module: "math".into(),
3557 params: vec![
3558 FunctionParamDescriptor {
3559 name: "a".into(),
3560 pg_type: "int8".into(),
3561 },
3562 FunctionParamDescriptor {
3563 name: "b".into(),
3564 pg_type: "int8".into(),
3565 },
3566 ],
3567 return_pg_type: "int8".into(),
3568 return_is_object: false,
3569 return_is_set: false,
3570 return_is_polymorphic: false,
3571 volatility: "immutable".into(),
3572 body: "a + b".into(),
3573 };
3574 let schema = minimal_schema(vec![fd.clone()]);
3575 let ddl = emit_one_function(&fd, &schema).unwrap();
3576 assert!(
3577 ddl.contains("CREATE OR REPLACE FUNCTION \"math\".\"mysum\""),
3578 "got:\n{}",
3579 ddl
3580 );
3581 assert!(ddl.contains("\"a\" int8, \"b\" int8"), "got:\n{}", ddl);
3582 assert!(ddl.contains("RETURNS int8"), "got:\n{}", ddl);
3583 assert!(ddl.contains("IMMUTABLE"), "got:\n{}", ddl);
3584 assert!(ddl.contains("SELECT"), "got:\n{}", ddl);
3585 }
3586
3587 #[test]
3588 fn test_emit_setof_function_ddl() {
3589 let fd = FunctionDescriptor {
3590 name: "counters".into(),
3591 module: "default".into(),
3592 params: vec![],
3593 return_pg_type: "int8".into(),
3594 return_is_object: false,
3595 return_is_set: true,
3596 return_is_polymorphic: false,
3597 volatility: "stable".into(),
3598 body: "1".into(),
3599 };
3600 let schema = minimal_schema(vec![fd.clone()]);
3601 let ddl = emit_one_function(&fd, &schema).unwrap();
3602 assert!(ddl.contains("RETURNS SETOF int8"), "got:\n{}", ddl);
3603 assert!(ddl.contains("STABLE"), "got:\n{}", ddl);
3604 }
3605
3606 #[test]
3607 fn test_emit_object_function_ddl() {
3608 let fd = FunctionDescriptor {
3609 name: "adults".into(),
3610 module: "default".into(),
3611 params: vec![],
3612 return_pg_type: "default::Person".into(),
3613 return_is_object: true,
3614 return_is_set: true,
3615 return_is_polymorphic: false,
3616 volatility: "stable".into(),
3617 body: "select Person filter .age > 18".into(),
3618 };
3619 let schema = minimal_schema(vec![fd.clone()]);
3620 let ddl = emit_one_function(&fd, &schema).unwrap();
3621 assert!(
3622 ddl.contains("CREATE OR REPLACE FUNCTION \"public\".\"adults\"()"),
3623 "got:\n{}",
3624 ddl
3625 );
3626 assert!(ddl.contains("RETURNS TABLE("), "got:\n{}", ddl);
3627 assert!(ddl.contains("\"id\" uuid"), "got:\n{}", ddl);
3628 assert!(ddl.contains("\"age\" int8"), "got:\n{}", ddl);
3629 assert!(ddl.contains("STABLE"), "got:\n{}", ddl);
3630 assert!(ddl.contains("SELECT * FROM"), "got:\n{}", ddl);
3631 }
3632
3633 #[test]
3634 fn test_object_function_can_reference_its_own_parameter() {
3635 let fd = FunctionDescriptor {
3640 name: "older_than".into(),
3641 module: "default".into(),
3642 params: vec![FunctionParamDescriptor {
3643 name: "min_age".into(),
3644 pg_type: "int8".into(),
3645 }],
3646 return_pg_type: "default::Person".into(),
3647 return_is_object: true,
3648 return_is_set: true,
3649 return_is_polymorphic: false,
3650 volatility: "stable".into(),
3651 body: "select Person filter .age > min_age".into(),
3652 };
3653 let schema = minimal_schema(vec![fd.clone()]);
3654 let ddl = emit_one_function(&fd, &schema).unwrap();
3655 assert!(
3656 ddl.contains("\"min_age\" int8"),
3657 "parameter missing from the signature, got:\n{}",
3658 ddl
3659 );
3660 assert!(
3661 ddl.contains("\"min_age\")") || ddl.contains("= \"min_age\"") || ddl.contains("> \"min_age\""),
3662 "parameter not referenced in the body, got:\n{}",
3663 ddl
3664 );
3665 }
3666
3667 #[test]
3668 fn test_a_function_reading_a_global_takes_the_globals_argument() {
3669 use crate::schema::GlobalDescriptor;
3672 let reader = FunctionDescriptor {
3673 name: "cutoff".into(),
3674 module: "default".into(),
3675 params: vec![],
3676 return_pg_type: "timestamptz".into(),
3677 return_is_object: false,
3678 return_is_set: false,
3679 return_is_polymorphic: false,
3680 volatility: "stable".into(),
3681 body: "global snapshot_at ?? datetime_of_transaction()".into(),
3682 };
3683 let plain = FunctionDescriptor {
3684 name: "bump".into(),
3685 module: "default".into(),
3686 params: vec![FunctionParamDescriptor {
3687 name: "n".into(),
3688 pg_type: "int8".into(),
3689 }],
3690 return_pg_type: "int8".into(),
3691 return_is_object: false,
3692 return_is_set: false,
3693 return_is_polymorphic: false,
3694 volatility: "immutable".into(),
3695 body: "n + 1".into(),
3696 };
3697 let mut schema = minimal_schema(vec![reader.clone(), plain.clone()]);
3698 schema.globals.push(GlobalDescriptor {
3699 name: "snapshot_at".into(),
3700 module: "default".into(),
3701 scalar_type: "datetime".into(),
3702 required: false,
3703 default_expr: None,
3704 computed_expr: None,
3705 });
3706
3707 let reader_ddl = emit_one_function(&reader, &schema).unwrap();
3708 assert!(
3709 reader_ddl.contains("\"__pylon_json_globals__\" jsonb"),
3710 "the global reader should take the argument, got:\n{}",
3711 reader_ddl
3712 );
3713 assert!(
3714 reader_ddl.contains("__pylon_json_globals__ ->> 'default::snapshot_at'"),
3715 "the body should read the global out of it, got:\n{}",
3716 reader_ddl
3717 );
3718 assert!(
3719 !reader_ddl.contains("$1"),
3720 "no parameter placeholder should survive, got:\n{}",
3721 reader_ddl
3722 );
3723
3724 let plain_ddl = emit_one_function(&plain, &schema).unwrap();
3728 assert!(
3729 !plain_ddl.contains("__pylon_json_globals__"),
3730 "a function that cannot reach a global should be untouched, got:\n{}",
3731 plain_ddl
3732 );
3733 }
3734
3735 #[test]
3736 fn test_the_globals_argument_is_forwarded_to_a_callee_that_needs_it() {
3737 use crate::schema::GlobalDescriptor;
3738 let reader = FunctionDescriptor {
3739 name: "cutoff".into(),
3740 module: "default".into(),
3741 params: vec![],
3742 return_pg_type: "timestamptz".into(),
3743 return_is_object: false,
3744 return_is_set: false,
3745 return_is_polymorphic: false,
3746 volatility: "stable".into(),
3747 body: "global snapshot_at ?? datetime_of_transaction()".into(),
3748 };
3749 let caller = FunctionDescriptor {
3750 name: "is_past".into(),
3751 module: "default".into(),
3752 params: vec![],
3753 return_pg_type: "bool".into(),
3754 return_is_object: false,
3755 return_is_set: false,
3756 return_is_polymorphic: false,
3757 volatility: "stable".into(),
3758 body: "cutoff() < datetime_of_transaction()".into(),
3759 };
3760 let mut schema = minimal_schema(vec![reader, caller.clone()]);
3761 schema.globals.push(GlobalDescriptor {
3762 name: "snapshot_at".into(),
3763 module: "default".into(),
3764 scalar_type: "datetime".into(),
3765 required: false,
3766 default_expr: None,
3767 computed_expr: None,
3768 });
3769
3770 let ddl = emit_one_function(&caller, &schema).unwrap();
3773 assert!(
3774 ddl.contains("\"__pylon_json_globals__\" jsonb"),
3775 "a caller of a global reader needs the argument too, got:\n{}",
3776 ddl
3777 );
3778 assert!(
3779 ddl.contains("cutoff\"((__pylon_json_globals__))") || ddl.contains("__pylon_json_globals__)"),
3780 "it should forward the argument, got:\n{}",
3781 ddl
3782 );
3783 }
3784
3785 #[test]
3786 fn test_calling_an_object_function_in_an_expression_says_why() {
3787 let fd = FunctionDescriptor {
3792 name: "adults".into(),
3793 module: "default".into(),
3794 params: vec![],
3795 return_pg_type: "default::Person".into(),
3796 return_is_object: true,
3797 return_is_set: true,
3798 return_is_polymorphic: false,
3799 volatility: "stable".into(),
3800 body: "select Person filter .age > 18".into(),
3801 };
3802 let schema = minimal_schema(vec![fd]);
3803 let ast = crate::parse::parse("select Person { x := adults() }").unwrap();
3808 assert!(
3809 crate::ir::compile(&ast, &schema).is_ok(),
3810 "an object-returning call should stand as a pointer's own value"
3811 );
3812 let ast = crate::parse::parse("select Person { x := count(adults()) }").unwrap();
3813 let err = match crate::ir::compile(&ast, &schema) {
3814 Err(e) => e.to_string(),
3815 Ok(_) => panic!("expected the call to be rejected"),
3816 };
3817 assert!(
3818 err.contains("returns objects") && err.contains("subject of a select"),
3819 "expected an explanation of the restriction, got: {}",
3820 err
3821 );
3822 assert!(
3823 !err.contains("does not exist"),
3824 "the function does exist; the message should not claim otherwise: {}",
3825 err
3826 );
3827 }
3828
3829 #[test]
3830 fn test_emit_sequence_scalar_ddl() {
3831 use crate::schema::ScalarDescriptor;
3832 let schema = SchemaDescriptor {
3833 types: vec![],
3834 scalars: vec![ScalarDescriptor {
3835 name: "OrderNumber".into(),
3836 module: "default".into(),
3837 base: "Sequence".into(),
3838 pg_type: "int8".into(),
3839 check_constraints: vec![],
3840 is_sequence: true,
3841 }],
3842 enums: vec![],
3843 named_tuples: vec![],
3844 globals: vec![],
3845 functions: vec![],
3846 aliases: vec![],
3847 channels: vec![],
3848 ..Default::default()
3849 };
3850 let ddl = export_schema(&schema).unwrap();
3851 assert!(
3852 ddl.contains("CREATE SEQUENCE \"public\".\"OrderNumber_seq\""),
3853 "got:\n{}",
3854 ddl
3855 );
3856 assert!(
3857 ddl.contains("CREATE DOMAIN \"public\".\"OrderNumber\" AS int8"),
3858 "got:\n{}",
3859 ddl
3860 );
3861 let seq_pos = ddl.find("CREATE SEQUENCE").unwrap();
3863 let dom_pos = ddl.find("CREATE DOMAIN").unwrap();
3864 assert!(seq_pos < dom_pos, "sequence must appear before domain");
3865 }
3866
3867 #[test]
3868 fn test_registered_scalar_domain_is_used_as_the_column_type() {
3869 use crate::schema::ScalarDescriptor;
3875 let schema = SchemaDescriptor {
3876 types: vec![TypeDescriptor {
3877 name: "Contact".into(),
3878 module: "default".into(),
3879 table: "Contact".into(),
3880 abstract_: false,
3881 materialized: true,
3882 description: None,
3883 parents: vec![],
3884 interfaces: vec![],
3885 bases: vec![],
3886 properties: vec![PropertyDescriptor {
3887 name: "email".into(),
3888 pg_type: "text".into(),
3889 nullable: false,
3890 default_sql: None,
3891 default_pyql: None,
3892 description: None,
3893 check_constraints: vec![],
3894 is_exclusive: false,
3895 is_pk: false,
3896 is_readonly: false,
3897 rewrites: vec![],
3898 tuple_members: None,
3899 column_type: Some("\"public\".\"EmailStr\"".into()),
3900 }],
3901 links: vec![],
3902 multilinks: vec![],
3903 computed: vec![],
3904 constraints: vec![],
3905 indexes: vec![],
3906 partition: None,
3907 vector_indexes: vec![],
3908 search_indexes: vec![],
3909 triggers: vec![],
3910 junction: false,
3911 signals: vec![],
3912 }],
3913 scalars: vec![ScalarDescriptor {
3914 name: "EmailStr".into(),
3915 module: "default".into(),
3916 base: "Str".into(),
3917 pg_type: "text".into(),
3918 check_constraints: vec!["value ~ '^[^@]+@[^@]+\\.[^@]+$'".into()],
3919 is_sequence: false,
3920 }],
3921 enums: vec![],
3922 named_tuples: vec![],
3923 globals: vec![],
3924 functions: vec![],
3925 aliases: vec![],
3926 channels: vec![],
3927 ..Default::default()
3928 };
3929 let ddl = export_schema(&schema).unwrap();
3930 assert!(
3931 ddl.contains("CREATE DOMAIN \"public\".\"EmailStr\" AS text\n CONSTRAINT ")
3932 && ddl.contains("CHECK (value ~ '^[^@]+@[^@]+\\.[^@]+$')"),
3933 "got:\n{}",
3934 ddl
3935 );
3936 assert!(
3937 ddl.contains("\"email\" \"public\".\"EmailStr\" NOT NULL"),
3938 "column must use the domain type, not the plain base type — got:\n{}",
3939 ddl
3940 );
3941 }
3942
3943 fn account_interface_schema() -> SchemaDescriptor {
3946 let mut account = TypeDescriptor {
3947 name: "Account".into(),
3948 module: "default".into(),
3949 table: "Account".into(),
3950 abstract_: true,
3951 materialized: true,
3952 description: None,
3953 parents: vec![],
3954 interfaces: vec![],
3955 bases: vec![],
3956 properties: vec![PropertyDescriptor {
3957 name: "email".into(),
3958 pg_type: "text".into(),
3959 nullable: false,
3960 default_sql: None,
3961 default_pyql: None,
3962 description: None,
3963 check_constraints: vec![],
3964 is_exclusive: true,
3965 is_pk: false,
3966 is_readonly: false,
3967 rewrites: vec![],
3968 tuple_members: None,
3969 column_type: None,
3970 }],
3971 links: vec![],
3972 multilinks: vec![],
3973 computed: vec![],
3974 constraints: vec![],
3975 indexes: vec![],
3976 partition: None,
3977 vector_indexes: vec![],
3978 search_indexes: vec![],
3979 triggers: vec![],
3980 junction: false,
3981 signals: vec![],
3982 };
3983 let mut individual = account.clone();
3984 individual.name = "Individual".into();
3985 individual.table = "Individual".into();
3986 individual.abstract_ = false;
3987 individual.materialized = true;
3988 individual.interfaces = vec!["default::Account".into()];
3989 let mut organization = individual.clone();
3990 organization.name = "Organization".into();
3991 organization.table = "Organization".into();
3992 account.constraints = vec![]; SchemaDescriptor {
3994 types: vec![account, individual, organization],
3995 scalars: vec![],
3996 enums: vec![],
3997 named_tuples: vec![],
3998 globals: vec![],
3999 functions: vec![],
4000 aliases: vec![],
4001 channels: vec![],
4002 ..Default::default()
4003 }
4004 }
4005
4006 #[test]
4007 fn test_interface_exclusive_property_gets_per_table_index_and_cross_table_trigger() {
4008 let schema = account_interface_schema();
4009 let ddl = export_schema(&schema).unwrap();
4010
4011 assert!(
4012 ddl.contains("CREATE UNIQUE INDEX ON \"public\".\"Individual\" (\"email\")"),
4013 "got:\n{ddl}"
4014 );
4015 assert!(
4016 ddl.contains("CREATE UNIQUE INDEX ON \"public\".\"Organization\" (\"email\")"),
4017 "got:\n{ddl}"
4018 );
4019
4020 assert_eq!(
4022 ddl.matches("CREATE OR REPLACE FUNCTION \"public\".\"_excl_Account_email\"")
4023 .count(),
4024 1,
4025 "the trigger function must be emitted exactly once, shared by every implementor; got:\n{ddl}"
4026 );
4027 assert!(ddl.contains("SELECT 1 FROM \"public\".\"Account\""), "got:\n{ddl}");
4029
4030 assert!(
4032 ddl.contains(
4033 "CREATE CONSTRAINT TRIGGER \"_excl_Account_email_ins\"\nAFTER INSERT ON \"public\".\"Individual\""
4034 ),
4035 "got:\n{ddl}"
4036 );
4037 assert!(
4038 ddl.contains(
4039 "CREATE CONSTRAINT TRIGGER \"_excl_Account_email_ins\"\nAFTER INSERT ON \"public\".\"Organization\""
4040 ),
4041 "got:\n{ddl}"
4042 );
4043 assert!(ddl.contains("DEFERRABLE INITIALLY DEFERRED"), "got:\n{ddl}");
4044 assert!(
4045 ddl.contains("CREATE CONSTRAINT TRIGGER \"_excl_Account_email_upd\"\nAFTER UPDATE OF \"email\""),
4046 "the UPDATE trigger must only fire when the exclusive column itself changes; got:\n{ddl}"
4047 );
4048 }
4049
4050 #[test]
4051 fn test_junction_backed_exclusive_link_gets_a_cross_implementor_helper_view_and_trigger() {
4052 use crate::schema::LinkDescriptor;
4059 let mut schema = account_interface_schema();
4060 for t in &mut schema.types {
4061 t.properties.retain(|p| p.name != "email");
4062 if t.name == "Account" || t.name == "Individual" || t.name == "Organization" {
4063 t.links.push(LinkDescriptor {
4064 name: "owner".into(),
4065 target: "default::Person".into(),
4066 nullable: true,
4067 through: Some("default::AccountOwner".into()),
4068 description: None,
4069 default_pyql: None,
4070 is_exclusive: true,
4071 is_readonly: false,
4072 rewrites: vec![],
4073 on_delete: vec![],
4074 });
4075 }
4076 }
4077
4078 let views = interface_junction_view_ddl_with_names(&schema);
4079 assert_eq!(views.len(), 1, "expected exactly one helper view, got: {views:?}");
4080 let (view_module, view_name, view_ddl) = &views[0];
4081 assert_eq!(view_module, "default");
4082 assert_eq!(view_name, "Account.owner");
4083 assert!(
4084 view_ddl.contains("SELECT source, target FROM \"public\".\"Individual.owner\""),
4085 "got:\n{view_ddl}"
4086 );
4087 assert!(
4088 view_ddl.contains("SELECT source, target FROM \"public\".\"Organization.owner\""),
4089 "got:\n{view_ddl}"
4090 );
4091
4092 let infos = interface_exclusive_trigger_infos(&schema);
4093 assert_eq!(
4094 infos.len(),
4095 2,
4096 "expected one entry per implementor, got: {}",
4097 infos.len()
4098 );
4099 assert!(
4100 infos.iter().any(|i| i.impl_table == "Individual.owner"
4101 && i.ins_ddl.contains("AFTER INSERT ON \"public\".\"Individual.owner\"")),
4102 "the constraint trigger must attach to the implementor's own *junction* table, not the owner table"
4103 );
4104 assert!(
4105 infos
4106 .iter()
4107 .any(|i| i.fn_ddl.contains("SELECT 1 FROM \"public\".\"Account.owner\"")),
4108 "the trigger function must query the helper view"
4109 );
4110 assert!(infos.iter().all(|i| i.upd_ddl.contains("AFTER UPDATE OF \"target\"")));
4111 }
4112
4113 #[test]
4116 fn test_target_fk_suffix_forces_deferrable_when_source_side_deletes_target() {
4117 let policies = vec![OnDeletePolicy {
4124 side: DeleteSide::Source,
4125 action: DeleteAction::DeleteTarget,
4126 }];
4127 assert_eq!(target_fk_suffix(&policies), " DEFERRABLE INITIALLY DEFERRED");
4128
4129 let policies = vec![OnDeletePolicy {
4130 side: DeleteSide::Source,
4131 action: DeleteAction::DeleteTargetIfOrphan,
4132 }];
4133 assert_eq!(target_fk_suffix(&policies), " DEFERRABLE INITIALLY DEFERRED");
4134 }
4135
4136 #[test]
4137 fn test_target_fk_suffix_unaffected_when_no_source_side_policy() {
4138 assert_eq!(target_fk_suffix(&[]), " ON DELETE RESTRICT");
4140 let policies = vec![OnDeletePolicy {
4141 side: DeleteSide::Target,
4142 action: DeleteAction::Allow,
4143 }];
4144 assert_eq!(target_fk_suffix(&policies), " ON DELETE SET NULL");
4145 }
4146
4147 #[test]
4148 fn test_target_jt_fk_suffix_forces_deferrable_when_source_side_deletes_target() {
4149 let policies = vec![OnDeletePolicy {
4151 side: DeleteSide::Source,
4152 action: DeleteAction::DeleteTargetIfOrphan,
4153 }];
4154 assert_eq!(target_jt_fk_suffix(&policies), " DEFERRABLE INITIALLY DEFERRED");
4155 }
4156
4157 fn org_type(module: &str) -> TypeDescriptor {
4158 TypeDescriptor {
4159 name: "Org".into(),
4160 module: module.into(),
4161 table: "Org".into(),
4162 abstract_: false,
4163 materialized: true,
4164 description: None,
4165 parents: vec![],
4166 interfaces: vec![],
4167 bases: vec![],
4168 properties: vec![PropertyDescriptor {
4169 name: "id".into(),
4170 pg_type: "uuid".into(),
4171 nullable: false,
4172 default_sql: Some("gen_random_uuid()".into()),
4173 default_pyql: None,
4174 description: None,
4175 check_constraints: vec![],
4176 is_exclusive: true,
4177 is_pk: true,
4178 is_readonly: true,
4179 rewrites: vec![],
4180 tuple_members: None,
4181 column_type: None,
4182 }],
4183 links: vec![],
4184 multilinks: vec![],
4185 computed: vec![],
4186 constraints: vec![],
4187 indexes: vec![],
4188 partition: None,
4189 vector_indexes: vec![],
4190 search_indexes: vec![],
4191 triggers: vec![],
4192 junction: false,
4193 signals: vec![],
4194 }
4195 }
4196
4197 #[test]
4198 fn test_multilink_target_delete_source_trigger_is_after_not_before() {
4199 let module = "default";
4209 let mut owner = org_type(module);
4210 owner.name = "Product".into();
4211 owner.table = "Product".into();
4212 owner.multilinks = vec![MultiLinkDescriptor {
4213 name: "tags".into(),
4214 target: format!("{module}::Org"),
4215 through: None,
4216 nullable: false,
4217 description: None,
4218 default_pyql: None,
4219 on_delete: vec![OnDeletePolicy {
4220 side: DeleteSide::Target,
4221 action: DeleteAction::DeleteSource,
4222 }],
4223 is_exclusive: false,
4224 }];
4225 let schema = SchemaDescriptor {
4226 types: vec![org_type(module), owner],
4227 scalars: vec![],
4228 enums: vec![],
4229 named_tuples: vec![],
4230 globals: vec![],
4231 functions: vec![],
4232 aliases: vec![],
4233 channels: vec![],
4234 ..Default::default()
4235 };
4236 let ddl = export_schema(&schema).unwrap();
4237 assert!(
4238 ddl.contains("AFTER DELETE ON \"public\".\"Product.tags\""),
4239 "got:\n{ddl}"
4240 );
4241 assert!(
4242 !ddl.contains("BEFORE DELETE ON \"public\".\"Product.tags\""),
4243 "got:\n{ddl}"
4244 );
4245 }
4246
4247 fn referenced_pylon_functions(ddl: &str) -> std::collections::BTreeSet<String> {
4253 const DEFINITION: &str = "CREATE OR REPLACE FUNCTION _pylon.";
4254 let mut found = std::collections::BTreeSet::new();
4255 let mut rest = ddl;
4256 while let Some(pos) = rest.find("_pylon.") {
4257 let is_definition = rest[..pos]
4258 .len()
4259 .checked_sub(DEFINITION.len() - "_pylon.".len())
4260 .is_some_and(|start| rest[start..pos + "_pylon.".len()].ends_with(DEFINITION));
4261 let after = &rest[pos + "_pylon.".len()..];
4262 let name_len = after
4263 .find(|c: char| !c.is_alphanumeric() && c != '_')
4264 .unwrap_or(after.len());
4265 let name = &after[..name_len];
4266 if !is_definition && !name.is_empty() && after[name_len..].starts_with('(') {
4269 found.insert(name.to_string());
4270 }
4271 rest = &rest[pos + "_pylon.".len()..];
4272 }
4273 found
4274 }
4275
4276 #[test]
4277 fn only_known_stdlib_functions_reach_persisted_ddl() {
4278 let mut person = person_type();
4294 person.materialized = true;
4295 person.properties.push(PropertyDescriptor {
4296 name: "tags".into(),
4297 pg_type: "text[]".into(),
4298 nullable: true,
4299 default_sql: None,
4300 default_pyql: None,
4301 description: None,
4302 check_constraints: vec![],
4303 is_exclusive: false,
4304 is_pk: false,
4305 is_readonly: false,
4306 rewrites: vec![],
4307 tuple_members: None,
4308 column_type: None,
4309 });
4310 person.triggers = vec![trig(
4311 1,
4312 "After",
4313 "update Person set { age := <std::int64>__new__.tags[0] }",
4314 )];
4315 person.signals = vec![crate::schema::SignalEntry { on: 1 }];
4316 let schema = SchemaDescriptor {
4317 types: vec![person],
4318 scalars: vec![],
4319 enums: vec![],
4320 named_tuples: vec![],
4321 globals: vec![],
4322 functions: vec![],
4323 aliases: vec![],
4324 channels: vec![],
4325 ..Default::default()
4326 };
4327
4328 let ddl = export_schema(&schema).unwrap();
4329 let referenced = referenced_pylon_functions(&ddl);
4330 let allowed: std::collections::BTreeSet<String> = ["array_subscript", "notify_cache_invalidate"]
4331 .into_iter()
4332 .map(String::from)
4333 .collect();
4334 assert!(
4335 referenced.is_subset(&allowed),
4336 "new stdlib functions reached persisted DDL: {:?}\n\
4337 see this test's comment before widening the allowlist",
4338 referenced.difference(&allowed).collect::<Vec<_>>(),
4339 );
4340 assert!(
4342 referenced.contains("array_subscript"),
4343 "fixture no longer bakes a stdlib call; it is not testing anything"
4344 );
4345 }
4346
4347 fn multilink_trigger_schema(owner_table: &str) -> SchemaDescriptor {
4353 let module = "default";
4354 let mut owner = org_type(module);
4355 owner.name = owner_table.into();
4356 owner.table = owner_table.into();
4357 owner.multilinks = vec![MultiLinkDescriptor {
4358 name: "tags".into(),
4359 target: format!("{module}::Org"),
4360 through: None,
4361 nullable: false,
4362 description: None,
4363 default_pyql: None,
4364 on_delete: vec![OnDeletePolicy {
4365 side: DeleteSide::Target,
4366 action: DeleteAction::DeleteSource,
4367 }],
4368 is_exclusive: false,
4369 }];
4370 SchemaDescriptor {
4371 types: vec![org_type(module), owner],
4372 scalars: vec![],
4373 enums: vec![],
4374 named_tuples: vec![],
4375 globals: vec![],
4376 functions: vec![],
4377 aliases: vec![],
4378 channels: vec![],
4379 ..Default::default()
4380 }
4381 }
4382
4383 fn trigger_names_of(schema: &SchemaDescriptor) -> Vec<String> {
4384 let type_map: HashMap<String, (&str, &str)> = schema
4385 .types
4386 .iter()
4387 .map(|t| {
4388 (
4389 format!("{}::{}", t.module, t.name),
4390 (t.module.as_str(), t.table.as_str()),
4391 )
4392 })
4393 .collect();
4394 let mut names: Vec<String> = deletion_policy_trigger_infos(schema, &type_map)
4395 .into_iter()
4396 .map(|i| i.trigger_name)
4397 .collect();
4398 names.sort();
4399 names
4400 }
4401
4402 #[test]
4403 fn a_trigger_name_is_stable_for_an_unchanged_schema() {
4404 let schema = multilink_trigger_schema("Product");
4408 assert_eq!(trigger_names_of(&schema), trigger_names_of(&schema));
4409 }
4410
4411 #[test]
4412 fn a_trigger_name_changes_when_its_body_changes() {
4413 let before = trigger_names_of(&multilink_trigger_schema("Product"));
4419 let after = trigger_names_of(&multilink_trigger_schema("Widget"));
4420 assert_ne!(before, after, "trigger name did not follow its body");
4421 }
4422
4423 #[test]
4424 fn a_signal_trigger_name_follows_its_event_list() {
4425 let names_for = |on: u8| {
4429 let module = "default";
4430 let mut t = org_type(module);
4431 t.signals = vec![crate::schema::SignalEntry { on }];
4432 let schema = SchemaDescriptor {
4433 types: vec![t],
4434 scalars: vec![],
4435 enums: vec![],
4436 named_tuples: vec![],
4437 globals: vec![],
4438 functions: vec![],
4439 aliases: vec![],
4440 channels: vec![],
4441 ..Default::default()
4442 };
4443 signal_trigger_infos(&schema)
4444 .into_iter()
4445 .map(|i| i.trigger_name)
4446 .collect::<Vec<_>>()
4447 };
4448 assert_ne!(names_for(1), names_for(1 | 4));
4449 assert_eq!(names_for(1), names_for(1));
4450 }
4451
4452 #[test]
4455 fn test_junction_backed_single_link_gets_a_source_pk_junction_table_no_fk_column() {
4456 let module = "default";
4457 let mut owner = org_type(module);
4458 owner.name = "Person".into();
4459 owner.table = "Person".into();
4460 owner.links = vec![LinkDescriptor {
4461 name: "spouse".into(),
4462 target: format!("{module}::Org"),
4463 nullable: true,
4464 through: Some(format!("{module}::Marriage")),
4465 description: None,
4466 default_pyql: None,
4467 is_exclusive: true,
4468 is_readonly: false,
4469 rewrites: vec![],
4470 on_delete: vec![],
4471 }];
4472 let schema = SchemaDescriptor {
4473 types: vec![org_type(module), owner],
4474 scalars: vec![],
4475 enums: vec![],
4476 named_tuples: vec![],
4477 globals: vec![],
4478 functions: vec![],
4479 aliases: vec![],
4480 channels: vec![],
4481 ..Default::default()
4482 };
4483 let ddl = export_schema(&schema).unwrap();
4484
4485 assert!(
4486 !ddl.contains("spouse_id"),
4487 "no {{name}}_id column/FK for a junction-backed link, got:\n{ddl}"
4488 );
4489 assert!(ddl.contains("CREATE TABLE \"public\".\"Person.spouse\""), "got:\n{ddl}");
4490 assert!(
4491 ddl.contains("PRIMARY KEY (source)"),
4492 "single-link junction table must be capped to one row per source, got:\n{ddl}"
4493 );
4494 assert!(
4495 ddl.contains("UNIQUE (target)"),
4496 "exclusive single link must also be unique on the target side, got:\n{ddl}"
4497 );
4498 assert!(
4499 !ddl.contains("CREATE UNIQUE INDEX ON \"public\".\"Person\" (\"spouse_id\")"),
4500 "got:\n{ddl}"
4501 );
4502 }
4503
4504 #[test]
4507 fn test_type_with_a_signal_gets_a_capture_trigger() {
4508 use crate::schema::SignalEntry;
4509 let mut with_signal = person_type();
4510 with_signal.signals = vec![SignalEntry { on: 5 }]; let schema = SchemaDescriptor {
4512 types: vec![with_signal],
4513 scalars: vec![],
4514 enums: vec![],
4515 named_tuples: vec![],
4516 globals: vec![],
4517 functions: vec![],
4518 aliases: vec![],
4519 channels: vec![],
4520 ..Default::default()
4521 };
4522 let ddl = export_schema(&schema).unwrap();
4523 assert!(
4524 ddl.contains("AFTER INSERT OR UPDATE OR DELETE ON \"public\".\"Person\""),
4525 "got:\n{ddl}"
4526 );
4527 assert!(ddl.contains("_pylon.\"SignalOutbox\""), "got:\n{ddl}");
4528 assert!(ddl.contains("'default::Person'"), "got:\n{ddl}");
4529 }
4530
4531 #[test]
4532 fn test_type_without_a_signal_gets_no_capture_trigger() {
4533 let schema = minimal_schema(vec![]);
4536 let ddl = export_schema(&schema).unwrap();
4537 assert!(!ddl.contains("SignalOutbox"), "got:\n{ddl}");
4538 }
4539
4540 #[test]
4541 fn test_signal_update_capture_skips_index_maintenance_only_changes() {
4542 use crate::schema::{SignalEntry, VectorIndexDescriptor};
4549 let mut with_signal = person_type();
4550 with_signal.signals = vec![SignalEntry { on: 2 }]; with_signal.vector_indexes = vec![VectorIndexDescriptor {
4552 index_name: None,
4553 pointers: vec!["name".into()],
4554 model: "mistral-embed".into(),
4555 metric: "cosine".into(),
4556 dimensions: 1024,
4557 }];
4558 let schema = SchemaDescriptor {
4559 types: vec![with_signal],
4560 scalars: vec![],
4561 enums: vec![],
4562 named_tuples: vec![],
4563 globals: vec![],
4564 functions: vec![],
4565 aliases: vec![],
4566 channels: vec![],
4567 ..Default::default()
4568 };
4569 let ddl = export_schema(&schema).unwrap();
4570 assert!(
4571 ddl.contains(
4572 "IF TG_OP = 'UPDATE' AND (to_jsonb(OLD) - '__vector__') = (to_jsonb(NEW) - '__vector__') THEN"
4573 ),
4574 "got:\n{ddl}"
4575 );
4576 assert!(ddl.contains("RETURN NULL;\n END IF;"), "got:\n{ddl}");
4577 }
4578
4579 #[test]
4580 fn test_signal_update_capture_has_no_guard_without_an_index() {
4581 use crate::schema::SignalEntry;
4584 let mut with_signal = person_type();
4585 with_signal.signals = vec![SignalEntry { on: 2 }]; let schema = SchemaDescriptor {
4587 types: vec![with_signal],
4588 scalars: vec![],
4589 enums: vec![],
4590 named_tuples: vec![],
4591 globals: vec![],
4592 functions: vec![],
4593 aliases: vec![],
4594 channels: vec![],
4595 ..Default::default()
4596 };
4597 let ddl = export_schema(&schema).unwrap();
4598 assert!(!ddl.contains("IF TG_OP = 'UPDATE'"), "got:\n{ddl}");
4599 }
4600}
4601
4602#[cfg(test)]
4603mod partition_tests {
4604 use super::*;
4605 use crate::schema::{PartitionDescriptor, PartitionInterval, PropertyDescriptor};
4606
4607 fn prop(name: &str, pg_type: &str, is_pk: bool) -> PropertyDescriptor {
4608 PropertyDescriptor {
4609 name: name.into(),
4610 pg_type: pg_type.into(),
4611 nullable: false,
4612 default_sql: None,
4613 default_pyql: None,
4614 description: None,
4615 check_constraints: vec![],
4616 is_exclusive: false,
4617 is_pk,
4618 is_readonly: false,
4619 rewrites: vec![],
4620 tuple_members: None,
4621 column_type: None,
4622 }
4623 }
4624
4625 fn event_schema(retention: Option<u32>) -> SchemaDescriptor {
4626 let mut td = TypeDescriptor {
4627 name: "Event".into(),
4628 module: "default".into(),
4629 table: "Event".into(),
4630 abstract_: false,
4631 materialized: true,
4632 description: None,
4633 parents: vec![],
4634 interfaces: vec![],
4635 bases: vec![],
4636 properties: vec![prop("id", "uuid", true), prop("occurred_at", "timestamptz", false)],
4637 links: vec![],
4638 multilinks: vec![],
4639 computed: vec![],
4640 constraints: vec![],
4641 indexes: vec![],
4642 partition: None,
4643 vector_indexes: vec![],
4644 search_indexes: vec![],
4645 triggers: vec![],
4646 junction: false,
4647 signals: vec![],
4648 };
4649 td.partition = Some(PartitionDescriptor {
4650 pointer: "occurred_at".into(),
4651 interval: PartitionInterval::Monthly,
4652 premake: 4,
4653 retention,
4654 });
4655 SchemaDescriptor {
4656 types: vec![td],
4657 ..Default::default()
4658 }
4659 }
4660
4661 #[test]
4662 fn a_partitioned_table_declares_its_range_key() {
4663 let ddl = export_schema(&event_schema(None)).unwrap();
4664 assert!(ddl.contains("PARTITION BY RANGE (\"occurred_at\")"), "got:\n{ddl}");
4665 }
4666
4667 #[test]
4668 fn the_partition_key_is_added_to_the_primary_key() {
4669 let ddl = export_schema(&event_schema(None)).unwrap();
4672 assert!(ddl.contains("PRIMARY KEY (\"id\", \"occurred_at\")"), "got:\n{ddl}");
4673 }
4674
4675 #[test]
4676 fn an_unpartitioned_table_is_untouched() {
4677 let mut schema = event_schema(None);
4678 schema.types[0].partition = None;
4679 let ddl = export_schema(&schema).unwrap();
4680 assert!(!ddl.contains("PARTITION BY"), "got:\n{ddl}");
4681 assert!(!ddl.contains("partman"), "got:\n{ddl}");
4682 assert!(ddl.contains("PRIMARY KEY (\"id\")"), "got:\n{ddl}");
4683 }
4684
4685 #[test]
4686 fn registration_is_guarded_so_re_export_is_idempotent() {
4687 let ddl = export_schema(&event_schema(None)).unwrap();
4690 assert!(
4691 ddl.contains("IF NOT EXISTS (SELECT 1 FROM partman.part_config"),
4692 "got:\n{ddl}"
4693 );
4694 assert!(ddl.contains("partman.create_parent("), "got:\n{ddl}");
4695 assert!(ddl.contains("p_interval := '1 month'"), "got:\n{ddl}");
4696 assert!(ddl.contains("p_premake := 4"), "got:\n{ddl}");
4697 }
4698
4699 #[test]
4700 fn retention_is_applied_when_declared() {
4701 let ddl = export_schema(&event_schema(Some(12))).unwrap();
4702 assert!(ddl.contains("SET retention = '12 months'"), "got:\n{ddl}");
4703 assert!(ddl.contains("retention_keep_table = false"), "got:\n{ddl}");
4704 }
4705
4706 #[test]
4707 fn retention_is_cleared_when_not_declared() {
4708 let ddl = export_schema(&event_schema(None)).unwrap();
4711 assert!(ddl.contains("SET retention = NULL"), "got:\n{ddl}");
4712 }
4713
4714 #[test]
4715 fn a_partitioned_schema_requires_the_partman_extension() {
4716 let exts = crate::diff::required_extensions(&event_schema(None));
4717 assert!(exts.contains(&"pg_partman"), "got: {exts:?}");
4718 }
4719}