Skip to content

Commit 1f6e611

Browse files
authored
Merge pull request #100 from constructive-io/feat/sync-pgpm-modules-from-db
feat: sync pgpm modules from constructive-db
2 parents 7195fba + abf866f commit 1f6e611

135 files changed

Lines changed: 2956 additions & 267 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

packages/database-jobs/deploy/schemas/app_jobs/procedures/add_scheduled_job.sql

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,8 @@ CREATE FUNCTION app_jobs.add_scheduled_job(
1515
queue_name text DEFAULT NULL,
1616
max_attempts integer DEFAULT 25,
1717
priority integer DEFAULT 0,
18-
entity_id uuid DEFAULT NULL
18+
entity_id uuid DEFAULT NULL,
19+
db_id uuid DEFAULT NULL
1920
)
2021
RETURNS app_jobs.scheduled_jobs
2122
AS $$
@@ -24,7 +25,9 @@ DECLARE
2425
v_database_id uuid;
2526
v_actor_id uuid;
2627
BEGIN
27-
v_database_id := jwt_private.current_database_id();
28+
-- Callers that run outside a JWT context (e.g. provisioning triggers) pass
29+
-- db_id explicitly; everyone else keeps the JWT-derived default.
30+
v_database_id := coalesce(db_id, jwt_private.current_database_id());
2831
v_actor_id := jwt_public.current_user_id();
2932

3033
IF job_key IS NOT NULL THEN
Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,10 @@
1+
-- Deploy schemas/app_jobs/procedures/grants/grant_execute_add_scheduled_job_to_authenticated to pg
2+
3+
-- requires: schemas/app_jobs/schema
4+
-- requires: schemas/app_jobs/procedures/add_scheduled_job
5+
6+
BEGIN;
7+
8+
GRANT EXECUTE ON FUNCTION app_jobs.add_scheduled_job(text, json, json, text, text, integer, integer, uuid, uuid) TO authenticated;
9+
10+
COMMIT;

packages/database-jobs/deploy/schemas/app_jobs/procedures/run_scheduled_job.sql

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -67,10 +67,10 @@ BEGIN
6767
* INTO j;
6868
-- update the scheduled job
6969
UPDATE
70-
app_jobs.scheduled_jobs s
70+
app_jobs.scheduled_jobs s
7171
SET
7272
last_scheduled = NOW(),
73-
last_scheduled_id = j.id
73+
last_scheduled_id = COALESCE(j.id, s.last_scheduled_id)
7474
WHERE
7575
s.id = run_scheduled_job.id;
7676
RETURN j;
@@ -79,4 +79,3 @@ $$
7979
LANGUAGE 'plpgsql'
8080
VOLATILE;
8181
COMMIT;
82-
Lines changed: 95 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,95 @@
1+
-- Deploy schemas/app_jobs/procedures/schedule_min_interval_seconds to pg
2+
3+
-- requires: schemas/app_jobs/schema
4+
5+
BEGIN;
6+
7+
-- Conservative lower-bound estimate (in seconds) of how often a
8+
-- scheduled_jobs.schedule_info spec can fire. Used by generated schedule
9+
-- limit triggers to enforce a minimum schedule interval cap.
10+
--
11+
-- Supports the node-schedule shapes stored in schedule_info:
12+
-- { "rule": "*/5 * * * *" } 5-field cron (minute resolution)
13+
-- { "rule": "0 */5 * * * *" } 6-field cron (second resolution)
14+
-- { "minute": [...], "hour": [...] } recurrence-object form
15+
--
16+
-- Returns NULL when no lower bound can be derived (callers should treat
17+
-- NULL as "unknown" and skip enforcement).
18+
CREATE FUNCTION app_jobs.schedule_min_interval_seconds (schedule_info json)
19+
RETURNS integer
20+
AS $$
21+
DECLARE
22+
v_rule text;
23+
v_fields text[];
24+
v_seconds_field text;
25+
v_minutes_field text;
26+
v_step int;
27+
BEGIN
28+
IF schedule_info IS NULL THEN
29+
RETURN NULL;
30+
END IF;
31+
32+
v_rule := schedule_info ->> 'rule';
33+
34+
IF v_rule IS NULL THEN
35+
-- Recurrence-object form: second resolution when a "second" key is
36+
-- present, otherwise minute resolution.
37+
IF schedule_info -> 'second' IS NOT NULL THEN
38+
RETURN 1;
39+
ELSIF schedule_info -> 'minute' IS NOT NULL
40+
OR schedule_info -> 'hour' IS NOT NULL
41+
OR schedule_info -> 'dayOfWeek' IS NOT NULL
42+
OR schedule_info -> 'date' IS NOT NULL
43+
OR schedule_info -> 'month' IS NOT NULL THEN
44+
RETURN 60;
45+
END IF;
46+
RETURN NULL;
47+
END IF;
48+
49+
v_fields := regexp_split_to_array(trim(v_rule), '\s+');
50+
51+
IF array_length(v_fields, 1) = 6 THEN
52+
v_seconds_field := v_fields[1];
53+
v_minutes_field := v_fields[2];
54+
ELSIF array_length(v_fields, 1) = 5 THEN
55+
v_seconds_field := '0';
56+
v_minutes_field := v_fields[1];
57+
ELSE
58+
RETURN NULL;
59+
END IF;
60+
61+
-- Second-resolution rules
62+
IF v_seconds_field = '*' THEN
63+
RETURN 1;
64+
END IF;
65+
IF v_seconds_field ~ '^\*/\d+$' THEN
66+
v_step := substring(v_seconds_field FROM 3)::int;
67+
RETURN GREATEST(v_step, 1);
68+
END IF;
69+
IF position(',' IN v_seconds_field) > 0 OR position('-' IN v_seconds_field) > 0 THEN
70+
-- Multiple seconds within each matching minute: could fire every second
71+
RETURN 1;
72+
END IF;
73+
74+
-- Fixed second: at most once per matching minute
75+
IF v_minutes_field = '*' THEN
76+
RETURN 60;
77+
END IF;
78+
IF v_minutes_field ~ '^\*/\d+$' THEN
79+
v_step := substring(v_minutes_field FROM 3)::int;
80+
RETURN GREATEST(v_step, 1) * 60;
81+
END IF;
82+
IF position(',' IN v_minutes_field) > 0 OR position('-' IN v_minutes_field) > 0 THEN
83+
-- Multiple minutes within each matching hour: adjacent values can be
84+
-- one minute apart
85+
RETURN 60;
86+
END IF;
87+
88+
-- Fixed minute: at most once per hour
89+
RETURN 3600;
90+
END;
91+
$$
92+
LANGUAGE 'plpgsql'
93+
IMMUTABLE;
94+
95+
COMMIT;

packages/database-jobs/pgpm.plan

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,8 @@ schemas/app_jobs/procedures/fail_job [schemas/app_jobs/schema schemas/app_jobs/t
3535
schemas/app_jobs/procedures/complete_jobs [schemas/app_jobs/schema schemas/app_jobs/tables/job_queues/table schemas/app_jobs/tables/jobs/table] 2025-08-26T23:57:41Z pgpm <pgpm@5b0c196eeb62> # add schemas/app_jobs/procedures/complete_jobs
3636
schemas/app_jobs/procedures/complete_job [schemas/app_jobs/schema schemas/app_jobs/tables/jobs/table schemas/app_jobs/tables/job_queues/table] 2025-08-26T23:57:41Z pgpm <pgpm@5b0c196eeb62> # add schemas/app_jobs/procedures/complete_job
3737
schemas/app_jobs/procedures/add_scheduled_job [schemas/app_jobs/schema schemas/app_jobs/tables/scheduled_jobs/table pgpm-jwt-claims:schemas/jwt_private/procedures/current_database_id] 2025-08-26T23:57:41Z pgpm <pgpm@5b0c196eeb62> # add schemas/app_jobs/procedures/add_scheduled_job
38+
schemas/app_jobs/procedures/grants/grant_execute_add_scheduled_job_to_authenticated [schemas/app_jobs/schema schemas/app_jobs/procedures/add_scheduled_job] 2026-07-15T06:07:00Z pgpm <pgpm@localhost> # grant authenticated EXECUTE on add_scheduled_job for INVOKER trigger support
39+
schemas/app_jobs/procedures/schedule_min_interval_seconds [schemas/app_jobs/schema] 2026-07-12T13:30:00Z pgpm <pgpm@localhost> # add schedule_min_interval_seconds estimator for schedule interval caps
3840
schemas/app_jobs/procedures/add_job [schemas/app_jobs/schema schemas/app_jobs/tables/jobs/table schemas/app_jobs/tables/job_queues/table pgpm-jwt-claims:schemas/jwt_private/procedures/current_database_id pgpm-jwt-claims:schemas/jwt_public/procedures/current_user_id] 2025-08-26T23:57:41Z pgpm <pgpm@5b0c196eeb62> # add schemas/app_jobs/procedures/add_job
3941
schemas/app_jobs/procedures/grants/grant_execute_add_job_to_authenticated [schemas/app_jobs/schema schemas/app_jobs/procedures/add_job] 2026-06-03T01:15:00Z pgpm <pgpm@localhost> # grant authenticated EXECUTE on add_job for INVOKER trigger support
4042
schemas/app_jobs/procedures/remove_job [schemas/app_jobs/schema schemas/app_jobs/tables/jobs/table] 2025-08-26T23:57:41Z pgpm <pgpm@5b0c196eeb62> # add schemas/app_jobs/procedures/remove_job
Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,7 @@
1+
-- Revert schemas/app_jobs/procedures/grants/grant_execute_add_scheduled_job_to_authenticated from pg
2+
3+
BEGIN;
4+
5+
REVOKE EXECUTE ON FUNCTION app_jobs.add_scheduled_job(text, json, json, text, text, integer, integer, uuid, uuid) FROM authenticated;
6+
7+
COMMIT;
Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,7 @@
1+
-- Revert schemas/app_jobs/procedures/schedule_min_interval_seconds from pg
2+
3+
BEGIN;
4+
5+
DROP FUNCTION app_jobs.schedule_min_interval_seconds;
6+
7+
COMMIT;

packages/database-jobs/sql/pgpm-database-jobs--0.22.0.sql

Lines changed: 84 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -449,10 +449,10 @@ BEGIN
449449
* INTO j;
450450
-- update the scheduled job
451451
UPDATE
452-
app_jobs.scheduled_jobs s
452+
app_jobs.scheduled_jobs s
453453
SET
454454
last_scheduled = NOW(),
455-
last_scheduled_id = j.id
455+
last_scheduled_id = COALESCE(j.id, s.last_scheduled_id)
456456
WHERE
457457
s.id = run_scheduled_job.id;
458458
RETURN j;
@@ -735,14 +735,17 @@ CREATE FUNCTION app_jobs.add_scheduled_job(
735735
queue_name text DEFAULT NULL,
736736
max_attempts int DEFAULT 25,
737737
priority int DEFAULT 0,
738-
entity_id uuid DEFAULT NULL
738+
entity_id uuid DEFAULT NULL,
739+
db_id uuid DEFAULT NULL
739740
) RETURNS app_jobs.scheduled_jobs AS $EOFCODE$
740741
DECLARE
741742
v_job app_jobs.scheduled_jobs;
742743
v_database_id uuid;
743744
v_actor_id uuid;
744745
BEGIN
745-
v_database_id := jwt_private.current_database_id();
746+
-- Callers that run outside a JWT context (e.g. provisioning triggers) pass
747+
-- db_id explicitly; everyone else keeps the JWT-derived default.
748+
v_database_id := coalesce(db_id, jwt_private.current_database_id());
746749
v_actor_id := jwt_public.current_user_id();
747750

748751
IF job_key IS NOT NULL THEN
@@ -824,6 +827,83 @@ BEGIN
824827
END;
825828
$EOFCODE$ LANGUAGE plpgsql VOLATILE SECURITY DEFINER;
826829

830+
GRANT EXECUTE ON FUNCTION app_jobs.add_scheduled_job(text, pg_catalog.json, pg_catalog.json, text, text, int, int, uuid, uuid) TO authenticated;
831+
832+
CREATE FUNCTION app_jobs.schedule_min_interval_seconds(
833+
schedule_info pg_catalog.json
834+
) RETURNS int AS $EOFCODE$
835+
DECLARE
836+
v_rule text;
837+
v_fields text[];
838+
v_seconds_field text;
839+
v_minutes_field text;
840+
v_step int;
841+
BEGIN
842+
IF schedule_info IS NULL THEN
843+
RETURN NULL;
844+
END IF;
845+
846+
v_rule := schedule_info ->> 'rule';
847+
848+
IF v_rule IS NULL THEN
849+
-- Recurrence-object form: second resolution when a "second" key is
850+
-- present, otherwise minute resolution.
851+
IF schedule_info -> 'second' IS NOT NULL THEN
852+
RETURN 1;
853+
ELSIF schedule_info -> 'minute' IS NOT NULL
854+
OR schedule_info -> 'hour' IS NOT NULL
855+
OR schedule_info -> 'dayOfWeek' IS NOT NULL
856+
OR schedule_info -> 'date' IS NOT NULL
857+
OR schedule_info -> 'month' IS NOT NULL THEN
858+
RETURN 60;
859+
END IF;
860+
RETURN NULL;
861+
END IF;
862+
863+
v_fields := regexp_split_to_array(trim(v_rule), '\s+');
864+
865+
IF array_length(v_fields, 1) = 6 THEN
866+
v_seconds_field := v_fields[1];
867+
v_minutes_field := v_fields[2];
868+
ELSIF array_length(v_fields, 1) = 5 THEN
869+
v_seconds_field := '0';
870+
v_minutes_field := v_fields[1];
871+
ELSE
872+
RETURN NULL;
873+
END IF;
874+
875+
-- Second-resolution rules
876+
IF v_seconds_field = '*' THEN
877+
RETURN 1;
878+
END IF;
879+
IF v_seconds_field ~ '^\*/\d+$' THEN
880+
v_step := substring(v_seconds_field FROM 3)::int;
881+
RETURN GREATEST(v_step, 1);
882+
END IF;
883+
IF position(',' IN v_seconds_field) > 0 OR position('-' IN v_seconds_field) > 0 THEN
884+
-- Multiple seconds within each matching minute: could fire every second
885+
RETURN 1;
886+
END IF;
887+
888+
-- Fixed second: at most once per matching minute
889+
IF v_minutes_field = '*' THEN
890+
RETURN 60;
891+
END IF;
892+
IF v_minutes_field ~ '^\*/\d+$' THEN
893+
v_step := substring(v_minutes_field FROM 3)::int;
894+
RETURN GREATEST(v_step, 1) * 60;
895+
END IF;
896+
IF position(',' IN v_minutes_field) > 0 OR position('-' IN v_minutes_field) > 0 THEN
897+
-- Multiple minutes within each matching hour: adjacent values can be
898+
-- one minute apart
899+
RETURN 60;
900+
END IF;
901+
902+
-- Fixed minute: at most once per hour
903+
RETURN 3600;
904+
END;
905+
$EOFCODE$ LANGUAGE plpgsql IMMUTABLE;
906+
827907
CREATE FUNCTION app_jobs.add_job(
828908
identifier text,
829909
payload pg_catalog.json DEFAULT '{}'::json,
Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,7 @@
1+
-- Verify schemas/app_jobs/procedures/grants/grant_execute_add_scheduled_job_to_authenticated on pg
2+
3+
BEGIN;
4+
5+
SELECT has_function_privilege('authenticated', 'app_jobs.add_scheduled_job(text, json, json, text, text, integer, integer, uuid, uuid)', 'EXECUTE');
6+
7+
ROLLBACK;
Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,7 @@
1+
-- Verify schemas/app_jobs/procedures/schedule_min_interval_seconds on pg
2+
3+
BEGIN;
4+
5+
SELECT verify_function ('app_jobs.schedule_min_interval_seconds');
6+
7+
ROLLBACK;

0 commit comments

Comments
 (0)