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';