LCOV - code coverage report
Current view: top level - exchangedb - iterate_aggregation_wtids_above_serial_id.c (source / functions) Coverage Total Hit
Test: coverage.info Lines: 85.7 % 28 24
Test Date: 2026-09-11 18:55:36 Functions: 100.0 % 2 2

            Line data    Source code
       1              : /*
       2              :    This file is part of TALER
       3              :    Copyright (C) 2026 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              :  * @file src/exchangedb/iterate_aggregation_wtids_above_serial_id.c
      18              :  * @brief Implementation of the iterate_aggregation_wtids_above_serial_id function for Postgres
      19              :  * @author Christian Grothoff
      20              :  */
      21              : #include "taler/taler_pq_lib.h"
      22              : #include "helper.h"
      23              : #include "exchange-database/iterate_aggregation_wtids_above_serial_id.h"
      24              : 
      25              : 
      26              : /**
      27              :  * Closure for #aggregation_wtid_cb().
      28              :  */
      29              : struct AggregationWtidContext
      30              : {
      31              :   /**
      32              :    * Function to call for each aggregation.
      33              :    */
      34              :   TALER_EXCHANGEDB_AggregationWtidCallback cb;
      35              : 
      36              :   /**
      37              :    * Closure for @e cb.
      38              :    */
      39              :   void *cb_cls;
      40              : 
      41              :   /**
      42              :    * Query status to return.
      43              :    */
      44              :   enum GNUNET_DB_QueryStatus qs;
      45              : };
      46              : 
      47              : 
      48              : /**
      49              :  * Helper function for
      50              :  * #TALER_EXCHANGEDB_iterate_aggregation_wtids_above_serial_id().
      51              :  * To be called with the results of a SELECT statement
      52              :  * that has returned @a num_results results.
      53              :  *
      54              :  * @param cls closure of type `struct AggregationWtidContext *`
      55              :  * @param result the postgres result
      56              :  * @param num_results the number of results in @a result
      57              :  */
      58              : static void
      59            3 : aggregation_wtid_cb (void *cls,
      60              :                      PGresult *result,
      61              :                      unsigned int num_results)
      62              : {
      63            3 :   struct AggregationWtidContext *awc = cls;
      64              : 
      65            6 :   for (unsigned int i = 0; i < num_results; i++)
      66              :   {
      67              :     uint64_t rowid;
      68              :     struct TALER_WireTransferIdentifierRawP wtid;
      69              :     struct TALER_FullPaytoHashP wire_target_h_payto;
      70              :     bool pending;
      71            3 :     struct GNUNET_PQ_ResultSpec rs[] = {
      72            3 :       GNUNET_PQ_result_spec_uint64 ("aggregation_serial_id",
      73              :                                     &rowid),
      74            3 :       GNUNET_PQ_result_spec_auto_from_type ("wtid_raw",
      75              :                                             &wtid),
      76            3 :       GNUNET_PQ_result_spec_auto_from_type ("wire_target_h_payto",
      77              :                                             &wire_target_h_payto),
      78            3 :       GNUNET_PQ_result_spec_bool ("pending",
      79              :                                   &pending),
      80              :       GNUNET_PQ_result_spec_end
      81              :     };
      82              :     enum GNUNET_GenericReturnValue rval;
      83              : 
      84            3 :     if (GNUNET_OK !=
      85            3 :         GNUNET_PQ_extract_result (result,
      86              :                                   rs,
      87              :                                   i))
      88              :     {
      89            0 :       GNUNET_break (0);
      90            0 :       awc->qs = GNUNET_DB_STATUS_HARD_ERROR;
      91            0 :       return;
      92              :     }
      93            3 :     awc->qs = i + 1;
      94            3 :     rval = awc->cb (awc->cb_cls,
      95              :                     rowid,
      96              :                     &wtid,
      97              :                     &wire_target_h_payto,
      98              :                     pending);
      99            3 :     GNUNET_PQ_cleanup_result (rs);
     100            3 :     if (GNUNET_OK != rval)
     101            0 :       break;
     102              :   }
     103              : }
     104              : 
     105              : 
     106              : enum GNUNET_DB_QueryStatus
     107            3 : TALER_EXCHANGEDB_iterate_aggregation_wtids_above_serial_id (
     108              :   struct TALER_EXCHANGEDB_PostgresContext *pg,
     109              :   uint64_t serial_id,
     110              :   TALER_EXCHANGEDB_AggregationWtidCallback cb,
     111              :   TALER_EXCHANGEDB_AGGREGATION_WTID_RESULT_CLOSURE *cb_cls)
     112              : {
     113            3 :   struct GNUNET_PQ_QueryParam params[] = {
     114            3 :     GNUNET_PQ_query_param_uint64 (&serial_id),
     115              :     GNUNET_PQ_query_param_end
     116              :   };
     117            3 :   struct AggregationWtidContext awc = {
     118              :     .cb = cb,
     119              :     .cb_cls = cb_cls
     120              :   };
     121              :   enum GNUNET_DB_QueryStatus qs;
     122              : 
     123              :   /* Rows that already have their `wire_out' are returned as well, flagged as
     124              :      not pending: the caller needs every serial ID to be able to advance its
     125              :      cursor past them, but knowing here that the transfer happened saves it
     126              :      from asking #TALER_EXCHANGEDB_get_pending_aggregation() -- a full
     127              :      recomputation of the transfer -- once per wire transfer the exchange ever
     128              :      made, which is what the first run after an upgrade would otherwise do. */
     129            3 :   PREPARE (pg,
     130              :            "iterate_aggregation_wtids_above_serial_id",
     131              :            "SELECT"
     132              :            " atr.aggregation_serial_id"
     133              :            ",atr.wtid_raw"
     134              :            ",bdep.wire_target_h_payto"
     135              :            ",NOT EXISTS ("
     136              :            "   SELECT 1"
     137              :            "     FROM wire_out wo"
     138              :            "    WHERE wo.wtid_raw=atr.wtid_raw"
     139              :            " ) AS pending"
     140              :            " FROM aggregation_tracking atr"
     141              :            " JOIN batch_deposits bdep"
     142              :            "   USING (batch_deposit_serial_id)"
     143              :            " WHERE atr.aggregation_serial_id>=$1"
     144              :            " ORDER BY atr.aggregation_serial_id ASC;");
     145            3 :   qs = GNUNET_PQ_eval_prepared_multi_select (
     146              :     pg->conn,
     147              :     "iterate_aggregation_wtids_above_serial_id",
     148              :     params,
     149              :     &aggregation_wtid_cb,
     150              :     &awc);
     151            3 :   if (qs > 0)
     152            2 :     return awc.qs;
     153            1 :   GNUNET_break (GNUNET_DB_STATUS_HARD_ERROR != qs);
     154            1 :   return qs;
     155              : }
        

Generated by: LCOV version 2.0-1