commit fe79c64eaf74efe0fdf3308ed7c797b1cbe56db1
parent 7bfaf2c5e3f7ee89c29cc9175da77ee7cb582b18
Author: Antoine A <>
Date: Fri, 2 Oct 2026 12:57:00 +0200
common: fix history id not being ordered correctly
Diffstat:
13 files changed, 218 insertions(+), 33 deletions(-)
diff --git a/adapters/taler-cyclos/db/cyclos-0003.sql b/adapters/taler-cyclos/db/cyclos-0003.sql
@@ -0,0 +1,31 @@
+--
+-- This file is part of TALER
+-- Copyright (C) 2025-2026 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/>
+
+SELECT _v.register_patch('cyclos-0003', NULL, NULL);
+
+SET search_path TO cyclos;
+
+-- Preserve existing incoming-history IDs and persisted consumer cursors.
+-- Allocate subsequent IDs when a transaction becomes eligible for history.
+ALTER TABLE taler_in ADD COLUMN taler_in_id INT8 GENERATED BY DEFAULT AS IDENTITY;
+UPDATE taler_in SET taler_in_id = tx_in_id;
+ALTER TABLE taler_in
+ ALTER COLUMN taler_in_id SET GENERATED ALWAYS,
+ DROP CONSTRAINT taler_in_pkey,
+ ADD UNIQUE (tx_in_id),
+ ADD PRIMARY KEY (taler_in_id);
+
+SELECT setval(pg_get_serial_sequence('taler_in', 'taler_in_id'),
+ (SELECT COALESCE(MAX(tx_in_id), 0) + 1 FROM tx_in), false);
diff --git a/adapters/taler-cyclos/db/cyclos-procedures.sql b/adapters/taler-cyclos/db/cyclos-procedures.sql
@@ -63,6 +63,7 @@ LANGUAGE plpgsql AS $$
DECLARE
local_authorization_pub BYTEA;
local_authorization_sig BYTEA;
+local_taler_in_id INT8;
BEGIN
out_pending=false;
-- Check for idempotence
@@ -144,9 +145,9 @@ ELSIF in_type IS NOT NULL THEN
in_metadata,
local_authorization_pub,
local_authorization_sig
- );
+ ) RETURNING taler_in_id INTO local_taler_in_id;
-- Notify new incoming talerable transaction registration
- PERFORM pg_notify('taler_in', out_tx_row_id || '');
+ PERFORM pg_notify('taler_in', local_taler_in_id::text);
END IF;
END $$;
COMMENT ON FUNCTION register_tx_in IS 'Register an incoming transaction idempotently';
@@ -425,6 +426,7 @@ CREATE FUNCTION register_prepared_transfers (
LANGUAGE plpgsql AS $$
DECLARE
talerable_tx INT8;
+ local_taler_in_id INT8;
idempotent BOOLEAN;
BEGIN
@@ -470,9 +472,9 @@ IF in_recurrent THEN
INSERT INTO taler_in (tx_in_id, type, metadata, authorization_pub, authorization_sig)
SELECT moved_tx.tx_in_id, in_type, in_account_pub, in_authorization_pub, in_authorization_sig
FROM moved_tx
- RETURNING tx_in_id INTO talerable_tx;
+ RETURNING tx_in_id, taler_in_id INTO talerable_tx, local_taler_in_id;
IF talerable_tx IS NOT NULL THEN
- PERFORM pg_notify('taler_in', talerable_tx::text);
+ PERFORM pg_notify('taler_in', local_taler_in_id::text);
END IF;
ELSE
-- Bounce all pending
diff --git a/adapters/taler-cyclos/src/db.rs b/adapters/taler-cyclos/src/db.rs
@@ -562,7 +562,7 @@ pub async fn incoming_history(
) -> sqlx::Result<Vec<IncomingBankTransaction>> {
history(
db,
- "tx_in_id",
+ "taler_in_id",
params,
listen,
|| {
@@ -570,7 +570,7 @@ pub async fn incoming_history(
"
SELECT
type,
- tx_in_id,
+ taler_in_id,
amount,
debit_account,
debit_name,
diff --git a/adapters/taler-magnet-bank/db/magnet-bank-0003.sql b/adapters/taler-magnet-bank/db/magnet-bank-0003.sql
@@ -0,0 +1,31 @@
+--
+-- This file is part of TALER
+-- Copyright (C) 2026 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/>
+
+SELECT _v.register_patch('magnet-bank-0003', NULL, NULL);
+
+SET search_path TO magnet_bank;
+
+-- Preserve existing incoming-history IDs and persisted consumer cursors.
+-- Allocate subsequent IDs when a transaction becomes eligible for history.
+ALTER TABLE taler_in ADD COLUMN taler_in_id INT8 GENERATED BY DEFAULT AS IDENTITY;
+UPDATE taler_in SET taler_in_id = tx_in_id;
+ALTER TABLE taler_in
+ ALTER COLUMN taler_in_id SET GENERATED ALWAYS,
+ DROP CONSTRAINT taler_in_pkey,
+ ADD UNIQUE (tx_in_id),
+ ADD PRIMARY KEY (taler_in_id);
+
+SELECT setval(pg_get_serial_sequence('taler_in', 'taler_in_id'),
+ (SELECT COALESCE(MAX(tx_in_id), 0) + 1 FROM tx_in), false);
diff --git a/adapters/taler-magnet-bank/db/magnet-bank-procedures.sql b/adapters/taler-magnet-bank/db/magnet-bank-procedures.sql
@@ -62,6 +62,7 @@ LANGUAGE plpgsql AS $$
DECLARE
local_authorization_pub BYTEA;
local_authorization_sig BYTEA;
+local_taler_in_id INT8;
BEGIN
out_pending=false;
-- Check for idempotence
@@ -140,9 +141,9 @@ ELSIF in_type IS NOT NULL THEN
in_metadata,
local_authorization_pub,
local_authorization_sig
- );
+ ) RETURNING taler_in_id INTO local_taler_in_id;
-- Notify new incoming talerable transaction registration
- PERFORM pg_notify('taler_in', out_tx_row_id || '');
+ PERFORM pg_notify('taler_in', local_taler_in_id::text);
END IF;
END $$;
COMMENT ON FUNCTION register_tx_in IS 'Register an incoming transaction idempotently';
@@ -477,6 +478,7 @@ CREATE FUNCTION register_prepared_transfers (
LANGUAGE plpgsql AS $$
DECLARE
talerable_tx INT8;
+ local_taler_in_id INT8;
idempotent BOOLEAN;
BEGIN
@@ -522,9 +524,9 @@ IF in_recurrent THEN
INSERT INTO taler_in (tx_in_id, type, metadata, authorization_pub, authorization_sig)
SELECT moved_tx.tx_in_id, in_type, in_account_pub, in_authorization_pub, in_authorization_sig
FROM moved_tx
- RETURNING tx_in_id INTO talerable_tx;
+ RETURNING tx_in_id, taler_in_id INTO talerable_tx, local_taler_in_id;
IF talerable_tx IS NOT NULL THEN
- PERFORM pg_notify('taler_in', talerable_tx::text);
+ PERFORM pg_notify('taler_in', local_taler_in_id::text);
END IF;
ELSE
-- Bounce all pending
diff --git a/adapters/taler-magnet-bank/src/db.rs b/adapters/taler-magnet-bank/src/db.rs
@@ -573,7 +573,7 @@ pub async fn incoming_history(
) -> sqlx::Result<Vec<IncomingBankTransaction>> {
history(
db,
- "tx_in_id",
+ "taler_in_id",
params,
listen,
|| {
@@ -581,7 +581,7 @@ pub async fn incoming_history(
"
SELECT
type,
- tx_in_id,
+ taler_in_id,
amount,
debit_account,
debit_name,
diff --git a/adapters/taler-wise/db/wise-0001.sql b/adapters/taler-wise/db/wise-0001.sql
@@ -48,7 +48,8 @@ CREATE TYPE incoming_type AS ENUM
COMMENT ON TYPE incoming_type IS 'Types of incoming talerable transactions';
CREATE TABLE taler_in (
- tx_in_id INT8 PRIMARY KEY REFERENCES tx_in(tx_in_id) ON DELETE CASCADE,
+ taler_in_id INT8 GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
+ tx_in_id INT8 NOT NULL UNIQUE REFERENCES tx_in(tx_in_id) ON DELETE CASCADE,
type incoming_type NOT NULL,
metadata BYTEA NOT NULL CHECK (LENGTH(metadata)=32),
authorization_pub BYTEA CHECK (LENGTH(authorization_pub)=32),
diff --git a/adapters/taler-wise/db/wise-procedures.sql b/adapters/taler-wise/db/wise-procedures.sql
@@ -64,6 +64,7 @@ LANGUAGE plpgsql AS $$
DECLARE
local_authorization_pub BYTEA;
local_authorization_sig BYTEA;
+local_taler_in_id INT8;
BEGIN
out_pending=false;
-- Check for idempotence
@@ -145,9 +146,9 @@ ELSIF in_type IS NOT NULL AND in_debit_payto IS NOT NULL THEN
in_metadata,
local_authorization_pub,
local_authorization_sig
- );
+ ) RETURNING taler_in_id INTO local_taler_in_id;
-- Notify new incoming talerable transaction registration
- PERFORM pg_notify('taler_in', in_balance_id || ' ' || out_tx_row_id);
+ PERFORM pg_notify('taler_in', in_balance_id || ' ' || local_taler_in_id);
END IF;
END $$;
COMMENT ON FUNCTION register_tx_in IS 'Register an incoming transaction idempotently';
@@ -166,6 +167,7 @@ CREATE FUNCTION register_prepared_transfers (
LANGUAGE plpgsql AS $$
DECLARE
talerable_tx INT8;
+ local_taler_in_id INT8;
local_balance_id INT8;
idempotent BOOLEAN;
BEGIN
@@ -214,14 +216,14 @@ IF in_recurrent THEN
INSERT INTO taler_in (tx_in_id, type, metadata, authorization_pub, authorization_sig)
SELECT tx_in_id, in_type, in_account_pub, in_authorization_pub, in_authorization_sig
FROM moved_tx
- RETURNING tx_in_id
+ RETURNING tx_in_id, taler_in_id
)
- SELECT inserted.tx_in_id, moved_tx.balance_id
- INTO talerable_tx, local_balance_id
+ SELECT inserted.tx_in_id, inserted.taler_in_id, moved_tx.balance_id
+ INTO talerable_tx, local_taler_in_id, local_balance_id
FROM inserted
JOIN moved_tx USING (tx_in_id);
IF talerable_tx IS NOT NULL THEN
- PERFORM pg_notify('taler_in', local_balance_id || ' ' || talerable_tx);
+ PERFORM pg_notify('taler_in', local_balance_id || ' ' || local_taler_in_id);
END IF;
END IF;
diff --git a/adapters/taler-wise/src/db.rs b/adapters/taler-wise/src/db.rs
@@ -183,7 +183,7 @@ pub async fn incoming_history(
) -> sqlx::Result<Vec<IncomingBankTransaction>> {
history(
db,
- "tx_in_id",
+ "taler_in_id",
params,
listen,
|| {
@@ -191,7 +191,7 @@ pub async fn incoming_history(
"
SELECT
type,
- tx_in_id,
+ taler_in_id,
amount,
debit_payto,
debit_name,
diff --git a/common/taler-api/db/taler-api-0001.sql b/common/taler-api/db/taler-api-0001.sql
@@ -44,7 +44,8 @@ CREATE TYPE incoming_type AS ENUM
COMMENT ON TYPE incoming_type IS 'Types of incoming talerable transactions';
CREATE TABLE taler_in (
- tx_in_id INT8 PRIMARY KEY REFERENCES tx_in(tx_in_id) ON DELETE CASCADE,
+ taler_in_id INT8 GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
+ tx_in_id INT8 NOT NULL UNIQUE REFERENCES tx_in(tx_in_id) ON DELETE CASCADE,
type incoming_type NOT NULL,
account_pub BYTEA NOT NULL CHECK (LENGTH(account_pub)=32),
authorization_pub BYTEA CHECK (LENGTH(authorization_pub)=32),
diff --git a/common/taler-api/db/taler-api-procedures.sql b/common/taler-api/db/taler-api-procedures.sql
@@ -132,6 +132,7 @@ DECLARE
local_pending BOOLEAN;
local_authorization_pub BYTEA;
local_authorization_sig BYTEA;
+local_taler_in_id INT8;
BEGIN
local_pending=false;
@@ -190,9 +191,9 @@ ELSE
in_account_pub,
local_authorization_pub,
local_authorization_sig
- );
+ ) RETURNING taler_in_id INTO local_taler_in_id;
-- Notify new incoming transaction
- PERFORM pg_notify('incoming_tx', out_tx_row_id || '');
+ PERFORM pg_notify('incoming_tx', local_taler_in_id::text);
END IF;
END $$;
@@ -211,6 +212,7 @@ CREATE FUNCTION register_prepared_transfers (
LANGUAGE plpgsql AS $$
DECLARE
talerable_tx INT8;
+ local_taler_in_id INT8;
idempotent BOOLEAN;
BEGIN
@@ -256,9 +258,9 @@ IF in_recurrent THEN
INSERT INTO taler_in (tx_in_id, type, account_pub, authorization_pub, authorization_sig)
SELECT moved_tx.tx_in_id, in_type, in_account_pub, in_authorization_pub, in_authorization_sig
FROM moved_tx
- RETURNING tx_in_id INTO talerable_tx;
+ RETURNING tx_in_id, taler_in_id INTO talerable_tx, local_taler_in_id;
IF talerable_tx IS NOT NULL THEN
- PERFORM pg_notify('incoming_tx', talerable_tx::text);
+ PERFORM pg_notify('incoming_tx', local_taler_in_id::text);
END IF;
ELSE
-- Bounce all pending
diff --git a/common/taler-api/src/test/db.rs b/common/taler-api/src/test/db.rs
@@ -273,7 +273,7 @@ pub async fn incoming_history(
) -> sqlx::Result<Vec<IncomingBankTransaction>> {
history(
db,
- "tx_in_id",
+ "taler_in_id",
params,
listen,
|| {
@@ -281,7 +281,7 @@ pub async fn incoming_history(
"
SELECT
type,
- tx_in_id,
+ taler_in_id,
amount,
created_at,
debit_payto,
@@ -297,7 +297,7 @@ pub async fn incoming_history(
|r: PgRow| {
Ok(match r.try_get("type")? {
IncomingType::reserve => IncomingBankTransaction::Reserve {
- row_id: r.try_get_u64("tx_in_id")?,
+ row_id: r.try_get_u64("taler_in_id")?,
date: r.try_get_timestamp("created_at")?.into(),
amount: r.try_get_amount("amount", currency)?,
credit_fee: None,
@@ -307,7 +307,7 @@ pub async fn incoming_history(
authorization_sig: r.try_get("authorization_sig")?,
},
IncomingType::kyc => IncomingBankTransaction::Kyc {
- row_id: r.try_get_u64("tx_in_id")?,
+ row_id: r.try_get_u64("taler_in_id")?,
date: r.try_get_timestamp("created_at")?.into(),
amount: r.try_get_amount("amount", currency)?,
credit_fee: None,
diff --git a/common/taler-test-utils/src/routine.rs b/common/taler-test-utils/src/routine.rs
@@ -146,19 +146,20 @@ pub async fn routine_pagination<T: Page>(
assert_history(&format!("limit=-10&{}", id - 4), 10).await;
}
-pub async fn assert_time<R: Debug>(range: std::ops::Range<u128>, task: impl Future<Output = R>) {
+pub async fn assert_time<R>(range: std::ops::Range<u128>, task: impl Future<Output = R>) -> R {
let start = Instant::now();
let limit = Duration::from_millis(range.end.try_into().expect("timing bound exceeds u64"));
- if tokio::time::timeout(limit, task).await.is_err() {
+ let Ok(res) = tokio::time::timeout(limit, task).await else {
panic!(
"Expected to last {range:?} ms, timed out after {} ms",
start.elapsed().as_millis()
);
- }
+ };
let elapsed = start.elapsed().as_millis();
if !range.contains(&elapsed) {
panic!("Expected to last {range:?} got {elapsed:?}")
}
+ res
}
pub async fn routine_history<T: Page>(
@@ -1025,6 +1026,118 @@ pub async fn in_history_routine(
ignored,
)
.await;
+
+ // Release an older bank payment after a consumer advances its cursor
+ let history = wire_gateway.suffix("/history/incoming");
+ let incoming = async |cursor| history.get(format!("?limit=10&offset={cursor}")).await;
+ let register = async |request: &RegistrationRequest| {
+ prepared_transfer
+ .post("/registration")
+ .json(request)
+ .await
+ .assert_ok_json::<RegistrationResponse>()
+ };
+ let initial_cursor = latest_id::<IncomingHistory>(&history).await;
+ let pair = Ed25519KeyPair::generate().unwrap();
+ let authorization_pub = EddsaPublicKey::try_from(pair.public_key().as_ref()).unwrap();
+ let first_reserve = EddsaPublicKey::rand();
+ let unrelated_reserve = EddsaPublicKey::rand();
+ let amount = Amount::new(currency, 42, 0);
+ let registration = RegistrationRequest {
+ credit_account: credit_account.clone(),
+ r#type: TransferType::reserve,
+ recurrent: true,
+ credit_amount: amount,
+ alg: PublicKeyAlg::EdDSA,
+ account_pub: first_reserve,
+ authorization_pub,
+ authorization_sig: EddsaSignature::ZEROED,
+ }
+ .signed(&pair);
+ register(®istration).await;
+ for _ in 0..2 {
+ wire_gateway
+ .post("/admin/add-mapped")
+ .json(json!({
+ "amount": amount,
+ "authorization_pub": authorization_pub,
+ "debit_account": debit_account,
+ }))
+ .await
+ .assert_ok_json::<TransferResponse>();
+ }
+ wire_gateway
+ .post("/admin/add-incoming")
+ .json(json!({
+ "amount": amount,
+ "reserve_pub": unrelated_reserve,
+ "debit_account": debit_account,
+ }))
+ .await
+ .assert_ok_json::<TransferResponse>();
+ let before = incoming(initial_cursor)
+ .await
+ .assert_ok_json::<IncomingHistory>();
+ assert_eq!(before.incoming_transactions.len(), 2);
+ for (tx, expected) in before
+ .incoming_transactions
+ .iter()
+ .zip([first_reserve, unrelated_reserve])
+ {
+ assert!(
+ matches!(tx, IncomingBankTransaction::Reserve { reserve_pub, .. }
+ if *reserve_pub == expected)
+ );
+ }
+ let cursor = *before.ids().last().unwrap();
+ incoming(cursor).await.assert_no_content();
+ let release = RegistrationRequest {
+ account_pub: EddsaPublicKey::rand(),
+ ..registration
+ }
+ .signed(&pair);
+ let (response, ()) = tokio::join!(
+ assert_time(0..5000, async {
+ history
+ .get(format!("?limit=10&offset={cursor}&timeout_ms=10000"))
+ .await
+ }),
+ async {
+ sleep(Duration::from_millis(100)).await;
+ register(&release).await;
+ }
+ );
+ let resumed = response.assert_ok_json::<IncomingHistory>();
+ let [
+ IncomingBankTransaction::Reserve {
+ row_id,
+ reserve_pub,
+ amount: credited,
+ authorization_pub: Some(auth_pub),
+ authorization_sig: Some(auth_sig),
+ ..
+ },
+ ] = resumed.incoming_transactions.as_slice()
+ else {
+ panic!("expected exactly one released reserve payment");
+ };
+ assert_eq!(*reserve_pub, release.account_pub);
+ assert_eq!(*credited, amount);
+ assert_eq!(*auth_pub, authorization_pub);
+ assert!(release.verify(auth_pub, auth_sig));
+ let next_cursor = *row_id as i64;
+ assert!(
+ next_cursor > cursor,
+ "released payment must advance the history cursor"
+ );
+ assert_eq!(
+ incoming(cursor)
+ .await
+ .assert_ok_json::<IncomingHistory>()
+ .incoming_transactions,
+ resumed.incoming_transactions
+ );
+ incoming(next_cursor).await.assert_no_content();
}
/// Test standard behavior of the admin add incoming endpoints