mte/include/sql/exchange_statistics_helpers.sql

1043 lines
31 KiB
PL/PgSQL

--
-- 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 <http://www.gnu.org/licenses/>
--
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';