use super::*;
pub(crate) fn render_materialized_view(view: &ManifestMaterializedView) -> String {
let with_data = if view.with_data {
"WITH DATA"
} else {
"WITH NO DATA"
};
format!(
"\nCREATE SCHEMA IF NOT EXISTS {};\nCREATE MATERIALIZED VIEW IF NOT EXISTS {}.{} AS\n{}\n{};\n",
qi(&view.schema),
qi(&view.schema),
qi(&view.name),
view.query,
with_data
)
}
pub(crate) fn render_sql_artifacts(table: &ManifestTable, phase: &str) -> String {
let mut out = String::new();
for artifact in &table.sql_artifacts {
if !artifact_applies_to_phase(artifact, phase) {
continue;
}
if !artifact.backend.trim().is_empty() && artifact.backend != "postgres" {
continue;
}
if !artifact.sql.trim().is_empty() {
out.push_str(&format!(
"\n-- UDB:sql_artifact={}\n-- UDB:sql_artifact_phase={}\n",
artifact.name, phase
));
if !artifact.checksum_sha256.trim().is_empty() {
out.push_str(&format!(
"-- UDB:sql_artifact_sha256={}\n",
artifact.checksum_sha256
));
}
out.push_str(artifact.sql.trim());
out.push_str("\n\n");
} else if !artifact.file.trim().is_empty() {
out.push_str(&format!(
"\n-- UDB:sql_artifact={} file={} phase={} sha256={}\n",
artifact.name, artifact.file, phase, artifact.checksum_sha256
));
}
}
out
}
pub(crate) fn artifact_applies_to_phase(artifact: &ManifestSqlArtifact, phase: &str) -> bool {
let artifact_phase = artifact.phase.trim();
if artifact_phase.is_empty() {
return phase == "before_triggers";
}
artifact_phase == phase
}
pub(crate) fn render_trigger(trigger: &ManifestTrigger) -> String {
let timing = if trigger.timing.trim().is_empty() {
"AFTER"
} else {
trigger.timing.as_str()
};
let event = if trigger.event.trim().is_empty() {
"INSERT"
} else {
trigger.event.as_str()
};
let for_each = if trigger.for_each.trim().eq_ignore_ascii_case("STATEMENT") {
"STATEMENT"
} else {
"ROW"
};
let when_clause = if trigger.when_clause.trim().is_empty() {
String::new()
} else {
format!("\n WHEN ({})", trigger.when_clause)
};
format!(
"\nDROP TRIGGER IF EXISTS {} ON {}.{};\nCREATE TRIGGER {}\n {} {} ON {}.{}\n FOR EACH {}{}\n EXECUTE FUNCTION {};\n",
qi(&trigger.name),
qi(&trigger.schema),
qi(&trigger.table),
qi(&trigger.name),
timing,
event,
qi(&trigger.schema),
qi(&trigger.table),
for_each,
when_clause,
trigger.function
)
}
pub(crate) fn render_add_fk(
schema: &str,
table: &str,
fk: &ManifestForeignKey,
table_is_partitioned: bool,
) -> String {
let name = derive_fk_name(table, fk);
let mut inner = format!(
"ALTER TABLE {}.{}\n ADD CONSTRAINT {}\n FOREIGN KEY ({}) REFERENCES {}.{} ({})",
qi(schema),
qi(table),
qi(&name),
quote_list(&fk.columns),
qi(&fk.ref_schema),
qi(&fk.ref_table),
quote_list(&fk.ref_columns)
);
if !fk.on_delete.trim().is_empty() && fk.on_delete != "NO ACTION" {
inner.push_str(&format!("\n ON DELETE {}", fk.on_delete));
}
if !fk.on_update.trim().is_empty() && fk.on_update != "NO ACTION" {
inner.push_str(&format!("\n ON UPDATE {}", fk.on_update));
}
if fk.deferrable {
inner.push_str("\n DEFERRABLE");
if fk.initially_deferred {
inner.push_str(" INITIALLY DEFERRED");
} else {
inner.push_str(" INITIALLY IMMEDIATE");
}
}
if fk.not_valid && !table_is_partitioned {
inner.push_str("\n NOT VALID");
}
inner.push(';');
format!(
"DO $$\n\
DECLARE\n\
\t_missing_partition_cols TEXT[];\n\
BEGIN\n\
\tSELECT COALESCE(array_agg(a.attname ORDER BY key.ord), ARRAY[]::TEXT[])\n\
\tINTO _missing_partition_cols\n\
\tFROM pg_partitioned_table p\n\
\tJOIN LATERAL unnest(p.partattrs) WITH ORDINALITY AS key(attnum, ord) ON TRUE\n\
\tJOIN pg_attribute a ON a.attrelid = p.partrelid AND a.attnum = key.attnum\n\
\tWHERE p.partrelid = to_regclass({ref_relation})\n\
\t AND NOT (a.attname = ANY({ref_columns}));\n\
\n\
\tIF COALESCE(array_length(_missing_partition_cols, 1), 0) > 0 THEN\n\
\t\tRAISE NOTICE 'Skipping FK {fk_name}: referenced partition key columns % are not present in ref_columns {ref_columns_notice}', _missing_partition_cols;\n\
\t\tRETURN;\n\
\tEND IF;\n\
\n\
\t{inner}\n\
EXCEPTION WHEN duplicate_object THEN NULL;\n\
END$$;",
ref_relation = ql(&format!("{}.{}", fk.ref_schema, fk.ref_table)),
ref_columns = sql_text_array(&fk.ref_columns),
fk_name = name.replace('\'', "''"),
ref_columns_notice = fk.ref_columns.join(", ").replace('\'', "''"),
inner = inner
)
}
pub(crate) fn render_add_check(schema: &str, table: &str, check: &ManifestCheck) -> String {
let name = if check.name.trim().is_empty() {
format!("chk_{}_auto", table)
} else {
check.name.clone()
};
format!(
"DO $$BEGIN\n ALTER TABLE {}.{} ADD CONSTRAINT {} CHECK ({});\nEXCEPTION WHEN duplicate_object THEN NULL;\nEND$$;",
qi(schema),
qi(table),
qi(&name),
check.expression
)
}
pub(crate) fn render_policy(schema: &str, table: &str, policy: &ManifestPolicy) -> String {
let mode = if policy.permissive {
"PERMISSIVE"
} else {
"RESTRICTIVE"
};
let command = if policy.command.trim().is_empty() {
"ALL"
} else {
policy.command.as_str()
};
let using = policy.using_expression.trim();
let check = policy.with_check.trim();
let using_clause = if using.is_empty() {
String::new()
} else {
format!("\n USING ({})", using)
};
let check_clause = if check.is_empty() {
String::new()
} else {
format!("\n WITH CHECK ({})", check)
};
format!(
"DO $$BEGIN\n\
\tALTER POLICY {name} ON {schema}.{table}{using_clause}{check_clause};\n\
EXCEPTION WHEN undefined_object THEN\n\
\tCREATE POLICY {name} ON {schema}.{table}\n\
\t AS {mode} FOR {command}{using_clause}{check_clause};\n\
END$$;",
name = qi(&policy.name),
schema = qi(schema),
table = qi(table),
mode = mode,
command = command,
using_clause = using_clause,
check_clause = check_clause,
)
}
pub(crate) fn find_column<'a>(table: &'a ManifestTable, name: &str) -> Option<&'a ManifestColumn> {
table
.columns
.iter()
.find(|column| column.column_name == name)
}
pub(crate) fn derive_index_name(table: &ManifestTable, index: &ManifestIndex) -> String {
if index.name.trim().is_empty() {
format!(
"idx_{}_{}_{}",
table.schema,
table.table,
index.columns.join("_")
)
} else {
index.name.clone()
}
}
pub(crate) fn derive_fk_name(table: &str, fk: &ManifestForeignKey) -> String {
if fk.name.trim().is_empty() {
format!("fk_{}_{}", table, fk.columns.join("_"))
} else {
fk.name.clone()
}
}
pub(crate) fn find_index<'a>(table: &'a ManifestTable, name: &str) -> Option<&'a ManifestIndex> {
table
.indexes
.iter()
.find(|index| derive_index_name(table, index) == name)
}
pub(crate) fn find_fk<'a>(table: &'a ManifestTable, name: &str) -> Option<&'a ManifestForeignKey> {
table
.foreign_keys
.iter()
.find(|fk| derive_fk_name(&table.table, fk) == name)
}
pub(crate) fn find_check<'a>(table: &'a ManifestTable, name: &str) -> Option<&'a ManifestCheck> {
table.checks.iter().find(|check| {
if check.name.trim().is_empty() {
check.expression == name
} else {
check.name == name
}
})
}
pub(crate) fn find_policy<'a>(table: &'a ManifestTable, name: &str) -> Option<&'a ManifestPolicy> {
table.rls_policies.iter().find(|policy| policy.name == name)
}
pub(crate) fn find_extension<'a>(
manifest: &'a CatalogManifest,
schema: &str,
name: &str,
) -> Option<&'a ManifestExtension> {
manifest
.tables
.iter()
.flat_map(|table| table.extensions.iter())
.find(|extension| extension.schema == schema && extension.name == name)
}
pub(crate) fn find_materialized_view<'a>(
table: &'a ManifestTable,
object_name: &str,
) -> Option<&'a ManifestMaterializedView> {
table
.materialized_views
.iter()
.find(|view| format!("{}.{}", view.schema, view.name) == object_name)
}
pub(crate) fn find_trigger<'a>(
table: &'a ManifestTable,
name: &str,
) -> Option<&'a ManifestTrigger> {
table.triggers.iter().find(|trigger| trigger.name == name)
}
pub(crate) fn render_index_standalone(table: &ManifestTable, index: &ManifestIndex) -> String {
render_index_impl(table, index, false)
}
#[cfg(test)]
pub(crate) fn render_index(table: &ManifestTable, index: &ManifestIndex) -> String {
render_index_impl(table, index, false)
}
pub(crate) fn render_index_in_tx(table: &ManifestTable, index: &ManifestIndex) -> String {
render_index_impl(table, index, true)
}
pub(crate) fn render_index_impl(
table: &ManifestTable,
index: &ManifestIndex,
in_transaction: bool,
) -> String {
if index.columns.is_empty() {
return String::new();
}
let unique = if index.unique { "UNIQUE " } else { "" };
let concurrently = if index.concurrent && !in_transaction && !index.unique {
"CONCURRENTLY "
} else {
""
};
let name = derive_index_name(table, index);
let raw_method = if index.method.trim().is_empty() {
"BTREE".to_string()
} else {
index.method.to_ascii_uppercase()
};
let method = if index.unique && raw_method != "BTREE" {
"BTREE".to_string()
} else {
raw_method
};
let effective_columns = if index.unique {
partition_aware_unique_columns(table, &index.columns)
} else {
index.columns.clone()
};
let column_sql_parts = effective_columns
.iter()
.map(|column| {
let op_class = if !index.operator_class.trim().is_empty() {
index.operator_class.as_str().to_string()
} else if matches!(method.as_str(), "HNSW" | "IVFFLAT") {
"vector_cosine_ops".to_string()
} else {
String::new()
};
if op_class.is_empty() {
col_or_expr(column)
} else {
format!("{} {}", col_or_expr(column), op_class)
}
})
.collect::<Vec<_>>();
let include = if index.include_columns.is_empty() {
String::new()
} else {
format!(" INCLUDE ({})", quote_list(&index.include_columns))
};
let params = if index.index_params.is_empty() {
String::new()
} else {
format!(
" WITH ({})",
index
.index_params
.iter()
.map(|param| format!("{} = {}", param.key, param.value))
.collect::<Vec<_>>()
.join(", ")
)
};
let where_clause = if index.where_clause.trim().is_empty() {
String::new()
} else {
format!("\n WHERE {}", index.where_clause)
};
if index.unique {
return render_partitioned_unique_index_create(
table,
&name,
&effective_columns,
&column_sql_parts,
&method,
&format!("{params}{include}"),
&where_clause,
);
}
format!(
"CREATE {}INDEX {}IF NOT EXISTS {}\n ON {}.{} USING {} ({}){}{}{};\n\n",
unique,
concurrently,
qi(&name),
qi(&table.schema),
qi(&table.table),
method,
column_sql_parts.join(", "),
params,
include,
where_clause
)
}
pub(crate) fn validate_enum_value(v: &str) -> Result<(), String> {
if v.is_empty() {
return Err("enum value cannot be empty".to_string());
}
if v.len() > 128 {
return Err(format!(
"enum value exceeds 128 characters (got {})",
v.len()
));
}
if v.contains(['\'', '\\', '$', '\0', '\n', '\r']) {
return Err(
"enum value contains an unsafe character (quote, backslash, dollar, null, or newline)"
.to_string(),
);
}
Ok(())
}
pub(crate) fn render_enum_types(table: &ManifestTable) -> String {
let mut out = String::new();
for column in &table.columns {
if column.enum_values.is_empty() {
continue;
}
let mut values = column.enum_values.clone();
values.sort();
values.dedup();
let type_name = format!(
"{}.{}_{}_enum",
qi(&table.schema),
table.table,
column.column_name
);
let value_list = values
.iter()
.filter(|v| {
if let Err(reason) = validate_enum_value(v) {
tracing::warn!(
enum_type = %type_name,
value = %v,
reason = %reason,
"skipping unsafe enum value in DDL generation"
);
return false;
}
true
})
.map(|v| format!("'{}'", v.replace('\'', "''")))
.collect::<Vec<_>>()
.join(", ");
out.push_str(&format!(
"DO $$BEGIN\n CREATE TYPE {type_name} AS ENUM ({value_list});\nEXCEPTION WHEN duplicate_object THEN NULL;\nEND$$;\n\n"
));
}
out
}
pub(crate) fn render_jsonb_gin_indexes(table: &ManifestTable) -> String {
let mut out = String::new();
for column in &table.columns {
if !column.is_jsonb {
continue;
}
let op_class = if column.json_path_ops {
" jsonb_path_ops".to_string()
} else {
String::new()
};
let name = format!(
"gin_{}_{}_{}",
table.schema, table.table, column.column_name
);
out.push_str(&format!(
"CREATE INDEX IF NOT EXISTS {}\n ON {}.{} USING GIN ({}{});\n\n",
qi(&name),
qi(&table.schema),
qi(&table.table),
qi(&column.column_name),
op_class
));
}
out
}
pub(crate) fn render_tsvector_indexes(table: &ManifestTable) -> String {
let mut out = String::new();
for column in &table.columns {
if column.is_tsvector && !column.tsvector_source_columns.is_empty() {
let raw_lang = if column.tsvector_language.trim().is_empty() {
"simple"
} else {
column.tsvector_language.as_str()
};
let lang = if raw_lang
.chars()
.all(|c| c.is_ascii_alphanumeric() || c == '_')
{
raw_lang
} else {
tracing::warn!(
schema = %table.schema,
table = %table.table,
column = %column.column_name,
lang = %raw_lang,
"tsvector_language contains unsafe characters — falling back to 'simple'"
);
"simple"
};
let concat = column
.tsvector_source_columns
.iter()
.map(|c| format!("coalesce({}, '')", qi(c)))
.collect::<Vec<_>>()
.join(" || ' ' || ");
let idx_name = format!(
"idx_fts_{}_{}_{}",
table.schema, table.table, column.column_name
);
out.push_str(&format!(
"CREATE INDEX IF NOT EXISTS {}\n ON {}.{} USING GIN (to_tsvector('{}', {}));\n\n",
qi(&idx_name),
qi(&table.schema),
qi(&table.table),
lang,
concat,
));
}
if column.trigram_index {
let idx_name = format!(
"idx_trgm_{}_{}_{}",
table.schema, table.table, column.column_name
);
out.push_str(&format!(
"CREATE INDEX IF NOT EXISTS {}\n ON {}.{} USING GIN ({} gin_trgm_ops);\n\n",
qi(&idx_name),
qi(&table.schema),
qi(&table.table),
qi(&column.column_name),
));
}
}
out
}
pub(crate) fn render_partition_setup(table: &ManifestTable) -> String {
if !is_partitioned(table) {
return String::new();
}
let interval = partition_interval_for_table(table);
if interval.trim().is_empty() {
return String::new();
}
let premake = if table.partition_premake > 0 {
table.partition_premake
} else {
4
};
let mut out = String::new();
out.push_str(&render_partition_unique_constraint_cleanup(table));
out.push_str(&format!(
"DO $$\n\
DECLARE\n\
\t_rel REGCLASS := to_regclass({parent_table});\n\
\t_control_col TEXT := {part_col};\n\
BEGIN\n\
IF _rel IS NOT NULL THEN\n\
\tSELECT a.attname\n\
\tINTO _control_col\n\
\tFROM pg_partitioned_table p\n\
\tJOIN LATERAL unnest(p.partattrs) WITH ORDINALITY AS key(attnum, ord) ON TRUE\n\
\tJOIN pg_attribute a ON a.attrelid = p.partrelid AND a.attnum = key.attnum\n\
\tWHERE p.partrelid = _rel\n\
\tORDER BY key.ord\n\
\tLIMIT 1;\n\
\t_control_col := COALESCE(_control_col, {part_col});\n\
END IF;\n\
\n\
IF NOT EXISTS (\n\
\t SELECT 1 FROM partman.part_config WHERE parent_table = {parent_table}\n\
) THEN\n\
\tPERFORM partman.create_parent(\n\
\t\tp_parent_table := {parent_table},\n\
\t\tp_control := _control_col,\n\
\t\tp_interval := {interval},\n\
\t\tp_premake := {premake},\n\
\t\tp_start_partition := to_char(now(), 'YYYY-MM-DD'),\n\
\t\tp_date_trunc_interval := {date_trunc_interval}\n\
\t);\n\
END IF;\n\
END$$;\n\n",
parent_table = ql(&format!("{}.{}", table.schema, table.table)),
part_col = ql(&table.partition_column),
interval = ql(&normalize_partition_interval(&interval)),
premake = premake,
date_trunc_interval = partition_date_trunc_interval(&interval)
.map(|value| ql(&value))
.unwrap_or_else(|| "NULL".to_string()),
));
let parent_ident = format!("{}.{}", qi(&table.schema), qi(&table.table));
if table.partition_default {
let default_ident = format!(
"{}.{}",
qi(&table.schema),
qi(&format!("{}_default", table.table))
);
out.push_str(&format!(
"CREATE TABLE IF NOT EXISTS {default_ident} PARTITION OF {parent_ident} DEFAULT;\n\n"
));
}
if table.retention_days > 0 {
out.push_str(&format!(
"UPDATE partman.part_config\n\
\tSET retention = {retention}, retention_keep_table = false\n\
\tWHERE parent_table = {parent_lit};\n\n",
retention = ql(&format!("{} days", table.retention_days)),
parent_lit = ql(&format!("{}.{}", table.schema, table.table)),
));
}
out
}
pub(crate) fn is_partitioned(table: &ManifestTable) -> bool {
!table.partition_strategy.trim().is_empty()
&& !table.partition_column.trim().is_empty()
&& !table.partition_strategy.ends_with("NONE")
&& !table.partition_strategy.ends_with("UNSPECIFIED")
}
pub(crate) fn normalize_partition_interval(interval: &str) -> String {
match interval.trim().to_ascii_uppercase().as_str() {
"MONTHLY" | "MONTH" => "1 month".to_string(),
"WEEKLY" | "WEEK" => "1 week".to_string(),
"DAILY" | "DAY" => "1 day".to_string(),
"HOURLY" | "HOUR" => "1 hour".to_string(),
"YEARLY" | "YEAR" => "1 year".to_string(),
_ => interval.to_string(), }
}
fn partition_interval_for_table(table: &ManifestTable) -> String {
let explicit = table.partition_interval.trim();
if !explicit.is_empty() {
return explicit.to_string();
}
match table
.partition_strategy
.trim()
.to_ascii_uppercase()
.as_str()
{
value if value.contains("RANGE_YEAR") => "YEARLY".to_string(),
value if value.contains("RANGE_MONTH") => "MONTHLY".to_string(),
value if value.contains("RANGE_WEEK") => "WEEKLY".to_string(),
value if value.contains("RANGE_DAY") => "DAILY".to_string(),
value if value.contains("RANGE_HOUR") => "HOURLY".to_string(),
_ => String::new(),
}
}
fn partition_date_trunc_interval(interval: &str) -> Option<String> {
match interval.trim().to_ascii_uppercase().as_str() {
"YEARLY" | "YEAR" | "1 YEAR" => Some("year".to_string()),
"MONTHLY" | "MONTH" | "1 MONTH" => Some("month".to_string()),
"WEEKLY" | "WEEK" | "1 WEEK" => Some("week".to_string()),
"DAILY" | "DAY" | "1 DAY" => Some("day".to_string()),
"HOURLY" | "HOUR" | "1 HOUR" => Some("hour".to_string()),
_ => None,
}
}
pub(crate) fn normalize_partition_strategy(value: &str) -> &'static str {
let value = value.to_ascii_uppercase();
if value.contains("LIST") {
"LIST"
} else if value.contains("HASH") {
"HASH"
} else {
"RANGE"
}
}
pub(crate) fn quote_list(values: &[String]) -> String {
values
.iter()
.map(|value| qi(value))
.collect::<Vec<_>>()
.join(", ")
}
pub(crate) fn first_non_empty<'a>(a: &'a str, b: &'a str, fallback: &'a str) -> &'a str {
if !a.trim().is_empty() {
a
} else if !b.trim().is_empty() {
b
} else {
fallback
}
}
pub(crate) fn delta_slug(ops: &[&ChangeOperation]) -> String {
let mut parts = ops
.iter()
.map(|op| format!("{:?}", op.kind).to_ascii_lowercase())
.collect::<Vec<_>>();
parts.sort();
parts.dedup();
let slug = parts.join("_");
if slug.len() > 64 {
slug[..64].to_string()
} else {
slug
}
}
pub(crate) fn split_qualified_name(value: &str, fallback_schema: &str) -> (String, String) {
if let Some((schema, name)) = value.split_once('.') {
(schema.to_string(), name.to_string())
} else {
(fallback_schema.to_string(), value.to_string())
}
}
pub(crate) fn qi(value: &str) -> String {
format!("\"{}\"", value.replace('"', "\"\""))
}
pub(crate) fn col_or_expr(col: &str) -> String {
if col.contains('(') || col.contains('\'') {
col.to_string()
} else {
qi(col)
}
}
pub(crate) fn ql(value: &str) -> String {
format!("'{}'", value.replace('\'', "''"))
}