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