commit 832f4cc8dbdc9ab000ab281f74b39cce42586e86
parent f69521e9ef7422acbc8f8e175cc64c28a32f9f0b
Author: Florian Dold <dold@taler.net>
Date: Mon, 31 Aug 2026 20:43:30 +0200
wallet-core: bound shepherd task concurrency
Diffstat:
3 files changed, 456 insertions(+), 17 deletions(-)
diff --git a/packages/taler-wallet-core/src/shepherd.ts b/packages/taler-wallet-core/src/shepherd.ts
@@ -113,6 +113,7 @@ import {
processWithdrawalGroup,
} from "./withdraw.js";
import { WalletDbTransaction } from "./db/transaction.js";
+import { TaskRunLimiter } from "./task-run-limiter.js";
const logger = new Logger("shepherd.ts");
@@ -212,6 +213,8 @@ export class TaskSchedulerImpl implements TaskScheduler {
private throttler = new TaskThrottler();
+ private runLimiter = new TaskRunLimiter();
+
isRunning: boolean = false;
private shepCounter = 1;
@@ -539,25 +542,54 @@ export class TaskSchedulerImpl implements TaskScheduler {
Duration.fromSpec({ seconds: 60 }),
);
}
- const wex = getWalletExecutionContextForTask(
- this.ws,
- taskId,
- info.cts.token,
- );
- const startTime = AbsoluteTime.now();
- // Reset before running the handler so that any wakeup arriving while the
- // handler (and its result storage) runs is captured and honored below.
- info.rerunRequested = false;
- logger.trace(`Shepherd for ${taskId} will call handler`);
- let res: TaskRunResult;
+ let taskExecution: {
+ endTime: AbsoluteTime;
+ res: TaskRunResult;
+ startTime: AbsoluteTime;
+ wex: WalletExecutionContext;
+ };
try {
- res = await callOperationHandlerForTaskId(wex, taskId);
+ taskExecution = await this.runLimiter.run(
+ taskId,
+ parseTaskIdentifier(taskId).tag === PendingTaskType.Refresh,
+ info.cts.token,
+ async () => {
+ const wex = getWalletExecutionContextForTask(
+ this.ws,
+ taskId,
+ info.cts.token,
+ );
+ const startTime = AbsoluteTime.now();
+ // Reset immediately before running the handler so that a wakeup
+ // while queued is satisfied by this run, while a wakeup during
+ // the handler is captured and honored below.
+ info.rerunRequested = false;
+ logger.trace(`Shepherd for ${taskId} will call handler`);
+ let res: TaskRunResult;
+ try {
+ res = await callOperationHandlerForTaskId(wex, taskId);
+ } catch (e) {
+ res = {
+ type: TaskRunResultType.Error,
+ errorDetail: getErrorDetailFromException(e),
+ };
+ }
+ return {
+ endTime: AbsoluteTime.now(),
+ res,
+ startTime,
+ wex,
+ };
+ },
+ );
} catch (e) {
- res = {
- type: TaskRunResultType.Error,
- errorDetail: getErrorDetailFromException(e),
- };
+ if (e instanceof CancellationToken.CancellationError) {
+ logger.trace(`task ${taskId} cancelled while waiting to run`);
+ return;
+ }
+ throw e;
}
+ const { endTime, res, startTime, wex } = taskExecution;
if (this.ws.stopped) {
logger.trace("wallet stopped, not processing result");
return;
@@ -566,7 +598,6 @@ export class TaskSchedulerImpl implements TaskScheduler {
return;
}
logger.trace(`Shepherd for ${taskId} got result ${res.type}`);
- const endTime = AbsoluteTime.now();
const taskDuration = AbsoluteTime.difference(endTime, startTime);
if (taskDuration.d_ms === "forever") {
throw Error("assertion failed");
diff --git a/packages/taler-wallet-core/src/task-run-limiter.test.ts b/packages/taler-wallet-core/src/task-run-limiter.test.ts
@@ -0,0 +1,254 @@
+/*
+ This file is part of GNU Taler
+ (C) 2026 Taler Systems S.A.
+
+ GNU 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.
+ */
+
+import assert from "node:assert/strict";
+import test from "node:test";
+import { CancellationToken } from "@gnu-taler/taler-util";
+import { TaskRunLimiter } from "./task-run-limiter.js";
+
+function deferred(): { promise: Promise<void>; resolve: () => void } {
+ let resolve!: () => void;
+ return { promise: new Promise<void>((r) => (resolve = r)), resolve };
+}
+
+async function nextTurn(): Promise<void> {
+ await Promise.resolve();
+ await Promise.resolve();
+}
+
+test("task run limiter bounds mixed global concurrency", async () => {
+ const limiter = new TaskRunLimiter({
+ maxConcurrentTasks: 2,
+ maxConcurrentRefreshTasks: 1,
+ });
+ const release = deferred();
+ let running = 0;
+ let maximum = 0;
+ const run = (id: string, refresh: boolean) =>
+ limiter.run(id, refresh, CancellationToken.CONTINUE, async () => {
+ running++;
+ maximum = Math.max(maximum, running);
+ await release.promise;
+ running--;
+ });
+
+ const tasks = [run("ordinary-1", false), run("refresh-1", true)];
+ const queued = run("ordinary-2", false);
+ await nextTurn();
+ assert.equal(running, 2);
+ assert.equal(maximum, 2);
+
+ release.resolve();
+ await Promise.all([...tasks, queued]);
+ assert.equal(maximum, 2);
+});
+
+test("refresh backlog leaves global capacity for ordinary tasks", async () => {
+ const limiter = new TaskRunLimiter({
+ maxConcurrentTasks: 3,
+ maxConcurrentRefreshTasks: 1,
+ });
+ const releaseRefresh = deferred();
+ const releaseOrdinary = deferred();
+ const started: string[] = [];
+ const firstRefresh = limiter.run(
+ "refresh-1",
+ true,
+ CancellationToken.CONTINUE,
+ async () => {
+ started.push("refresh-1");
+ await releaseRefresh.promise;
+ },
+ );
+ const secondRefresh = limiter.run(
+ "refresh-2",
+ true,
+ CancellationToken.CONTINUE,
+ async () => {
+ started.push("refresh-2");
+ },
+ );
+ const ordinary = limiter.run(
+ "ordinary",
+ false,
+ CancellationToken.CONTINUE,
+ async () => {
+ started.push("ordinary");
+ await releaseOrdinary.promise;
+ },
+ );
+
+ await nextTurn();
+ assert.deepEqual(new Set(started), new Set(["refresh-1", "ordinary"]));
+ releaseOrdinary.resolve();
+ await ordinary;
+ releaseRefresh.resolve();
+ await Promise.all([firstRefresh, secondRefresh]);
+ assert.equal(started.at(-1), "refresh-2");
+});
+
+test("task run limiter admits global waiters in FIFO order", async () => {
+ const limiter = new TaskRunLimiter({
+ maxConcurrentTasks: 1,
+ maxConcurrentRefreshTasks: 1,
+ });
+ const firstRelease = deferred();
+ const order: string[] = [];
+ const first = limiter.run(
+ "first",
+ false,
+ CancellationToken.CONTINUE,
+ async () => {
+ order.push("first");
+ await firstRelease.promise;
+ },
+ );
+ const second = limiter.run(
+ "second",
+ false,
+ CancellationToken.CONTINUE,
+ async () => {
+ order.push("second");
+ },
+ );
+ const third = limiter.run(
+ "third",
+ false,
+ CancellationToken.CONTINUE,
+ async () => {
+ order.push("third");
+ },
+ );
+
+ await nextTurn();
+ assert.deepEqual(order, ["first"]);
+ firstRelease.resolve();
+ await Promise.all([first, second, third]);
+ assert.deepEqual(order, ["first", "second", "third"]);
+});
+
+test("cancelling queued work does not run it or leak capacity", async () => {
+ const limiter = new TaskRunLimiter({
+ maxConcurrentTasks: 1,
+ maxConcurrentRefreshTasks: 1,
+ });
+ const releaseBlocker = deferred();
+ const blocker = limiter.run(
+ "blocker",
+ false,
+ CancellationToken.CONTINUE,
+ async () => releaseBlocker.promise,
+ );
+ const cts = CancellationToken.create();
+ let cancelledRan = false;
+ const cancelled = limiter.run("cancelled", false, cts.token, async () => {
+ cancelledRan = true;
+ });
+ await nextTurn();
+ cts.cancel("test cancellation");
+ await assert.rejects(cancelled, CancellationToken.CancellationError);
+ assert.equal(cancelledRan, false);
+
+ releaseBlocker.resolve();
+ await blocker;
+ let successorRan = false;
+ await limiter.run(
+ "successor",
+ false,
+ CancellationToken.CONTINUE,
+ async () => {
+ successorRan = true;
+ },
+ );
+ assert.equal(successorRan, true);
+});
+
+test("cancelling a queued refresh releases both pools", async () => {
+ const limiter = new TaskRunLimiter({
+ maxConcurrentTasks: 1,
+ maxConcurrentRefreshTasks: 1,
+ });
+ const releaseBlocker = deferred();
+ const blocker = limiter.run(
+ "global-blocker",
+ false,
+ CancellationToken.CONTINUE,
+ async () => releaseBlocker.promise,
+ );
+ const cts = CancellationToken.create();
+ const cancelled = limiter.run(
+ "cancelled-refresh",
+ true,
+ cts.token,
+ async () => assert.fail("cancelled refresh was run"),
+ );
+ await nextTurn();
+ cts.cancel("test cancellation");
+ await assert.rejects(cancelled, CancellationToken.CancellationError);
+
+ releaseBlocker.resolve();
+ await blocker;
+ let successorRan = false;
+ await limiter.run(
+ "successor-refresh",
+ true,
+ CancellationToken.CONTINUE,
+ async () => {
+ successorRan = true;
+ },
+ );
+ assert.equal(successorRan, true);
+});
+
+test("failure releases global and refresh permits", async () => {
+ const limiter = new TaskRunLimiter({
+ maxConcurrentTasks: 1,
+ maxConcurrentRefreshTasks: 1,
+ });
+ await assert.rejects(
+ limiter.run(
+ "failing-refresh",
+ true,
+ CancellationToken.CONTINUE,
+ async () => {
+ throw Error("handler failed");
+ },
+ ),
+ /handler failed/,
+ );
+ let successorRan = false;
+ await limiter.run(
+ "successor-refresh",
+ true,
+ CancellationToken.CONTINUE,
+ async () => {
+ successorRan = true;
+ },
+ );
+ assert.equal(successorRan, true);
+});
+
+test("task run limits must be positive and nested", () => {
+ assert.throws(
+ () =>
+ new TaskRunLimiter({
+ maxConcurrentTasks: 0,
+ maxConcurrentRefreshTasks: 0,
+ }),
+ /positive integer/,
+ );
+ assert.throws(
+ () =>
+ new TaskRunLimiter({
+ maxConcurrentTasks: 1,
+ maxConcurrentRefreshTasks: 2,
+ }),
+ /must not exceed/,
+ );
+});
diff --git a/packages/taler-wallet-core/src/task-run-limiter.ts b/packages/taler-wallet-core/src/task-run-limiter.ts
@@ -0,0 +1,154 @@
+/*
+ This file is part of GNU Taler
+ (C) 2026 Taler Systems S.A.
+
+ GNU 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.
+
+ GNU 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
+ GNU Taler; see the file COPYING. If not, see <http://www.gnu.org/licenses/>
+ */
+
+import { CancellationToken, Logger, openPromise } from "@gnu-taler/taler-util";
+
+const logger = new Logger("task-run-limiter.ts");
+
+export interface TaskRunLimits {
+ maxConcurrentTasks: number;
+ maxConcurrentRefreshTasks: number;
+}
+
+const defaultTaskRunLimits: TaskRunLimits = {
+ maxConcurrentTasks: 32,
+ maxConcurrentRefreshTasks: 4,
+};
+
+interface Waiter {
+ resolve: () => void;
+ taskId: string;
+}
+
+/** A small cancellation-aware FIFO semaphore. */
+class FifoSemaphore {
+ private available: number;
+ private waiters: Waiter[] = [];
+
+ constructor(
+ private readonly name: string,
+ private readonly capacity: number,
+ ) {
+ if (!Number.isSafeInteger(capacity) || capacity <= 0) {
+ throw Error(`${name} concurrency limit must be a positive integer`);
+ }
+ this.available = capacity;
+ }
+
+ private releasePermit(): void {
+ const next = this.waiters.shift();
+ if (next) {
+ logger.trace(
+ `admitting queued task ${next.taskId} to ${this.name} pool (${this.waiters.length} waiting)`,
+ );
+ next.resolve();
+ return;
+ }
+ this.available++;
+ if (this.available > this.capacity) {
+ throw Error(`${this.name} concurrency permit released more than once`);
+ }
+ }
+
+ async acquire(
+ taskId: string,
+ cancellationToken: CancellationToken,
+ ): Promise<() => void> {
+ cancellationToken.throwIfCancelled();
+ if (this.available > 0) {
+ this.available--;
+ } else {
+ const promCap = openPromise<void>();
+ const waiter: Waiter = { resolve: promCap.resolve, taskId };
+ this.waiters.push(waiter);
+ logger.trace(
+ `queueing task ${taskId} for ${this.name} pool (${this.waiters.length} waiting)`,
+ );
+ try {
+ await cancellationToken.racePromise(promCap.promise);
+ } catch (e) {
+ const idx = this.waiters.indexOf(waiter);
+ if (idx >= 0) {
+ this.waiters.splice(idx, 1);
+ } else {
+ // The waiter was already handed a permit when cancellation won the
+ // promise race. Pass that permit on instead of leaking it.
+ this.releasePermit();
+ }
+ throw e;
+ }
+ }
+
+ let released = false;
+ return () => {
+ if (released) {
+ return;
+ }
+ released = true;
+ this.releasePermit();
+ };
+ }
+
+ async run<T>(
+ taskId: string,
+ cancellationToken: CancellationToken,
+ f: () => Promise<T>,
+ ): Promise<T> {
+ const release = await this.acquire(taskId, cancellationToken);
+ try {
+ // Cancellation can happen after the permit was handed to us but before
+ // the awaiting continuation runs.
+ cancellationToken.throwIfCancelled();
+ return await f();
+ } finally {
+ release();
+ }
+ }
+}
+
+/**
+ * Bounds active shepherd handlers. Refresh tasks first enter their smaller
+ * pool so that a large refresh backlog cannot fill the global queue.
+ */
+export class TaskRunLimiter {
+ private readonly allTasks: FifoSemaphore;
+ private readonly refreshTasks: FifoSemaphore;
+
+ constructor(limits: TaskRunLimits = defaultTaskRunLimits) {
+ if (limits.maxConcurrentRefreshTasks > limits.maxConcurrentTasks) {
+ throw Error("refresh concurrency limit must not exceed global limit");
+ }
+ this.allTasks = new FifoSemaphore("global", limits.maxConcurrentTasks);
+ this.refreshTasks = new FifoSemaphore(
+ "refresh",
+ limits.maxConcurrentRefreshTasks,
+ );
+ }
+
+ async run<T>(
+ taskId: string,
+ isRefreshTask: boolean,
+ cancellationToken: CancellationToken,
+ f: () => Promise<T>,
+ ): Promise<T> {
+ if (isRefreshTask) {
+ return await this.refreshTasks.run(taskId, cancellationToken, async () =>
+ this.allTasks.run(taskId, cancellationToken, f),
+ );
+ }
+ return await this.allTasks.run(taskId, cancellationToken, f);
+ }
+}