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 CACHE_INVALIDATE_DDL: &str = concat!(
251 "CREATE OR REPLACE FUNCTION _pylon.notify_cache_invalidate()\n",
252 " RETURNS trigger LANGUAGE plpgsql AS $$\n",
253 "BEGIN\n",
254 " PERFORM pg_notify('pylon_cache_invalidate', TG_TABLE_SCHEMA || '.' || TG_TABLE_NAME);\n",
255 " RETURN NULL;\n",
256 "END\n",
257 "$$;\n",
258);
259
260pub const MIGRATION_TRACKING_DDL: &str = concat!(
269 "CREATE TABLE IF NOT EXISTS _pylon.\"Migrations\" (\n",
270 " id text PRIMARY KEY,\n",
271 " onto text NOT NULL,\n",
272 " filename text NOT NULL,\n",
273 " db_state jsonb NULL,\n",
274 " applied_at timestamptz NULL\n",
275 ");\n",
276 "ALTER TABLE _pylon.\"Migrations\" ADD COLUMN IF NOT EXISTS schema_state jsonb NULL;\n\n",
281 "CREATE TABLE IF NOT EXISTS _pylon.\"Progress\" (\n",
282 " id text PRIMARY KEY,\n",
283 " step_index integer NOT NULL,\n",
284 " updated_at timestamptz NOT NULL DEFAULT now()\n",
285 ");\n\n",
286 "CREATE TABLE IF NOT EXISTS _pylon.\"Schema\" (\n",
287 " singleton boolean PRIMARY KEY DEFAULT true CHECK (singleton),\n",
288 " snapshot jsonb NOT NULL,\n",
289 " updated_at timestamptz NOT NULL DEFAULT now()\n",
290 ");\n\n",
291 "CREATE TABLE IF NOT EXISTS _pylon.\"Internal\" (\n",
293 " singleton boolean PRIMARY KEY DEFAULT true CHECK (singleton),\n",
294 " version integer NOT NULL,\n",
295 " updated_at timestamptz NOT NULL DEFAULT now()\n",
296 ");\n",
297 "INSERT INTO _pylon.\"Internal\" (singleton, version) VALUES (true, ",
298 internal_schema_version_literal!(),
299 ")\n",
300 " ON CONFLICT (singleton) DO UPDATE SET version = ",
301 internal_schema_version_literal!(),
302 ", updated_at = now();\n",
303);
304
305pub const INTERNAL_SCHEMA_VERSION: i32 = 1;
320
321pub const MIN_SUPPORTED_INTERNAL_VERSION: i32 = 1;
340
341macro_rules! internal_schema_version_literal {
346 () => {
347 "1"
348 };
349}
350use internal_schema_version_literal;
351
352pub fn export_stdlib() -> String {
359 let mut out = String::from("CREATE SCHEMA IF NOT EXISTS _pylon;\n\n");
360
361 out.push_str(INDEX_OUTBOX_DDL);
362 out.push('\n');
363 out.push_str(SIGNAL_OUTBOX_DDL);
364 out.push('\n');
365 out.push_str(MIGRATION_TRACKING_DDL);
366 out.push('\n');
367 out.push_str(CACHE_INVALIDATE_DDL);
368 out.push('\n');
369
370 out.push_str(concat!(
372 "CREATE OR REPLACE FUNCTION _pylon.array_subscript(arr anyarray, idx bigint)\n",
373 "\tRETURNS anyelement\n",
374 "\tLANGUAGE plpgsql STABLE PARALLEL SAFE\n",
375 "AS $$\n",
376 "DECLARE\n",
377 " element_index bigint := CASE WHEN idx < 0 THEN idx + cardinality(arr) ELSE idx END;\n",
378 "BEGIN\n",
379 " IF element_index < 0 OR element_index >= cardinality(arr) THEN\n",
380 " RAISE EXCEPTION 'array index % is out of bounds', idx\n",
381 " USING ERRCODE = 'array_subscript_error';\n",
382 " END IF;\n",
383 " RETURN arr[element_index + 1];\n",
384 "END\n",
385 "$$;\n\n",
386 "CREATE OR REPLACE FUNCTION _pylon.raise_invalid_parameter(msg text)\n",
390 "\tRETURNS text\n",
391 "\tLANGUAGE plpgsql IMMUTABLE PARALLEL SAFE STRICT\n",
392 "AS $$\n",
393 "BEGIN\n",
394 " RAISE EXCEPTION '%', msg USING ERRCODE = 'invalid_parameter_value';\n",
395 "END\n",
396 "$$;\n\n",
397 "CREATE OR REPLACE FUNCTION _pylon.duration_in(val text)\n",
401 "\tRETURNS interval\n",
402 "\tLANGUAGE plpgsql IMMUTABLE PARALLEL SAFE STRICT\n",
403 "AS $$\n",
404 "DECLARE\n",
405 " parsed interval := val::interval;\n",
406 "BEGIN\n",
407 " IF date_part('year', parsed) != 0 OR date_part('month', parsed) != 0\n",
408 " OR date_part('day', parsed) != 0 THEN\n",
409 " RAISE EXCEPTION 'invalid input syntax for type duration: %', quote_literal(val)\n",
410 " USING ERRCODE = 'invalid_datetime_format',\n",
411 " HINT = 'Day, month and year units cannot be used for duration.';\n",
412 " END IF;\n",
413 " RETURN parsed;\n",
414 "END\n",
415 "$$;\n\n",
416 "CREATE OR REPLACE FUNCTION _pylon.date_duration_in(val text)\n",
417 "\tRETURNS interval\n",
418 "\tLANGUAGE plpgsql IMMUTABLE PARALLEL SAFE STRICT\n",
419 "AS $$\n",
420 "DECLARE\n",
421 " parsed interval := val::interval;\n",
422 "BEGIN\n",
423 " IF date_part('epoch', parsed - date_trunc('day', parsed)) != 0 THEN\n",
424 " RAISE EXCEPTION 'invalid input syntax for type cal::date_duration: %', quote_literal(val)\n",
425 " USING ERRCODE = 'invalid_datetime_format',\n",
426 " HINT = 'Units smaller than days cannot be used for cal::date_duration.';\n",
427 " END IF;\n",
428 " RETURN parsed;\n",
429 "END\n",
430 "$$;\n\n",
431 "CREATE OR REPLACE FUNCTION _pylon.str_subscript(s text, idx bigint)\n",
432 "\tRETURNS text\n",
433 "\tLANGUAGE plpgsql STABLE PARALLEL SAFE\n",
434 "AS $$\n",
435 "DECLARE\n",
436 " element_index bigint := CASE WHEN idx < 0 THEN idx + char_length(s) ELSE idx END;\n",
437 "BEGIN\n",
438 " IF element_index < 0 OR element_index >= char_length(s) THEN\n",
439 " RAISE EXCEPTION 'string index % is out of bounds', idx\n",
440 " USING ERRCODE = 'array_subscript_error';\n",
441 " END IF;\n",
442 " RETURN substr(s, (element_index + 1)::int, 1);\n",
443 "END\n",
444 "$$;\n\n",
445 "CREATE OR REPLACE FUNCTION _pylon.str_subscript(s bytea, idx bigint)\n",
446 "\tRETURNS bytea\n",
447 "\tLANGUAGE plpgsql STABLE PARALLEL SAFE\n",
448 "AS $$\n",
449 "DECLARE\n",
450 " element_index bigint := CASE WHEN idx < 0 THEN idx + length(s) ELSE idx END;\n",
451 "BEGIN\n",
452 " IF element_index < 0 OR element_index >= length(s) THEN\n",
453 " RAISE EXCEPTION 'byte string index % is out of bounds', idx\n",
454 " USING ERRCODE = 'array_subscript_error';\n",
455 " END IF;\n",
456 " RETURN substr(s, (element_index + 1)::int, 1);\n",
457 "END\n",
458 "$$;\n\n",
459 ));
460
461 for desc in super::registry() {
462 if let ImplStrategy::PylonFunction(def) = &desc.impl_strategy {
463 out.push_str(&render_function(desc, def));
464 out.push('\n');
465 }
466 }
467 out
468}
469
470#[cfg(test)]
471mod tests {
472 use super::{INTERNAL_SCHEMA_VERSION, export_stdlib};
473
474 #[test]
475 fn ddl_smoke() {
476 let ddl = export_stdlib();
477 let fn_count = ddl.matches("CREATE OR REPLACE FUNCTION").count();
478 assert!(ddl.starts_with("CREATE SCHEMA IF NOT EXISTS _pylon;"));
479 assert!(fn_count > 0, "no functions generated");
480 assert!(ddl.contains("_pylon.to_bool"), "to_bool missing");
481 assert!(ddl.contains("_pylon.enumerate"), "enumerate missing");
482 assert!(ddl.contains("_pylon.datetime_get"), "datetime_get missing");
483 assert!(
484 !ddl.contains("_pylon.range("),
485 "range must not be installed (TranspilerIntrinsic)"
486 );
487 assert!(!ddl.contains("_pylon.multirange("), "multirange must not be installed");
488 eprintln!("export_stdlib: {} PylonFunction overloads installed", fn_count);
489 }
490
491 #[test]
492 fn ddl_to_bool_has_three_overloads() {
493 let ddl = export_stdlib();
494 let count = ddl.matches("_pylon.to_bool(").count();
495 assert_eq!(count, 3, "expected int2/int4/int8 overloads; got {count}");
496 }
497
498 #[test]
499 fn ddl_enumerate_returns_table() {
500 let ddl = export_stdlib();
501 assert!(ddl.contains("RETURNS TABLE(index bigint, value anyelement)"));
502 }
503
504 #[test]
505 fn ddl_json_get_uses_variadic() {
506 let ddl = export_stdlib();
507 assert!(ddl.contains("VARIADIC path text[]"), "json_get must use VARIADIC");
508 }
509
510 #[test]
511 fn ddl_installs_cache_invalidate_notify_function() {
512 let ddl = export_stdlib();
513 assert!(ddl.contains("CREATE OR REPLACE FUNCTION _pylon.notify_cache_invalidate()"));
514 assert!(ddl.contains("pg_notify('pylon_cache_invalidate', TG_TABLE_SCHEMA || '.' || TG_TABLE_NAME)"));
515 }
516
517 #[test]
518 fn internal_schema_version_literal_matches_the_constant() {
519 assert_eq!(internal_schema_version_literal!(), INTERNAL_SCHEMA_VERSION.to_string(),);
524 }
525
526 #[test]
527 fn the_internal_version_marker_is_created_and_upserted() {
528 let ddl = export_stdlib();
529 assert!(ddl.contains("CREATE TABLE IF NOT EXISTS _pylon.\"Internal\""));
530 assert!(
533 ddl.contains("ON CONFLICT (singleton) DO UPDATE SET version = 1"),
534 "got:\n{ddl}"
535 );
536 }
537}