commit 45a0bbb0118051f70b95f6b8d49fd934c0eda2d1
parent 053bedf6a2460b1a75ecc680a85f6229c128991f
Author: Florian Dold <dold@taler.net>
Date: Thu, 24 Sep 2026 12:46:53 +0200
wallet-core: coalesce duplicate balance notifications
Merge balance invalidations for the same transaction when publishing a
database commit, retaining public visibility and all state events.
Notify for changed renewal reports and restored checking results instead
of cache generation or scheduling changes alone.
Diffstat:
7 files changed, 313 insertions(+), 31 deletions(-)
diff --git a/packages/taler-wallet-core/src/db/indexeddb/handle.ts b/packages/taler-wallet-core/src/db/indexeddb/handle.ts
@@ -58,6 +58,7 @@ import {
import { IdbWalletTransaction } from "./transaction.js";
import { WalletDbTransaction } from "../transaction.js";
import { DbAccess, DbAccessImpl } from "../query.js";
+import { coalesceWalletNotifications } from "../shared.js";
const logger = new Logger("db/indexeddb/handle.ts");
@@ -236,7 +237,7 @@ export class IdbWalletDbHandle implements WalletDbHandle {
WalletIndexedDbStoresV1,
CancellationToken.CONTINUE,
(notifs: WalletNotification[]) => {
- for (const n of notifs) {
+ for (const n of coalesceWalletNotifications(notifs)) {
notificationSink(n);
}
},
diff --git a/packages/taler-wallet-core/src/db/shared.ts b/packages/taler-wallet-core/src/db/shared.ts
@@ -30,6 +30,9 @@ import {
ScopeType,
encodeCrock,
getRandomBytes,
+ BalanceChangeNotification,
+ NotificationType,
+ WalletNotification,
} from "@gnu-taler/taler-util";
import { WalletDbTransaction } from "./transaction.js";
import {
@@ -41,6 +44,32 @@ import {
} from "./records.js";
import { auditorProvidesVerifiedTrust } from "../auditorTrust.js";
+/** Merge balance invalidations from one commit without dropping state events. */
+export function coalesceWalletNotifications(
+ notifications: readonly WalletNotification[],
+): WalletNotification[] {
+ const result: WalletNotification[] = [];
+ const balances = new Map<string, BalanceChangeNotification>();
+ for (const notification of notifications) {
+ if (notification.type !== NotificationType.BalanceChange) {
+ result.push(notification);
+ continue;
+ }
+ const previous = balances.get(notification.hintTransactionId);
+ if (previous) {
+ // An unspecified flag is public too. An internal event must never
+ // hide a public invalidation emitted by the same transaction.
+ previous.isInternal =
+ previous.isInternal === true && notification.isInternal === true;
+ } else {
+ const copy = { ...notification };
+ balances.set(copy.hintTransactionId, copy);
+ result.push(copy);
+ }
+ }
+ return result;
+}
+
/**
* Does the exchange fall within the given scope?
*/
diff --git a/packages/taler-wallet-core/src/db/sqlite/database.ts b/packages/taler-wallet-core/src/db/sqlite/database.ts
@@ -28,6 +28,7 @@ import {
WalletNotification,
} from "@gnu-taler/taler-util";
import { WalletDbTransaction } from "../transaction.js";
+import { coalesceWalletNotifications } from "../shared.js";
import {
DATA_TABLES_CONDITION,
SchemaMigration,
@@ -470,7 +471,7 @@ async function runNativeSqliteWalletTxLocked<T>(
for (const h of tx.afterCommitHandlers) {
h();
}
- for (const notif of tx.pendingNotifications) {
+ for (const notif of coalesceWalletNotifications(tx.pendingNotifications)) {
notifyFn(notif);
}
return res;
diff --git a/packages/taler-wallet-core/src/db/testing/conformance.test.ts b/packages/taler-wallet-core/src/db/testing/conformance.test.ts
@@ -22,13 +22,19 @@
* themselves do not change.
*/
-import { CancellationToken } from "@gnu-taler/taler-util";
+import {
+ CancellationToken,
+ NotificationType,
+ WalletNotification,
+ TransactionMajorState,
+} from "@gnu-taler/taler-util";
import assert from "node:assert";
import { test } from "node:test";
import { conformanceCases } from "./conformance-cases.js";
import { ConformanceAsserts } from "./conformance.js";
import { runnerFactories } from "./runners.js";
import { ConfigRecordKey } from "../records.js";
+import { WalletDbTransaction } from "../transaction.js";
const asserts: ConformanceAsserts = {
equal: (a, e, m) => assert.strictEqual(a, e, m),
@@ -52,6 +58,78 @@ function isNotImplemented(e: unknown): boolean {
}
for (const makeRunner of runnerFactories) {
+ test(`dbtx ${makeRunner.name}: coalesces balance notifications only within a committed transaction`, async () => {
+ const runner = await makeRunner();
+ const notifications: WalletNotification[] = [];
+ runner.setNotificationSink((n) => notifications.push(n));
+ const state: WalletNotification = {
+ type: NotificationType.TransactionStateTransition,
+ transactionId: "txn:refresh:one",
+ oldTxState: { major: TransactionMajorState.Pending },
+ newTxState: { major: TransactionMajorState.Done },
+ newStId: 83886080,
+ causeHint: "refresh-done",
+ };
+ const emit = (tx: WalletDbTransaction) => {
+ tx.notify({
+ type: NotificationType.BalanceChange,
+ hintTransactionId: "txn:refresh:one",
+ isInternal: true,
+ });
+ tx.notify(state);
+ tx.notify({
+ type: NotificationType.BalanceChange,
+ hintTransactionId: "txn:refresh:one",
+ });
+ tx.notify({
+ type: NotificationType.BalanceChange,
+ hintTransactionId: "txn:refresh:one",
+ isInternal: true,
+ });
+ tx.notify({
+ type: NotificationType.BalanceChange,
+ hintTransactionId: "txn:refresh:two",
+ isInternal: true,
+ });
+ tx.notify({
+ type: NotificationType.BalanceChange,
+ hintTransactionId: "txn:refresh:two",
+ isInternal: true,
+ });
+ };
+ try {
+ await assert.rejects(
+ runner.runReadWriteTx(async (tx) => {
+ emit(tx);
+ throw Error("rollback");
+ }),
+ /rollback/,
+ );
+ assert.deepStrictEqual(notifications, []);
+ await runner.runReadWriteTx(async (tx) => {
+ emit(tx);
+ assert.deepStrictEqual(notifications, []);
+ });
+ assert.deepStrictEqual(notifications, [
+ {
+ type: NotificationType.BalanceChange,
+ hintTransactionId: "txn:refresh:one",
+ isInternal: false,
+ },
+ state,
+ {
+ type: NotificationType.BalanceChange,
+ hintTransactionId: "txn:refresh:two",
+ isInternal: true,
+ },
+ ]);
+ await runner.runReadWriteTx(async (tx) => emit(tx));
+ assert.strictEqual(notifications.length, 6);
+ } finally {
+ await runner.close();
+ }
+ });
+
for (const c of conformanceCases) {
test(`dbtx conformance: ${c.name}`, async (t) => {
const runner = await makeRunner();
diff --git a/packages/taler-wallet-core/src/refresh.test.ts b/packages/taler-wallet-core/src/refresh.test.ts
@@ -23,6 +23,9 @@ import {
TalerError,
TalerErrorCode,
TransactionAction,
+ NotificationType,
+ RefreshReason,
+ WalletNotification,
} from "@gnu-taler/taler-util";
import assert from "node:assert";
import { test } from "node:test";
@@ -114,6 +117,50 @@ test("refresh costs exclude outputs with a different age mask", () => {
});
for (const makeRunner of runnerFactories) {
+ test(`${makeRunner.name}: refresh completion emits one balance invalidation`, async () => {
+ const runner = await makeRunner();
+ const notifications: WalletNotification[] = [];
+ runner.setNotificationSink((n) => notifications.push(n));
+ try {
+ const group: WalletRefreshGroup = {
+ refreshGroupId: "notification-test",
+ currency: "TESTKUDOS",
+ reason: RefreshReason.Manual,
+ operationStatus: RefreshOperationStatus.Pending,
+ oldCoinPubs: [],
+ inputPerCoin: [],
+ expectedOutputPerCoin: [],
+ statusPerCoin: [],
+ refundRequests: {},
+ timestampCreated: 1 as WalletRefreshGroup["timestampCreated"],
+ timestampFinished: undefined,
+ };
+ await runner.runReadWriteTx((tx) => tx.upsertRefreshGroup(group));
+ const ctx = new RefreshTransactionContext(
+ {} as WalletExecutionContext,
+ group.refreshGroupId,
+ );
+ await runner.runReadWriteTx(async (tx) => {
+ const [record, handle] = await ctx.getRecordHandle(tx);
+ record!.operationStatus = RefreshOperationStatus.Finished;
+ await handle.update(record, "refresh-done");
+ });
+ const balances = notifications.filter(
+ (n) => n.type === NotificationType.BalanceChange,
+ );
+ assert.strictEqual(balances.length, 1);
+ assert.strictEqual(balances[0].isInternal, false);
+ assert.strictEqual(
+ notifications.filter(
+ (n) => n.type === NotificationType.TransactionStateTransition,
+ ).length,
+ 1,
+ );
+ } finally {
+ await runner.close();
+ }
+ });
+
test(`${makeRunner.name}: refresh finds a matching age mask within a denomination family`, async () => {
const runner = await makeRunner();
try {
diff --git a/packages/taler-wallet-core/src/refreshBalance.test.ts b/packages/taler-wallet-core/src/refreshBalance.test.ts
@@ -24,6 +24,8 @@ import {
stringifyScopeInfo,
encodeCrock,
DenomKeyType,
+ NotificationType,
+ WalletNotification,
} from "@gnu-taler/taler-util";
import {
ConfigRecordKey,
@@ -101,6 +103,103 @@ function fixture() {
}
for (const makeRunner of runnerFactories) {
+ test(`${makeRunner.name}: renewal publication notifies only for visible changes`, async () => {
+ const runner = await makeRunner();
+ const notifications: WalletNotification[] = [];
+ runner.setNotificationSink((n) => notifications.push(n));
+ try {
+ const { saved, response } = fixture();
+ const scope = stringifyScopeInfo(response.balances[0].scopeInfo);
+ // A refreshed generation with no balances has no visible effect.
+ await runner.runReadWriteTx((tx) =>
+ publishRefreshBalance(tx, { ...saved, scopes: {} }),
+ );
+ notifications.length = 0;
+ await runner.runReadWriteTx(async (tx) => {
+ await tx.upsertConfig({
+ key: ConfigRecordKey.RefreshBalanceGeneration,
+ value: "next",
+ });
+ await publishRefreshBalance(tx, {
+ ...saved,
+ generation: "next",
+ scopes: {},
+ });
+ });
+ assert.equal(notifications.length, 0);
+
+ saved.generation = "next";
+ await runner.runReadWriteTx((tx) => publishRefreshBalance(tx, saved));
+ assert.equal(notifications.length, 1);
+ notifications.length = 0;
+ // Advancing the cache's scheduling timestamps is not a balance change.
+ await runner.runReadWriteTx((tx) =>
+ publishRefreshBalance(tx, {
+ ...saved,
+ nextCheck: saved.nextCheck + 60000,
+ }),
+ );
+ assert.equal(notifications.length, 0);
+
+ // Invalidating and restoring the same cost does change the visible
+ // "checking" status, so the frontend still needs an invalidation.
+ await runner.runReadWriteTx((tx) =>
+ tx.upsertConfig({
+ key: ConfigRecordKey.RefreshBalanceGeneration,
+ value: "newer",
+ }),
+ );
+ await runner.runReadWriteTx((tx) => readRefreshBalanceInfo(tx, response));
+ assert.equal(
+ response.balances[0].refreshInfo?.annualCostBound.status,
+ "unavailable",
+ );
+ saved.generation = "newer";
+ await runner.runReadWriteTx((tx) => publishRefreshBalance(tx, saved));
+ assert.equal(notifications.length, 1);
+ notifications.length = 0;
+ await runner.runReadWriteTx((tx) => readRefreshBalanceInfo(tx, response));
+ assert.equal(
+ response.balances[0].refreshInfo?.annualCostBound.status,
+ "available",
+ );
+
+ // Warnings can change without any monetary amount changing.
+ saved.scopes[scope].info.risks = [];
+ await runner.runReadWriteTx((tx) => publishRefreshBalance(tx, saved));
+ assert.equal(notifications.length, 1);
+ assert.equal(notifications[0].type, NotificationType.BalanceChange);
+ notifications.length = 0;
+
+ // A calculation that still reports "checking" must not notify merely
+ // because its inputs were invalidated again.
+ saved.scopes[scope].info.annualCostBound = {
+ status: "unavailable",
+ horizonDays: 365,
+ projection: "stable-current-offerings",
+ reasons: ["checking"],
+ };
+ await runner.runReadWriteTx((tx) => publishRefreshBalance(tx, saved));
+ notifications.length = 0;
+ saved.generation = "again";
+ await runner.runReadWriteTx(async (tx) => {
+ await tx.upsertConfig({
+ key: ConfigRecordKey.RefreshBalanceGeneration,
+ value: saved.generation,
+ });
+ await publishRefreshBalance(tx, saved);
+ });
+ assert.equal(notifications.length, 0);
+ // Removing the last reported scope also invalidates an existing view.
+ await runner.runReadWriteTx((tx) =>
+ publishRefreshBalance(tx, { ...saved, scopes: {} }),
+ );
+ assert.equal(notifications.length, 1);
+ } finally {
+ await runner.close();
+ }
+ });
+
test(`${makeRunner.name}: renewal cache survives reopening; invalidation is atomic and rejects stale publication`, async () => {
const dir = mkdtempSync(join(tmpdir(), "refresh-balance-"));
const file = join(dir, "wallet.sqlite3");
diff --git a/packages/taler-wallet-core/src/refreshBalance.ts b/packages/taler-wallet-core/src/refreshBalance.ts
@@ -15,8 +15,10 @@
import {
AbsoluteTime,
BalancesResponse,
+ canonicalJson,
NotificationType,
stringifyScopeInfo,
+ WalletRefreshInfo,
} from "@gnu-taler/taler-util";
import {
ConfigRecordKey,
@@ -45,6 +47,38 @@ import { TaskIdentifiers, TaskRunResult } from "./common.js";
export const REFRESH_BALANCE_NOTIFICATION = "refresh-balance";
+/** Project the cache exactly as a balance reader would see it. */
+function visibleRefreshInfo(
+ saved: WalletRefreshBalance | undefined,
+ generation: string,
+ now: number,
+ scopeKey: string,
+ available: string,
+): WalletRefreshInfo {
+ const current =
+ saved?.version === 1 &&
+ saved.generation === generation &&
+ saved.computedAt <= now &&
+ saved.nextCheck > now;
+ const cached = saved?.version === 1 ? saved.scopes[scopeKey] : undefined;
+ return {
+ // Retain known warnings during recomputation. Expired coins use the
+ // ordinary expiry indication, even if the maintenance task is delayed.
+ risks:
+ cached?.info.risks.filter(
+ (r) =>
+ AbsoluteTime.toStampMs(
+ AbsoluteTime.fromProtocolTimestamp(r.earliestDepositExpiration),
+ ) > now,
+ ) ?? [],
+ recoveries: cached?.info.recoveries ?? [],
+ annualCostBound:
+ current && cached?.available === available
+ ? cached.info.annualCostBound
+ : unavailableAnnualBound("checking"),
+ };
+}
+
/** Balance reads do no denomination selection or history scans. */
export async function readRefreshBalanceInfo(
tx: WalletDbTransaction,
@@ -54,32 +88,14 @@ export async function readRefreshBalanceInfo(
const generation =
(await tx.getConfig(ConfigRecordKey.RefreshBalanceGeneration))?.value ?? "";
const now = AbsoluteTime.toStampMs(AbsoluteTime.now());
- const current =
- saved?.version === 1 &&
- saved.generation === generation &&
- saved.computedAt <= now &&
- saved.nextCheck > now;
for (const balance of response.balances) {
- const cached =
- saved?.version === 1
- ? saved.scopes[stringifyScopeInfo(balance.scopeInfo)]
- : undefined;
- balance.refreshInfo = {
- // Retain known warnings during recomputation. Expired coins use the
- // ordinary expiry indication, even if the maintenance task is delayed.
- risks:
- cached?.info.risks.filter(
- (r) =>
- AbsoluteTime.toStampMs(
- AbsoluteTime.fromProtocolTimestamp(r.earliestDepositExpiration),
- ) > now,
- ) ?? [],
- recoveries: cached?.info.recoveries ?? [],
- annualCostBound:
- current && cached?.available === balance.available
- ? cached.info.annualCostBound
- : unavailableAnnualBound("checking"),
- };
+ balance.refreshInfo = visibleRefreshInfo(
+ saved,
+ generation,
+ now,
+ stringifyScopeInfo(balance.scopeInfo),
+ balance.available,
+ );
}
}
@@ -247,10 +263,21 @@ export async function publishRefreshBalance(
if (generation !== value.generation) return false;
const previous = (await tx.getConfig(ConfigRecordKey.RefreshBalance))?.value;
await tx.upsertConfig({ key: ConfigRecordKey.RefreshBalance, value });
+ const now = AbsoluteTime.toStampMs(AbsoluteTime.now());
+ // Generation and scheduling changes alone are not balance changes. Still
+ // notify when publication replaces a reader's "checking" status, even if
+ // the newly computed cost equals the previously cached cost.
if (
- JSON.stringify(previous?.scopes) !== JSON.stringify(value.scopes) ||
- previous?.generation !== value.generation ||
- (previous?.nextCheck ?? 0) <= value.computedAt
+ canonicalJson(previous?.scopes ?? {}) !== canonicalJson(value.scopes) ||
+ Object.entries(value.scopes).some(
+ ([key, scope]) =>
+ canonicalJson(
+ visibleRefreshInfo(previous, generation, now, key, scope.available),
+ ) !==
+ canonicalJson(
+ visibleRefreshInfo(value, generation, now, key, scope.available),
+ ),
+ )
) {
tx.notify({
type: NotificationType.BalanceChange,