exchange

Base system with REST service to issue digital coins, run by the payment service provider
Log | Files | Refs | Submodules | README | LICENSE

exchange_statistics_helpers.sql (32004B)


      1 --
      2 -- This file is part of TALER
      3 -- Copyright (C) 2025 Taler Systems SA
      4 --
      5 -- TALER is free software; you can redistribute it and/or modify it under the
      6 -- terms of the GNU General Public License as published by the Free Software
      7 -- Foundation; either version 3, or (at your option) any later version.
      8 --
      9 -- TALER is distributed in the hope that it will be useful, but WITHOUT ANY
     10 -- WARRANTY; without even the implied warranty of MERCHANTABILITY or FITNESS FOR
     11 -- A PARTICULAR PURPOSE.  See the GNU General Public License for more details.
     12 --
     13 -- You should have received a copy of the GNU General Public License along with
     14 -- TALER; see the file COPYING.  If not, see <http://www.gnu.org/licenses/>
     15 --
     16 
     17 SET search_path TO exchange;
     18 DROP FUNCTION IF EXISTS interval_to_start;
     19 CREATE OR REPLACE FUNCTION interval_to_start (
     20   IN in_timestamp TIMESTAMP,
     21   IN in_range statistic_range,
     22   OUT out_bucket_start INT8
     23 )
     24 LANGUAGE plpgsql
     25 AS $$
     26 BEGIN
     27   out_bucket_start = EXTRACT(EPOCH FROM DATE_TRUNC(in_range::text, in_timestamp));
     28 END $$;
     29 COMMENT ON FUNCTION interval_to_start
     30  IS 'computes the start time of the bucket for an event at the current time given the desired bucket range';
     31 
     32 
     33 DROP PROCEDURE IF EXISTS exchange_do_bump_number_bucket_stat;
     34 CREATE OR REPLACE PROCEDURE exchange_do_bump_number_bucket_stat(
     35   in_slug TEXT,
     36   in_h_payto BYTEA,
     37   in_timestamp TIMESTAMP,
     38   in_delta INT8
     39 )
     40 LANGUAGE plpgsql
     41 AS $$
     42 DECLARE
     43   my_meta INT8;
     44   my_range statistic_range;
     45   my_bucket_start INT8;
     46   my_curs CURSOR (arg_slug TEXT)
     47    FOR SELECT UNNEST(ranges)
     48          FROM exchange_statistic_bucket_meta
     49         WHERE slug=arg_slug;
     50 BEGIN
     51   SELECT bmeta_serial_id
     52     INTO my_meta
     53     FROM exchange_statistic_bucket_meta
     54    WHERE slug=in_slug
     55      AND stype='number';
     56   IF NOT FOUND
     57   THEN
     58     RETURN;
     59   END IF;
     60   OPEN my_curs (arg_slug:=in_slug);
     61   LOOP
     62     FETCH NEXT
     63       FROM my_curs
     64       INTO my_range;
     65     EXIT WHEN NOT FOUND;
     66     SELECT *
     67       INTO my_bucket_start
     68       FROM interval_to_start (in_timestamp, my_range);
     69 
     70     UPDATE exchange_statistic_bucket_counter
     71        SET cumulative_number = cumulative_number + in_delta
     72      WHERE bmeta_serial_id=my_meta
     73        AND h_payto=in_h_payto
     74        AND bucket_start=my_bucket_start
     75        AND bucket_range=my_range;
     76     IF NOT FOUND
     77     THEN
     78       INSERT INTO exchange_statistic_bucket_counter
     79         (bmeta_serial_id
     80         ,h_payto
     81         ,bucket_start
     82         ,bucket_range
     83         ,cumulative_number
     84         ) VALUES (
     85          my_meta
     86         ,in_h_payto
     87         ,my_bucket_start
     88         ,my_range
     89         ,in_delta);
     90     END IF;
     91   END LOOP;
     92   CLOSE my_curs;
     93 END $$;
     94 
     95 
     96 DROP PROCEDURE IF EXISTS exchange_do_bump_amount_bucket_stat;
     97 CREATE OR REPLACE PROCEDURE exchange_do_bump_amount_bucket_stat(
     98   in_slug TEXT,
     99   in_h_payto BYTEA,
    100   in_timestamp TIMESTAMP,
    101   in_delta taler_amount
    102 )
    103 LANGUAGE plpgsql
    104 AS $$
    105 DECLARE
    106   my_meta INT8;
    107   my_range statistic_range;
    108   my_bucket_start INT8;
    109   my_curs CURSOR (arg_slug TEXT)
    110    FOR SELECT UNNEST(ranges)
    111          FROM exchange_statistic_bucket_meta
    112         WHERE slug=arg_slug;
    113 BEGIN
    114   SELECT bmeta_serial_id
    115     INTO my_meta
    116     FROM exchange_statistic_bucket_meta
    117    WHERE slug=in_slug
    118      AND stype='amount';
    119   IF NOT FOUND
    120   THEN
    121     RETURN;
    122   END IF;
    123   OPEN my_curs (arg_slug:=in_slug);
    124   LOOP
    125     FETCH NEXT
    126       FROM my_curs
    127       INTO my_range;
    128     EXIT WHEN NOT FOUND;
    129     SELECT *
    130       INTO my_bucket_start
    131       FROM interval_to_start (in_timestamp, my_range);
    132 
    133     UPDATE exchange_statistic_bucket_amount
    134       SET
    135         cumulative_value.val = (cumulative_value).val + (in_delta).val
    136         + CASE
    137             WHEN (in_delta).frac + (cumulative_value).frac >= 100000000
    138             THEN 1
    139             ELSE 0
    140           END,
    141         cumulative_value.frac = (cumulative_value).frac + (in_delta).frac
    142         - CASE
    143             WHEN (in_delta).frac + (cumulative_value).frac >= 100000000
    144             THEN 100000000
    145             ELSE 0
    146           END
    147      WHERE bmeta_serial_id=my_meta
    148        AND h_payto=in_h_payto
    149        AND bucket_start=my_bucket_start
    150        AND bucket_range=my_range;
    151     IF NOT FOUND
    152     THEN
    153       INSERT INTO exchange_statistic_bucket_amount
    154         (bmeta_serial_id
    155         ,h_payto
    156         ,bucket_start
    157         ,bucket_range
    158         ,cumulative_value
    159         ) VALUES (
    160          my_meta
    161         ,in_h_payto
    162         ,my_bucket_start
    163         ,my_range
    164         ,in_delta);
    165     END IF;
    166   END LOOP;
    167   CLOSE my_curs;
    168 END $$;
    169 
    170 COMMENT ON PROCEDURE exchange_do_bump_amount_bucket_stat
    171   IS 'Updates an amount statistic tracked over buckets';
    172 
    173 
    174 DROP PROCEDURE IF EXISTS exchange_do_bump_number_interval_stat;
    175 CREATE OR REPLACE PROCEDURE exchange_do_bump_number_interval_stat(
    176   in_slug TEXT,
    177   in_h_payto BYTEA,
    178   in_timestamp TIMESTAMP,
    179   in_delta INT8
    180 )
    181 LANGUAGE plpgsql
    182 AS $$
    183 DECLARE
    184   my_now INT8;
    185   my_record RECORD;
    186   my_meta INT8;
    187   my_ranges INT8[];
    188   my_precisions INT8[];
    189   my_rangex INT8;
    190   my_precisionx INT8;
    191   my_start INT8;
    192   my_event INT8;
    193 BEGIN
    194   my_now = ROUND(EXTRACT(epoch FROM exchange_now()) * 1000000)::INT8 / 1000 / 1000;
    195   SELECT imeta_serial_id
    196         ,ranges AS ranges
    197         ,precisions AS precisions
    198     INTO my_record
    199     FROM exchange_statistic_interval_meta
    200    WHERE slug=in_slug
    201      AND stype='number';
    202   IF NOT FOUND
    203   THEN
    204     RETURN;
    205   END IF;
    206 
    207   my_start = ROUND(EXTRACT(epoch FROM in_timestamp) * 1000000)::INT8 / 1000 / 1000; -- convert to seconds
    208   my_precisions = my_record.precisions;
    209   my_ranges = my_record.ranges;
    210   my_rangex = NULL;
    211   FOR my_x IN 1..COALESCE(array_length(my_ranges,1),0)
    212   LOOP
    213     IF my_now - my_ranges[my_x] < my_start
    214     THEN
    215       my_rangex = my_ranges[my_x];
    216       my_precisionx = my_precisions[my_x];
    217       EXIT;
    218     END IF;
    219   END LOOP;
    220   IF my_rangex IS NULL
    221   THEN
    222     -- event is beyond the ranges we care about
    223     RETURN;
    224   END IF;
    225 
    226   my_meta = my_record.imeta_serial_id;
    227   my_start = my_start - my_start % my_precisionx; -- round down
    228 
    229   INSERT INTO exchange_statistic_counter_event AS msce
    230     (imeta_serial_id
    231     ,h_payto
    232     ,slot
    233     ,delta)
    234    VALUES
    235     (my_meta
    236     ,in_h_payto
    237     ,my_start
    238     ,in_delta)
    239    ON CONFLICT (imeta_serial_id, h_payto, slot)
    240    DO UPDATE SET
    241      delta = msce.delta + in_delta
    242    RETURNING nevent_serial_id
    243         INTO my_event;
    244 
    245   UPDATE exchange_statistic_interval_counter
    246      SET cumulative_number = cumulative_number + in_delta
    247    WHERE imeta_serial_id = my_meta
    248      AND h_payto = in_h_payto
    249      AND range=my_rangex;
    250   IF NOT FOUND
    251   THEN
    252     INSERT INTO exchange_statistic_interval_counter
    253       (imeta_serial_id
    254       ,h_payto
    255       ,range
    256       ,event_delimiter
    257       ,cumulative_number
    258      ) VALUES (
    259        my_meta
    260       ,in_h_payto
    261       ,my_rangex
    262       ,my_event
    263       ,in_delta);
    264   END IF;
    265 END $$;
    266 
    267 COMMENT ON PROCEDURE exchange_do_bump_number_interval_stat
    268   IS 'Updates a numeric statistic tracked over an interval';
    269 
    270 
    271 DROP PROCEDURE IF EXISTS exchange_do_bump_amount_interval_stat;
    272 CREATE OR REPLACE PROCEDURE exchange_do_bump_amount_interval_stat(
    273   in_slug TEXT,
    274   in_h_payto BYTEA,
    275   in_timestamp TIMESTAMP,
    276   in_delta taler_amount
    277 )
    278 LANGUAGE plpgsql
    279 AS $$
    280 DECLARE
    281   my_now INT8;
    282   my_record RECORD;
    283   my_meta INT8;
    284   my_ranges INT8[];
    285   my_precisions INT8[];
    286   my_x INT;
    287   my_rangex INT8;
    288   my_precisionx INT8;
    289   my_start INT8;
    290   my_event INT8;
    291 BEGIN
    292   my_now = ROUND(EXTRACT(epoch FROM exchange_now()) * 1000000)::INT8 / 1000 / 1000;
    293   SELECT imeta_serial_id
    294         ,ranges
    295         ,precisions
    296     INTO my_record
    297     FROM exchange_statistic_interval_meta
    298    WHERE slug=in_slug
    299      AND stype='amount';
    300   IF NOT FOUND
    301   THEN
    302     RETURN;
    303   END IF;
    304 
    305   my_start = ROUND(EXTRACT(epoch FROM in_timestamp) * 1000000)::INT8 / 1000 / 1000; -- convert to seconds since epoch
    306   my_precisions = my_record.precisions;
    307   my_ranges = my_record.ranges;
    308   my_rangex = NULL;
    309   FOR my_x IN 1..COALESCE(array_length(my_ranges,1),0)
    310   LOOP
    311     IF my_now - my_ranges[my_x] < my_start
    312     THEN
    313       my_rangex = my_ranges[my_x];
    314       my_precisionx = my_precisions[my_x];
    315       EXIT;
    316     END IF;
    317   END LOOP;
    318   IF my_rangex IS NULL
    319   THEN
    320     -- event is beyond the ranges we care about
    321     RETURN;
    322   END IF;
    323   my_start = my_start - my_start % my_precisionx; -- round down
    324   my_meta = my_record.imeta_serial_id;
    325 
    326   INSERT INTO exchange_statistic_amount_event AS msae
    327     (imeta_serial_id
    328     ,h_payto
    329     ,slot
    330     ,delta
    331     ) VALUES (
    332      my_meta
    333     ,in_h_payto
    334     ,my_start
    335     ,in_delta
    336     )
    337     ON CONFLICT (imeta_serial_id, h_payto, slot)
    338     DO UPDATE SET
    339       delta.val = (msae.delta).val + (in_delta).val
    340         + CASE
    341           WHEN (in_delta).frac + (msae.delta).frac >= 100000000
    342           THEN 1
    343           ELSE 0
    344         END,
    345       delta.frac = (msae.delta).frac + (in_delta).frac
    346         - CASE
    347           WHEN (in_delta).frac + (msae.delta).frac >= 100000000
    348           THEN 100000000
    349           ELSE 0
    350         END
    351     RETURNING aevent_serial_id
    352          INTO my_event;
    353 
    354   UPDATE exchange_statistic_interval_amount
    355     SET
    356       cumulative_value.val = (cumulative_value).val + (in_delta).val
    357       + CASE
    358           WHEN (in_delta).frac + (cumulative_value).frac >= 100000000
    359           THEN 1
    360           ELSE 0
    361         END,
    362       cumulative_value.frac = (cumulative_value).frac + (in_delta).frac
    363       - CASE
    364           WHEN (in_delta).frac + (cumulative_value).frac >= 100000000
    365           THEN 100000000
    366           ELSE 0
    367         END
    368    WHERE imeta_serial_id=my_meta
    369      AND h_payto=in_h_payto
    370      AND range=my_rangex;
    371   IF NOT FOUND
    372   THEN
    373     INSERT INTO exchange_statistic_interval_amount
    374       (imeta_serial_id
    375       ,h_payto
    376       ,range
    377       ,event_delimiter
    378       ,cumulative_value
    379       ) VALUES (
    380        my_meta
    381       ,in_h_payto
    382       ,my_rangex
    383       ,my_event
    384       ,in_delta);
    385   END IF;
    386 END $$;
    387 COMMENT ON PROCEDURE exchange_do_bump_amount_interval_stat
    388   IS 'Updates an amount statistic tracked over an interval';
    389 
    390 
    391 DROP PROCEDURE IF EXISTS exchange_do_bump_number_stat;
    392 CREATE OR REPLACE PROCEDURE exchange_do_bump_number_stat(
    393   in_slug TEXT,
    394   in_h_payto BYTEA,
    395   in_timestamp TIMESTAMP,
    396   in_delta INT8
    397 )
    398 LANGUAGE plpgsql
    399 AS $$
    400 BEGIN
    401   CALL exchange_do_bump_number_bucket_stat (in_slug, in_h_payto, in_timestamp, in_delta);
    402   CALL exchange_do_bump_number_interval_stat (in_slug, in_h_payto, in_timestamp, in_delta);
    403 END $$;
    404 COMMENT ON PROCEDURE exchange_do_bump_number_stat
    405   IS 'Updates a numeric statistic (bucket or interval)';
    406 
    407 
    408 DROP PROCEDURE IF EXISTS exchange_do_bump_amount_stat;
    409 CREATE OR REPLACE PROCEDURE exchange_do_bump_amount_stat(
    410   in_slug TEXT,
    411   in_h_payto BYTEA,
    412   in_timestamp TIMESTAMP,
    413   in_delta taler_amount
    414 )
    415 LANGUAGE plpgsql
    416 AS $$
    417 BEGIN
    418   CALL exchange_do_bump_amount_bucket_stat (in_slug, in_h_payto, in_timestamp, in_delta);
    419   CALL exchange_do_bump_amount_interval_stat (in_slug, in_h_payto, in_timestamp, in_delta);
    420 END $$;
    421 COMMENT ON PROCEDURE exchange_do_bump_amount_stat
    422   IS 'Updates an amount statistic (bucket or interval)';
    423 
    424 
    425 DROP FUNCTION IF EXISTS exchange_statistic_interval_number_get;
    426 CREATE OR REPLACE FUNCTION exchange_statistic_interval_number_get (
    427   IN in_slug TEXT,
    428   IN in_h_payto BYTEA
    429 )
    430 RETURNS SETOF exchange_statistic_interval_number_get_return_value
    431 LANGUAGE plpgsql
    432 AS $$
    433 DECLARE
    434   my_time INT8 DEFAULT ROUND(EXTRACT(epoch FROM exchange_now()) * 1000000)::INT8 / 1000 / 1000;
    435   my_ranges INT8[];
    436   my_range INT8;
    437   my_delta INT8;
    438   my_meta INT8;
    439   my_next_max_serial INT8;
    440   my_rec RECORD;
    441   my_irec RECORD;
    442   my_i INT;
    443   my_min_serial INT8 DEFAULT NULL;
    444   my_rval exchange_statistic_interval_number_get_return_value;
    445 BEGIN
    446   SELECT imeta_serial_id
    447         ,ranges
    448         ,precisions
    449     INTO my_rec
    450     FROM exchange_statistic_interval_meta
    451    WHERE slug=in_slug;
    452   IF NOT FOUND
    453   THEN
    454     RETURN;
    455   END IF;
    456   my_rval.rvalue = 0;
    457   my_ranges = my_rec.ranges;
    458   my_meta = my_rec.imeta_serial_id;
    459 
    460   FOR my_i IN 1..COALESCE(array_length(my_ranges,1),0)
    461   LOOP
    462     my_range = my_ranges[my_i];
    463     SELECT event_delimiter
    464           ,cumulative_number
    465       INTO my_irec
    466       FROM exchange_statistic_interval_counter
    467      WHERE imeta_serial_id = my_meta
    468        AND range = my_range
    469        AND h_payto = in_h_payto;
    470     IF FOUND
    471     THEN
    472       my_min_serial = my_irec.event_delimiter;
    473       my_rval.rvalue = my_rval.rvalue + my_irec.cumulative_number;
    474 
    475       -- Check if we have events that left the applicable range
    476       SELECT SUM(delta) AS delta_sum
    477         INTO my_irec
    478         FROM exchange_statistic_counter_event
    479        WHERE imeta_serial_id = my_meta
    480          AND h_payto = in_h_payto
    481          AND slot < my_time - my_range
    482          AND nevent_serial_id >= my_min_serial;
    483 
    484       IF FOUND AND my_irec.delta_sum IS NOT NULL
    485       THEN
    486         my_delta = my_irec.delta_sum;
    487         my_rval.rvalue = my_rval.rvalue - my_delta;
    488 
    489         -- First find out the next event delimiter value
    490         SELECT nevent_serial_id
    491           INTO my_next_max_serial
    492           FROM exchange_statistic_counter_event
    493          WHERE imeta_serial_id = my_meta
    494            AND h_payto = in_h_payto
    495            AND slot >= my_time - my_range
    496            AND nevent_serial_id >= my_min_serial
    497          ORDER BY slot ASC
    498          LIMIT 1;
    499 
    500         IF FOUND
    501         THEN
    502           -- remove expired events from the sum of the current slot
    503 
    504           UPDATE exchange_statistic_interval_counter
    505              SET cumulative_number = cumulative_number - my_delta,
    506                  event_delimiter = my_next_max_serial
    507            WHERE imeta_serial_id = my_meta
    508              AND h_payto = in_h_payto
    509              AND range = my_range;
    510         ELSE
    511           -- actually, slot is now empty, remove it entirely
    512           DELETE FROM exchange_statistic_interval_counter
    513            WHERE imeta_serial_id = my_meta
    514              AND h_payto = in_h_payto
    515              AND range = my_range;
    516         END IF;
    517         IF (my_i < array_length(my_ranges,1))
    518         THEN
    519           -- carry over all events into the next slot
    520           UPDATE exchange_statistic_interval_counter AS usic SET
    521             cumulative_number = cumulative_number + my_delta,
    522             event_delimiter = LEAST(usic.event_delimiter,my_min_serial)
    523            WHERE imeta_serial_id = my_meta
    524              AND h_payto = in_h_payto
    525              AND range=my_ranges[my_i+1];
    526           IF NOT FOUND
    527           THEN
    528             INSERT INTO exchange_statistic_interval_counter
    529               (imeta_serial_id
    530               ,h_payto
    531               ,range
    532               ,event_delimiter
    533               ,cumulative_number
    534               ) VALUES (
    535                my_meta
    536               ,in_h_payto
    537               ,my_ranges[my_i+1]
    538               ,my_min_serial
    539               ,my_delta);
    540           END IF;
    541         ELSE
    542           -- events are obsolete, delete them
    543           DELETE FROM exchange_statistic_counter_event
    544                 WHERE imeta_serial_id = my_meta
    545                   AND h_payto = in_h_payto
    546                   AND slot < my_time - my_range;
    547         END IF;
    548       END IF;
    549 
    550     END IF;
    551     -- Smaller ranges contribute to this interval even when this range has
    552     -- no stored events of its own (for example the 52-week TOPS total).
    553     my_rval.range = my_range;
    554     RETURN NEXT my_rval;
    555   END LOOP;
    556 END $$;
    557 
    558 COMMENT ON FUNCTION exchange_statistic_interval_number_get
    559   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';
    560 
    561 
    562 DROP FUNCTION IF EXISTS exchange_statistic_interval_amount_get;
    563 CREATE OR REPLACE FUNCTION exchange_statistic_interval_amount_get (
    564   IN in_slug TEXT,
    565   IN in_h_payto BYTEA
    566 )
    567 RETURNS SETOF exchange_statistic_interval_amount_get_return_value
    568 LANGUAGE plpgsql
    569 AS $$
    570 DECLARE
    571   my_time INT8 DEFAULT ROUND(EXTRACT(epoch FROM exchange_now()) * 1000000)::INT8 / 1000 / 1000;
    572   my_ranges INT8[];
    573   my_range INT8;
    574   my_delta_value INT8;
    575   my_delta_frac INT8;
    576   my_delta taler_amount;
    577   my_meta INT8;
    578   my_next_max_serial INT8;
    579   my_rec RECORD;
    580   my_irec RECORD;
    581   my_jrec RECORD;
    582   my_i INT;
    583   my_min_serial INT8 DEFAULT NULL;
    584   my_rval exchange_statistic_interval_amount_get_return_value;
    585 BEGIN
    586   SELECT imeta_serial_id
    587         ,ranges
    588         ,precisions
    589     INTO my_rec
    590     FROM exchange_statistic_interval_meta
    591    WHERE slug=in_slug;
    592   IF NOT FOUND
    593   THEN
    594     RETURN;
    595   END IF;
    596 
    597   my_meta = my_rec.imeta_serial_id;
    598   my_ranges = my_rec.ranges;
    599 
    600   my_rval.rvalue.val = 0;
    601   my_rval.rvalue.frac = 0;
    602 
    603   FOR my_i IN 1..COALESCE(array_length(my_ranges,1),0)
    604   LOOP
    605     my_range = my_ranges[my_i];
    606     SELECT event_delimiter
    607           ,cumulative_value
    608       INTO my_irec
    609       FROM exchange_statistic_interval_amount
    610      WHERE imeta_serial_id = my_meta
    611        AND h_payto = in_h_payto
    612        AND range = my_range;
    613 
    614     IF FOUND
    615     THEN
    616       my_min_serial = my_irec.event_delimiter;
    617       my_rval.rvalue.val = (my_rval.rvalue).val + (my_irec.cumulative_value).val + (my_irec.cumulative_value).frac / 100000000;
    618       my_rval.rvalue.frac = (my_rval.rvalue).frac + (my_irec.cumulative_value).frac % 100000000;
    619       IF (my_rval.rvalue).frac >= 100000000
    620       THEN
    621         my_rval.rvalue.frac = (my_rval.rvalue).frac - 100000000;
    622         my_rval.rvalue.val = (my_rval.rvalue).val + 1;
    623       END IF;
    624 
    625       -- Check if we have events that left the applicable range
    626       SELECT SUM((esae.delta).val) AS value_sum
    627             ,SUM((esae.delta).frac) AS frac_sum
    628         INTO my_jrec
    629         FROM exchange_statistic_amount_event esae
    630        WHERE imeta_serial_id = my_meta
    631          AND h_payto = in_h_payto
    632          AND slot < my_time - my_range
    633          AND aevent_serial_id >= my_min_serial;
    634 
    635       IF FOUND AND my_jrec.value_sum IS NOT NULL
    636       THEN
    637         -- Normalize sum
    638         my_delta_value = my_jrec.value_sum + my_jrec.frac_sum / 100000000;
    639         my_delta_frac = my_jrec.frac_sum % 100000000;
    640         my_rval.rvalue.val = (my_rval.rvalue).val - my_delta_value;
    641         IF ((my_rval.rvalue).frac >= my_delta_frac)
    642         THEN
    643           my_rval.rvalue.frac = (my_rval.rvalue).frac - my_delta_frac;
    644         ELSE
    645           my_rval.rvalue.frac = 100000000 + (my_rval.rvalue).frac - my_delta_frac;
    646           my_rval.rvalue.val = (my_rval.rvalue).val - 1;
    647         END IF;
    648 
    649         -- First find out the next event delimiter value
    650         SELECT aevent_serial_id
    651           INTO my_next_max_serial
    652           FROM exchange_statistic_amount_event
    653          WHERE imeta_serial_id = my_meta
    654            AND h_payto = in_h_payto
    655            AND slot >= my_time - my_range
    656            AND aevent_serial_id >= my_min_serial
    657          ORDER BY slot ASC
    658          LIMIT 1;
    659         IF FOUND
    660         THEN
    661           -- remove expired events from the sum of the current slot
    662           UPDATE exchange_statistic_interval_amount SET
    663              cumulative_value.val = (cumulative_value).val - my_delta_value
    664               - CASE
    665                   WHEN (cumulative_value).frac < my_delta_frac
    666                   THEN 1
    667                   ELSE 0
    668                 END,
    669              cumulative_value.frac = (cumulative_value).frac - my_delta_frac
    670              + CASE
    671                  WHEN (cumulative_value).frac < my_delta_frac
    672                  THEN 100000000
    673                  ELSE 0
    674                END,
    675              event_delimiter = my_next_max_serial
    676            WHERE imeta_serial_id = my_meta
    677              AND h_payto = in_h_payto
    678              AND range = my_range;
    679         ELSE
    680           -- actually, slot is now empty, remove it entirely
    681           DELETE FROM exchange_statistic_interval_amount
    682            WHERE imeta_serial_id = my_meta
    683              AND h_payto = in_h_payto
    684              AND range = my_range;
    685         END IF;
    686         IF (my_i < array_length(my_ranges,1))
    687         THEN
    688           -- carry over all events into the next (larger) slot
    689           UPDATE exchange_statistic_interval_amount AS msia SET
    690             cumulative_value.val = (cumulative_value).val + my_delta_value
    691               + CASE
    692                  WHEN (cumulative_value).frac + my_delta_frac >= 100000000
    693                  THEN 1
    694                  ELSE 0
    695                END,
    696             -- Note: the fraction, not the value; this used to add
    697             -- my_delta_value here and so replaced the fraction of the
    698             -- carried-over amount with its integer part.
    699             cumulative_value.frac = (cumulative_value).frac + my_delta_frac
    700               - CASE
    701                  WHEN (cumulative_value).frac + my_delta_frac >= 100000000
    702                  THEN 100000000
    703                  ELSE 0
    704                END,
    705             event_delimiter = LEAST (msia.event_delimiter,my_min_serial)
    706            WHERE imeta_serial_id = my_meta
    707              AND h_payto = in_h_payto
    708              AND range=my_ranges[my_i+1];
    709           IF NOT FOUND
    710           THEN
    711             my_delta.val = my_delta_value;
    712             my_delta.frac = my_delta_frac;
    713             INSERT INTO exchange_statistic_interval_amount
    714               (imeta_serial_id
    715               ,h_payto
    716               ,event_delimiter
    717               ,range
    718               ,cumulative_value
    719               ) VALUES (
    720                my_meta
    721               ,in_h_payto
    722               ,my_min_serial
    723               ,my_ranges[my_i+1]
    724               ,my_delta);
    725           END IF;
    726         ELSE
    727           -- events are obsolete, delete them
    728           DELETE FROM exchange_statistic_amount_event
    729                 WHERE imeta_serial_id = my_meta
    730                   AND h_payto = in_h_payto
    731                   AND slot < my_time - my_range;
    732         END IF;
    733       END IF;
    734 
    735     END IF;
    736     -- Smaller ranges contribute to this interval even when this range has
    737     -- no stored events of its own (for example the 52-week TOPS total).
    738     my_rval.range = my_range;
    739     RETURN NEXT my_rval;
    740   END LOOP; -- over my_ranges
    741 END $$;
    742 
    743 COMMENT ON FUNCTION exchange_statistic_interval_amount_get
    744   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';
    745 
    746 
    747 
    748 
    749 
    750 DROP PROCEDURE IF EXISTS exchange_statistic_counter_gc;
    751 CREATE OR REPLACE PROCEDURE exchange_statistic_counter_gc ()
    752 LANGUAGE plpgsql
    753 AS $$
    754 DECLARE
    755   my_time INT8 DEFAULT ROUND(EXTRACT(epoch FROM exchange_now()) * 1000000)::INT8 / 1000 / 1000;
    756   my_h_payto BYTEA;
    757   my_rec RECORD;
    758   my_sum RECORD;
    759   my_meta INT8;
    760   my_ranges INT8[];
    761   my_precisions INT8[];
    762   my_precision INT4;
    763   my_i INT4;
    764   min_slot INT8;
    765   max_slot INT8;
    766   end_slot INT8;
    767   my_total INT8;
    768 BEGIN
    769   -- GC for all instances
    770   FOR my_h_payto IN
    771     SELECT DISTINCT h_payto
    772       FROM exchange_statistic_counter_event
    773   LOOP
    774   -- Do combination work for all numeric statistic events
    775   FOR my_rec IN
    776     SELECT imeta_serial_id
    777           ,ranges
    778           ,precisions
    779           ,slug
    780       FROM exchange_statistic_interval_meta
    781   LOOP
    782     -- First, we query the current interval statistic to update its counters
    783     PERFORM FROM exchange_statistic_interval_number_get (my_rec.slug, my_h_payto);
    784 
    785     my_meta = my_rec.imeta_serial_id;
    786     my_ranges = my_rec.ranges;
    787     my_precisions = my_rec.precisions;
    788 
    789     FOR my_i IN 1..COALESCE(array_length(my_ranges,1),0)
    790     LOOP
    791       my_precision = my_precisions[my_i];
    792       IF 1 >= my_precision
    793       THEN
    794         -- Cannot coarsen in this case
    795         CONTINUE;
    796       END IF;
    797 
    798       IF 1 = my_i
    799       THEN
    800         min_slot = 0;
    801       ELSE
    802         min_slot = my_ranges[my_i - 1];
    803       END IF;
    804       end_slot = my_ranges[my_i];
    805 --    RAISE NOTICE 'Coarsening from [%,%) at %', my_time - end_slot, my_time - min_slot, my_precision;
    806 
    807       LOOP
    808         EXIT WHEN min_slot >= end_slot;
    809         max_slot = min_slot + my_precision;
    810         SELECT SUM(delta) AS total,
    811                COUNT(*)   AS matches,
    812                MIN(nevent_serial_id) AS rep_serial_id
    813           INTO my_sum
    814           FROM exchange_statistic_counter_event
    815          WHERE h_payto=my_h_payto
    816            AND imeta_serial_id=my_meta
    817            AND slot >= my_time - max_slot
    818            AND slot  < my_time - min_slot;
    819 
    820 --      RAISE NOTICE 'Found % entries between [%,%)', my_sum.matches, my_time - max_slot, my_time - min_slot;
    821         -- we only proceed if we had more then one match (optimization)
    822         IF FOUND AND my_sum.matches > 1
    823         THEN
    824           my_total = my_sum.total;
    825 
    826 --        RAISE NOTICE 'combining % entries to representative % for slots [%-%)', my_sum.matches, my_sum.rep_serial_id, my_time - max_slot, my_time - min_slot;
    827 
    828           -- combine entries
    829           DELETE FROM exchange_statistic_counter_event
    830            WHERE h_payto=my_h_payto
    831              AND imeta_serial_id=my_meta
    832              AND slot >= my_time - max_slot
    833              AND slot  < my_time - min_slot
    834              AND nevent_serial_id > my_sum.rep_serial_id;
    835            -- Now update the representative to the sum
    836           UPDATE exchange_statistic_counter_event SET
    837             delta = my_total
    838            WHERE imeta_serial_id = my_meta
    839              AND h_payto = my_h_payto
    840              AND nevent_serial_id = my_sum.rep_serial_id;
    841         END IF;
    842         min_slot = min_slot + my_precision;
    843       END LOOP; -- min_slot to end_slot by precision loop
    844     END LOOP; -- my_i loop
    845     -- Finally, delete all events beyond the range we care about
    846 
    847 --  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)];
    848     DELETE FROM exchange_statistic_counter_event
    849      WHERE h_payto=my_h_payto
    850        AND imeta_serial_id=my_meta
    851        AND slot < my_time - my_ranges[array_length(my_ranges,1)];
    852   END LOOP; -- my_rec loop
    853   END LOOP; -- my_h_payto loop
    854 END $$;
    855 COMMENT ON PROCEDURE exchange_statistic_counter_gc
    856   IS 'Performs garbage collection and compaction of the exchange_statistic_counter_event table';
    857 
    858 
    859 
    860 DROP PROCEDURE IF EXISTS exchange_statistic_amount_gc;
    861 CREATE OR REPLACE PROCEDURE exchange_statistic_amount_gc ()
    862 LANGUAGE plpgsql
    863 AS $$
    864 DECLARE
    865   my_time INT8 DEFAULT ROUND(EXTRACT(epoch FROM exchange_now()) * 1000000)::INT8 / 1000 / 1000;
    866   my_h_payto BYTEA;
    867   my_rec RECORD;
    868   my_sum RECORD;
    869   my_meta INT8;
    870   my_ranges INT8[];
    871   my_precisions INT8[];
    872   my_precision INT4;
    873   my_i INT4;
    874   min_slot INT8;
    875   max_slot INT8;
    876   end_slot INT8;
    877   my_total_val INT8;
    878   my_total_frac INT8;
    879 BEGIN
    880   -- GC for all accounts
    881   FOR my_h_payto IN
    882     SELECT DISTINCT h_payto
    883       FROM exchange_statistic_counter_event
    884   LOOP
    885   -- Do combination work for all numeric statistic events
    886   FOR my_rec IN
    887     SELECT imeta_serial_id
    888           ,ranges
    889           ,precisions
    890           ,slug
    891       FROM exchange_statistic_interval_meta
    892   LOOP
    893 
    894     -- First, we query the current interval statistic to update its counters
    895     PERFORM FROM exchange_statistic_interval_amount_get (my_rec.slug, my_h_payto);
    896 
    897     my_meta = my_rec.imeta_serial_id;
    898     my_ranges = my_rec.ranges;
    899     my_precisions = my_rec.precisions;
    900 
    901     FOR my_i IN 1..COALESCE(array_length(my_ranges,1),0)
    902     LOOP
    903       my_precision = my_precisions[my_i];
    904       IF 1 >= my_precision
    905       THEN
    906         -- Cannot coarsen in this case
    907         CONTINUE;
    908       END IF;
    909 
    910       IF 1 = my_i
    911       THEN
    912         min_slot = 0;
    913       ELSE
    914         min_slot = my_ranges[my_i - 1];
    915       END IF;
    916       end_slot = my_ranges[my_i];
    917 
    918 --    RAISE NOTICE 'Coarsening from [%,%) at %', my_time - end_slot, my_time - min_slot, my_precision;
    919       LOOP
    920         EXIT WHEN min_slot >= end_slot;
    921         max_slot = min_slot + my_precision;
    922         SELECT SUM((delta).val)  AS total_val,
    923                SUM((delta).frac) AS total_frac,
    924                COUNT(*)          AS matches,
    925                MIN(aevent_serial_id) AS rep_serial_id
    926           INTO my_sum
    927           FROM exchange_statistic_amount_event
    928          WHERE imeta_serial_id=my_meta
    929            AND h_payto=my_h_payto
    930            AND slot >= my_time - max_slot
    931            AND slot  < my_time - min_slot;
    932         -- we only proceed if we had more then one match (optimization)
    933         IF FOUND AND my_sum.matches > 1
    934         THEN
    935           -- normalize new total
    936           my_total_frac = my_sum.total_frac % 100000000;
    937           my_total_val = my_sum.total_val + my_sum.total_frac / 100000000;
    938 
    939           -- combine entries
    940           DELETE FROM exchange_statistic_amount_event
    941            WHERE imeta_serial_id=my_meta
    942              AND h_payto=my_h_payto
    943              AND slot >= my_time - max_slot
    944              AND slot  < my_time - min_slot
    945              AND aevent_serial_id > my_sum.rep_serial_id;
    946           -- Now update the representative to the sum
    947           UPDATE exchange_statistic_amount_event SET
    948              delta.val = my_total_val
    949             ,delta.frac = my_total_frac
    950            WHERE imeta_serial_id = my_meta
    951              AND h_payto = my_h_payto
    952              AND aevent_serial_id = my_sum.rep_serial_id;
    953         END IF;
    954         min_slot = min_slot + my_precision;
    955       END LOOP; -- min_slot to end_slot by precision loop
    956     END LOOP; -- my_i loop
    957     -- Finally, delete all events beyond the range we care about
    958 
    959 --  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)];
    960     DELETE FROM exchange_statistic_amount_event
    961      WHERE h_payto=my_h_payto
    962        AND imeta_serial_id=my_meta
    963        AND slot < my_time - my_ranges[array_length(my_ranges,1)];
    964     END LOOP; -- my_rec loop
    965   END LOOP; -- my_h_payto loop
    966 END $$;
    967 COMMENT ON PROCEDURE exchange_statistic_amount_gc
    968   IS 'Performs garbage collection and compaction of the exchange_statistic_amount_event table';
    969 
    970 
    971 
    972 DROP PROCEDURE IF EXISTS exchange_statistic_bucket_gc;
    973 CREATE OR REPLACE PROCEDURE exchange_statistic_bucket_gc ()
    974 LANGUAGE plpgsql
    975 AS $$
    976 DECLARE
    977   my_rec RECORD;
    978   my_range TEXT;
    979   my_now INT8;
    980   my_end INT8;
    981 BEGIN
    982   my_now = EXTRACT(EPOCH FROM exchange_now()); -- seconds since epoch
    983   FOR my_rec IN
    984     SELECT bmeta_serial_id
    985           ,stype
    986           ,ranges[array_length(ranges,1)] AS range
    987           ,ages[array_length(ages,1)] AS age
    988       FROM exchange_statistic_bucket_meta
    989   LOOP
    990     my_range = '1 ' || my_rec.range::TEXT;
    991     my_end = my_now - my_rec.age * EXTRACT(SECONDS FROM (SELECT my_range::INTERVAL)); -- age is given in multiples of the range (in seconds)
    992     IF my_rec.stype = 'amount'
    993     THEN
    994       DELETE
    995         FROM exchange_statistic_bucket_amount
    996        WHERE bmeta_serial_id = my_rec.bmeta_serial_id
    997          AND bucket_start < my_end;
    998     ELSE
    999       DELETE
   1000         FROM exchange_statistic_bucket_counter
   1001        WHERE bmeta_serial_id = my_rec.bmeta_serial_id
   1002          AND bucket_start < my_end;
   1003     END IF;
   1004   END LOOP;
   1005 END $$;
   1006 COMMENT ON PROCEDURE exchange_statistic_bucket_gc
   1007   IS 'Performs garbage collection of the exchange_statistic_bucket_counter and exchange_statistic_bucket_amount tables';
   1008 
   1009 
   1010 
   1011 DROP FUNCTION IF EXISTS exchange_drop_customization;
   1012 CREATE OR REPLACE FUNCTION exchange_drop_customization (
   1013   IN in_schema TEXT,
   1014   OUT out_found BOOLEAN
   1015 )
   1016 LANGUAGE plpgsql
   1017 AS $$
   1018 DECLARE
   1019   my_xpatches TEXT;
   1020 BEGIN
   1021   -- Update DB versioning table.
   1022   out_found = FALSE;
   1023   FOR my_xpatches IN
   1024     SELECT patch_name
   1025       FROM _v.patches
   1026      WHERE starts_with(patch_name, in_schema || '-')
   1027   LOOP
   1028     PERFORM _v.unregister_patch(my_xpatches);
   1029     out_found = TRUE;
   1030   END LOOP;
   1031 
   1032   IF out_found
   1033   THEN
   1034     -- Drop the schema with all stored procedures/functions.
   1035     -- This also removes all associated triggers, hence CASCADE.
   1036     EXECUTE FORMAT('DROP SCHEMA %s CASCADE'
   1037       ,in_schema);
   1038   END IF;
   1039 
   1040   -- Finally, need to also remove entries from the statistics meta-tables.
   1041   -- Doing so also DELETEs the associated statistics, hence CASCADE.
   1042   DELETE
   1043      FROM exchange_statistic_interval_meta
   1044     WHERE origin=in_schema;
   1045   DELETE
   1046      FROM exchange_statistic_bucket_meta
   1047     WHERE origin=in_schema;
   1048 END $$;
   1049 COMMENT ON FUNCTION exchange_drop_customization
   1050   IS 'Removes all entries related to a particular exchange customization schema';