use crate::error::{PyQLError, PyQLFragmentError};
use crate::schema::{
DeleteAction, DeleteSide, FunctionDescriptor, OnDeletePolicy, SchemaDescriptor, TypeConstraint, TypeDescriptor,
};
use std::collections::{BTreeSet, HashMap, HashSet};
pub mod python_snippet;
pub fn export_schema(schema: &SchemaDescriptor) -> Result<String, PyQLError> {
let mut out = String::new();
let type_map: HashMap<String, (&str, &str)> = schema
.types
.iter()
.map(|t| {
(
format!("{}::{}", t.module, t.name),
(t.module.as_str(), t.table.as_str()),
)
})
.collect();
emit_schemas(schema, &mut out);
emit_enums(schema, &mut out);
emit_scalars(schema, &mut out);
emit_scalar_functions(schema, &mut out)?;
emit_tables(schema, &mut out);
emit_fk_constraints(schema, &type_map, &mut out);
emit_link_source_triggers(schema, &type_map, &mut out);
emit_junction_tables(schema, &mut out);
emit_junction_fk_constraints(schema, &type_map, &mut out);
emit_multilink_deletion_triggers(schema, &type_map, &mut out);
emit_interface_link_triggers(schema, &mut out);
emit_signal_triggers(schema, &mut out);
emit_unique_indexes(schema, &mut out);
emit_check_constraints(schema, &mut out)?;
emit_plain_indexes(schema, &mut out);
emit_triggers(schema, &mut out)?;
emit_interface_views(schema, &mut out);
emit_interface_junction_views(schema, &mut out);
emit_interface_exclusive_triggers(schema, &mut out);
emit_object_functions(schema, &mut out)?;
emit_vector_columns(schema, &mut out);
emit_vector_indexes(schema, &mut out);
emit_search_columns(schema, &mut out);
emit_search_indexes(schema, &mut out);
Ok(out)
}
fn qi(s: &str) -> String {
format!("\"{}\"", s.replace('"', "\"\""))
}
fn pg_schema(module: &str) -> String {
if module == "default" {
"\"public\"".into()
} else {
qi(module)
}
}
fn qn(module: &str, name: &str) -> String {
format!("{}.{}", pg_schema(module), qi(name))
}
pub(crate) fn fnv(parts: &[&str]) -> String {
let mut h: u64 = 0xcbf29ce484222325;
for p in parts {
for b in p.bytes() {
h ^= b as u64;
h = h.wrapping_mul(0x100000001b3);
}
h ^= b'|' as u64;
h = h.wrapping_mul(0x100000001b3);
}
format!("{:016x}", h)
}
fn trigger_names(module: &str, table: &str, pointer: &str, suffix: &str, body: &str) -> (String, String) {
let hash = fnv(&[table, pointer, suffix, body]);
let fname = format!("{}_{}_{}", table, pointer, &hash[..8]);
let fn_qname = qn(module, &fname);
(fname, fn_qname)
}
fn trigger_events(on: u8) -> String {
let mut events = Vec::new();
if on & 1 != 0 {
events.push("INSERT");
}
if on & 2 != 0 {
events.push("UPDATE");
}
if on & 4 != 0 {
events.push("DELETE");
}
events.join(" OR ")
}
fn trigger_timing(timing: &str) -> &str {
match timing {
"Before" => "BEFORE",
"After" => "AFTER",
"InsteadOf" => "INSTEAD OF",
t => t,
}
}
fn emit_schemas(schema: &SchemaDescriptor, out: &mut String) {
let mut modules: BTreeSet<&str> = BTreeSet::new();
for t in &schema.types {
modules.insert(&t.module);
}
for s in &schema.scalars {
modules.insert(&s.module);
}
for e in &schema.enums {
modules.insert(&e.module);
}
for f in &schema.functions {
modules.insert(&f.module);
}
for g in &schema.globals {
modules.insert(&g.module);
}
for a in &schema.aliases {
modules.insert(&a.module);
}
let non_default: Vec<&str> = modules.into_iter().filter(|m| *m != "default").collect();
for module in &non_default {
out.push_str(&format!("CREATE SCHEMA IF NOT EXISTS {};\n", pg_schema(module)));
}
if !non_default.is_empty() {
out.push('\n');
}
}
fn emit_enums(schema: &SchemaDescriptor, out: &mut String) {
for e in &schema.enums {
let members: Vec<String> = e
.members
.iter()
.map(|m| format!("'{}'", m.replace('\'', "''")))
.collect();
out.push_str(&format!(
"DO $$ BEGIN CREATE TYPE {}.{} AS ENUM ({}); EXCEPTION WHEN duplicate_object THEN NULL; END $$;\n",
pg_schema(&e.module),
qi(&e.name),
members.join(", "),
));
}
if !schema.enums.is_empty() {
out.push('\n');
}
}
fn emit_scalars(schema: &SchemaDescriptor, out: &mut String) {
for s in &schema.scalars {
if s.is_sequence {
out.push_str(&format!(
"CREATE SEQUENCE {}.{};\n",
pg_schema(&s.module),
qi(&format!("{}_seq", s.name)),
));
}
let check_clause = scalar_check_clauses(schema, &s.module, &s.name);
out.push_str(&format!(
"CREATE DOMAIN {}.{} AS {}{};\n",
pg_schema(&s.module),
qi(&s.name),
s.pg_type,
check_clause,
));
}
if !schema.scalars.is_empty() {
out.push('\n');
}
}
fn emit_tables(schema: &SchemaDescriptor, out: &mut String) {
for t in &schema.types {
if t.abstract_ || t.junction {
continue;
}
emit_one_table(t, Some(schema), out);
}
}
fn column_default(p: &crate::schema::PropertyDescriptor, schema: Option<&SchemaDescriptor>) -> Option<String> {
if let Some(sql) = &p.default_sql {
return Some(sql.clone());
}
let pyql = p.default_pyql.as_deref()?;
crate::ir::compile_scalar_default(pyql, schema?).ok()
}
fn emit_one_table(t: &TypeDescriptor, schema: Option<&SchemaDescriptor>, out: &mut String) {
out.push_str(&format!("CREATE TABLE {} (\n", qn(&t.module, &t.table)));
let mut lines: Vec<String> = Vec::new();
for p in &t.properties {
let not_null = if p.nullable { "" } else { " NOT NULL" };
let default = column_default(p, schema)
.map(|d| format!(" DEFAULT {}", d))
.unwrap_or_default();
let col_type = p
.column_type
.as_deref()
.unwrap_or_else(|| p.pg_type.strip_prefix("__nt__:").map(|_| "jsonb").unwrap_or(&p.pg_type));
lines.push(format!(" {} {}{}{}", qi(&p.name), col_type, not_null, default));
}
for l in &t.links {
if l.is_junction_backed() {
continue;
}
let not_null = if l.nullable { "" } else { " NOT NULL" };
lines.push(format!(" {} uuid{}", qi(&format!("{}_id", l.name)), not_null));
}
let mut pk_cols: Vec<String> = t.properties.iter().filter(|p| p.is_pk).map(|p| qi(&p.name)).collect();
if let Some(part) = &t.partition {
let key = qi(&part.pointer);
if !pk_cols.contains(&key) {
pk_cols.push(key);
}
}
if !pk_cols.is_empty() {
lines.push(format!(" PRIMARY KEY ({})", pk_cols.join(", ")));
}
out.push_str(&lines.join(",\n"));
out.push_str("\n)");
if let Some(part) = &t.partition {
out.push_str(&format!(" PARTITION BY RANGE ({})", qi(&part.pointer)));
}
out.push_str(";\n\n");
if let Some(part) = &t.partition {
out.push_str(&partman_setup_sql(&t.module, &t.table, part));
out.push_str("\n\n");
}
out.push_str(&cache_invalidate_trigger_sql(&qn(&t.module, &t.table)));
out.push_str("\n\n");
}
fn partman_setup_sql(module: &str, table: &str, part: &crate::schema::PartitionDescriptor) -> String {
let pg_schema = if module == "default" { "public" } else { module };
let parent = format!("{pg_schema}.{table}").replace('\'', "''");
let mut out = String::new();
out.push_str("DO $$ BEGIN\n");
out.push_str(&format!(
" IF NOT EXISTS (SELECT 1 FROM partman.part_config WHERE parent_table = '{parent}') THEN\n"
));
out.push_str(&format!(
" PERFORM partman.create_parent(\n\
\x20 p_parent_table := '{parent}',\n\
\x20 p_control := '{control}',\n\
\x20 p_interval := '{interval}',\n\
\x20 p_premake := {premake}\n\
\x20 );\n",
control = part.pointer.replace('\'', "''"),
interval = part.interval.as_pg_interval(),
premake = part.premake,
));
out.push_str(" END IF;\nEND $$;\n");
match part.retention_interval() {
Some(retention) => {
out.push_str(&format!(
"UPDATE partman.part_config\n\
\x20 SET retention = '{retention}', retention_keep_table = false\n\
\x20 WHERE parent_table = '{parent}';",
));
}
None => {
out.push_str(&format!(
"UPDATE partman.part_config\n SET retention = NULL\n WHERE parent_table = '{parent}';",
));
}
}
out
}
fn cache_invalidate_trigger_sql(qualified_table: &str) -> String {
format!(
"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();",
qualified_table
)
}
pub(crate) fn polymorphic_types(schema: &SchemaDescriptor) -> HashSet<String> {
let extended: HashSet<&str> = schema
.types
.iter()
.flat_map(|t| t.bases.iter().map(String::as_str))
.collect();
schema
.types
.iter()
.map(|t| format!("{}::{}", t.module, t.name))
.zip(&schema.types)
.filter(|(qname, t)| (t.abstract_ && t.materialized) || extended.contains(qname.as_str()))
.map(|(qname, _)| qname)
.collect()
}
pub(crate) fn interface_implementors(schema: &SchemaDescriptor) -> HashMap<String, Vec<&TypeDescriptor>> {
let mut implementors: HashMap<String, Vec<&TypeDescriptor>> = HashMap::new();
for t in &schema.types {
if !t.abstract_ {
for iface in &t.interfaces {
implementors.entry(iface.clone()).or_default().push(t);
}
}
}
for t in &schema.types {
for base in &t.bases {
if let Some(base_td) = schema
.types
.iter()
.find(|b| format!("{}::{}", b.module, b.name) == *base)
{
let entry = implementors.entry(base.clone()).or_default();
if entry.is_empty() {
entry.push(base_td);
}
entry.push(t);
}
}
}
implementors
}
fn policy_for<'a>(policies: &'a [OnDeletePolicy], side: &DeleteSide) -> Option<&'a DeleteAction> {
policies.iter().find(|p| &p.side == side).map(|p| &p.action)
}
pub(crate) fn needs_deferred_target_fk(policies: &[OnDeletePolicy]) -> bool {
policies.iter().any(|p| {
p.side == DeleteSide::Source
&& matches!(
p.action,
DeleteAction::DeleteTarget | DeleteAction::DeleteTargetIfOrphan
)
})
}
fn target_fk_suffix(policies: &[OnDeletePolicy]) -> String {
match policy_for(policies, &DeleteSide::Target).unwrap_or(&DeleteAction::Restrict) {
DeleteAction::Restrict if needs_deferred_target_fk(policies) => " DEFERRABLE INITIALLY DEFERRED".into(),
DeleteAction::Restrict => " ON DELETE RESTRICT".into(),
DeleteAction::DeferredRestrict => " DEFERRABLE INITIALLY DEFERRED".into(),
DeleteAction::DeleteSource => " ON DELETE CASCADE".into(),
DeleteAction::Allow => " ON DELETE SET NULL".into(),
_ => " ON DELETE RESTRICT".into(),
}
}
fn source_jt_fk_suffix(_policies: &[OnDeletePolicy]) -> &'static str {
" ON DELETE CASCADE"
}
fn target_jt_fk_suffix(policies: &[OnDeletePolicy]) -> String {
match policy_for(policies, &DeleteSide::Target).unwrap_or(&DeleteAction::Restrict) {
DeleteAction::Restrict if needs_deferred_target_fk(policies) => " DEFERRABLE INITIALLY DEFERRED".into(),
DeleteAction::Restrict => " ON DELETE RESTRICT".into(),
DeleteAction::DeferredRestrict => " DEFERRABLE INITIALLY DEFERRED".into(),
DeleteAction::Allow => " ON DELETE CASCADE".into(),
DeleteAction::DeleteSource => " ON DELETE CASCADE".into(),
_ => " ON DELETE RESTRICT".into(),
}
}
fn emit_fk_constraints(schema: &SchemaDescriptor, type_map: &HashMap<String, (&str, &str)>, out: &mut String) {
let polymorphic = polymorphic_types(schema);
let mut emitted = false;
for t in &schema.types {
if t.abstract_ || t.junction {
continue;
}
for l in &t.links {
if l.is_junction_backed() {
continue;
}
if polymorphic.contains(&l.target) {
continue;
}
let Some((tgt_module, tgt_table)) = type_map.get(&l.target) else {
continue;
};
let cname = qi(&format!("{}_{}_fkey", t.table, l.name));
let suffix = target_fk_suffix(&l.on_delete);
out.push_str(&format!(
"ALTER TABLE {} ADD CONSTRAINT {} FOREIGN KEY ({}) REFERENCES {}(id){};\n",
qn(&t.module, &t.table),
cname,
qi(&format!("{}_id", l.name)),
qn(tgt_module, tgt_table),
suffix,
));
emitted = true;
}
}
if emitted {
out.push('\n');
}
}
pub fn junction_fk_constraints(
schema: &SchemaDescriptor,
type_map: &HashMap<String, (&str, &str)>,
) -> Vec<(String, String, String, String)> {
let polymorphic = polymorphic_types(schema);
let mut result = Vec::new();
for t in &schema.types {
if t.abstract_ || t.junction {
continue;
}
let pointers = t
.multilinks
.iter()
.map(|ml| (ml.name.as_str(), &ml.target, &ml.on_delete))
.chain(
t.links
.iter()
.filter(|l| l.is_junction_backed())
.map(|l| (l.name.as_str(), &l.target, &l.on_delete)),
);
for (name, target, on_delete) in pointers {
if polymorphic.contains(target) {
continue;
}
let Some((tgt_module, tgt_table)) = type_map.get(target) else {
continue;
};
let jt_name = format!("{}.{}", t.table, name);
let cname = format!("{}_{}_target_fkey", t.table, name);
let ddl = format!(
"ALTER TABLE {} ADD CONSTRAINT {} FOREIGN KEY (target) REFERENCES {}(id){};",
qn(&t.module, &jt_name),
qi(&cname),
qn(tgt_module, tgt_table),
target_jt_fk_suffix(on_delete),
);
result.push((t.module.clone(), jt_name, cname, ddl));
}
}
result
}
fn emit_junction_fk_constraints(schema: &SchemaDescriptor, type_map: &HashMap<String, (&str, &str)>, out: &mut String) {
let constraints = junction_fk_constraints(schema, type_map);
for (_, _, _, ddl) in &constraints {
out.push_str(ddl);
out.push('\n');
}
if !constraints.is_empty() {
out.push('\n');
}
}
fn emit_before_delete_trigger(fn_qname: &str, trigger_name: &str, table_qname: &str, body: &str, out: &mut String) {
out.push_str(&format!(
"CREATE OR REPLACE FUNCTION {fn_qname}()\n\
RETURNS trigger LANGUAGE plpgsql AS $$\n\
BEGIN\n"
));
out.push_str(body);
out.push_str(&format!(
"\n RETURN OLD;\nEND;\n$$;\n\n\
CREATE TRIGGER {trigger_name}\n\
BEFORE DELETE ON {table_qname}\n\
FOR EACH ROW EXECUTE FUNCTION {fn_qname}();\n\n"
));
}
fn emit_after_delete_trigger(fn_qname: &str, trigger_name: &str, table_qname: &str, body: &str, out: &mut String) {
out.push_str(&format!(
"CREATE OR REPLACE FUNCTION {fn_qname}()\n\
RETURNS trigger LANGUAGE plpgsql AS $$\n\
BEGIN\n"
));
out.push_str(body);
out.push_str(&format!(
"\n RETURN NULL;\nEND;\n$$;\n\n\
CREATE TRIGGER {trigger_name}\n\
AFTER DELETE ON {table_qname}\n\
FOR EACH ROW EXECUTE FUNCTION {fn_qname}();\n\n"
));
}
fn emit_after_mutation_trigger(
fn_qname: &str,
trigger_name: &str,
table_qname: &str,
events: &str,
body: &str,
out: &mut String,
) {
out.push_str(&format!(
"CREATE OR REPLACE FUNCTION {fn_qname}()\n\
RETURNS trigger LANGUAGE plpgsql AS $$\n\
BEGIN\n"
));
out.push_str(body);
out.push_str(&format!(
"\n RETURN NULL;\nEND;\n$$;\n\n\
CREATE OR REPLACE TRIGGER {trigger_name}\n\
AFTER {events} ON {table_qname}\n\
FOR EACH ROW EXECUTE FUNCTION {fn_qname}();\n\n"
));
}
fn target_table_qnames(
target: &str,
type_map: &HashMap<String, (&str, &str)>,
polymorphic: &HashSet<String>,
implementors: &HashMap<String, Vec<&TypeDescriptor>>,
) -> Option<Vec<String>> {
if polymorphic.contains(target) {
let impls = implementors.get(target)?;
if impls.is_empty() {
return None;
}
return Some(impls.iter().map(|i| qn(&i.module, &i.table)).collect());
}
let (tgt_module, tgt_table) = type_map.get(target)?;
Some(vec![qn(tgt_module, tgt_table)])
}
pub struct DeletionTriggerInfo {
pub table_module: String,
pub table_name: String,
pub trigger_name: String,
pub ddl: String,
}
fn link_source_trigger_infos(
schema: &SchemaDescriptor,
type_map: &HashMap<String, (&str, &str)>,
) -> Vec<DeletionTriggerInfo> {
let polymorphic = polymorphic_types(schema);
let implementors = interface_implementors(schema);
let mut result = Vec::new();
for t in &schema.types {
if t.abstract_ || t.junction {
continue;
}
for l in &t.links {
if l.is_junction_backed() {
continue;
}
let src_action = policy_for(&l.on_delete, &DeleteSide::Source).unwrap_or(&DeleteAction::Allow);
match src_action {
DeleteAction::Allow => continue,
DeleteAction::DeleteTarget | DeleteAction::DeleteTargetIfOrphan => {}
_ => continue,
}
let suffix = if matches!(src_action, DeleteAction::DeleteTargetIfOrphan) {
"del_orphan"
} else {
"del_target"
};
let tbl_qname = qn(&t.module, &t.table);
let col = qi(&format!("{}_id", l.name));
let Some(tgt_qnames) = target_table_qnames(&l.target, type_map, &polymorphic, &implementors) else {
continue;
};
let deletes = tgt_qnames
.iter()
.map(|tgt_qname| format!("DELETE FROM {tgt_qname} WHERE id = OLD.{col};"))
.collect::<Vec<_>>();
let body = if matches!(src_action, DeleteAction::DeleteTargetIfOrphan) {
let indented = deletes
.iter()
.map(|d| format!(" {d}"))
.collect::<Vec<_>>()
.join("\n");
format!(
" IF NOT EXISTS (\n SELECT 1 FROM {tbl_qname} WHERE {col} = OLD.{col} AND id != OLD.id\n ) THEN\n{indented}\n END IF;"
)
} else {
deletes
.iter()
.map(|d| format!(" {d}"))
.collect::<Vec<_>>()
.join("\n")
};
let (fname, fn_qname) = trigger_names(&t.module, &t.table, &l.name, suffix, &body);
let mut ddl = String::new();
emit_before_delete_trigger(&fn_qname, &qi(&fname), &tbl_qname, &body, &mut ddl);
result.push(DeletionTriggerInfo {
table_module: t.module.clone(),
table_name: t.table.clone(),
trigger_name: fname,
ddl,
});
}
}
result
}
fn emit_link_source_triggers(schema: &SchemaDescriptor, type_map: &HashMap<String, (&str, &str)>, out: &mut String) {
for info in link_source_trigger_infos(schema, type_map) {
out.push_str(&info.ddl);
}
}
fn emit_junction_tables(schema: &SchemaDescriptor, out: &mut String) {
for t in &schema.types {
if t.abstract_ || t.junction {
continue;
}
for ml in &t.multilinks {
emit_one_junction_table(
schema,
t,
&ml.name,
&ml.on_delete,
ml.through.as_deref(),
false,
ml.is_exclusive,
out,
);
}
for l in &t.links {
if !l.is_junction_backed() {
continue;
}
emit_one_junction_table(
schema,
t,
&l.name,
&l.on_delete,
l.through.as_deref(),
true,
l.is_exclusive,
out,
);
}
}
}
#[allow(clippy::too_many_arguments)]
fn emit_one_junction_table(
schema: &SchemaDescriptor,
t: &TypeDescriptor,
name: &str,
on_delete: &[OnDeletePolicy],
through: Option<&str>,
single: bool,
exclusive: bool,
out: &mut String,
) {
let jt_name = format!("{}.{}", t.table, name);
let src_suffix = source_jt_fk_suffix(on_delete);
out.push_str(&format!(
"CREATE TABLE {} (\n source uuid NOT NULL REFERENCES {}(id){},\n",
qn(&t.module, &jt_name),
qn(&t.module, &t.table),
src_suffix,
));
out.push_str(" target uuid NOT NULL,\n");
if let Some(through_qname) = through {
let through_td = schema
.types
.iter()
.find(|td| format!("{}::{}", td.module, td.name) == *through_qname);
if let Some(td) = through_td
&& td.junction
{
for p in &td.properties {
if p.name == "id" {
continue;
}
let not_null = if p.nullable { "" } else { " NOT NULL" };
out.push_str(&format!(" {} {}{},\n", qi(&p.name), p.pg_type, not_null));
}
}
}
if single {
out.push_str(" PRIMARY KEY (source)");
} else {
out.push_str(" PRIMARY KEY (source, target)");
}
if exclusive {
out.push_str(",\n UNIQUE (target)\n);\n\n");
} else {
out.push_str("\n);\n\n");
}
out.push_str(&cache_invalidate_trigger_sql(&qn(&t.module, &jt_name)));
out.push_str("\n\n");
}
fn multilink_deletion_trigger_infos(
schema: &SchemaDescriptor,
type_map: &HashMap<String, (&str, &str)>,
) -> Vec<DeletionTriggerInfo> {
let polymorphic = polymorphic_types(schema);
let implementors = interface_implementors(schema);
let mut result = Vec::new();
for t in &schema.types {
if t.abstract_ || t.junction {
continue;
}
for ml in &t.multilinks {
push_junction_deletion_triggers(
t,
&ml.name,
&ml.target,
&ml.on_delete,
type_map,
&polymorphic,
&implementors,
&mut result,
);
}
for l in &t.links {
if !l.is_junction_backed() {
continue;
}
push_junction_deletion_triggers(
t,
&l.name,
&l.target,
&l.on_delete,
type_map,
&polymorphic,
&implementors,
&mut result,
);
}
}
result
}
#[allow(clippy::too_many_arguments)]
fn push_junction_deletion_triggers(
t: &TypeDescriptor,
name: &str,
target: &str,
on_delete: &[OnDeletePolicy],
type_map: &HashMap<String, (&str, &str)>,
polymorphic: &HashSet<String>,
implementors: &HashMap<String, Vec<&TypeDescriptor>>,
result: &mut Vec<DeletionTriggerInfo>,
) {
let jt_name = format!("{}.{}", t.table, name);
let jt_qname = qn(&t.module, &jt_name);
let src_action = policy_for(on_delete, &DeleteSide::Source).unwrap_or(&DeleteAction::Allow);
if matches!(
src_action,
DeleteAction::DeleteTarget | DeleteAction::DeleteTargetIfOrphan
) && let Some(tgt_qnames) = target_table_qnames(target, type_map, polymorphic, implementors)
{
let suffix = if matches!(src_action, DeleteAction::DeleteTargetIfOrphan) {
"del_orphan"
} else {
"del_target"
};
let owner_qname = qn(&t.module, &t.table);
let orphan = matches!(src_action, DeleteAction::DeleteTargetIfOrphan);
let body = tgt_qnames
.iter()
.map(|tgt_qname| {
let mut delete = format!(
" DELETE FROM {tgt_qname} AS _tgt\n WHERE _tgt.id IN (SELECT target FROM {jt_qname} WHERE source = OLD.id)"
);
if orphan {
delete.push_str(&format!(
"\n AND NOT EXISTS (SELECT 1 FROM {jt_qname} WHERE target = _tgt.id AND source != OLD.id)"
));
}
delete.push(';');
delete
})
.collect::<Vec<_>>()
.join("\n");
let hash = fnv(&[&t.table, name, suffix, &body]);
let fname = format!("{}_{}_{}_{}", t.table, name, suffix, &hash[..8]);
let fn_qname = qn(&t.module, &fname);
let mut ddl = String::new();
emit_before_delete_trigger(&fn_qname, &qi(&fname), &owner_qname, &body, &mut ddl);
result.push(DeletionTriggerInfo {
table_module: t.module.clone(),
table_name: t.table.clone(),
trigger_name: fname,
ddl,
});
}
let tgt_action = policy_for(on_delete, &DeleteSide::Target).unwrap_or(&DeleteAction::Restrict);
if matches!(tgt_action, DeleteAction::DeleteSource) {
let src_qname = qn(&t.module, &t.table);
let body = format!(" DELETE FROM {src_qname} WHERE id = OLD.source;");
let (fname, fn_qname) = trigger_names(&t.module, &t.table, name, "del_source", &body);
let mut ddl = String::new();
emit_after_delete_trigger(&fn_qname, &qi(&fname), &jt_qname, &body, &mut ddl);
result.push(DeletionTriggerInfo {
table_module: t.module.clone(),
table_name: jt_name.clone(),
trigger_name: fname,
ddl,
});
}
}
#[allow(clippy::too_many_arguments)]
fn push_interface_link_triggers(
module: &str,
referencing: &str,
column: &str,
pointer: &str,
src_table: &str,
via_junction: bool,
on_delete: &[OnDeletePolicy],
impls: &[&TypeDescriptor],
result: &mut Vec<DeletionTriggerInfo>,
) {
let action = policy_for(on_delete, &DeleteSide::Target).unwrap_or(&DeleteAction::Restrict);
let col = qi(column);
let deferred = matches!(action, DeleteAction::DeferredRestrict) || needs_deferred_target_fk(on_delete);
let body = match action {
DeleteAction::Allow if via_junction => {
format!(" DELETE FROM {referencing} WHERE {col} = OLD.id;")
}
DeleteAction::Allow => format!(" UPDATE {referencing} SET {col} = NULL WHERE {col} = OLD.id;"),
DeleteAction::DeleteSource => format!(" DELETE FROM {referencing} WHERE {col} = OLD.id;"),
_ => format!(
" IF EXISTS (SELECT 1 FROM {referencing} WHERE {col} = OLD.id) THEN\n\
\x20 RAISE foreign_key_violation\n\
\x20 USING MESSAGE = 'update or delete on table \"' || TG_TABLE_NAME || '\" violates foreign key constraint on table {src_table}',\n\
\x20 DETAIL = format('Key (id)=(%s) is still referenced from table \"{src_table}\".', OLD.id);\n\
\x20 END IF;"
),
};
let suffix = match action {
DeleteAction::Allow => "ifl_allow",
DeleteAction::DeleteSource => "ifl_del_source",
_ => "ifl_restrict",
};
for impl_t in impls {
let hash = fnv(&[&impl_t.table, src_table, pointer, suffix, &body]);
let fname = format!("_ifl_{}_{}_{}", src_table, pointer, &hash[..8]);
let fn_qname = qn(module, &fname);
let impl_qname = qn(&impl_t.module, &impl_t.table);
let mut ddl = String::new();
if deferred {
ddl.push_str(&format!(
"CREATE OR REPLACE FUNCTION {fn_qname}()\n\
RETURNS trigger LANGUAGE plpgsql AS $$\n\
BEGIN\n{body}\n RETURN NULL;\nEND;\n$$;\n\n\
CREATE CONSTRAINT TRIGGER {}\n\
AFTER DELETE ON {impl_qname}\n\
DEFERRABLE INITIALLY DEFERRED\n\
FOR EACH ROW EXECUTE FUNCTION {fn_qname}();\n\n",
qi(&fname),
));
} else {
emit_before_delete_trigger(&fn_qname, &qi(&fname), &impl_qname, &body, &mut ddl);
}
result.push(DeletionTriggerInfo {
table_module: impl_t.module.clone(),
table_name: impl_t.table.clone(),
trigger_name: fname,
ddl,
});
}
}
fn interface_link_trigger_infos(schema: &SchemaDescriptor) -> Vec<DeletionTriggerInfo> {
let polymorphic = polymorphic_types(schema);
let implementors = interface_implementors(schema);
let mut result = Vec::new();
for t in &schema.types {
if t.abstract_ || t.junction {
continue;
}
for l in &t.links {
if !polymorphic.contains(&l.target) {
continue;
}
let Some(impls) = implementors.get(&l.target) else {
continue;
};
let via_junction = l.is_junction_backed();
let (referencing, column) = if via_junction {
(qn(&t.module, &format!("{}.{}", t.table, l.name)), "target".to_string())
} else {
(qn(&t.module, &t.table), format!("{}_id", l.name))
};
push_interface_link_triggers(
&t.module,
&referencing,
&column,
&l.name,
&t.table,
via_junction,
&l.on_delete,
impls,
&mut result,
);
}
for ml in &t.multilinks {
if !polymorphic.contains(&ml.target) {
continue;
}
let Some(impls) = implementors.get(&ml.target) else {
continue;
};
let referencing = qn(&t.module, &format!("{}.{}", t.table, ml.name));
push_interface_link_triggers(
&t.module,
&referencing,
"target",
&ml.name,
&t.table,
true,
&ml.on_delete,
impls,
&mut result,
);
}
}
result
}
fn emit_interface_link_triggers(schema: &SchemaDescriptor, out: &mut String) {
for info in interface_link_trigger_infos(schema) {
out.push_str(&info.ddl);
}
}
fn emit_multilink_deletion_triggers(
schema: &SchemaDescriptor,
type_map: &HashMap<String, (&str, &str)>,
out: &mut String,
) {
for info in multilink_deletion_trigger_infos(schema, type_map) {
out.push_str(&info.ddl);
}
}
pub fn deletion_policy_trigger_infos(
schema: &SchemaDescriptor,
type_map: &HashMap<String, (&str, &str)>,
) -> Vec<DeletionTriggerInfo> {
let mut result = link_source_trigger_infos(schema, type_map);
result.extend(multilink_deletion_trigger_infos(schema, type_map));
result.extend(interface_link_trigger_infos(schema));
result
}
pub fn signal_trigger_infos(schema: &SchemaDescriptor) -> Vec<DeletionTriggerInfo> {
use crate::schema::SearchBackend;
let mut result = Vec::new();
for t in &schema.types {
if t.abstract_ || t.junction || t.signals.is_empty() {
continue;
}
let qname = format!("{}::{}", t.module, t.name);
let qname_literal = format!("'{}'", qname.replace('\'', "''"));
let tbl_qname = qn(&t.module, &t.table);
let combined_on = t.signals.iter().fold(0u8, |acc, s| acc | s.on);
let mut events = Vec::new();
if combined_on & 1 != 0 {
events.push("INSERT");
}
if combined_on & 2 != 0 {
events.push("UPDATE");
}
if combined_on & 4 != 0 {
events.push("DELETE");
}
let events_str = events.join(" OR ");
let index_cols: Vec<String> = t
.vector_indexes
.iter()
.map(|vi| vi.column_name())
.chain(
t.search_indexes
.iter()
.filter(|si| si.backend == SearchBackend::Postgres)
.map(|si| si.column_name()),
)
.collect();
let update_guard = if combined_on & 2 != 0 && !index_cols.is_empty() {
let strip: String = index_cols
.iter()
.map(|c| format!(" - '{}'", c.replace('\'', "''")))
.collect();
format!(
" IF TG_OP = 'UPDATE' AND (to_jsonb(OLD){strip}) = (to_jsonb(NEW){strip}) THEN\n \
RETURN NULL;\n \
END IF;\n"
)
} else {
String::new()
};
let body = format!(
"{update_guard} INSERT INTO _pylon.\"SignalOutbox\" (type_name, operation, old_row, new_row)\n \
VALUES (\n \
{qname_literal},\n \
TG_OP,\n \
CASE WHEN TG_OP IN ('UPDATE', 'DELETE') THEN to_jsonb(OLD) ELSE NULL END,\n \
CASE WHEN TG_OP IN ('INSERT', 'UPDATE') THEN to_jsonb(NEW) ELSE NULL END\n \
);"
);
let hash = fnv(&[&t.table, "signal", &events_str, &body]);
let fname = format!("{}_signal_{}", t.table, &hash[..8]);
let fn_qname = qn(&t.module, &fname);
let mut ddl = String::new();
emit_after_mutation_trigger(&fn_qname, &qi(&fname), &tbl_qname, &events_str, &body, &mut ddl);
result.push(DeletionTriggerInfo {
table_module: t.module.clone(),
table_name: t.table.clone(),
trigger_name: fname,
ddl,
});
}
result
}
fn emit_signal_triggers(schema: &SchemaDescriptor, out: &mut String) {
for info in signal_trigger_infos(schema) {
out.push_str(&info.ddl);
}
}
fn constraint_column(t: &TypeDescriptor, pointer: &str) -> String {
if t.links.iter().any(|l| l.name == pointer && !l.is_junction_backed()) {
format!("{}_id", pointer)
} else {
pointer.to_string()
}
}
fn emit_unique_indexes(schema: &SchemaDescriptor, out: &mut String) {
let mut emitted = false;
for t in &schema.types {
if t.abstract_ || t.junction {
continue;
}
let qname = qn(&t.module, &t.table);
for p in &t.properties {
if p.is_exclusive && !p.is_pk {
out.push_str(&format!("CREATE UNIQUE INDEX ON {} ({});\n", qname, qi(&p.name),));
emitted = true;
}
}
for l in &t.links {
if l.is_exclusive && !l.is_junction_backed() {
out.push_str(&format!(
"CREATE UNIQUE INDEX ON {} ({});\n",
qname,
qi(&format!("{}_id", l.name)),
));
emitted = true;
}
}
for c in &t.constraints {
if let TypeConstraint::Exclusive {
pointers: fields,
unless,
} = c
{
let cols: Vec<String> = fields.iter().map(|f| qi(&constraint_column(t, f))).collect();
let where_clause = match unless.as_deref() {
Some(u) => {
let qualified = format!("{}::{}", t.module, t.name);
match crate::ir::compile_constraint_expr(u, &qualified, schema) {
Ok(compiled) => format!(" WHERE NOT ({compiled})"),
Err(_) => continue,
}
}
None => String::new(),
};
out.push_str(&format!(
"CREATE UNIQUE INDEX ON {} ({}){};\n",
qname,
cols.join(", "),
where_clause,
));
emitted = true;
}
}
}
if emitted {
out.push('\n');
}
}
pub fn scalar_check_constraints(schema: &SchemaDescriptor) -> Vec<(String, String, String, String)> {
let mut result = Vec::new();
for s in &schema.scalars {
for expr in &s.check_constraints {
let hash = fnv(&[&s.name, expr.as_str()]);
result.push((
s.module.clone(),
s.name.clone(),
format!("{}_{}_check", s.name, &hash[..8]),
expr.clone(),
));
}
}
result
}
pub fn scalar_check_clauses(schema: &SchemaDescriptor, module: &str, name: &str) -> String {
let clauses: Vec<String> = scalar_check_constraints(schema)
.into_iter()
.filter(|(m, n, _, _)| m == module && n == name)
.map(|(_, _, cname, expr)| format!(" CONSTRAINT {} CHECK ({})", qi(&cname), expr))
.collect();
if clauses.is_empty() {
String::new()
} else {
format!("\n{}", clauses.join("\n"))
}
}
pub fn check_constraints(schema: &SchemaDescriptor) -> Result<Vec<(String, String, String, String)>, PyQLError> {
let mut result = Vec::new();
for t in &schema.types {
if t.abstract_ || t.junction {
continue;
}
let qname = qn(&t.module, &t.table);
for p in &t.properties {
for expr in &p.check_constraints {
let hash = fnv(&[&t.table, &p.name, expr.as_str()]);
let cname = format!("{}_{}_{}_check", t.table, p.name, &hash[..8]);
result.push((
t.module.clone(),
t.table.clone(),
cname.clone(),
format!("ALTER TABLE {} ADD CONSTRAINT {} CHECK ({});", qname, qi(&cname), expr),
));
}
}
for c in &t.constraints {
if let TypeConstraint::Expression { expr } = c {
let qualified = format!("{}::{}", t.module, t.name);
let sql = crate::ir::compile_constraint_expr(expr, &qualified, schema).map_err(|e| {
PyQLError::Fragment(PyQLFragmentError {
message: format!("error in constraint on '{}': {}", qualified, e),
context: qualified.clone(),
position: crate::error::Position { line: 0, col: 0 },
})
})?;
let hash = fnv(&[&t.table, expr.as_str()]);
let cname = format!("{}_{}_check", t.table, &hash[..8]);
result.push((
t.module.clone(),
t.table.clone(),
cname.clone(),
format!("ALTER TABLE {} ADD CONSTRAINT {} CHECK ({});", qname, qi(&cname), sql),
));
}
}
}
Ok(result)
}
fn emit_check_constraints(schema: &SchemaDescriptor, out: &mut String) -> Result<(), PyQLError> {
let constraints = check_constraints(schema)?;
for (_, _, _, ddl) in &constraints {
out.push_str(ddl);
out.push('\n');
}
if !constraints.is_empty() {
out.push('\n');
}
Ok(())
}
pub(crate) fn index_body_and_predicate(
t: &TypeDescriptor,
pointers: &[String],
expression: Option<&str>,
unless: Option<&str>,
schema: &SchemaDescriptor,
) -> Result<(String, String), PyQLError> {
let qualified = format!("{}::{}", t.module, t.name);
let body = match expression {
Some(expr) => format!("({})", crate::ir::compile_constraint_expr(expr, &qualified, schema)?),
None => {
let mut cols = Vec::with_capacity(pointers.len());
for pointer in pointers {
cols.push(index_pointer_sql(t, pointer, &qualified, schema)?);
}
format!("({})", cols.join(", "))
}
};
let predicate = match unless {
Some(u) => format!(
" WHERE NOT ({})",
crate::ir::compile_constraint_expr(u, &qualified, schema)?
),
None => String::new(),
};
Ok((body, predicate))
}
pub(crate) fn index_pointer_sql(
t: &TypeDescriptor,
pointer: &str,
qualified: &str,
schema: &SchemaDescriptor,
) -> Result<String, PyQLError> {
if t.properties.iter().any(|p| p.name == pointer) {
return Ok(qi(pointer));
}
if t.links.iter().any(|l| l.name == pointer && !l.is_junction_backed()) {
return Ok(qi(&format!("{pointer}_id")));
}
Ok(format!(
"({})",
crate::ir::compile_constraint_expr(&format!(".{pointer}"), qualified, schema)?
))
}
fn emit_plain_indexes(schema: &SchemaDescriptor, out: &mut String) {
let mut emitted = false;
for t in &schema.types {
if t.abstract_ || t.junction {
continue;
}
let qname = qn(&t.module, &t.table);
for idx in &t.indexes {
let unique = if idx.unique { "UNIQUE " } else { "" };
let Ok((body, where_clause)) = index_body_and_predicate(
t,
&idx.pointers,
idx.expression.as_deref(),
idx.unless.as_deref(),
schema,
) else {
continue;
};
out.push_str(&format!(
"CREATE {}INDEX ON {} {}{};\n",
unique, qname, body, where_clause,
));
emitted = true;
}
}
if emitted {
out.push('\n');
}
}
fn emit_triggers(schema: &SchemaDescriptor, out: &mut String) -> Result<(), PyQLError> {
for info in user_trigger_infos(schema)? {
out.push_str(&info.ddl);
}
Ok(())
}
fn trigger_return_statement(timing: &str, on: u8) -> &'static str {
if timing == "After" {
return "RETURN NULL;";
}
let has_delete = on & 4 != 0;
let has_insert_or_update = on & (1 | 2) != 0;
match (has_delete, has_insert_or_update) {
(true, true) => "IF TG_OP = 'DELETE' THEN RETURN OLD; ELSE RETURN NEW; END IF;",
(true, false) => "RETURN OLD;",
_ => "RETURN NEW;",
}
}
fn trigger_ddl_name(table: &str, trig: &crate::schema::TriggerDescriptor, body_sql: &str) -> String {
let hash = fnv(&[
table,
&trig.on.to_string(),
trig.timing.as_str(),
trig.handler.as_str(),
body_sql,
]);
format!("{table}_{}", &hash[..12])
}
fn trigger_name_body(trig: &crate::schema::TriggerDescriptor, type_name: &str, schema: &SchemaDescriptor) -> String {
crate::query::compile_trigger_handler(&trig.handler, type_name, trig.on, schema)
.unwrap_or_else(|_| trig.handler.clone())
}
pub fn user_trigger_names(schema: &SchemaDescriptor) -> Vec<(String, String, String)> {
let mut result = Vec::new();
for t in &schema.types {
if t.abstract_ || t.junction {
continue;
}
let type_name = format!("{}::{}", t.module, t.name);
for trig in &t.triggers {
let body = trigger_name_body(trig, &type_name, schema);
result.push((
t.module.clone(),
t.table.clone(),
trigger_ddl_name(&t.table, trig, &body),
));
}
for (on, _) in REWRITE_EVENTS {
for deferred in [false, true] {
if let Some(name) = rewrite_trigger_name(t, on, schema, deferred) {
result.push((t.module.clone(), t.table.clone(), name));
}
}
}
}
result
}
const REWRITE_EVENTS: [(u8, &str); 2] = [(1, "INSERT"), (2, "UPDATE")];
fn rewrite_trigger_name(t: &TypeDescriptor, on: u8, schema: &SchemaDescriptor, deferred: bool) -> Option<String> {
let handlers: Vec<String> = t
.properties
.iter()
.map(|p| (&p.name, &p.rewrites))
.chain(t.links.iter().map(|l| (&l.name, &l.rewrites)))
.flat_map(|(name, rewrites)| {
rewrites
.iter()
.filter(move |rw| rw.on & on != 0)
.map(move |rw| format!("{name}:{}", rw.handler))
})
.collect();
if handlers.is_empty() {
return None;
}
let event = if on == 1 { "ins" } else { "upd" };
let assignments =
crate::ir::compile_rewrite_assignments(&format!("{}::{}", t.module, t.name), on, schema).unwrap_or_default();
let assignments: Vec<&crate::ir::RewriteAssignment> = assignments
.iter()
.filter(|a| a.reads_a_multi_link == deferred)
.collect();
if assignments.is_empty() {
return None;
}
let compiled = assignments.iter().map(|a| a.sql.clone()).collect::<Vec<_>>().join("\n");
let timing = if deferred { "rwa" } else { "rw" };
let hash = fnv(&[&t.table, event, &handlers.join("\n"), &compiled]);
Some(format!("{}_{timing}_{event}_{}", t.table, &hash[..12]))
}
fn rewrite_trigger_infos(schema: &SchemaDescriptor) -> Result<Vec<DeletionTriggerInfo>, PyQLError> {
let mut result = Vec::new();
for t in &schema.types {
if t.abstract_ || t.junction {
continue;
}
let type_name = format!("{}::{}", t.module, t.name);
for (on, event) in REWRITE_EVENTS {
let compiled = crate::ir::compile_rewrite_assignments(&type_name, on, schema).map_err(|e| {
PyQLError::Fragment(PyQLFragmentError {
message: format!("error in a rewrite of '{type_name}': {e}"),
context: type_name.clone(),
position: crate::error::Position { line: 0, col: 0 },
})
})?;
for deferred in [false, true] {
let Some(fname) = rewrite_trigger_name(t, on, schema, deferred) else {
continue;
};
let assignments: Vec<&crate::ir::RewriteAssignment> = compiled
.iter()
.filter(|assignment| assignment.reads_a_multi_link == deferred)
.collect();
if assignments.is_empty() {
continue;
}
let values = assignments
.iter()
.enumerate()
.map(|(i, assignment)| format!("{} AS \"v{i}\"", assignment.sql))
.collect::<Vec<_>>()
.join(",\n\t\t");
let fn_qname = qn(&t.module, &fname);
let table_qname = qn(&t.module, &t.table);
let globals_arg = qi(crate::ir::GLOBALS_ARG);
let body = if deferred {
let sets = assignments
.iter()
.enumerate()
.map(|(i, assignment)| format!("\t\t{} = _pylon_rewrites.\"v{i}\"", qi(&assignment.column)))
.collect::<Vec<_>>()
.join(",\n");
let stored = assignments
.iter()
.map(|assignment| qi(&assignment.column))
.collect::<Vec<_>>()
.join(", ");
let computed = (0..assignments.len())
.map(|i| format!("_pylon_rewrites.\"v{i}\""))
.collect::<Vec<_>>()
.join(", ");
format!(
"\tSELECT {values} INTO _pylon_rewrites;\n\
\tUPDATE {table_qname} SET\n\
{sets}\n\
\tWHERE \"id\" = NEW.\"id\"\n\
\t AND ({stored}) IS DISTINCT FROM ({computed});\n\
\tRETURN NULL;"
)
} else {
let sets = assignments
.iter()
.enumerate()
.map(|(i, assignment)| format!("\tNEW.{} := _pylon_rewrites.\"v{i}\";", qi(&assignment.column)))
.collect::<Vec<_>>()
.join("\n");
format!("\tSELECT {values} INTO _pylon_rewrites;\n{sets}\n\tRETURN NEW;")
};
let timing = if deferred { "AFTER" } else { "BEFORE" };
let ddl = format!(
"CREATE OR REPLACE FUNCTION {fn_qname}()\n\
RETURNS trigger LANGUAGE plpgsql AS $$\n\
DECLARE\n\
\t_pylon_rewrites record;\n\
\t{globals_arg} jsonb := nullif(current_setting('pylon.globals', true), '')::jsonb;\n\
BEGIN\n\
{body}\n\
END;\n\
$$;\n\n\
CREATE OR REPLACE TRIGGER {}\n\
{timing} {event} ON {table_qname}\n\
FOR EACH ROW EXECUTE FUNCTION {fn_qname}();\n\n",
qi(&fname),
);
result.push(DeletionTriggerInfo {
table_module: t.module.clone(),
table_name: t.table.clone(),
trigger_name: fname,
ddl,
});
}
}
}
Ok(result)
}
pub fn user_trigger_infos(schema: &SchemaDescriptor) -> Result<Vec<DeletionTriggerInfo>, PyQLError> {
let mut result = rewrite_trigger_infos(schema)?;
for t in &schema.types {
if t.abstract_ || t.junction {
continue;
}
let table_qname = qn(&t.module, &t.table);
let type_name = format!("{}::{}", t.module, t.name);
for trig in &t.triggers {
let fname = trigger_ddl_name(&t.table, trig, &trigger_name_body(trig, &type_name, schema));
let fn_qname = qn(&t.module, &fname);
let events = trigger_events(trig.on);
let timing = trigger_timing(&trig.timing);
let return_stmt = trigger_return_statement(&trig.timing, trig.on);
let body_sql =
crate::query::compile_trigger_handler(&trig.handler, &type_name, trig.on, schema).map_err(|e| {
let msg = format!("error in trigger handler for '{type_name}': {e}");
PyQLError::Fragment(PyQLFragmentError {
message: msg,
context: type_name.clone(),
position: crate::error::Position { line: 0, col: 0 },
})
})?;
let body_sql = body_sql.trim_end_matches(';');
let globals_arg = qi(crate::ir::GLOBALS_ARG);
let ddl = format!(
"CREATE OR REPLACE FUNCTION {fn_qname}()\n\
RETURNS trigger LANGUAGE plpgsql AS $$\n\
DECLARE\n\
\t_pylon_trigger_result record;\n\
\t{globals_arg} jsonb := nullif(current_setting('pylon.globals', true), '')::jsonb;\n\
BEGIN\n\
{body_sql} INTO _pylon_trigger_result;\n\
{return_stmt}\n\
END;\n\
$$;\n\n\
CREATE OR REPLACE TRIGGER {}\n\
{timing} {events} ON {table_qname}\n\
FOR EACH ROW EXECUTE FUNCTION {fn_qname}();\n\n",
qi(&fname),
);
result.push(DeletionTriggerInfo {
table_module: t.module.clone(),
table_name: t.table.clone(),
trigger_name: fname,
ddl,
});
}
}
Ok(result)
}
pub fn interface_view_ddl(schema: &SchemaDescriptor) -> Vec<String> {
interface_view_ddl_with_names(schema)
.into_iter()
.map(|(_, _, ddl)| ddl)
.collect()
}
pub fn interface_view_ddl_with_names(schema: &SchemaDescriptor) -> Vec<(String, String, String)> {
let mut implementors: HashMap<String, Vec<&TypeDescriptor>> = HashMap::new();
for t in &schema.types {
if !t.abstract_ {
for iface in &t.interfaces {
implementors.entry(iface.clone()).or_default().push(t);
}
}
}
let mut result = Vec::new();
for t in &schema.types {
if !(t.abstract_ && t.materialized) {
continue;
}
let key = format!("{}::{}", t.module, t.name);
let Some(impls) = implementors.get(&key) else { continue };
if impls.is_empty() {
continue;
}
let mut ddl = String::new();
emit_one_interface_view(t, impls, &mut ddl);
let ddl = ddl.trim().to_string();
if !ddl.is_empty() {
result.push((t.module.clone(), t.name.clone(), ddl));
}
}
result
}
fn interface_junction_view_name(iface_table: &str, link_name: &str) -> String {
format!("{}.{}", iface_table, link_name)
}
pub fn interface_junction_view_ddl_with_names(schema: &SchemaDescriptor) -> Vec<(String, String, String)> {
let implementors = interface_implementors(schema);
let mut result = Vec::new();
for t in &schema.types {
if !(t.abstract_ && t.materialized) {
continue;
}
let key = format!("{}::{}", t.module, t.name);
let Some(impls) = implementors.get(&key) else { continue };
if impls.is_empty() {
continue;
}
let pointers: Vec<(&str, Option<&str>)> = t
.multilinks
.iter()
.map(|ml| (ml.name.as_str(), ml.through.as_deref()))
.chain(
t.links
.iter()
.filter(|l| l.is_junction_backed())
.map(|l| (l.name.as_str(), l.through.as_deref())),
)
.collect();
for (link_name, through) in pointers {
let mut columns = vec!["source".to_string(), "target".to_string()];
if let Some(through_qname) = through
&& let Some(td) = schema
.types
.iter()
.find(|td| format!("{}::{}", td.module, td.name) == through_qname && td.junction)
{
columns.extend(td.properties.iter().filter(|p| p.name != "id").map(|p| qi(&p.name)));
}
let column_list = columns.join(", ");
let view_name = interface_junction_view_name(&t.table, link_name);
let selects: Vec<String> = impls
.iter()
.map(|impl_t| {
let jt_name = format!("{}.{}", impl_t.table, link_name);
format!(" SELECT {} FROM {}", column_list, qn(&impl_t.module, &jt_name))
})
.collect();
let ddl = format!(
"CREATE VIEW {} AS\n{};",
qn(&t.module, &view_name),
selects.join("\n UNION ALL\n"),
);
result.push((t.module.clone(), view_name, ddl));
}
}
result
}
fn emit_interface_junction_views(schema: &SchemaDescriptor, out: &mut String) {
for (_, _, ddl) in interface_junction_view_ddl_with_names(schema) {
out.push_str(&ddl);
out.push_str("\n\n");
}
}
pub fn function_ddl(schema: &SchemaDescriptor) -> Result<Vec<String>, crate::error::PyQLError> {
function_ddl_with_names(schema).map(|v| v.into_iter().map(|(_, _, ddl)| ddl).collect())
}
pub fn function_ddl_with_names(
schema: &SchemaDescriptor,
) -> Result<Vec<(String, String, String)>, crate::error::PyQLError> {
schema
.functions
.iter()
.map(|fd| emit_one_function(fd, schema).map(|ddl| (fd.module.clone(), fd.name.clone(), ddl)))
.collect()
}
pub fn scalar_function_ddl_with_names(
schema: &SchemaDescriptor,
) -> Result<Vec<(String, String, String)>, crate::error::PyQLError> {
schema
.functions
.iter()
.filter(|fd| !fd.return_is_object)
.map(|fd| emit_one_function(fd, schema).map(|ddl| (fd.module.clone(), fd.name.clone(), ddl)))
.collect()
}
pub fn object_function_ddl_with_names(
schema: &SchemaDescriptor,
) -> Result<Vec<(String, String, String)>, crate::error::PyQLError> {
schema
.functions
.iter()
.filter(|fd| fd.return_is_object)
.map(|fd| emit_one_function(fd, schema).map(|ddl| (fd.module.clone(), fd.name.clone(), ddl)))
.collect()
}
fn emit_one_interface_view(t: &TypeDescriptor, impls: &[&TypeDescriptor], out: &mut String) {
let cols: Vec<String> = t
.properties
.iter()
.map(|p| qi(&p.name))
.chain(
t.links
.iter()
.filter(|l| !l.is_junction_backed())
.map(|l| qi(&format!("{}_id", l.name))),
)
.collect();
let col_list = cols.join(", ");
let selects: Vec<String> = impls
.iter()
.map(|impl_t| format!(" SELECT {} FROM {}", col_list, qn(&impl_t.module, &impl_t.table)))
.collect();
out.push_str(&format!("CREATE VIEW {} AS\n", qn(&t.module, &t.table)));
out.push_str(&selects.join("\n UNION ALL\n"));
out.push_str(";\n\n");
}
fn emit_interface_views(schema: &SchemaDescriptor, out: &mut String) {
let implementors = interface_implementors(schema);
for t in &schema.types {
if !(t.abstract_ && t.materialized) {
continue;
}
let key = format!("{}::{}", t.module, t.name);
let Some(impls) = implementors.get(&key) else { continue };
if impls.is_empty() {
continue;
}
emit_one_interface_view(t, impls, out);
}
}
pub struct ExclTriggerInfo {
pub fn_module: String,
pub fn_name: String,
pub fn_ddl: String,
pub impl_module: String,
pub impl_table: String,
pub ins_trigger_name: String,
pub ins_ddl: String,
pub upd_trigger_name: String,
pub upd_ddl: String,
}
fn excl_fn_name(iface_table: &str, fields: &[String]) -> String {
format!("_excl_{}_{}", iface_table, fields.join("_"))
}
fn make_excl_info(
iface: &TypeDescriptor,
fields: &[String],
columns: &[String],
impl_t: &TypeDescriptor,
impls: &[&TypeDescriptor],
) -> ExclTriggerInfo {
let fn_name = excl_fn_name(&iface.table, fields);
let fn_qname = format!("{}.{}", pg_schema(&iface.module), qi(&fn_name));
let view_qname = if iface.abstract_ {
qn(&iface.module, &iface.table)
} else {
let spanned_columns = std::iter::once("id".to_string())
.chain(columns.iter().cloned())
.map(|c| qi(&c))
.collect::<Vec<_>>()
.join(", ");
let branches = impls
.iter()
.map(|t| format!("SELECT {spanned_columns} FROM {}", qn(&t.module, &t.table)))
.collect::<Vec<_>>()
.join(" UNION ALL ");
format!("({branches}) AS \"_spanned\"")
};
let tbl_qname = qn(&impl_t.module, &impl_t.table);
let field_conds: Vec<String> = columns.iter().map(|c| format!("{} = NEW.{}", qi(c), qi(c))).collect();
let where_clause = format!("{} AND \"id\" <> NEW.\"id\"", field_conds.join(" AND "));
let detail_keys = columns.join(", ");
let detail_vals = columns
.iter()
.map(|c| format!("NEW.{}::text", qi(c)))
.collect::<Vec<_>>()
.join(" || ', ' || ");
let fn_ddl = format!(
"CREATE OR REPLACE FUNCTION {}()\n\
RETURNS trigger LANGUAGE plpgsql AS $$\n\
BEGIN\n\
IF EXISTS (\n\
SELECT 1 FROM {}\n\
WHERE {}\n\
) THEN\n\
RAISE unique_violation\n\
USING CONSTRAINT = '{}',\n\
DETAIL = format('Key ({})=(%s) already exists.', {});\n\
END IF;\n\
RETURN NEW;\n\
END;\n\
$$;",
fn_qname, view_qname, where_clause, fn_name, detail_keys, detail_vals,
);
let ins_trigger_name = format!("{}_ins", fn_name);
let upd_trigger_name = format!("{}_upd", fn_name);
let of_cols = columns.iter().map(|c| qi(c)).collect::<Vec<_>>().join(", ");
let when_clause = columns
.iter()
.map(|c| format!("OLD.{} IS DISTINCT FROM NEW.{}", qi(c), qi(c)))
.collect::<Vec<_>>()
.join(" OR ");
let ins_ddl = format!(
"CREATE CONSTRAINT TRIGGER {}\n\
AFTER INSERT ON {}\n\
DEFERRABLE INITIALLY DEFERRED\n\
FOR EACH ROW EXECUTE FUNCTION {}();",
qi(&ins_trigger_name),
tbl_qname,
fn_qname,
);
let upd_ddl = format!(
"CREATE CONSTRAINT TRIGGER {}\n\
AFTER UPDATE OF {} ON {}\n\
DEFERRABLE INITIALLY DEFERRED\n\
FOR EACH ROW WHEN ({})\n\
EXECUTE FUNCTION {}();",
qi(&upd_trigger_name),
of_cols,
tbl_qname,
when_clause,
fn_qname,
);
ExclTriggerInfo {
fn_module: iface.module.clone(),
fn_name,
fn_ddl,
impl_module: impl_t.module.clone(),
impl_table: impl_t.table.clone(),
ins_trigger_name,
ins_ddl,
upd_trigger_name,
upd_ddl,
}
}
fn make_excl_junction_info(
iface: &TypeDescriptor,
link_name: &str,
impl_t: &TypeDescriptor,
impls: &[&TypeDescriptor],
) -> ExclTriggerInfo {
let fn_name = excl_fn_name(&iface.table, std::slice::from_ref(&link_name.to_string()));
let fn_qname = format!("{}.{}", pg_schema(&iface.module), qi(&fn_name));
let view_qname = if iface.abstract_ {
qn(&iface.module, &interface_junction_view_name(&iface.table, link_name))
} else {
let branches = impls
.iter()
.map(|t| {
format!(
"SELECT \"source\", \"target\" FROM {}",
qn(&t.module, &format!("{}.{}", t.table, link_name))
)
})
.collect::<Vec<_>>()
.join(" UNION ALL ");
format!("({branches}) AS \"_spanned\"")
};
let jt_name = format!("{}.{}", impl_t.table, link_name);
let jt_qname = qn(&impl_t.module, &jt_name);
let where_clause = "\"target\" = NEW.\"target\" AND \"source\" <> NEW.\"source\"";
let fn_ddl = format!(
"CREATE OR REPLACE FUNCTION {}()\n\
RETURNS trigger LANGUAGE plpgsql AS $$\n\
BEGIN\n\
IF EXISTS (\n\
SELECT 1 FROM {}\n\
WHERE {}\n\
) THEN\n\
RAISE unique_violation\n\
USING CONSTRAINT = '{}',\n\
DETAIL = format('Key (target)=(%s) already exists.', NEW.\"target\"::text);\n\
END IF;\n\
RETURN NEW;\n\
END;\n\
$$;",
fn_qname, view_qname, where_clause, fn_name,
);
let ins_trigger_name = format!("{}_ins", fn_name);
let upd_trigger_name = format!("{}_upd", fn_name);
let ins_ddl = format!(
"CREATE CONSTRAINT TRIGGER {}\n\
AFTER INSERT ON {}\n\
DEFERRABLE INITIALLY DEFERRED\n\
FOR EACH ROW EXECUTE FUNCTION {}();",
qi(&ins_trigger_name),
jt_qname,
fn_qname,
);
let upd_ddl = format!(
"CREATE CONSTRAINT TRIGGER {}\n\
AFTER UPDATE OF \"target\" ON {}\n\
DEFERRABLE INITIALLY DEFERRED\n\
FOR EACH ROW WHEN (OLD.\"target\" IS DISTINCT FROM NEW.\"target\")\n\
EXECUTE FUNCTION {}();",
qi(&upd_trigger_name),
jt_qname,
fn_qname,
);
ExclTriggerInfo {
fn_module: iface.module.clone(),
fn_name,
fn_ddl,
impl_module: impl_t.module.clone(),
impl_table: jt_name,
ins_trigger_name,
ins_ddl,
upd_trigger_name,
upd_ddl,
}
}
pub fn interface_exclusive_trigger_infos(schema: &SchemaDescriptor) -> Vec<ExclTriggerInfo> {
let implementors = interface_implementors(schema);
let mut result = Vec::new();
for t in &schema.types {
let key = format!("{}::{}", t.module, t.name);
let spans_subtypes = !t.abstract_ && schema.types.iter().any(|sub| sub.bases.contains(&key));
if !(spans_subtypes || t.abstract_ && t.materialized) {
continue;
}
let Some(impls) = implementors.get(&key) else { continue };
if impls.is_empty() {
continue;
}
let interfaces: Vec<&TypeDescriptor> = if t.abstract_ {
vec![]
} else {
schema
.types
.iter()
.filter(|i| {
i.abstract_ && i.materialized && t.interfaces.contains(&format!("{}::{}", i.module, i.name))
})
.collect()
};
let declared_by_interface = |name: &str| {
interfaces.iter().any(|i| {
i.properties.iter().any(|p| p.name == name && p.is_exclusive)
|| i.links.iter().any(|l| l.name == name && l.is_exclusive)
|| i.multilinks.iter().any(|ml| ml.name == name && ml.is_exclusive)
})
};
for p in &t.properties {
if !p.is_exclusive || p.is_pk || declared_by_interface(&p.name) {
continue;
}
let fields = vec![p.name.clone()];
for impl_t in impls {
result.push(make_excl_info(t, &fields, &fields, impl_t, impls));
}
}
for l in &t.links {
if !l.is_exclusive || declared_by_interface(&l.name) {
continue;
}
if l.is_junction_backed() {
for impl_t in impls {
result.push(make_excl_junction_info(t, &l.name, impl_t, impls));
}
continue;
}
let fields = vec![format!("{}_id", l.name)];
for impl_t in impls {
result.push(make_excl_info(t, &fields, &fields, impl_t, impls));
}
}
for ml in &t.multilinks {
if !ml.is_exclusive || declared_by_interface(&ml.name) {
continue;
}
for impl_t in impls {
result.push(make_excl_junction_info(t, &ml.name, impl_t, impls));
}
}
for c in &t.constraints {
if let TypeConstraint::Exclusive { pointers: fields, .. } = c {
let from_interface = interfaces.iter().any(|i| {
i.constraints
.iter()
.any(|ic| matches!(ic, TypeConstraint::Exclusive { pointers, .. } if pointers == fields))
});
if from_interface {
continue;
}
let columns: Vec<String> = fields.iter().map(|f| constraint_column(t, f)).collect();
for impl_t in impls {
result.push(make_excl_info(t, fields, &columns, impl_t, impls));
}
}
}
}
result
}
fn emit_interface_exclusive_triggers(schema: &SchemaDescriptor, out: &mut String) {
use std::collections::HashSet;
let mut fn_emitted: HashSet<String> = HashSet::new();
for info in interface_exclusive_trigger_infos(schema) {
if fn_emitted.insert(info.fn_name.clone()) {
out.push_str(&info.fn_ddl);
out.push_str("\n\n");
}
out.push_str(&info.ins_ddl);
out.push('\n');
out.push_str(&info.upd_ddl);
out.push_str("\n\n");
}
}
fn emit_scalar_functions(schema: &SchemaDescriptor, out: &mut String) -> Result<(), PyQLError> {
for fd in schema.functions.iter().filter(|fd| !fd.return_is_object) {
let ddl = emit_one_function(fd, schema)?;
out.push_str(&ddl);
out.push('\n');
}
Ok(())
}
fn emit_object_functions(schema: &SchemaDescriptor, out: &mut String) -> Result<(), PyQLError> {
for fd in schema.functions.iter().filter(|fd| fd.return_is_object) {
let ddl = emit_one_function(fd, schema)?;
out.push_str(&ddl);
out.push('\n');
}
Ok(())
}
fn emit_one_function(fd: &FunctionDescriptor, schema: &SchemaDescriptor) -> Result<String, PyQLError> {
use crate::ir::compile_fn_body;
use crate::sql::emit_fn_body;
let ir_output = compile_fn_body(fd, schema).map_err(|e| {
let msg = format!("error in function '{}::{}' body: {}", fd.module, fd.name, e);
PyQLError::Fragment(PyQLFragmentError {
message: msg,
context: format!("{}::{}", fd.module, fd.name),
position: crate::error::Position { line: 0, col: 0 },
})
})?;
let mut body_sql = emit_fn_body(&ir_output);
if fd.return_is_object
&& !matches!(
ir_output.stmt,
crate::ir::IrStmt::Insert(_) | crate::ir::IrStmt::Update(_) | crate::ir::IrStmt::Delete(_)
)
&& let Some(columns) = fn_return_column_names(fd, schema)
{
body_sql = format!(
"SELECT {} FROM (\n {}\n ) AS \"_returned\"",
columns.join(", "),
body_sql
);
}
let mut param_parts: Vec<String> = Vec::with_capacity(fd.params.len() + 1);
if ir_output.uses_globals_arg {
param_parts.push(format!("{} jsonb", qi(crate::ir::GLOBALS_ARG)));
}
param_parts.extend(fd.params.iter().map(|p| format!("{} {}", qi(&p.name), p.pg_type)));
let params_sql = param_parts.join(", ");
let returns_sql = if fd.return_is_object {
let type_columns = emit_fn_return_table(fd, schema);
if fd.return_is_set {
format!("TABLE({})", type_columns)
} else {
format!("TABLE({})", type_columns)
}
} else if fd.return_is_set {
format!("SETOF {}", fd.return_pg_type)
} else {
fd.return_pg_type.clone()
};
let volatility_kw = match fd.volatility.as_str() {
"immutable" => "IMMUTABLE",
"stable" => "STABLE",
_ => "VOLATILE",
};
Ok(format!(
"CREATE OR REPLACE FUNCTION {fn_name}({params})\nRETURNS {returns}\nLANGUAGE SQL {vol}\nAS $$\n {body}\n$$;\n",
fn_name = qn(&fd.module, &fd.name),
params = params_sql,
returns = returns_sql,
vol = volatility_kw,
body = body_sql,
))
}
fn fn_return_column_names(fd: &FunctionDescriptor, schema: &SchemaDescriptor) -> Option<Vec<String>> {
let td = schema
.types
.iter()
.find(|t| format!("{}::{}", t.module, t.name) == fd.return_pg_type)?;
let mut names: Vec<String> = Vec::new();
if fd.return_is_polymorphic {
names.push(qi("__type__"));
}
names.extend(td.properties.iter().map(|p| qi(&p.name)));
names.extend(
td.links
.iter()
.filter(|l| !l.is_junction_backed())
.map(|l| qi(&format!("{}_id", l.name))),
);
Some(names)
}
fn emit_fn_return_table(fd: &FunctionDescriptor, schema: &SchemaDescriptor) -> String {
let type_name = &fd.return_pg_type; let td = schema
.types
.iter()
.find(|t| format!("{}::{}", t.module, t.name) == *type_name);
let Some(td) = td else {
return "__type__ text, id uuid".to_string();
};
let mut cols: Vec<String> = Vec::new();
if fd.return_is_polymorphic {
cols.push("__type__ text".to_string());
}
for p in &td.properties {
let pg_type = p.pg_type.strip_prefix("__nt__:").map(|_| "jsonb").unwrap_or(&p.pg_type);
cols.push(format!("{} {}", qi(&p.name), pg_type));
}
for l in &td.links {
if l.is_junction_backed() {
continue;
}
cols.push(format!("{} uuid", qi(&format!("{}_id", l.name))));
}
cols.join(", ")
}
fn emit_vector_columns(schema: &SchemaDescriptor, out: &mut String) {
for td in &schema.types {
if td.abstract_ || td.vector_indexes.is_empty() {
continue;
}
for vi in &td.vector_indexes {
out.push_str(&format!(
"ALTER TABLE {} ADD COLUMN IF NOT EXISTS {} vector({});\n",
qn(&td.module, &td.table),
qi(&vi.column_name()),
vi.dimensions,
));
}
}
if schema
.types
.iter()
.any(|t| !t.abstract_ && !t.vector_indexes.is_empty())
{
out.push('\n');
}
}
fn emit_vector_indexes(schema: &SchemaDescriptor, out: &mut String) {
for td in &schema.types {
if td.abstract_ || td.vector_indexes.is_empty() {
continue;
}
for vi in &td.vector_indexes {
let index_name = match &vi.index_name {
None => format!("{}__vector__", td.table),
Some(name) => format!("{}__vector_{}__", td.table, name),
};
out.push_str(&format!(
"CREATE INDEX IF NOT EXISTS {} ON {} USING hnsw ({} {});\n",
qi(&index_name),
qn(&td.module, &td.table),
qi(&vi.column_name()),
vi.ops_class(),
));
}
}
}
fn emit_search_columns(schema: &SchemaDescriptor, out: &mut String) {
use crate::schema::SearchBackend;
let mut emitted = false;
for td in &schema.types {
if td.abstract_ {
continue;
}
for si in &td.search_indexes {
if si.backend != SearchBackend::Postgres {
continue;
}
let parts: Vec<String> = si
.pointers
.iter()
.map(|sf| {
let col = qi(&sf.name);
let w = sf.weight.as_str();
format!("setweight(to_tsvector('english', coalesce({col}, '')), '{w}')")
})
.collect();
let expr = if parts.len() == 1 {
parts.into_iter().next().unwrap()
} else {
parts.join(" || ")
};
out.push_str(&format!(
"ALTER TABLE {} ADD COLUMN IF NOT EXISTS {} tsvector GENERATED ALWAYS AS ({}) STORED;\n",
qn(&td.module, &td.table),
qi(&si.column_name()),
expr,
));
emitted = true;
}
}
if emitted {
out.push('\n');
}
}
fn emit_search_indexes(schema: &SchemaDescriptor, out: &mut String) {
use crate::schema::SearchBackend;
for td in &schema.types {
if td.abstract_ {
continue;
}
for si in &td.search_indexes {
if si.backend != SearchBackend::Postgres {
continue;
}
let col = si.column_name();
let index_name = match &si.index_name {
None => format!("{}__search__", td.table),
Some(name) => format!("{}__search_{}__", td.table, name),
};
out.push_str(&format!(
"CREATE INDEX IF NOT EXISTS {} ON {} USING gin ({});\n",
qi(&index_name),
qn(&td.module, &td.table),
qi(&col),
));
}
}
}
pub fn compile_index_fetch(
type_name: &str,
index_name: Option<&str>,
schema: &SchemaDescriptor,
) -> Result<String, PyQLError> {
let td = schema
.types
.iter()
.find(|t| format!("{}::{}", t.module, t.name) == type_name)
.ok_or_else(|| {
PyQLError::Fragment(PyQLFragmentError {
message: format!("compile_index_fetch: unknown type '{}'", type_name),
context: type_name.to_string(),
position: crate::error::Position { line: 0, col: 0 },
})
})?;
let vi = td
.vector_indexes
.iter()
.find(|vi| vi.index_name.as_deref() == index_name)
.ok_or_else(|| {
let key = index_name.unwrap_or("<default>");
PyQLError::Fragment(PyQLFragmentError {
message: format!("compile_index_fetch: no vector index '{}' on type '{}'", key, type_name),
context: type_name.to_string(),
position: crate::error::Position { line: 0, col: 0 },
})
})?;
let field_exprs = vi
.pointers
.iter()
.map(|f| {
let pg_type = td
.properties
.iter()
.find(|p| p.name == *f)
.map(|p| p.pg_type.as_str())
.unwrap_or("text");
let col = qi(f);
if pg_type == "text" {
col
} else {
format!("{}::text", col)
}
})
.collect::<Vec<_>>();
let concat = if field_exprs.len() == 1 {
field_exprs.into_iter().next().unwrap()
} else {
format!("concat_ws(E'\\n', {})", field_exprs.join(", "))
};
Ok(format!(
"SELECT \"id\", {} AS source_text\nFROM {}\nWHERE \"id\" = ANY($1::uuid[])",
concat,
qn(&td.module, &td.table),
))
}
pub fn compile_search_index_fetch(
type_name: &str,
index_name: Option<&str>,
schema: &SchemaDescriptor,
) -> Result<String, PyQLError> {
let td = schema
.types
.iter()
.find(|t| format!("{}::{}", t.module, t.name) == type_name)
.ok_or_else(|| {
PyQLError::Fragment(PyQLFragmentError {
message: format!("compile_search_index_fetch: unknown type '{}'", type_name),
context: type_name.to_string(),
position: crate::error::Position { line: 0, col: 0 },
})
})?;
use crate::schema::SearchBackend;
let si = td
.search_indexes
.iter()
.find(|si| si.index_name.as_deref() == index_name && si.backend != SearchBackend::Postgres)
.ok_or_else(|| {
let key = index_name.unwrap_or("<default>");
PyQLError::Fragment(PyQLFragmentError {
message: format!(
"compile_search_index_fetch: no remote SearchIndex '{}' on type '{}'",
key, type_name
),
context: type_name.to_string(),
position: crate::error::Position { line: 0, col: 0 },
})
})?;
let field_exprs = si
.pointers
.iter()
.map(|sf| {
let pg_type = td
.properties
.iter()
.find(|p| p.name == sf.name)
.map(|p| p.pg_type.as_str())
.unwrap_or("text");
let col = qi(&sf.name);
if pg_type == "text" {
col
} else {
format!("{}::text", col)
}
})
.collect::<Vec<_>>();
let concat = if field_exprs.len() == 1 {
field_exprs.into_iter().next().unwrap()
} else {
format!("concat_ws(E'\\n', {})", field_exprs.join(", "))
};
Ok(format!(
"SELECT \"id\", {} AS source_text\nFROM {}\nWHERE \"id\" = ANY($1::uuid[])",
concat,
qn(&td.module, &td.table),
))
}
#[cfg(test)]
mod tests {
use super::*;
use crate::schema::{
DeleteAction, DeleteSide, FunctionDescriptor, FunctionParamDescriptor, LinkDescriptor, MultiLinkDescriptor,
OnDeletePolicy, PropertyDescriptor, SchemaDescriptor, TypeDescriptor,
};
fn person_type() -> TypeDescriptor {
TypeDescriptor {
name: "Person".into(),
module: "default".into(),
table: "Person".into(),
abstract_: false,
materialized: false,
description: None,
parents: vec![],
interfaces: vec![],
bases: vec![],
properties: vec![
PropertyDescriptor {
name: "id".into(),
pg_type: "uuid".into(),
nullable: false,
default_sql: Some("uuidv7()".into()),
default_pyql: None,
description: None,
check_constraints: vec![],
is_exclusive: true,
is_pk: true,
is_readonly: true,
rewrites: vec![],
tuple_members: None,
column_type: None,
},
PropertyDescriptor {
name: "age".into(),
pg_type: "int8".into(),
nullable: true,
default_sql: None,
default_pyql: None,
description: None,
check_constraints: vec![],
is_exclusive: false,
is_pk: false,
is_readonly: false,
rewrites: vec![],
tuple_members: None,
column_type: None,
},
],
links: vec![],
multilinks: vec![],
computed: vec![],
constraints: vec![],
indexes: vec![],
partition: None,
vector_indexes: vec![],
search_indexes: vec![],
triggers: vec![],
junction: false,
signals: vec![],
}
}
fn minimal_schema(fns: Vec<FunctionDescriptor>) -> SchemaDescriptor {
SchemaDescriptor {
types: vec![person_type()],
scalars: vec![],
enums: vec![],
named_tuples: vec![],
globals: vec![],
functions: fns,
aliases: vec![],
channels: vec![],
..Default::default()
}
}
fn interface_schema(referencing: Vec<TypeDescriptor>) -> SchemaDescriptor {
fn bare(name: &str, interfaces: Vec<String>, abstract_: bool, materialized: bool) -> TypeDescriptor {
TypeDescriptor {
name: name.into(),
module: "default".into(),
table: name.into(),
abstract_,
materialized,
description: None,
parents: vec![],
interfaces,
bases: vec![],
properties: vec![],
links: vec![],
multilinks: vec![],
computed: vec![],
constraints: vec![],
indexes: vec![],
partition: None,
vector_indexes: vec![],
search_indexes: vec![],
triggers: vec![],
junction: false,
signals: vec![],
}
}
let mut types = vec![
bare("Account", vec![], true, true),
bare("Individual", vec!["default::Account".into()], false, false),
bare("Organization", vec!["default::Account".into()], false, false),
];
types.extend(referencing);
SchemaDescriptor {
types,
scalars: vec![],
enums: vec![],
named_tuples: vec![],
globals: vec![],
functions: vec![],
aliases: vec![],
channels: vec![],
..Default::default()
}
}
fn link_to_account(name: &str, on_delete: Vec<OnDeletePolicy>, through: Option<String>) -> LinkDescriptor {
LinkDescriptor {
name: name.into(),
target: "default::Account".into(),
nullable: true,
description: None,
default_pyql: None,
is_exclusive: false,
is_readonly: false,
rewrites: vec![],
on_delete,
through,
}
}
fn referencing_type(
name: &str,
links: Vec<LinkDescriptor>,
multilinks: Vec<MultiLinkDescriptor>,
) -> TypeDescriptor {
TypeDescriptor {
name: name.into(),
module: "default".into(),
table: name.into(),
abstract_: false,
materialized: false,
description: None,
parents: vec![],
interfaces: vec![],
bases: vec![],
properties: vec![],
links,
multilinks,
computed: vec![],
constraints: vec![],
indexes: vec![],
partition: None,
vector_indexes: vec![],
search_indexes: vec![],
triggers: vec![],
junction: false,
signals: vec![],
}
}
#[test]
fn test_multilink_on_an_interface_gets_a_union_view_over_the_implementors() {
let mut schema = interface_schema(vec![referencing_type("Email", vec![], vec![])]);
for t in schema.types.iter_mut() {
if t.name == "Account" || t.name == "Individual" || t.name == "Organization" {
t.multilinks.push(MultiLinkDescriptor {
name: "emails".into(),
target: "default::Email".into(),
through: None,
nullable: true,
description: None,
default_pyql: None,
on_delete: vec![],
is_exclusive: false,
});
}
}
let ddl = export_schema(&schema).unwrap();
assert!(
ddl.contains("CREATE VIEW \"public\".\"Account.emails\""),
"no union view for the interface's multi-link, got:\n{}",
ddl
);
for implementor in ["Individual", "Organization"] {
assert!(
ddl.contains(&format!("FROM \"public\".\"{}.emails\"", implementor)),
"union view misses {}, got:\n{}",
implementor,
ddl
);
}
}
fn schema_with_a_link_to_a_type_with_subtypes() -> SchemaDescriptor {
let mut link = link_to_account("holder", vec![], None);
link.target = "default::Individual".into();
let mut schema = interface_schema(vec![referencing_type("Session", vec![link], vec![])]);
let mut staff = schema.types[1].clone();
staff.name = "Staff".into();
staff.table = "Staff".into();
staff.bases = vec!["default::Individual".into()];
schema.types.push(staff);
schema
}
#[test]
fn an_exclusive_multilink_on_a_type_with_subtypes_is_unique_across_them() {
let mut schema = schema_with_a_link_to_a_type_with_subtypes();
for t in schema
.types
.iter_mut()
.filter(|t| t.name == "Individual" || t.name == "Staff")
{
t.multilinks.push(MultiLinkDescriptor {
name: "keys".into(),
target: "default::Session".into(),
through: None,
nullable: false,
description: None,
default_pyql: None,
on_delete: vec![],
is_exclusive: true,
});
}
let ddl = export_schema(&schema).unwrap();
for table in ["Individual.keys", "Staff.keys"] {
assert!(
ddl.contains(&format!("CREATE TABLE \"public\".\"{table}\"")),
"got:\n{ddl}"
);
assert!(
ddl.contains(&format!("AFTER INSERT ON \"public\".\"{table}\"")),
"missing exclusive trigger on {table}, got:\n{ddl}"
);
}
assert!(ddl.contains("UNIQUE (target)"), "got:\n{ddl}");
assert!(
ddl.contains(
"SELECT \"source\", \"target\" FROM \"public\".\"Individual.keys\" UNION ALL \
SELECT \"source\", \"target\" FROM \"public\".\"Staff.keys\""
),
"the trigger must check every subtype's junction, got:\n{ddl}"
);
}
#[test]
fn a_link_to_a_type_with_subtypes_emits_no_foreign_key() {
let ddl = export_schema(&schema_with_a_link_to_a_type_with_subtypes()).unwrap();
assert!(!ddl.contains("Session_holder_fkey"), "got:\n{}", ddl);
for table in ["Individual", "Staff"] {
assert!(
ddl.contains(&format!("BEFORE DELETE ON \"public\".\"{table}\"")),
"missing enforcement trigger on {table}, got:\n{ddl}"
);
}
}
#[test]
fn test_link_to_interface_emits_no_foreign_key() {
let schema = interface_schema(vec![referencing_type(
"Session",
vec![link_to_account("account", vec![], None)],
vec![],
)]);
let ddl = export_schema(&schema).unwrap();
assert!(
!ddl.contains("Session_account_fkey"),
"a link targeting an interface must not get a FK, got:\n{}",
ddl
);
assert!(ddl.contains("CREATE VIEW \"public\".\"Account\""), "got:\n{}", ddl);
}
#[test]
fn test_link_to_interface_enforces_restrict_on_every_implementor() {
let schema = interface_schema(vec![referencing_type(
"Session",
vec![link_to_account("account", vec![], None)],
vec![],
)]);
let ddl = export_schema(&schema).unwrap();
for implementor in ["Individual", "Organization"] {
assert!(
ddl.contains(&format!("BEFORE DELETE ON \"public\".\"{}\"", implementor)),
"missing enforcement trigger on {}, got:\n{}",
implementor,
ddl
);
}
assert!(ddl.contains("RAISE foreign_key_violation"), "got:\n{}", ddl);
assert!(
ddl.contains("SELECT 1 FROM \"public\".\"Session\" WHERE \"account_id\" = OLD.id"),
"got:\n{}",
ddl
);
}
#[test]
fn test_interface_link_allow_nulls_the_column_but_clears_a_junction_row() {
let allow = vec![OnDeletePolicy {
side: DeleteSide::Target,
action: DeleteAction::Allow,
}];
let schema = interface_schema(vec![
referencing_type(
"AuditEntry",
vec![link_to_account("actor", allow.clone(), None)],
vec![],
),
referencing_type(
"Watchlist",
vec![],
vec![MultiLinkDescriptor {
name: "watched".into(),
target: "default::Account".into(),
through: None,
nullable: true,
description: None,
default_pyql: None,
on_delete: allow,
is_exclusive: false,
}],
),
]);
let ddl = export_schema(&schema).unwrap();
assert!(
ddl.contains("UPDATE \"public\".\"AuditEntry\" SET \"actor_id\" = NULL WHERE \"actor_id\" = OLD.id;"),
"single link Allow must null the column, got:\n{}",
ddl
);
assert!(
ddl.contains("DELETE FROM \"public\".\"Watchlist.watched\" WHERE \"target\" = OLD.id;"),
"multi-link Allow must drop the junction row, got:\n{}",
ddl
);
}
#[test]
fn test_multilink_to_interface_junction_target_has_no_reference() {
let schema = interface_schema(vec![referencing_type(
"Watchlist",
vec![],
vec![MultiLinkDescriptor {
name: "watched".into(),
target: "default::Account".into(),
through: None,
nullable: true,
description: None,
default_pyql: None,
on_delete: vec![],
is_exclusive: false,
}],
)]);
let ddl = export_schema(&schema).unwrap();
assert!(
ddl.contains(" target uuid NOT NULL,\n"),
"junction target must be a bare uuid column, got:\n{}",
ddl
);
assert!(
!ddl.contains("target uuid NOT NULL REFERENCES \"public\".\"Account\""),
"got:\n{}",
ddl
);
}
#[test]
fn test_source_side_cascade_deletes_from_implementors_not_the_view() {
let schema = interface_schema(vec![referencing_type(
"Session",
vec![link_to_account(
"account",
vec![OnDeletePolicy {
side: DeleteSide::Source,
action: DeleteAction::DeleteTarget,
}],
None,
)],
vec![],
)]);
let ddl = export_schema(&schema).unwrap();
assert!(
ddl.contains("DELETE FROM \"public\".\"Individual\" WHERE id = OLD.\"account_id\";"),
"got:\n{}",
ddl
);
assert!(
ddl.contains("DELETE FROM \"public\".\"Organization\" WHERE id = OLD.\"account_id\";"),
"got:\n{}",
ddl
);
assert!(
!ddl.contains("DELETE FROM \"public\".\"Account\" WHERE id ="),
"must never delete through the interface view, got:\n{}",
ddl
);
}
#[test]
fn test_link_to_concrete_type_still_gets_its_foreign_key() {
let mut schema = interface_schema(vec![referencing_type(
"Session",
vec![LinkDescriptor {
name: "owner".into(),
target: "default::Individual".into(),
nullable: true,
description: None,
default_pyql: None,
is_exclusive: false,
is_readonly: false,
rewrites: vec![],
on_delete: vec![],
through: None,
}],
vec![],
)]);
schema.types.retain(|t| t.name != "Organization");
let ddl = export_schema(&schema).unwrap();
assert!(
ddl.contains("FOREIGN KEY (\"owner_id\") REFERENCES \"public\".\"Individual\"(id)"),
"got:\n{}",
ddl
);
}
#[test]
fn test_emit_one_table_includes_cache_invalidate_trigger() {
let mut out = String::new();
emit_one_table(&person_type(), None, &mut out);
assert!(
out.contains("CREATE OR REPLACE TRIGGER pylon_cache_invalidate\n AFTER INSERT OR UPDATE OR DELETE ON \"public\".\"Person\""),
"got:\n{out}"
);
}
#[test]
fn test_emit_table_compiles_a_pyql_default() {
let mut td = person_type();
td.properties[1].default_sql = None;
td.properties[1].default_pyql = Some("std::uuid_generate_v7()".into());
td.properties[1].pg_type = "uuid".into();
let schema = SchemaDescriptor {
types: vec![td.clone()],
..minimal_schema(vec![])
};
let mut out = String::new();
emit_one_table(&td, Some(&schema), &mut out);
assert!(
out.contains("\"age\" uuid NULL DEFAULT uuidv7()") || out.contains("DEFAULT uuidv7()"),
"got:\n{out}"
);
}
#[test]
fn test_export_and_migration_agree_on_a_pyql_default() {
let mut td = person_type();
td.properties[1].default_sql = None;
td.properties[1].default_pyql = Some("std::uuid_generate_v7()".into());
td.properties[1].pg_type = "uuid".into();
let schema = SchemaDescriptor {
types: vec![td.clone()],
..minimal_schema(vec![])
};
let mut out = String::new();
emit_one_table(&td, Some(&schema), &mut out);
let from_migration = crate::diff::resolve_default_for_test(&td.properties[1], &schema);
assert_eq!(from_migration.as_deref(), Some("uuidv7()"));
assert!(
out.contains(&format!("DEFAULT {}", from_migration.unwrap())),
"export DDL disagrees with the migration path:\n{out}"
);
}
fn trig(on: u8, timing: &str, handler: &str) -> crate::schema::TriggerDescriptor {
crate::schema::TriggerDescriptor {
on,
timing: timing.into(),
handler: handler.into(),
}
}
fn schema_with_trigger(trigger: crate::schema::TriggerDescriptor) -> SchemaDescriptor {
let mut t = person_type();
t.triggers = vec![trigger];
SchemaDescriptor {
types: vec![t],
scalars: vec![],
enums: vec![],
named_tuples: vec![],
globals: vec![],
functions: vec![],
aliases: vec![],
channels: vec![],
..Default::default()
}
}
#[test]
fn test_trigger_new_anchor_resolves_to_new_alias() {
let schema = schema_with_trigger(trig(1, "After", "update Person set { age := __new__.age }"));
let ddl = export_schema(&schema).unwrap();
assert!(ddl.contains("NEW.\"age\""), "got:\n{ddl}");
}
#[test]
fn test_trigger_that_would_refire_itself_is_rejected() {
let schema = schema_with_trigger(trig(1, "After", "insert Person { age := __new__.age }"));
let err = export_schema(&schema).unwrap_err();
assert!(err.to_string().contains("is recursive"), "got: {err}");
}
#[test]
fn test_recursive_insert_wrapped_in_a_select_shape_is_still_caught() {
let schema = schema_with_trigger(trig(
1,
"After",
"select (insert Person { age := __new__.age }) { age }",
));
let err = export_schema(&schema).unwrap_err();
assert!(err.to_string().contains("is recursive"), "got: {err}");
}
#[test]
fn test_recursive_check_is_scoped_to_the_triggers_own_events() {
let schema = schema_with_trigger(trig(4, "After", "insert Person { age := __old__.age }"));
assert!(export_schema(&schema).is_ok());
}
#[test]
fn test_trigger_old_anchor_resolves_to_old_alias() {
let schema = schema_with_trigger(trig(4, "After", "update Person set { age := __old__.age }"));
let ddl = export_schema(&schema).unwrap();
assert!(ddl.contains("OLD.\"age\""), "got:\n{ddl}");
}
#[test]
fn test_trigger_update_can_reference_both_new_and_old() {
let schema = schema_with_trigger(trig(2, "After", "select Person filter (__new__.age = __old__.age)"));
let ddl = export_schema(&schema).unwrap();
assert!(ddl.contains("NEW.\"age\""), "got:\n{ddl}");
assert!(ddl.contains("OLD.\"age\""), "got:\n{ddl}");
}
#[test]
fn test_trigger_insert_only_cannot_reference_old() {
let schema = schema_with_trigger(trig(1, "After", "update Person set { age := __old__.age }"));
let err = export_schema(&schema).unwrap_err();
assert!(err.to_string().contains("__old__ cannot be used"), "got: {err}");
}
#[test]
fn test_trigger_delete_only_cannot_reference_new() {
let schema = schema_with_trigger(trig(4, "After", "update Person set { age := __new__.age }"));
let err = export_schema(&schema).unwrap_err();
assert!(err.to_string().contains("__new__ cannot be used"), "got: {err}");
}
#[test]
fn test_trigger_combined_insert_update_cannot_reference_old() {
let ok = schema_with_trigger(trig(3, "After", "select Person filter (__new__.age > 0)"));
assert!(export_schema(&ok).is_ok());
let bad = schema_with_trigger(trig(3, "After", "select Person filter (__old__.age > 0)"));
let err = export_schema(&bad).unwrap_err();
assert!(err.to_string().contains("__old__ cannot be used"), "got: {err}");
}
#[test]
fn test_trigger_after_timing_returns_null() {
let schema = schema_with_trigger(trig(1, "After", "update Person set { age := __new__.age }"));
let ddl = export_schema(&schema).unwrap();
assert!(ddl.contains("RETURN NULL;"), "got:\n{ddl}");
}
#[test]
fn test_trigger_before_timing_returns_new_or_old_appropriately() {
let insert_only = schema_with_trigger(trig(1, "Before", "update Person set { age := __new__.age }"));
let ddl = export_schema(&insert_only).unwrap();
assert!(
ddl.contains("RETURN NEW;"),
"insert-only Before should return NEW, got:\n{ddl}"
);
let delete_only = schema_with_trigger(trig(4, "Before", "update Person set { age := __old__.age }"));
let ddl = export_schema(&delete_only).unwrap();
assert!(
ddl.contains("RETURN OLD;"),
"delete-only Before should return OLD, got:\n{ddl}"
);
let combined = schema_with_trigger(trig(5, "Before", "update Person set { age := 1 }"));
let ddl = export_schema(&combined).unwrap();
assert!(
ddl.contains("IF TG_OP = 'DELETE' THEN RETURN OLD; ELSE RETURN NEW; END IF;"),
"got:\n{ddl}"
);
}
#[test]
fn test_multiple_triggers_get_separate_functions() {
let mut t = person_type();
t.triggers = vec![
trig(1, "After", "update Person set { age := __new__.age }"),
trig(4, "Before", "update Person set { age := __old__.age }"),
];
let schema = SchemaDescriptor {
types: vec![t],
scalars: vec![],
enums: vec![],
named_tuples: vec![],
globals: vec![],
functions: vec![],
aliases: vec![],
channels: vec![],
..Default::default()
};
let ddl = export_schema(&schema).unwrap();
let fn_count = ddl.matches("CREATE OR REPLACE FUNCTION \"public\".\"Person_").count();
assert_eq!(fn_count, 2, "expected one function per Trigger(...), got:\n{ddl}");
}
#[test]
fn test_emit_scalar_function_ddl() {
let fd = FunctionDescriptor {
name: "mysum".into(),
module: "math".into(),
params: vec![
FunctionParamDescriptor {
name: "a".into(),
pg_type: "int8".into(),
},
FunctionParamDescriptor {
name: "b".into(),
pg_type: "int8".into(),
},
],
return_pg_type: "int8".into(),
return_is_object: false,
return_is_set: false,
return_is_polymorphic: false,
volatility: "immutable".into(),
body: "a + b".into(),
};
let schema = minimal_schema(vec![fd.clone()]);
let ddl = emit_one_function(&fd, &schema).unwrap();
assert!(
ddl.contains("CREATE OR REPLACE FUNCTION \"math\".\"mysum\""),
"got:\n{}",
ddl
);
assert!(ddl.contains("\"a\" int8, \"b\" int8"), "got:\n{}", ddl);
assert!(ddl.contains("RETURNS int8"), "got:\n{}", ddl);
assert!(ddl.contains("IMMUTABLE"), "got:\n{}", ddl);
assert!(ddl.contains("SELECT"), "got:\n{}", ddl);
}
#[test]
fn test_emit_setof_function_ddl() {
let fd = FunctionDescriptor {
name: "counters".into(),
module: "default".into(),
params: vec![],
return_pg_type: "int8".into(),
return_is_object: false,
return_is_set: true,
return_is_polymorphic: false,
volatility: "stable".into(),
body: "1".into(),
};
let schema = minimal_schema(vec![fd.clone()]);
let ddl = emit_one_function(&fd, &schema).unwrap();
assert!(ddl.contains("RETURNS SETOF int8"), "got:\n{}", ddl);
assert!(ddl.contains("STABLE"), "got:\n{}", ddl);
}
#[test]
fn test_emit_object_function_ddl() {
let fd = FunctionDescriptor {
name: "adults".into(),
module: "default".into(),
params: vec![],
return_pg_type: "default::Person".into(),
return_is_object: true,
return_is_set: true,
return_is_polymorphic: false,
volatility: "stable".into(),
body: "select Person filter .age > 18".into(),
};
let schema = minimal_schema(vec![fd.clone()]);
let ddl = emit_one_function(&fd, &schema).unwrap();
assert!(
ddl.contains("CREATE OR REPLACE FUNCTION \"public\".\"adults\"()"),
"got:\n{}",
ddl
);
assert!(ddl.contains("RETURNS TABLE("), "got:\n{}", ddl);
assert!(ddl.contains("\"id\" uuid"), "got:\n{}", ddl);
assert!(ddl.contains("\"age\" int8"), "got:\n{}", ddl);
assert!(ddl.contains("STABLE"), "got:\n{}", ddl);
assert!(ddl.contains("SELECT * FROM"), "got:\n{}", ddl);
}
#[test]
fn test_object_function_can_reference_its_own_parameter() {
let fd = FunctionDescriptor {
name: "older_than".into(),
module: "default".into(),
params: vec![FunctionParamDescriptor {
name: "min_age".into(),
pg_type: "int8".into(),
}],
return_pg_type: "default::Person".into(),
return_is_object: true,
return_is_set: true,
return_is_polymorphic: false,
volatility: "stable".into(),
body: "select Person filter .age > min_age".into(),
};
let schema = minimal_schema(vec![fd.clone()]);
let ddl = emit_one_function(&fd, &schema).unwrap();
assert!(
ddl.contains("\"min_age\" int8"),
"parameter missing from the signature, got:\n{}",
ddl
);
assert!(
ddl.contains("\"min_age\")") || ddl.contains("= \"min_age\"") || ddl.contains("> \"min_age\""),
"parameter not referenced in the body, got:\n{}",
ddl
);
}
#[test]
fn test_a_function_reading_a_global_takes_the_globals_argument() {
use crate::schema::GlobalDescriptor;
let reader = FunctionDescriptor {
name: "cutoff".into(),
module: "default".into(),
params: vec![],
return_pg_type: "timestamptz".into(),
return_is_object: false,
return_is_set: false,
return_is_polymorphic: false,
volatility: "stable".into(),
body: "global snapshot_at ?? datetime_of_transaction()".into(),
};
let plain = FunctionDescriptor {
name: "bump".into(),
module: "default".into(),
params: vec![FunctionParamDescriptor {
name: "n".into(),
pg_type: "int8".into(),
}],
return_pg_type: "int8".into(),
return_is_object: false,
return_is_set: false,
return_is_polymorphic: false,
volatility: "immutable".into(),
body: "n + 1".into(),
};
let mut schema = minimal_schema(vec![reader.clone(), plain.clone()]);
schema.globals.push(GlobalDescriptor {
name: "snapshot_at".into(),
module: "default".into(),
scalar_type: "datetime".into(),
required: false,
default_expr: None,
computed_expr: None,
});
let reader_ddl = emit_one_function(&reader, &schema).unwrap();
assert!(
reader_ddl.contains("\"__pylon_json_globals__\" jsonb"),
"the global reader should take the argument, got:\n{}",
reader_ddl
);
assert!(
reader_ddl.contains("__pylon_json_globals__ ->> 'default::snapshot_at'"),
"the body should read the global out of it, got:\n{}",
reader_ddl
);
assert!(
!reader_ddl.contains("$1"),
"no parameter placeholder should survive, got:\n{}",
reader_ddl
);
let plain_ddl = emit_one_function(&plain, &schema).unwrap();
assert!(
!plain_ddl.contains("__pylon_json_globals__"),
"a function that cannot reach a global should be untouched, got:\n{}",
plain_ddl
);
}
#[test]
fn test_the_globals_argument_is_forwarded_to_a_callee_that_needs_it() {
use crate::schema::GlobalDescriptor;
let reader = FunctionDescriptor {
name: "cutoff".into(),
module: "default".into(),
params: vec![],
return_pg_type: "timestamptz".into(),
return_is_object: false,
return_is_set: false,
return_is_polymorphic: false,
volatility: "stable".into(),
body: "global snapshot_at ?? datetime_of_transaction()".into(),
};
let caller = FunctionDescriptor {
name: "is_past".into(),
module: "default".into(),
params: vec![],
return_pg_type: "bool".into(),
return_is_object: false,
return_is_set: false,
return_is_polymorphic: false,
volatility: "stable".into(),
body: "cutoff() < datetime_of_transaction()".into(),
};
let mut schema = minimal_schema(vec![reader, caller.clone()]);
schema.globals.push(GlobalDescriptor {
name: "snapshot_at".into(),
module: "default".into(),
scalar_type: "datetime".into(),
required: false,
default_expr: None,
computed_expr: None,
});
let ddl = emit_one_function(&caller, &schema).unwrap();
assert!(
ddl.contains("\"__pylon_json_globals__\" jsonb"),
"a caller of a global reader needs the argument too, got:\n{}",
ddl
);
assert!(
ddl.contains("cutoff\"((__pylon_json_globals__))") || ddl.contains("__pylon_json_globals__)"),
"it should forward the argument, got:\n{}",
ddl
);
}
#[test]
fn test_calling_an_object_function_in_an_expression_says_why() {
let fd = FunctionDescriptor {
name: "adults".into(),
module: "default".into(),
params: vec![],
return_pg_type: "default::Person".into(),
return_is_object: true,
return_is_set: true,
return_is_polymorphic: false,
volatility: "stable".into(),
body: "select Person filter .age > 18".into(),
};
let schema = minimal_schema(vec![fd]);
let ast = crate::parse::parse("select Person { x := adults() }").unwrap();
assert!(
crate::ir::compile(&ast, &schema).is_ok(),
"an object-returning call should stand as a pointer's own value"
);
let ast = crate::parse::parse("select Person { x := count(adults()) }").unwrap();
let err = match crate::ir::compile(&ast, &schema) {
Err(e) => e.to_string(),
Ok(_) => panic!("expected the call to be rejected"),
};
assert!(
err.contains("returns objects") && err.contains("subject of a select"),
"expected an explanation of the restriction, got: {}",
err
);
assert!(
!err.contains("does not exist"),
"the function does exist; the message should not claim otherwise: {}",
err
);
}
#[test]
fn test_emit_sequence_scalar_ddl() {
use crate::schema::ScalarDescriptor;
let schema = SchemaDescriptor {
types: vec![],
scalars: vec![ScalarDescriptor {
name: "OrderNumber".into(),
module: "default".into(),
base: "Sequence".into(),
pg_type: "int8".into(),
check_constraints: vec![],
is_sequence: true,
}],
enums: vec![],
named_tuples: vec![],
globals: vec![],
functions: vec![],
aliases: vec![],
channels: vec![],
..Default::default()
};
let ddl = export_schema(&schema).unwrap();
assert!(
ddl.contains("CREATE SEQUENCE \"public\".\"OrderNumber_seq\""),
"got:\n{}",
ddl
);
assert!(
ddl.contains("CREATE DOMAIN \"public\".\"OrderNumber\" AS int8"),
"got:\n{}",
ddl
);
let seq_pos = ddl.find("CREATE SEQUENCE").unwrap();
let dom_pos = ddl.find("CREATE DOMAIN").unwrap();
assert!(seq_pos < dom_pos, "sequence must appear before domain");
}
#[test]
fn test_registered_scalar_domain_is_used_as_the_column_type() {
use crate::schema::ScalarDescriptor;
let schema = SchemaDescriptor {
types: vec![TypeDescriptor {
name: "Contact".into(),
module: "default".into(),
table: "Contact".into(),
abstract_: false,
materialized: true,
description: None,
parents: vec![],
interfaces: vec![],
bases: vec![],
properties: vec![PropertyDescriptor {
name: "email".into(),
pg_type: "text".into(),
nullable: false,
default_sql: None,
default_pyql: None,
description: None,
check_constraints: vec![],
is_exclusive: false,
is_pk: false,
is_readonly: false,
rewrites: vec![],
tuple_members: None,
column_type: Some("\"public\".\"EmailStr\"".into()),
}],
links: vec![],
multilinks: vec![],
computed: vec![],
constraints: vec![],
indexes: vec![],
partition: None,
vector_indexes: vec![],
search_indexes: vec![],
triggers: vec![],
junction: false,
signals: vec![],
}],
scalars: vec![ScalarDescriptor {
name: "EmailStr".into(),
module: "default".into(),
base: "Str".into(),
pg_type: "text".into(),
check_constraints: vec!["value ~ '^[^@]+@[^@]+\\.[^@]+$'".into()],
is_sequence: false,
}],
enums: vec![],
named_tuples: vec![],
globals: vec![],
functions: vec![],
aliases: vec![],
channels: vec![],
..Default::default()
};
let ddl = export_schema(&schema).unwrap();
assert!(
ddl.contains("CREATE DOMAIN \"public\".\"EmailStr\" AS text\n CONSTRAINT ")
&& ddl.contains("CHECK (value ~ '^[^@]+@[^@]+\\.[^@]+$')"),
"got:\n{}",
ddl
);
assert!(
ddl.contains("\"email\" \"public\".\"EmailStr\" NOT NULL"),
"column must use the domain type, not the plain base type — got:\n{}",
ddl
);
}
fn account_interface_schema() -> SchemaDescriptor {
let mut account = TypeDescriptor {
name: "Account".into(),
module: "default".into(),
table: "Account".into(),
abstract_: true,
materialized: true,
description: None,
parents: vec![],
interfaces: vec![],
bases: vec![],
properties: vec![PropertyDescriptor {
name: "email".into(),
pg_type: "text".into(),
nullable: false,
default_sql: None,
default_pyql: None,
description: None,
check_constraints: vec![],
is_exclusive: true,
is_pk: false,
is_readonly: false,
rewrites: vec![],
tuple_members: None,
column_type: None,
}],
links: vec![],
multilinks: vec![],
computed: vec![],
constraints: vec![],
indexes: vec![],
partition: None,
vector_indexes: vec![],
search_indexes: vec![],
triggers: vec![],
junction: false,
signals: vec![],
};
let mut individual = account.clone();
individual.name = "Individual".into();
individual.table = "Individual".into();
individual.abstract_ = false;
individual.materialized = true;
individual.interfaces = vec!["default::Account".into()];
let mut organization = individual.clone();
organization.name = "Organization".into();
organization.table = "Organization".into();
account.constraints = vec![]; SchemaDescriptor {
types: vec![account, individual, organization],
scalars: vec![],
enums: vec![],
named_tuples: vec![],
globals: vec![],
functions: vec![],
aliases: vec![],
channels: vec![],
..Default::default()
}
}
#[test]
fn test_interface_exclusive_property_gets_per_table_index_and_cross_table_trigger() {
let schema = account_interface_schema();
let ddl = export_schema(&schema).unwrap();
assert!(
ddl.contains("CREATE UNIQUE INDEX ON \"public\".\"Individual\" (\"email\")"),
"got:\n{ddl}"
);
assert!(
ddl.contains("CREATE UNIQUE INDEX ON \"public\".\"Organization\" (\"email\")"),
"got:\n{ddl}"
);
assert_eq!(
ddl.matches("CREATE OR REPLACE FUNCTION \"public\".\"_excl_Account_email\"")
.count(),
1,
"the trigger function must be emitted exactly once, shared by every implementor; got:\n{ddl}"
);
assert!(ddl.contains("SELECT 1 FROM \"public\".\"Account\""), "got:\n{ddl}");
assert!(
ddl.contains(
"CREATE CONSTRAINT TRIGGER \"_excl_Account_email_ins\"\nAFTER INSERT ON \"public\".\"Individual\""
),
"got:\n{ddl}"
);
assert!(
ddl.contains(
"CREATE CONSTRAINT TRIGGER \"_excl_Account_email_ins\"\nAFTER INSERT ON \"public\".\"Organization\""
),
"got:\n{ddl}"
);
assert!(ddl.contains("DEFERRABLE INITIALLY DEFERRED"), "got:\n{ddl}");
assert!(
ddl.contains("CREATE CONSTRAINT TRIGGER \"_excl_Account_email_upd\"\nAFTER UPDATE OF \"email\""),
"the UPDATE trigger must only fire when the exclusive column itself changes; got:\n{ddl}"
);
}
#[test]
fn test_junction_backed_exclusive_link_gets_a_cross_implementor_helper_view_and_trigger() {
use crate::schema::LinkDescriptor;
let mut schema = account_interface_schema();
for t in &mut schema.types {
t.properties.retain(|p| p.name != "email");
if t.name == "Account" || t.name == "Individual" || t.name == "Organization" {
t.links.push(LinkDescriptor {
name: "owner".into(),
target: "default::Person".into(),
nullable: true,
through: Some("default::AccountOwner".into()),
description: None,
default_pyql: None,
is_exclusive: true,
is_readonly: false,
rewrites: vec![],
on_delete: vec![],
});
}
}
let views = interface_junction_view_ddl_with_names(&schema);
assert_eq!(views.len(), 1, "expected exactly one helper view, got: {views:?}");
let (view_module, view_name, view_ddl) = &views[0];
assert_eq!(view_module, "default");
assert_eq!(view_name, "Account.owner");
assert!(
view_ddl.contains("SELECT source, target FROM \"public\".\"Individual.owner\""),
"got:\n{view_ddl}"
);
assert!(
view_ddl.contains("SELECT source, target FROM \"public\".\"Organization.owner\""),
"got:\n{view_ddl}"
);
let infos = interface_exclusive_trigger_infos(&schema);
assert_eq!(
infos.len(),
2,
"expected one entry per implementor, got: {}",
infos.len()
);
assert!(
infos.iter().any(|i| i.impl_table == "Individual.owner"
&& i.ins_ddl.contains("AFTER INSERT ON \"public\".\"Individual.owner\"")),
"the constraint trigger must attach to the implementor's own *junction* table, not the owner table"
);
assert!(
infos
.iter()
.any(|i| i.fn_ddl.contains("SELECT 1 FROM \"public\".\"Account.owner\"")),
"the trigger function must query the helper view"
);
assert!(infos.iter().all(|i| i.upd_ddl.contains("AFTER UPDATE OF \"target\"")));
}
#[test]
fn test_target_fk_suffix_forces_deferrable_when_source_side_deletes_target() {
let policies = vec![OnDeletePolicy {
side: DeleteSide::Source,
action: DeleteAction::DeleteTarget,
}];
assert_eq!(target_fk_suffix(&policies), " DEFERRABLE INITIALLY DEFERRED");
let policies = vec![OnDeletePolicy {
side: DeleteSide::Source,
action: DeleteAction::DeleteTargetIfOrphan,
}];
assert_eq!(target_fk_suffix(&policies), " DEFERRABLE INITIALLY DEFERRED");
}
#[test]
fn test_target_fk_suffix_unaffected_when_no_source_side_policy() {
assert_eq!(target_fk_suffix(&[]), " ON DELETE RESTRICT");
let policies = vec![OnDeletePolicy {
side: DeleteSide::Target,
action: DeleteAction::Allow,
}];
assert_eq!(target_fk_suffix(&policies), " ON DELETE SET NULL");
}
#[test]
fn test_target_jt_fk_suffix_forces_deferrable_when_source_side_deletes_target() {
let policies = vec![OnDeletePolicy {
side: DeleteSide::Source,
action: DeleteAction::DeleteTargetIfOrphan,
}];
assert_eq!(target_jt_fk_suffix(&policies), " DEFERRABLE INITIALLY DEFERRED");
}
fn org_type(module: &str) -> TypeDescriptor {
TypeDescriptor {
name: "Org".into(),
module: module.into(),
table: "Org".into(),
abstract_: false,
materialized: true,
description: None,
parents: vec![],
interfaces: vec![],
bases: vec![],
properties: vec![PropertyDescriptor {
name: "id".into(),
pg_type: "uuid".into(),
nullable: false,
default_sql: Some("gen_random_uuid()".into()),
default_pyql: None,
description: None,
check_constraints: vec![],
is_exclusive: true,
is_pk: true,
is_readonly: true,
rewrites: vec![],
tuple_members: None,
column_type: None,
}],
links: vec![],
multilinks: vec![],
computed: vec![],
constraints: vec![],
indexes: vec![],
partition: None,
vector_indexes: vec![],
search_indexes: vec![],
triggers: vec![],
junction: false,
signals: vec![],
}
}
#[test]
fn test_multilink_target_delete_source_trigger_is_after_not_before() {
let module = "default";
let mut owner = org_type(module);
owner.name = "Product".into();
owner.table = "Product".into();
owner.multilinks = vec![MultiLinkDescriptor {
name: "tags".into(),
target: format!("{module}::Org"),
through: None,
nullable: false,
description: None,
default_pyql: None,
on_delete: vec![OnDeletePolicy {
side: DeleteSide::Target,
action: DeleteAction::DeleteSource,
}],
is_exclusive: false,
}];
let schema = SchemaDescriptor {
types: vec![org_type(module), owner],
scalars: vec![],
enums: vec![],
named_tuples: vec![],
globals: vec![],
functions: vec![],
aliases: vec![],
channels: vec![],
..Default::default()
};
let ddl = export_schema(&schema).unwrap();
assert!(
ddl.contains("AFTER DELETE ON \"public\".\"Product.tags\""),
"got:\n{ddl}"
);
assert!(
!ddl.contains("BEFORE DELETE ON \"public\".\"Product.tags\""),
"got:\n{ddl}"
);
}
fn referenced_pylon_functions(ddl: &str) -> std::collections::BTreeSet<String> {
const DEFINITION: &str = "CREATE OR REPLACE FUNCTION _pylon.";
let mut found = std::collections::BTreeSet::new();
let mut rest = ddl;
while let Some(pos) = rest.find("_pylon.") {
let is_definition = rest[..pos]
.len()
.checked_sub(DEFINITION.len() - "_pylon.".len())
.is_some_and(|start| rest[start..pos + "_pylon.".len()].ends_with(DEFINITION));
let after = &rest[pos + "_pylon.".len()..];
let name_len = after
.find(|c: char| !c.is_alphanumeric() && c != '_')
.unwrap_or(after.len());
let name = &after[..name_len];
if !is_definition && !name.is_empty() && after[name_len..].starts_with('(') {
found.insert(name.to_string());
}
rest = &rest[pos + "_pylon.".len()..];
}
found
}
#[test]
fn only_known_stdlib_functions_reach_persisted_ddl() {
let mut person = person_type();
person.materialized = true;
person.properties.push(PropertyDescriptor {
name: "tags".into(),
pg_type: "text[]".into(),
nullable: true,
default_sql: None,
default_pyql: None,
description: None,
check_constraints: vec![],
is_exclusive: false,
is_pk: false,
is_readonly: false,
rewrites: vec![],
tuple_members: None,
column_type: None,
});
person.triggers = vec![trig(
1,
"After",
"update Person set { age := <std::int64>__new__.tags[0] }",
)];
person.signals = vec![crate::schema::SignalEntry { on: 1 }];
let schema = SchemaDescriptor {
types: vec![person],
scalars: vec![],
enums: vec![],
named_tuples: vec![],
globals: vec![],
functions: vec![],
aliases: vec![],
channels: vec![],
..Default::default()
};
let ddl = export_schema(&schema).unwrap();
let referenced = referenced_pylon_functions(&ddl);
let allowed: std::collections::BTreeSet<String> = ["array_subscript", "notify_cache_invalidate"]
.into_iter()
.map(String::from)
.collect();
assert!(
referenced.is_subset(&allowed),
"new stdlib functions reached persisted DDL: {:?}\n\
see this test's comment before widening the allowlist",
referenced.difference(&allowed).collect::<Vec<_>>(),
);
assert!(
referenced.contains("array_subscript"),
"fixture no longer bakes a stdlib call; it is not testing anything"
);
}
fn multilink_trigger_schema(owner_table: &str) -> SchemaDescriptor {
let module = "default";
let mut owner = org_type(module);
owner.name = owner_table.into();
owner.table = owner_table.into();
owner.multilinks = vec![MultiLinkDescriptor {
name: "tags".into(),
target: format!("{module}::Org"),
through: None,
nullable: false,
description: None,
default_pyql: None,
on_delete: vec![OnDeletePolicy {
side: DeleteSide::Target,
action: DeleteAction::DeleteSource,
}],
is_exclusive: false,
}];
SchemaDescriptor {
types: vec![org_type(module), owner],
scalars: vec![],
enums: vec![],
named_tuples: vec![],
globals: vec![],
functions: vec![],
aliases: vec![],
channels: vec![],
..Default::default()
}
}
fn trigger_names_of(schema: &SchemaDescriptor) -> Vec<String> {
let type_map: HashMap<String, (&str, &str)> = schema
.types
.iter()
.map(|t| {
(
format!("{}::{}", t.module, t.name),
(t.module.as_str(), t.table.as_str()),
)
})
.collect();
let mut names: Vec<String> = deletion_policy_trigger_infos(schema, &type_map)
.into_iter()
.map(|i| i.trigger_name)
.collect();
names.sort();
names
}
#[test]
fn a_trigger_name_is_stable_for_an_unchanged_schema() {
let schema = multilink_trigger_schema("Product");
assert_eq!(trigger_names_of(&schema), trigger_names_of(&schema));
}
#[test]
fn a_trigger_name_changes_when_its_body_changes() {
let before = trigger_names_of(&multilink_trigger_schema("Product"));
let after = trigger_names_of(&multilink_trigger_schema("Widget"));
assert_ne!(before, after, "trigger name did not follow its body");
}
#[test]
fn a_signal_trigger_name_follows_its_event_list() {
let names_for = |on: u8| {
let module = "default";
let mut t = org_type(module);
t.signals = vec![crate::schema::SignalEntry { on }];
let schema = SchemaDescriptor {
types: vec![t],
scalars: vec![],
enums: vec![],
named_tuples: vec![],
globals: vec![],
functions: vec![],
aliases: vec![],
channels: vec![],
..Default::default()
};
signal_trigger_infos(&schema)
.into_iter()
.map(|i| i.trigger_name)
.collect::<Vec<_>>()
};
assert_ne!(names_for(1), names_for(1 | 4));
assert_eq!(names_for(1), names_for(1));
}
#[test]
fn test_junction_backed_single_link_gets_a_source_pk_junction_table_no_fk_column() {
let module = "default";
let mut owner = org_type(module);
owner.name = "Person".into();
owner.table = "Person".into();
owner.links = vec![LinkDescriptor {
name: "spouse".into(),
target: format!("{module}::Org"),
nullable: true,
through: Some(format!("{module}::Marriage")),
description: None,
default_pyql: None,
is_exclusive: true,
is_readonly: false,
rewrites: vec![],
on_delete: vec![],
}];
let schema = SchemaDescriptor {
types: vec![org_type(module), owner],
scalars: vec![],
enums: vec![],
named_tuples: vec![],
globals: vec![],
functions: vec![],
aliases: vec![],
channels: vec![],
..Default::default()
};
let ddl = export_schema(&schema).unwrap();
assert!(
!ddl.contains("spouse_id"),
"no {{name}}_id column/FK for a junction-backed link, got:\n{ddl}"
);
assert!(ddl.contains("CREATE TABLE \"public\".\"Person.spouse\""), "got:\n{ddl}");
assert!(
ddl.contains("PRIMARY KEY (source)"),
"single-link junction table must be capped to one row per source, got:\n{ddl}"
);
assert!(
ddl.contains("UNIQUE (target)"),
"exclusive single link must also be unique on the target side, got:\n{ddl}"
);
assert!(
!ddl.contains("CREATE UNIQUE INDEX ON \"public\".\"Person\" (\"spouse_id\")"),
"got:\n{ddl}"
);
}
#[test]
fn test_type_with_a_signal_gets_a_capture_trigger() {
use crate::schema::SignalEntry;
let mut with_signal = person_type();
with_signal.signals = vec![SignalEntry { on: 5 }]; let schema = SchemaDescriptor {
types: vec![with_signal],
scalars: vec![],
enums: vec![],
named_tuples: vec![],
globals: vec![],
functions: vec![],
aliases: vec![],
channels: vec![],
..Default::default()
};
let ddl = export_schema(&schema).unwrap();
assert!(
ddl.contains("AFTER INSERT OR UPDATE OR DELETE ON \"public\".\"Person\""),
"got:\n{ddl}"
);
assert!(ddl.contains("_pylon.\"SignalOutbox\""), "got:\n{ddl}");
assert!(ddl.contains("'default::Person'"), "got:\n{ddl}");
}
#[test]
fn test_type_without_a_signal_gets_no_capture_trigger() {
let schema = minimal_schema(vec![]);
let ddl = export_schema(&schema).unwrap();
assert!(!ddl.contains("SignalOutbox"), "got:\n{ddl}");
}
#[test]
fn test_signal_update_capture_skips_index_maintenance_only_changes() {
use crate::schema::{SignalEntry, VectorIndexDescriptor};
let mut with_signal = person_type();
with_signal.signals = vec![SignalEntry { on: 2 }]; with_signal.vector_indexes = vec![VectorIndexDescriptor {
index_name: None,
pointers: vec!["name".into()],
model: "mistral-embed".into(),
metric: "cosine".into(),
dimensions: 1024,
}];
let schema = SchemaDescriptor {
types: vec![with_signal],
scalars: vec![],
enums: vec![],
named_tuples: vec![],
globals: vec![],
functions: vec![],
aliases: vec![],
channels: vec![],
..Default::default()
};
let ddl = export_schema(&schema).unwrap();
assert!(
ddl.contains(
"IF TG_OP = 'UPDATE' AND (to_jsonb(OLD) - '__vector__') = (to_jsonb(NEW) - '__vector__') THEN"
),
"got:\n{ddl}"
);
assert!(ddl.contains("RETURN NULL;\n END IF;"), "got:\n{ddl}");
}
#[test]
fn test_signal_update_capture_has_no_guard_without_an_index() {
use crate::schema::SignalEntry;
let mut with_signal = person_type();
with_signal.signals = vec![SignalEntry { on: 2 }]; let schema = SchemaDescriptor {
types: vec![with_signal],
scalars: vec![],
enums: vec![],
named_tuples: vec![],
globals: vec![],
functions: vec![],
aliases: vec![],
channels: vec![],
..Default::default()
};
let ddl = export_schema(&schema).unwrap();
assert!(!ddl.contains("IF TG_OP = 'UPDATE'"), "got:\n{ddl}");
}
}
#[cfg(test)]
mod partition_tests {
use super::*;
use crate::schema::{PartitionDescriptor, PartitionInterval, PropertyDescriptor};
fn prop(name: &str, pg_type: &str, is_pk: bool) -> PropertyDescriptor {
PropertyDescriptor {
name: name.into(),
pg_type: pg_type.into(),
nullable: false,
default_sql: None,
default_pyql: None,
description: None,
check_constraints: vec![],
is_exclusive: false,
is_pk,
is_readonly: false,
rewrites: vec![],
tuple_members: None,
column_type: None,
}
}
fn event_schema(retention: Option<u32>) -> SchemaDescriptor {
let mut td = TypeDescriptor {
name: "Event".into(),
module: "default".into(),
table: "Event".into(),
abstract_: false,
materialized: true,
description: None,
parents: vec![],
interfaces: vec![],
bases: vec![],
properties: vec![prop("id", "uuid", true), prop("occurred_at", "timestamptz", false)],
links: vec![],
multilinks: vec![],
computed: vec![],
constraints: vec![],
indexes: vec![],
partition: None,
vector_indexes: vec![],
search_indexes: vec![],
triggers: vec![],
junction: false,
signals: vec![],
};
td.partition = Some(PartitionDescriptor {
pointer: "occurred_at".into(),
interval: PartitionInterval::Monthly,
premake: 4,
retention,
});
SchemaDescriptor {
types: vec![td],
..Default::default()
}
}
#[test]
fn a_partitioned_table_declares_its_range_key() {
let ddl = export_schema(&event_schema(None)).unwrap();
assert!(ddl.contains("PARTITION BY RANGE (\"occurred_at\")"), "got:\n{ddl}");
}
#[test]
fn the_partition_key_is_added_to_the_primary_key() {
let ddl = export_schema(&event_schema(None)).unwrap();
assert!(ddl.contains("PRIMARY KEY (\"id\", \"occurred_at\")"), "got:\n{ddl}");
}
#[test]
fn an_unpartitioned_table_is_untouched() {
let mut schema = event_schema(None);
schema.types[0].partition = None;
let ddl = export_schema(&schema).unwrap();
assert!(!ddl.contains("PARTITION BY"), "got:\n{ddl}");
assert!(!ddl.contains("partman"), "got:\n{ddl}");
assert!(ddl.contains("PRIMARY KEY (\"id\")"), "got:\n{ddl}");
}
#[test]
fn registration_is_guarded_so_re_export_is_idempotent() {
let ddl = export_schema(&event_schema(None)).unwrap();
assert!(
ddl.contains("IF NOT EXISTS (SELECT 1 FROM partman.part_config"),
"got:\n{ddl}"
);
assert!(ddl.contains("partman.create_parent("), "got:\n{ddl}");
assert!(ddl.contains("p_interval := '1 month'"), "got:\n{ddl}");
assert!(ddl.contains("p_premake := 4"), "got:\n{ddl}");
}
#[test]
fn retention_is_applied_when_declared() {
let ddl = export_schema(&event_schema(Some(12))).unwrap();
assert!(ddl.contains("SET retention = '12 months'"), "got:\n{ddl}");
assert!(ddl.contains("retention_keep_table = false"), "got:\n{ddl}");
}
#[test]
fn retention_is_cleared_when_not_declared() {
let ddl = export_schema(&event_schema(None)).unwrap();
assert!(ddl.contains("SET retention = NULL"), "got:\n{ddl}");
}
#[test]
fn a_partitioned_schema_requires_the_partman_extension() {
let exts = crate::diff::required_extensions(&event_schema(None));
assert!(exts.contains(&"pg_partman"), "got: {exts:?}");
}
}