1use super::{FnDescriptor, FnVolatility, ImplStrategy, Param, PylonFnDef, PylonType, SqlLanguage};
21
22fn pg_type(ty: &PylonType) -> String {
25 use PylonType::*;
26 match ty {
27 Str => "text".into(),
28 Bool => "bool".into(),
29 Int16 => "int2".into(),
30 Int32 => "int4".into(),
31 Int64 => "int8".into(),
32 Float32 => "float4".into(),
33 Float64 => "float8".into(),
34 Decimal | BigInt => "numeric".into(),
35 Uuid => "uuid".into(),
36 Json => "jsonb".into(),
37 Bytes => "bytea".into(),
38 Datetime => "timestamptz".into(),
39 Duration | RelativeDuration | DateDuration => "interval".into(),
40 LocalDatetime => "timestamp".into(),
41 LocalDate => "date".into(),
42 LocalTime => "time".into(),
43 Vector => "vector".into(),
44 Geometry => "geometry".into(),
45 Geography => "geography".into(),
46 Box2D => "box2d".into(),
47 Box3D => "box3d".into(),
48 Any | AnyOrderable | AnyPoint => "anyelement".into(),
49 Array(inner) => match inner.as_ref() {
50 Any | AnyOrderable | AnyPoint => "anyarray".into(),
51 other => format!("{}[]", pg_type(other)),
52 },
53 Set(inner) => match inner.as_ref() {
55 Any | AnyOrderable | AnyPoint => "anyarray".into(),
56 other => format!("{}[]", pg_type(other)),
57 },
58 Optional(inner) => pg_type(inner),
59 Range(inner) => match inner.as_ref() {
60 Any | AnyOrderable | AnyPoint => "anyrange".into(),
61 other => format!("{}range", pg_type(other)),
62 },
63 Multirange(inner) => match inner.as_ref() {
64 Any | AnyOrderable | AnyPoint => "anymultirange".into(),
65 other => format!("{}multirange", pg_type(other)),
66 },
67 Tuple(_) => panic!("Tuple cannot appear as a PG function parameter type"),
68 }
69}
70
71fn pg_returns(ty: &PylonType) -> String {
74 use PylonType::*;
75 match ty {
76 Set(inner) => match inner.as_ref() {
77 Tuple(_) => panic!("TABLE returns must use PylonFnDef::returns_override"),
78 other => format!("SETOF {}", pg_type(other)),
79 },
80 Optional(inner) => pg_type(inner),
81 other => pg_type(other),
82 }
83}
84
85fn pg_params(params: &[Param]) -> String {
91 if params.is_empty() {
92 return String::new();
93 }
94 let last = params.len() - 1;
95 params
96 .iter()
97 .enumerate()
98 .map(|(i, p)| {
99 let arr_ty = match &p.ty {
100 PylonType::Any | PylonType::AnyOrderable | PylonType::AnyPoint => "anyarray".into(),
101 other => format!("{}[]", pg_type(other)),
102 };
103 if p.variadic && i == last {
104 format!("VARIADIC {} {}", p.name, arr_ty)
106 } else if p.variadic {
107 format!("{} {}", p.name, arr_ty)
109 } else {
110 format!("{} {}", p.name, pg_type(&p.ty))
111 }
112 })
113 .collect::<Vec<_>>()
114 .join(", ")
115}
116
117fn render_function(desc: &FnDescriptor, def: &PylonFnDef) -> String {
120 let params = pg_params(&desc.params);
121 let returns = def
122 .returns_override
123 .map(|s| s.to_owned())
124 .unwrap_or_else(|| pg_returns(&desc.return_type));
125 let lang = match def.language {
126 SqlLanguage::Sql => "sql",
127 SqlLanguage::PlPgSql => "plpgsql",
128 };
129 let volatility = match def.volatility {
130 FnVolatility::Immutable => "IMMUTABLE",
131 FnVolatility::Stable => "STABLE",
132 FnVolatility::Volatile | FnVolatility::Modifying => "VOLATILE",
133 };
134 let parallel = match def.volatility {
137 FnVolatility::Modifying => "PARALLEL UNSAFE",
138 _ => "PARALLEL SAFE",
139 };
140 let strict = if def.strict { " STRICT" } else { "" };
141
142 format!(
143 "CREATE OR REPLACE FUNCTION _pylon.{name}({params})\n\
144 \tRETURNS {returns}\n\
145 \tLANGUAGE {lang} {volatility} {parallel}{strict}\n\
146 AS $$\n\
147 {body}\n\
148 $$;\n",
149 name = def.name,
150 body = def.body,
151 )
152}
153
154pub const INDEX_OUTBOX_DDL: &str = concat!(
161 "DO $$ BEGIN\n",
162 " CREATE TYPE _pylon.\"IndexKind\" AS ENUM ('Vector', 'OpenSearch', 'Meilisearch');\n",
163 "EXCEPTION WHEN duplicate_object THEN NULL; END $$;\n",
164 "ALTER TYPE _pylon.\"IndexKind\" ADD VALUE IF NOT EXISTS 'Meilisearch';\n",
165 "DO $$ BEGIN\n",
166 " CREATE TYPE _pylon.\"IndexOutboxStatus\" AS ENUM ('Pending', 'Processing', 'Failed');\n",
167 "EXCEPTION WHEN duplicate_object THEN NULL; END $$;\n\n",
168 "CREATE TABLE IF NOT EXISTS _pylon.\"IndexOutbox\" (\n",
169 " id uuid NOT NULL DEFAULT uuidv7(),\n",
170 " object_id uuid NOT NULL,\n",
171 " type_name text NOT NULL,\n",
172 " index_kind _pylon.\"IndexKind\" NOT NULL,\n",
173 " index_name text,\n",
174 " operation text NOT NULL DEFAULT 'index',\n",
175 " status _pylon.\"IndexOutboxStatus\" NOT NULL DEFAULT 'Pending',\n",
176 " attempts int NOT NULL DEFAULT 0,\n",
177 " enqueued_at timestamptz NOT NULL DEFAULT now(),\n",
178 " next_attempt timestamptz,\n",
179 " claimed_at timestamptz,\n",
182 " PRIMARY KEY (id),\n",
183 " UNIQUE NULLS NOT DISTINCT (object_id, index_kind, index_name)\n",
184 ");\n",
185 "ALTER TABLE _pylon.\"IndexOutbox\" ADD COLUMN IF NOT EXISTS\n",
186 " operation text NOT NULL DEFAULT 'index';\n",
187 "ALTER TABLE _pylon.\"IndexOutbox\" ADD COLUMN IF NOT EXISTS\n",
188 " claimed_at timestamptz;\n\n",
189 "CREATE INDEX IF NOT EXISTS \"IndexOutbox_status_next_attempt\" ON _pylon.\"IndexOutbox\" (status, next_attempt)\n",
190 " WHERE status IN ('Pending', 'Failed');\n\n",
191 "CREATE OR REPLACE FUNCTION _pylon.notify_index_queue()\n",
192 " RETURNS trigger LANGUAGE plpgsql AS $$\n",
193 "BEGIN\n",
194 " PERFORM pg_notify('pylon_index_queue', NEW.object_id::text);\n",
195 " RETURN NEW;\n",
196 "END\n",
197 "$$;\n\n",
198 "CREATE OR REPLACE TRIGGER notify_index_queue\n",
199 " AFTER INSERT OR UPDATE ON _pylon.\"IndexOutbox\"\n",
200 " FOR EACH ROW EXECUTE FUNCTION _pylon.notify_index_queue();\n",
201);
202
203pub const SIGNAL_OUTBOX_DDL: &str = concat!(
214 "DO $$ BEGIN\n",
215 " CREATE TYPE _pylon.\"SignalOutboxStatus\" AS ENUM ('Pending', 'Processing', 'Failed');\n",
216 "EXCEPTION WHEN duplicate_object THEN NULL; END $$;\n\n",
217 "CREATE TABLE IF NOT EXISTS _pylon.\"SignalOutbox\" (\n",
218 " id uuid NOT NULL DEFAULT uuidv7(),\n",
219 " type_name text NOT NULL,\n",
220 " operation text NOT NULL,\n",
221 " old_row jsonb,\n",
222 " new_row jsonb,\n",
223 " status _pylon.\"SignalOutboxStatus\" NOT NULL DEFAULT 'Pending',\n",
224 " attempts int NOT NULL DEFAULT 0,\n",
225 " enqueued_at timestamptz NOT NULL DEFAULT now(),\n",
226 " next_attempt timestamptz,\n",
227 " PRIMARY KEY (id)\n",
228 ");\n\n",
229 "CREATE INDEX IF NOT EXISTS \"SignalOutbox_status_next_attempt\" ON _pylon.\"SignalOutbox\" (status, next_attempt)\n",
230 " WHERE status IN ('Pending', 'Failed');\n\n",
231 "CREATE OR REPLACE FUNCTION _pylon.notify_signal_queue()\n",
232 " RETURNS trigger LANGUAGE plpgsql AS $$\n",
233 "BEGIN\n",
234 " PERFORM pg_notify('pylon_signal_queue', NEW.id::text);\n",
235 " RETURN NEW;\n",
236 "END\n",
237 "$$;\n\n",
238 "CREATE OR REPLACE TRIGGER notify_signal_queue\n",
239 " AFTER INSERT ON _pylon.\"SignalOutbox\"\n",
240 " FOR EACH ROW EXECUTE FUNCTION _pylon.notify_signal_queue();\n",
241);
242
243pub const TUPLE_CONVERSION_DDL: &str = concat!(
265 "CREATE OR REPLACE FUNCTION _pylon.tuple_json_object(src jsonb)\n",
269 "\tRETURNS jsonb\n",
270 "\tLANGUAGE sql IMMUTABLE PARALLEL SAFE\n",
271 "AS $$\n",
272 " SELECT CASE jsonb_typeof(src)\n",
273 " WHEN 'array' THEN (\n",
274 " SELECT COALESCE(jsonb_object_agg((ord - 1)::text, val), '{}'::jsonb)\n",
275 " FROM jsonb_array_elements(src) WITH ORDINALITY AS e(val, ord)\n",
276 " )\n",
277 " ELSE src\n",
278 " END\n",
279 "$$;\n\n",
280 "CREATE OR REPLACE FUNCTION _pylon.populate_tuple(base anyelement, src jsonb)\n",
281 "\tRETURNS anyelement\n",
282 "\tLANGUAGE sql IMMUTABLE PARALLEL SAFE\n",
283 "AS $$\n",
284 " SELECT CASE WHEN src IS NULL OR jsonb_typeof(src) = 'null' THEN base\n",
288 " ELSE jsonb_populate_record(base, _pylon.tuple_json_object(src)) END\n",
289 "$$;\n\n",
290 "CREATE OR REPLACE FUNCTION _pylon.populate_tuples(base anyelement, src jsonb)\n",
291 "\tRETURNS anyarray\n",
292 "\tLANGUAGE sql IMMUTABLE PARALLEL SAFE\n",
293 "AS $$\n",
294 " SELECT CASE WHEN src IS NULL OR jsonb_typeof(src) = 'null' THEN NULL\n",
300 " ELSE COALESCE(\n",
301 " (SELECT array_agg(_pylon.populate_tuple(base, e) ORDER BY ord)\n",
302 " FROM jsonb_array_elements(src) WITH ORDINALITY AS u(e, ord)),\n",
303 " (ARRAY[base])[1:0])\n",
304 " END\n",
305 "$$;\n",
306);
307
308pub const CACHE_INVALIDATE_DDL: &str = concat!(
316 "CREATE OR REPLACE FUNCTION _pylon.notify_cache_invalidate()\n",
317 " RETURNS trigger LANGUAGE plpgsql AS $$\n",
318 "BEGIN\n",
319 " PERFORM pg_notify('pylon_cache_invalidate', TG_TABLE_SCHEMA || '.' || TG_TABLE_NAME);\n",
320 " RETURN NULL;\n",
321 "END\n",
322 "$$;\n",
323);
324
325pub const MIGRATION_TRACKING_DDL: &str = concat!(
334 "CREATE TABLE IF NOT EXISTS _pylon.\"Migrations\" (\n",
335 " id text PRIMARY KEY,\n",
336 " onto text NOT NULL,\n",
337 " filename text NOT NULL,\n",
338 " db_state jsonb NULL,\n",
339 " applied_at timestamptz NULL\n",
340 ");\n",
341 "ALTER TABLE _pylon.\"Migrations\" ADD COLUMN IF NOT EXISTS schema_state jsonb NULL;\n\n",
346 "CREATE TABLE IF NOT EXISTS _pylon.\"Progress\" (\n",
347 " id text PRIMARY KEY,\n",
348 " step_index integer NOT NULL,\n",
349 " updated_at timestamptz NOT NULL DEFAULT now()\n",
350 ");\n\n",
351 "CREATE TABLE IF NOT EXISTS _pylon.\"Schema\" (\n",
352 " singleton boolean PRIMARY KEY DEFAULT true CHECK (singleton),\n",
353 " snapshot jsonb NOT NULL,\n",
354 " updated_at timestamptz NOT NULL DEFAULT now()\n",
355 ");\n\n",
356 "CREATE TABLE IF NOT EXISTS _pylon.\"Internal\" (\n",
358 " singleton boolean PRIMARY KEY DEFAULT true CHECK (singleton),\n",
359 " version integer NOT NULL,\n",
360 " updated_at timestamptz NOT NULL DEFAULT now()\n",
361 ");\n",
362 "INSERT INTO _pylon.\"Internal\" (singleton, version) VALUES (true, ",
363 internal_schema_version_literal!(),
364 ")\n",
365 " ON CONFLICT (singleton) DO UPDATE SET version = ",
366 internal_schema_version_literal!(),
367 ", updated_at = now();\n",
368);
369
370pub const INTERNAL_SCHEMA_VERSION: i32 = 1;
385
386pub const MIN_SUPPORTED_INTERNAL_VERSION: i32 = 1;
405
406macro_rules! internal_schema_version_literal {
411 () => {
412 "1"
413 };
414}
415use internal_schema_version_literal;
416
417pub fn export_stdlib() -> String {
424 let mut out = String::from("CREATE SCHEMA IF NOT EXISTS _pylon;\n\n");
425
426 out.push_str(INDEX_OUTBOX_DDL);
427 out.push('\n');
428 out.push_str(SIGNAL_OUTBOX_DDL);
429 out.push('\n');
430 out.push_str(MIGRATION_TRACKING_DDL);
431 out.push('\n');
432 out.push_str(CACHE_INVALIDATE_DDL);
433 out.push('\n');
434 out.push_str(TUPLE_CONVERSION_DDL);
435 out.push('\n');
436
437 out.push_str(concat!(
439 "CREATE OR REPLACE FUNCTION _pylon.array_subscript(arr anyarray, idx bigint)\n",
440 "\tRETURNS anyelement\n",
441 "\tLANGUAGE plpgsql STABLE PARALLEL SAFE\n",
442 "AS $$\n",
443 "DECLARE\n",
444 " element_index bigint := CASE WHEN idx < 0 THEN idx + cardinality(arr) ELSE idx END;\n",
445 "BEGIN\n",
446 " IF element_index < 0 OR element_index >= cardinality(arr) THEN\n",
447 " RAISE EXCEPTION 'array index % is out of bounds', idx\n",
448 " USING ERRCODE = 'array_subscript_error';\n",
449 " END IF;\n",
450 " RETURN arr[element_index + 1];\n",
451 "END\n",
452 "$$;\n\n",
453 "CREATE OR REPLACE FUNCTION _pylon.raise_invalid_parameter(msg text)\n",
457 "\tRETURNS text\n",
458 "\tLANGUAGE plpgsql IMMUTABLE PARALLEL SAFE STRICT\n",
459 "AS $$\n",
460 "BEGIN\n",
461 " RAISE EXCEPTION '%', msg USING ERRCODE = 'invalid_parameter_value';\n",
462 "END\n",
463 "$$;\n\n",
464 "CREATE OR REPLACE FUNCTION _pylon.duration_in(val text)\n",
468 "\tRETURNS interval\n",
469 "\tLANGUAGE plpgsql IMMUTABLE PARALLEL SAFE STRICT\n",
470 "AS $$\n",
471 "DECLARE\n",
472 " parsed interval := val::interval;\n",
473 "BEGIN\n",
474 " IF date_part('year', parsed) != 0 OR date_part('month', parsed) != 0\n",
475 " OR date_part('day', parsed) != 0 THEN\n",
476 " RAISE EXCEPTION 'invalid input syntax for type duration: %', quote_literal(val)\n",
477 " USING ERRCODE = 'invalid_datetime_format',\n",
478 " HINT = 'Day, month and year units cannot be used for duration.';\n",
479 " END IF;\n",
480 " RETURN parsed;\n",
481 "END\n",
482 "$$;\n\n",
483 "CREATE OR REPLACE FUNCTION _pylon.date_duration_in(val text)\n",
484 "\tRETURNS interval\n",
485 "\tLANGUAGE plpgsql IMMUTABLE PARALLEL SAFE STRICT\n",
486 "AS $$\n",
487 "DECLARE\n",
488 " parsed interval := val::interval;\n",
489 "BEGIN\n",
490 " IF date_part('epoch', parsed - date_trunc('day', parsed)) != 0 THEN\n",
491 " RAISE EXCEPTION 'invalid input syntax for type cal::date_duration: %', quote_literal(val)\n",
492 " USING ERRCODE = 'invalid_datetime_format',\n",
493 " HINT = 'Units smaller than days cannot be used for cal::date_duration.';\n",
494 " END IF;\n",
495 " RETURN parsed;\n",
496 "END\n",
497 "$$;\n\n",
498 "CREATE OR REPLACE FUNCTION _pylon.str_subscript(s text, idx bigint)\n",
499 "\tRETURNS text\n",
500 "\tLANGUAGE plpgsql STABLE PARALLEL SAFE\n",
501 "AS $$\n",
502 "DECLARE\n",
503 " element_index bigint := CASE WHEN idx < 0 THEN idx + char_length(s) ELSE idx END;\n",
504 "BEGIN\n",
505 " IF element_index < 0 OR element_index >= char_length(s) THEN\n",
506 " RAISE EXCEPTION 'string index % is out of bounds', idx\n",
507 " USING ERRCODE = 'array_subscript_error';\n",
508 " END IF;\n",
509 " RETURN substr(s, (element_index + 1)::int, 1);\n",
510 "END\n",
511 "$$;\n\n",
512 "CREATE OR REPLACE FUNCTION _pylon.str_subscript(s bytea, idx bigint)\n",
513 "\tRETURNS bytea\n",
514 "\tLANGUAGE plpgsql STABLE PARALLEL SAFE\n",
515 "AS $$\n",
516 "DECLARE\n",
517 " element_index bigint := CASE WHEN idx < 0 THEN idx + length(s) ELSE idx END;\n",
518 "BEGIN\n",
519 " IF element_index < 0 OR element_index >= length(s) THEN\n",
520 " RAISE EXCEPTION 'byte string index % is out of bounds', idx\n",
521 " USING ERRCODE = 'array_subscript_error';\n",
522 " END IF;\n",
523 " RETURN substr(s, (element_index + 1)::int, 1);\n",
524 "END\n",
525 "$$;\n\n",
526 ));
527
528 for desc in super::registry() {
529 if let ImplStrategy::PylonFunction(def) = &desc.impl_strategy {
530 out.push_str(&render_function(desc, def));
531 out.push('\n');
532 }
533 }
534 out
535}
536
537#[cfg(test)]
538mod tests {
539 use super::{INTERNAL_SCHEMA_VERSION, export_stdlib};
540
541 #[test]
542 fn ddl_smoke() {
543 let ddl = export_stdlib();
544 let fn_count = ddl.matches("CREATE OR REPLACE FUNCTION").count();
545 assert!(ddl.starts_with("CREATE SCHEMA IF NOT EXISTS _pylon;"));
546 assert!(fn_count > 0, "no functions generated");
547 assert!(ddl.contains("_pylon.to_bool"), "to_bool missing");
548 assert!(ddl.contains("_pylon.enumerate"), "enumerate missing");
549 assert!(ddl.contains("_pylon.datetime_get"), "datetime_get missing");
550 assert!(
551 !ddl.contains("_pylon.range("),
552 "range must not be installed (TranspilerIntrinsic)"
553 );
554 assert!(!ddl.contains("_pylon.multirange("), "multirange must not be installed");
555 eprintln!("export_stdlib: {} PylonFunction overloads installed", fn_count);
556 }
557
558 #[test]
559 fn ddl_to_bool_has_three_overloads() {
560 let ddl = export_stdlib();
561 let count = ddl.matches("_pylon.to_bool(").count();
562 assert_eq!(count, 3, "expected int2/int4/int8 overloads; got {count}");
563 }
564
565 #[test]
566 fn ddl_enumerate_returns_table() {
567 let ddl = export_stdlib();
568 assert!(ddl.contains("RETURNS TABLE(index bigint, value anyelement)"));
569 }
570
571 #[test]
572 fn ddl_json_get_uses_variadic() {
573 let ddl = export_stdlib();
574 assert!(ddl.contains("VARIADIC path text[]"), "json_get must use VARIADIC");
575 }
576
577 #[test]
578 fn ddl_installs_cache_invalidate_notify_function() {
579 let ddl = export_stdlib();
580 assert!(ddl.contains("CREATE OR REPLACE FUNCTION _pylon.notify_cache_invalidate()"));
581 assert!(ddl.contains("pg_notify('pylon_cache_invalidate', TG_TABLE_SCHEMA || '.' || TG_TABLE_NAME)"));
582 }
583
584 #[test]
585 fn internal_schema_version_literal_matches_the_constant() {
586 assert_eq!(internal_schema_version_literal!(), INTERNAL_SCHEMA_VERSION.to_string(),);
591 }
592
593 #[test]
594 fn the_internal_version_marker_is_created_and_upserted() {
595 let ddl = export_stdlib();
596 assert!(ddl.contains("CREATE TABLE IF NOT EXISTS _pylon.\"Internal\""));
597 assert!(
600 ddl.contains("ON CONFLICT (singleton) DO UPDATE SET version = 1"),
601 "got:\n{ddl}"
602 );
603 }
604}