import { ContinuousRunnerResult, PgClient } from '@bmd-studio/genstack-pg'; import environment from '@bmd-studio/genstack-environment'; const { APP_PREFIX, POSTGRES_HIDDEN_SCHEMA_NAME, POSTGRES_DEFAULT_SCHEMA_NAME, DATABASE_ID_COLUMN_NAME, } = environment.env; // adapted from https://gist.github.com/colophonemes/9701b906c5be572a40a84b08f4d2fa4e export default (pgClient: PgClient): ContinuousRunnerResult => { return pgClient.query(` -- Trigger notification for messaging to PG Notify CREATE OR REPLACE FUNCTION ${POSTGRES_HIDDEN_SCHEMA_NAME}.notify_row_operation() RETURNS trigger AS $trigger$ DECLARE row RECORD; payload TEXT; column_names TEXT[]; column_name TEXT; row_id TEXT; notify_id UUID; column_value TEXT; BEGIN -- Set record row depending on operation CASE TG_OP WHEN 'INSERT', 'UPDATE' THEN row := NEW; WHEN 'DELETE' THEN row := OLD; ELSE RAISE EXCEPTION 'Unknown TG_OP: "%". Should not occur!', TG_OP; END CASE; -- Generate an unique ID to identity the notification SELECT ${POSTGRES_DEFAULT_SCHEMA_NAME}.uuid_generate_v4() INTO notify_id; -- Get the record ID EXECUTE format('SELECT $1.%I::TEXT', '${DATABASE_ID_COLUMN_NAME}') INTO row_id USING row; -- By default use all the TG_ARGV entries column_names := TG_ARGV; -- Get required fields FOREACH column_name IN ARRAY column_names LOOP EXECUTE 'SELECT $1.' || quote_ident(column_name) || '' INTO column_value USING row; -- Notify per field PERFORM pg_notify('${APP_PREFIX}_row_operation', json_build_object( 'operation', lower(TG_OP), 'tableName', TG_TABLE_NAME, 'rowId', row_id, 'columnName', column_name, 'columnValue', column_value, 'notifyId', notify_id::text )::text ); END LOOP; RETURN row; END; $trigger$ LANGUAGE plpgsql security definer; `); };