1044 lines
31 KiB
MySQL
1044 lines
31 KiB
MySQL
|
|
--
|
||
|
|
-- 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';
|