-- -- This file is part of TALER -- Copyright (C) 2025 Taler Systems SA -- -- TALER is free software; you can redistribute it and/or modify it under the -- terms of the GNU General Public License as published by the Free Software -- Foundation; either version 3, or (at your option) any later version. -- -- TALER is distributed in the hope that it will be useful, but WITHOUT ANY -- WARRANTY; without even the implied warranty of MERCHANTABILITY or FITNESS FOR -- A PARTICULAR PURPOSE. See the GNU General Public License for more details. -- -- You should have received a copy of the GNU General Public License along with -- TALER; see the file COPYING. If not, see -- SET search_path TO exchange; DROP FUNCTION IF EXISTS interval_to_start; CREATE OR REPLACE FUNCTION interval_to_start ( IN in_timestamp TIMESTAMP, IN in_range statistic_range, OUT out_bucket_start INT8 ) LANGUAGE plpgsql AS $$ BEGIN out_bucket_start = EXTRACT(EPOCH FROM DATE_TRUNC(in_range::text, in_timestamp)); END $$; COMMENT ON FUNCTION interval_to_start IS 'computes the start time of the bucket for an event at the current time given the desired bucket range'; DROP PROCEDURE IF EXISTS exchange_do_bump_number_bucket_stat; CREATE OR REPLACE PROCEDURE exchange_do_bump_number_bucket_stat( in_slug TEXT, in_h_payto BYTEA, in_timestamp TIMESTAMP, in_delta INT8 ) LANGUAGE plpgsql AS $$ DECLARE my_meta INT8; my_range statistic_range; my_bucket_start INT8; my_curs CURSOR (arg_slug TEXT) FOR SELECT UNNEST(ranges) FROM exchange_statistic_bucket_meta WHERE slug=arg_slug; BEGIN SELECT bmeta_serial_id INTO my_meta FROM exchange_statistic_bucket_meta WHERE slug=in_slug AND stype='number'; IF NOT FOUND THEN RETURN; END IF; OPEN my_curs (arg_slug:=in_slug); LOOP FETCH NEXT FROM my_curs INTO my_range; EXIT WHEN NOT FOUND; SELECT * INTO my_bucket_start FROM interval_to_start (in_timestamp, my_range); UPDATE exchange_statistic_bucket_counter SET cumulative_number = cumulative_number + in_delta WHERE bmeta_serial_id=my_meta AND h_payto=in_h_payto AND bucket_start=my_bucket_start AND bucket_range=my_range; IF NOT FOUND THEN INSERT INTO exchange_statistic_bucket_counter (bmeta_serial_id ,h_payto ,bucket_start ,bucket_range ,cumulative_number ) VALUES ( my_meta ,in_h_payto ,my_bucket_start ,my_range ,in_delta); END IF; END LOOP; CLOSE my_curs; END $$; DROP PROCEDURE IF EXISTS exchange_do_bump_amount_bucket_stat; CREATE OR REPLACE PROCEDURE exchange_do_bump_amount_bucket_stat( in_slug TEXT, in_h_payto BYTEA, in_timestamp TIMESTAMP, in_delta taler_amount ) LANGUAGE plpgsql AS $$ DECLARE my_meta INT8; my_range statistic_range; my_bucket_start INT8; my_curs CURSOR (arg_slug TEXT) FOR SELECT UNNEST(ranges) FROM exchange_statistic_bucket_meta WHERE slug=arg_slug; BEGIN SELECT bmeta_serial_id INTO my_meta FROM exchange_statistic_bucket_meta WHERE slug=in_slug AND stype='amount'; IF NOT FOUND THEN RETURN; END IF; OPEN my_curs (arg_slug:=in_slug); LOOP FETCH NEXT FROM my_curs INTO my_range; EXIT WHEN NOT FOUND; SELECT * INTO my_bucket_start FROM interval_to_start (in_timestamp, my_range); UPDATE exchange_statistic_bucket_amount SET cumulative_value.val = (cumulative_value).val + (in_delta).val + CASE WHEN (in_delta).frac + (cumulative_value).frac >= 100000000 THEN 1 ELSE 0 END, cumulative_value.frac = (cumulative_value).frac + (in_delta).frac - CASE WHEN (in_delta).frac + (cumulative_value).frac >= 100000000 THEN 100000000 ELSE 0 END WHERE bmeta_serial_id=my_meta AND h_payto=in_h_payto AND bucket_start=my_bucket_start AND bucket_range=my_range; IF NOT FOUND THEN INSERT INTO exchange_statistic_bucket_amount (bmeta_serial_id ,h_payto ,bucket_start ,bucket_range ,cumulative_value ) VALUES ( my_meta ,in_h_payto ,my_bucket_start ,my_range ,in_delta); END IF; END LOOP; CLOSE my_curs; END $$; COMMENT ON PROCEDURE exchange_do_bump_amount_bucket_stat IS 'Updates an amount statistic tracked over buckets'; DROP PROCEDURE IF EXISTS exchange_do_bump_number_interval_stat; CREATE OR REPLACE PROCEDURE exchange_do_bump_number_interval_stat( in_slug TEXT, in_h_payto BYTEA, in_timestamp TIMESTAMP, in_delta INT8 ) LANGUAGE plpgsql AS $$ DECLARE my_now INT8; my_record RECORD; my_meta INT8; my_ranges INT8[]; my_precisions INT8[]; my_rangex INT8; my_precisionx INT8; my_start INT8; my_event INT8; BEGIN my_now = ROUND(EXTRACT(epoch FROM CURRENT_TIMESTAMP(0)::TIMESTAMP) * 1000000)::INT8 / 1000 / 1000; SELECT imeta_serial_id ,ranges AS ranges ,precisions AS precisions INTO my_record FROM exchange_statistic_interval_meta WHERE slug=in_slug AND stype='number'; IF NOT FOUND THEN RETURN; END IF; my_start = ROUND(EXTRACT(epoch FROM in_timestamp) * 1000000)::INT8 / 1000 / 1000; -- convert to seconds my_precisions = my_record.precisions; my_ranges = my_record.ranges; my_rangex = NULL; FOR my_x IN 1..COALESCE(array_length(my_ranges,1),0) LOOP IF my_now - my_ranges[my_x] < my_start THEN my_rangex = my_ranges[my_x]; my_precisionx = my_precisions[my_x]; EXIT; END IF; END LOOP; IF my_rangex IS NULL THEN -- event is beyond the ranges we care about RETURN; END IF; my_meta = my_record.imeta_serial_id; my_start = my_start - my_start % my_precisionx; -- round down INSERT INTO exchange_statistic_counter_event AS msce (imeta_serial_id ,h_payto ,slot ,delta) VALUES (my_meta ,in_h_payto ,my_start ,in_delta) ON CONFLICT (imeta_serial_id, h_payto, slot) DO UPDATE SET delta = msce.delta + in_delta RETURNING nevent_serial_id INTO my_event; UPDATE exchange_statistic_interval_counter SET cumulative_number = cumulative_number + in_delta WHERE imeta_serial_id = my_meta AND h_payto = in_h_payto AND range=my_rangex; IF NOT FOUND THEN INSERT INTO exchange_statistic_interval_counter (imeta_serial_id ,h_payto ,range ,event_delimiter ,cumulative_number ) VALUES ( my_meta ,in_h_payto ,my_rangex ,my_event ,in_delta); END IF; END $$; COMMENT ON PROCEDURE exchange_do_bump_number_interval_stat IS 'Updates a numeric statistic tracked over an interval'; DROP PROCEDURE IF EXISTS exchange_do_bump_amount_interval_stat; CREATE OR REPLACE PROCEDURE exchange_do_bump_amount_interval_stat( in_slug TEXT, in_h_payto BYTEA, in_timestamp TIMESTAMP, in_delta taler_amount ) LANGUAGE plpgsql AS $$ DECLARE my_now INT8; my_record RECORD; my_meta INT8; my_ranges INT8[]; my_precisions INT8[]; my_x INT; my_rangex INT8; my_precisionx INT8; my_start INT8; my_event INT8; BEGIN my_now = ROUND(EXTRACT(epoch FROM CURRENT_TIMESTAMP(0)::TIMESTAMP) * 1000000)::INT8 / 1000 / 1000; SELECT imeta_serial_id ,ranges ,precisions INTO my_record FROM exchange_statistic_interval_meta WHERE slug=in_slug AND stype='amount'; IF NOT FOUND THEN RETURN; END IF; my_start = ROUND(EXTRACT(epoch FROM in_timestamp) * 1000000)::INT8 / 1000 / 1000; -- convert to seconds since epoch my_precisions = my_record.precisions; my_ranges = my_record.ranges; my_rangex = NULL; FOR my_x IN 1..COALESCE(array_length(my_ranges,1),0) LOOP IF my_now - my_ranges[my_x] < my_start THEN my_rangex = my_ranges[my_x]; my_precisionx = my_precisions[my_x]; EXIT; END IF; END LOOP; IF my_rangex IS NULL THEN -- event is beyond the ranges we care about RETURN; END IF; my_start = my_start - my_start % my_precisionx; -- round down my_meta = my_record.imeta_serial_id; INSERT INTO exchange_statistic_amount_event AS msae (imeta_serial_id ,h_payto ,slot ,delta ) VALUES ( my_meta ,in_h_payto ,my_start ,in_delta ) ON CONFLICT (imeta_serial_id, h_payto, slot) DO UPDATE SET delta.val = (msae.delta).val + (in_delta).val + CASE WHEN (in_delta).frac + (msae.delta).frac >= 100000000 THEN 1 ELSE 0 END, delta.frac = (msae.delta).frac + (in_delta).frac - CASE WHEN (in_delta).frac + (msae.delta).frac >= 100000000 THEN 100000000 ELSE 0 END RETURNING aevent_serial_id INTO my_event; UPDATE exchange_statistic_interval_amount SET cumulative_value.val = (cumulative_value).val + (in_delta).val + CASE WHEN (in_delta).frac + (cumulative_value).frac >= 100000000 THEN 1 ELSE 0 END, cumulative_value.frac = (cumulative_value).frac + (in_delta).frac - CASE WHEN (in_delta).frac + (cumulative_value).frac >= 100000000 THEN 100000000 ELSE 0 END WHERE imeta_serial_id=my_meta AND h_payto=in_h_payto AND range=my_rangex; IF NOT FOUND THEN INSERT INTO exchange_statistic_interval_amount (imeta_serial_id ,h_payto ,range ,event_delimiter ,cumulative_value ) VALUES ( my_meta ,in_h_payto ,my_rangex ,my_event ,in_delta); END IF; END $$; COMMENT ON PROCEDURE exchange_do_bump_amount_interval_stat IS 'Updates an amount statistic tracked over an interval'; DROP PROCEDURE IF EXISTS exchange_do_bump_number_stat; CREATE OR REPLACE PROCEDURE exchange_do_bump_number_stat( in_slug TEXT, in_h_payto BYTEA, in_timestamp TIMESTAMP, in_delta INT8 ) LANGUAGE plpgsql AS $$ BEGIN CALL exchange_do_bump_number_bucket_stat (in_slug, in_h_payto, in_timestamp, in_delta); CALL exchange_do_bump_number_interval_stat (in_slug, in_h_payto, in_timestamp, in_delta); END $$; COMMENT ON PROCEDURE exchange_do_bump_number_stat IS 'Updates a numeric statistic (bucket or interval)'; DROP PROCEDURE IF EXISTS exchange_do_bump_amount_stat; CREATE OR REPLACE PROCEDURE exchange_do_bump_amount_stat( in_slug TEXT, in_h_payto BYTEA, in_timestamp TIMESTAMP, in_delta taler_amount ) LANGUAGE plpgsql AS $$ BEGIN CALL exchange_do_bump_amount_bucket_stat (in_slug, in_h_payto, in_timestamp, in_delta); CALL exchange_do_bump_amount_interval_stat (in_slug, in_h_payto, in_timestamp, in_delta); END $$; COMMENT ON PROCEDURE exchange_do_bump_amount_stat IS 'Updates an amount statistic (bucket or interval)'; DROP FUNCTION IF EXISTS exchange_statistic_interval_number_get; CREATE OR REPLACE FUNCTION exchange_statistic_interval_number_get ( IN in_slug TEXT, IN in_h_payto BYTEA ) RETURNS SETOF exchange_statistic_interval_number_get_return_value LANGUAGE plpgsql AS $$ DECLARE my_time INT8 DEFAULT ROUND(EXTRACT(epoch FROM CURRENT_TIMESTAMP(0)::TIMESTAMP) * 1000000)::INT8 / 1000 / 1000; my_ranges INT8[]; my_range INT8; my_delta INT8; my_meta INT8; my_next_max_serial INT8; my_rec RECORD; my_irec RECORD; my_i INT; my_min_serial INT8 DEFAULT NULL; my_rval exchange_statistic_interval_number_get_return_value; BEGIN SELECT imeta_serial_id ,ranges ,precisions INTO my_rec FROM exchange_statistic_interval_meta WHERE slug=in_slug; IF NOT FOUND THEN RETURN; END IF; my_rval.rvalue = 0; my_ranges = my_rec.ranges; my_meta = my_rec.imeta_serial_id; FOR my_i IN 1..COALESCE(array_length(my_ranges,1),0) LOOP my_range = my_ranges[my_i]; SELECT event_delimiter ,cumulative_number INTO my_irec FROM exchange_statistic_interval_counter WHERE imeta_serial_id = my_meta AND range = my_range AND h_payto = in_h_payto; IF FOUND THEN my_min_serial = my_irec.event_delimiter; my_rval.rvalue = my_rval.rvalue + my_irec.cumulative_number; -- Check if we have events that left the applicable range SELECT SUM(delta) AS delta_sum INTO my_irec FROM exchange_statistic_counter_event WHERE imeta_serial_id = my_meta AND h_payto = in_h_payto AND slot < my_time - my_range AND nevent_serial_id >= my_min_serial; IF FOUND AND my_irec.delta_sum IS NOT NULL THEN my_delta = my_irec.delta_sum; my_rval.rvalue = my_rval.rvalue - my_delta; -- First find out the next event delimiter value SELECT nevent_serial_id INTO my_next_max_serial FROM exchange_statistic_counter_event WHERE imeta_serial_id = my_meta AND h_payto = in_h_payto AND slot >= my_time - my_range AND nevent_serial_id >= my_min_serial ORDER BY slot ASC LIMIT 1; IF FOUND THEN -- remove expired events from the sum of the current slot UPDATE exchange_statistic_interval_counter SET cumulative_number = cumulative_number - my_delta, event_delimiter = my_next_max_serial WHERE imeta_serial_id = my_meta AND h_payto = in_h_payto AND range = my_range; ELSE -- actually, slot is now empty, remove it entirely DELETE FROM exchange_statistic_interval_counter WHERE imeta_serial_id = my_meta AND h_payto = in_h_payto AND range = my_range; END IF; IF (my_i < array_length(my_ranges,1)) THEN -- carry over all events into the next slot UPDATE exchange_statistic_interval_counter AS usic SET cumulative_number = cumulative_number + my_delta, event_delimiter = LEAST(usic.event_delimiter,my_min_serial) WHERE imeta_serial_id = my_meta AND h_payto = in_h_payto AND range=my_ranges[my_i+1]; IF NOT FOUND THEN INSERT INTO exchange_statistic_interval_counter (imeta_serial_id ,h_payto ,range ,event_delimiter ,cumulative_number ) VALUES ( my_meta ,in_h_payto ,my_ranges[my_i+1] ,my_min_serial ,my_delta); END IF; ELSE -- events are obsolete, delete them DELETE FROM exchange_statistic_counter_event WHERE imeta_serial_id = my_meta AND h_payto = in_h_payto AND slot < my_time - my_range; END IF; END IF; my_rval.range = my_range; RETURN NEXT my_rval; END IF; END LOOP; END $$; COMMENT ON FUNCTION exchange_statistic_interval_number_get IS 'Returns deposit statistic tracking deposited amounts over certain time intervals; we first trim the stored data to only track what is still in-range, and then return the remaining value for each range'; DROP FUNCTION IF EXISTS exchange_statistic_interval_amount_get; CREATE OR REPLACE FUNCTION exchange_statistic_interval_amount_get ( IN in_slug TEXT, IN in_h_payto BYTEA ) RETURNS SETOF exchange_statistic_interval_amount_get_return_value LANGUAGE plpgsql AS $$ DECLARE my_time INT8 DEFAULT ROUND(EXTRACT(epoch FROM CURRENT_TIMESTAMP(0)::TIMESTAMP) * 1000000)::INT8 / 1000 / 1000; my_ranges INT8[]; my_range INT8; my_delta_value INT8; my_delta_frac INT8; my_delta taler_amount; my_meta INT8; my_next_max_serial INT8; my_rec RECORD; my_irec RECORD; my_jrec RECORD; my_i INT; my_min_serial INT8 DEFAULT NULL; my_rval exchange_statistic_interval_amount_get_return_value; BEGIN SELECT imeta_serial_id ,ranges ,precisions INTO my_rec FROM exchange_statistic_interval_meta WHERE slug=in_slug; IF NOT FOUND THEN RETURN; END IF; my_meta = my_rec.imeta_serial_id; my_ranges = my_rec.ranges; my_rval.rvalue.val = 0; my_rval.rvalue.frac = 0; FOR my_i IN 1..COALESCE(array_length(my_ranges,1),0) LOOP my_range = my_ranges[my_i]; SELECT event_delimiter ,cumulative_value INTO my_irec FROM exchange_statistic_interval_amount WHERE imeta_serial_id = my_meta AND h_payto = in_h_payto AND range = my_range; IF FOUND THEN my_min_serial = my_irec.event_delimiter; my_rval.rvalue.val = (my_rval.rvalue).val + (my_irec.cumulative_value).val + (my_irec.cumulative_value).frac / 100000000; my_rval.rvalue.frac = (my_rval.rvalue).frac + (my_irec.cumulative_value).frac % 100000000; IF (my_rval.rvalue).frac > 100000000 THEN my_rval.rvalue.frac = (my_rval.rvalue).frac - 100000000; my_rval.rvalue.val = (my_rval.rvalue).val + 1; END IF; -- Check if we have events that left the applicable range SELECT SUM((esae.delta).val) AS value_sum ,SUM((esae.delta).frac) AS frac_sum INTO my_jrec FROM exchange_statistic_amount_event esae WHERE imeta_serial_id = my_meta AND h_payto = in_h_payto AND slot < my_time - my_range AND aevent_serial_id >= my_min_serial; IF FOUND AND my_jrec.value_sum IS NOT NULL THEN -- Normalize sum my_delta_value = my_jrec.value_sum + my_jrec.frac_sum / 100000000; my_delta_frac = my_jrec.frac_sum % 100000000; my_rval.rvalue.val = (my_rval.rvalue).val - my_delta_value; IF ((my_rval.rvalue).frac >= my_delta_frac) THEN my_rval.rvalue.frac = (my_rval.rvalue).frac - my_delta_frac; ELSE my_rval.rvalue.frac = 100000000 + (my_rval.rvalue).frac - my_delta_frac; my_rval.rvalue.val = (my_rval.rvalue).val - 1; END IF; -- First find out the next event delimiter value SELECT aevent_serial_id INTO my_next_max_serial FROM exchange_statistic_amount_event WHERE imeta_serial_id = my_meta AND h_payto = in_h_payto AND slot >= my_time - my_range AND aevent_serial_id >= my_min_serial ORDER BY slot ASC LIMIT 1; IF FOUND THEN -- remove expired events from the sum of the current slot UPDATE exchange_statistic_interval_amount SET cumulative_value.val = (cumulative_value).val - my_delta_value - CASE WHEN (cumulative_value).frac < my_delta_frac THEN 1 ELSE 0 END, cumulative_value.frac = (cumulative_value).frac - my_delta_frac + CASE WHEN (cumulative_value).frac < my_delta_frac THEN 100000000 ELSE 0 END, event_delimiter = my_next_max_serial WHERE imeta_serial_id = my_meta AND h_payto = in_h_payto AND range = my_range; ELSE -- actually, slot is now empty, remove it entirely DELETE FROM exchange_statistic_interval_amount WHERE imeta_serial_id = my_meta AND h_payto = in_h_payto AND range = my_range; END IF; IF (my_i < array_length(my_ranges,1)) THEN -- carry over all events into the next (larger) slot UPDATE exchange_statistic_interval_amount AS msia SET cumulative_value.val = (cumulative_value).val + my_delta_value + CASE WHEN (cumulative_value).frac + my_delta_frac > 100000000 THEN 1 ELSE 0 END, cumulative_value.frac = (cumulative_value).frac + my_delta_value - CASE WHEN (cumulative_value).frac + my_delta_frac > 100000000 THEN 100000000 ELSE 0 END, event_delimiter = LEAST (msia.event_delimiter,my_min_serial) WHERE imeta_serial_id = my_meta AND h_payto = in_h_payto AND range=my_ranges[my_i+1]; IF NOT FOUND THEN my_delta.val = my_delta_value; my_delta.frac = my_delta_frac; INSERT INTO exchange_statistic_interval_amount (imeta_serial_id ,h_payto ,event_delimiter ,range ,cumulative_value ) VALUES ( my_meta ,in_h_payto ,my_min_serial ,my_ranges[my_i+1] ,my_delta); END IF; ELSE -- events are obsolete, delete them DELETE FROM exchange_statistic_amount_event WHERE imeta_serial_id = my_meta AND h_payto = in_h_payto AND slot < my_time - my_range; END IF; END IF; my_rval.range = my_range; RETURN NEXT my_rval; END IF; END LOOP; -- over my_ranges END $$; COMMENT ON FUNCTION exchange_statistic_interval_amount_get IS 'Returns deposit statistic tracking deposited amounts over certain time intervals; we first trim the stored data to only track what is still in-range, and then return the remaining value; multiple values are returned, one per range'; DROP PROCEDURE IF EXISTS exchange_statistic_counter_gc; CREATE OR REPLACE PROCEDURE exchange_statistic_counter_gc () LANGUAGE plpgsql AS $$ DECLARE my_time INT8 DEFAULT ROUND(EXTRACT(epoch FROM CURRENT_TIMESTAMP(0)::TIMESTAMP) * 1000000)::INT8 / 1000 / 1000; my_h_payto BYTEA; my_rec RECORD; my_sum RECORD; my_meta INT8; my_ranges INT8[]; my_precisions INT8[]; my_precision INT4; my_i INT4; min_slot INT8; max_slot INT8; end_slot INT8; my_total INT8; BEGIN -- GC for all instances FOR my_h_payto IN SELECT DISTINCT h_payto FROM exchange_statistic_counter_event LOOP -- Do combination work for all numeric statistic events FOR my_rec IN SELECT imeta_serial_id ,ranges ,precisions ,slug FROM exchange_statistic_interval_meta LOOP -- First, we query the current interval statistic to update its counters PERFORM FROM exchange_statistic_interval_number_get (my_rec.slug, my_h_payto); my_meta = my_rec.imeta_serial_id; my_ranges = my_rec.ranges; my_precisions = my_rec.precisions; FOR my_i IN 1..COALESCE(array_length(my_ranges,1),0) LOOP my_precision = my_precisions[my_i]; IF 1 >= my_precision THEN -- Cannot coarsen in this case CONTINUE; END IF; IF 1 = my_i THEN min_slot = 0; ELSE min_slot = my_ranges[my_i - 1]; END IF; end_slot = my_ranges[my_i]; -- RAISE NOTICE 'Coarsening from [%,%) at %', my_time - end_slot, my_time - min_slot, my_precision; LOOP EXIT WHEN min_slot >= end_slot; max_slot = min_slot + my_precision; SELECT SUM(delta) AS total, COUNT(*) AS matches, MIN(nevent_serial_id) AS rep_serial_id INTO my_sum FROM exchange_statistic_counter_event WHERE h_payto=my_h_payto AND imeta_serial_id=my_meta AND slot >= my_time - max_slot AND slot < my_time - min_slot; -- RAISE NOTICE 'Found % entries between [%,%)', my_sum.matches, my_time - max_slot, my_time - min_slot; -- we only proceed if we had more then one match (optimization) IF FOUND AND my_sum.matches > 1 THEN my_total = my_sum.total; -- RAISE NOTICE 'combining % entries to representative % for slots [%-%)', my_sum.matches, my_sum.rep_serial_id, my_time - max_slot, my_time - min_slot; -- combine entries DELETE FROM exchange_statistic_counter_event WHERE h_payto=my_h_payto AND imeta_serial_id=my_meta AND slot >= my_time - max_slot AND slot < my_time - min_slot AND nevent_serial_id > my_sum.rep_serial_id; -- Now update the representative to the sum UPDATE exchange_statistic_counter_event SET delta = my_total WHERE imeta_serial_id = my_meta AND h_payto = my_h_payto AND nevent_serial_id = my_sum.rep_serial_id; END IF; min_slot = min_slot + my_precision; END LOOP; -- min_slot to end_slot by precision loop END LOOP; -- my_i loop -- Finally, delete all events beyond the range we care about -- RAISE NOTICE 'deleting entries of %/% before % - % = %', my_h_payto, my_meta, my_time, my_ranges[array_length(my_ranges,1)], my_time - my_ranges[array_length(my_ranges,1)]; DELETE FROM exchange_statistic_counter_event WHERE h_payto=my_h_payto AND imeta_serial_id=my_meta AND slot < my_time - my_ranges[array_length(my_ranges,1)]; END LOOP; -- my_rec loop END LOOP; -- my_h_payto loop END $$; COMMENT ON PROCEDURE exchange_statistic_counter_gc IS 'Performs garbage collection and compaction of the exchange_statistic_counter_event table'; DROP PROCEDURE IF EXISTS exchange_statistic_amount_gc; CREATE OR REPLACE PROCEDURE exchange_statistic_amount_gc () LANGUAGE plpgsql AS $$ DECLARE my_time INT8 DEFAULT ROUND(EXTRACT(epoch FROM CURRENT_TIMESTAMP(0)::TIMESTAMP) * 1000000)::INT8 / 1000 / 1000; my_h_payto BYTEA; my_rec RECORD; my_sum RECORD; my_meta INT8; my_ranges INT8[]; my_precisions INT8[]; my_precision INT4; my_i INT4; min_slot INT8; max_slot INT8; end_slot INT8; my_total_val INT8; my_total_frac INT8; BEGIN -- GC for all accounts FOR my_h_payto IN SELECT DISTINCT h_payto FROM exchange_statistic_counter_event LOOP -- Do combination work for all numeric statistic events FOR my_rec IN SELECT imeta_serial_id ,ranges ,precisions ,slug FROM exchange_statistic_interval_meta LOOP -- First, we query the current interval statistic to update its counters PERFORM FROM exchange_statistic_interval_amount_get (my_rec.slug, my_h_payto); my_meta = my_rec.imeta_serial_id; my_ranges = my_rec.ranges; my_precisions = my_rec.precisions; FOR my_i IN 1..COALESCE(array_length(my_ranges,1),0) LOOP my_precision = my_precisions[my_i]; IF 1 >= my_precision THEN -- Cannot coarsen in this case CONTINUE; END IF; IF 1 = my_i THEN min_slot = 0; ELSE min_slot = my_ranges[my_i - 1]; END IF; end_slot = my_ranges[my_i]; -- RAISE NOTICE 'Coarsening from [%,%) at %', my_time - end_slot, my_time - min_slot, my_precision; LOOP EXIT WHEN min_slot >= end_slot; max_slot = min_slot + my_precision; SELECT SUM((delta).val) AS total_val, SUM((delta).frac) AS total_frac, COUNT(*) AS matches, MIN(aevent_serial_id) AS rep_serial_id INTO my_sum FROM exchange_statistic_amount_event WHERE imeta_serial_id=my_meta AND h_payto=my_h_payto AND slot >= my_time - max_slot AND slot < my_time - max_slot; -- we only proceed if we had more then one match (optimization) IF FOUND AND my_sum.matches > 1 THEN -- normalize new total my_total_frac = my_sum.total_frac % 100000000; my_total_val = my_sum.total_val + my_sum.total_frac / 100000000; -- combine entries DELETE FROM exchange_statistic_amount_event WHERE imeta_serial_id=my_meta AND h_payto=my_h_payto AND slot >= my_time - max_slot AND slot < my_time - max_slot AND aevent_serial_id > my_sum.rep_serial_id; -- Now update the representative to the sum UPDATE exchange_statistic_amount_event SET delta.val = my_total_value ,delta.frac = my_total_frac WHERE imeta_serial_id = my_meta AND h_payto = my_h_payto AND aevent_serial_id = my_sum.rep_serial_id; END IF; min_slot = min_slot + my_precision; END LOOP; -- min_slot to end_slot by precision loop END LOOP; -- my_i loop -- Finally, delete all events beyond the range we care about -- RAISE NOTICE 'deleting entries of %/% before % - % = %', my_h_payto, my_meta, my_time, my_ranges[array_length(my_ranges,1)], my_time - my_ranges[array_length(my_ranges,1)]; DELETE FROM exchange_statistic_amount_event WHERE h_payto=my_h_payto AND imeta_serial_id=my_meta AND slot < my_time - my_ranges[array_length(my_ranges,1)]; END LOOP; -- my_rec loop END LOOP; -- my_h_payto loop END $$; COMMENT ON PROCEDURE exchange_statistic_amount_gc IS 'Performs garbage collection and compaction of the exchange_statistic_amount_event table'; DROP PROCEDURE IF EXISTS exchange_statistic_bucket_gc; CREATE OR REPLACE PROCEDURE exchange_statistic_bucket_gc () LANGUAGE plpgsql AS $$ DECLARE my_rec RECORD; my_range TEXT; my_now INT8; my_end INT8; BEGIN my_now = EXTRACT(EPOCH FROM CURRENT_TIMESTAMP(0)::TIMESTAMP); -- seconds since epoch FOR my_rec IN SELECT bmeta_serial_id ,stype ,ranges[array_length(ranges,1)] AS range ,ages[array_length(ages,1)] AS age FROM exchange_statistic_bucket_meta LOOP my_range = '1 ' || my_rec.range::TEXT; my_end = my_now - my_rec.age * EXTRACT(SECONDS FROM (SELECT my_range::INTERVAL)); -- age is given in multiples of the range (in seconds) IF my_rec.stype = 'amount' THEN DELETE FROM exchange_statistic_bucket_amount WHERE bmeta_serial_id = my_rec.bmeta_serial_id AND bucket_start >= my_end; ELSE DELETE FROM exchange_statistic_bucket_counter WHERE bmeta_serial_id = my_rec.bmeta_serial_id AND bucket_start >= my_end; END IF; END LOOP; END $$; COMMENT ON PROCEDURE exchange_statistic_bucket_gc IS 'Performs garbage collection of the exchange_statistic_bucket_counter and exchange_statistic_bucket_amount tables'; DROP FUNCTION IF EXISTS exchange_drop_customization; CREATE OR REPLACE FUNCTION exchange_drop_customization ( IN in_schema TEXT, OUT out_found BOOLEAN ) LANGUAGE plpgsql AS $$ DECLARE my_xpatches TEXT; BEGIN -- Update DB versioning table. out_found = FALSE; FOR my_xpatches IN SELECT patch_name FROM _v.patches WHERE starts_with(patch_name, in_schema || '-') LOOP PERFORM _v.unregister_patch(my_xpatches); out_found = TRUE; END LOOP; IF out_found THEN -- Drop the schema with all stored procedures/functions. -- This also removes all associated triggers, hence CASCADE. EXECUTE FORMAT('DROP SCHEMA %s CASCADE' ,in_schema); END IF; -- Finally, need to also remove entries from the statistics meta-tables. -- Doing so also DELETEs the associated statistics, hence CASCADE. DELETE FROM exchange_statistic_interval_meta WHERE origin=in_schema; DELETE FROM exchange_statistic_bucket_meta WHERE origin=in_schema; END $$; COMMENT ON FUNCTION exchange_drop_customization IS 'Removes all entries related to a particular exchange customization schema';