mirror of
https://github.com/Routstr/routstrd.git
synced 2026-10-05 12:28:23 +00:00
fix(wallet): settle paid mint quotes live instead of on restart (#94)
This commit is contained in:
@@ -1,4 +1,4 @@
|
||||
import { afterEach, describe, expect, it, mock } from "bun:test";
|
||||
import { afterEach, beforeEach, describe, expect, it, mock, spyOn } from "bun:test";
|
||||
import { gunzipSync } from "bun";
|
||||
import {
|
||||
existsSync,
|
||||
@@ -16,9 +16,14 @@ import {
|
||||
createCocoClient,
|
||||
isZombieProcess,
|
||||
settleExpiredMintQuotes,
|
||||
settlePendingMintQuotes,
|
||||
stopLegacyCocod,
|
||||
type ExpiredMintQuoteSource,
|
||||
type PendingMintQuoteSource,
|
||||
type PendingMintSweepState,
|
||||
} from "./coco-client";
|
||||
import { OperationInProgressError } from "@cashu/coco-core";
|
||||
import { logger } from "../../utils/logger";
|
||||
|
||||
type GuardOptions = NonNullable<
|
||||
Parameters<typeof assertLegacyCocodNotRunning>[0]
|
||||
@@ -667,3 +672,228 @@ describe("settleExpiredMintQuotes", () => {
|
||||
expect(failPendingOperation).not.toHaveBeenCalled();
|
||||
});
|
||||
});
|
||||
|
||||
describe("settlePendingMintQuotes", () => {
|
||||
const NOW_MS = 2_000_000_000_000;
|
||||
const UNEXPIRED_S = NOW_MS / 1000 + 600;
|
||||
const EXPIRED_S = NOW_MS / 1000 - 1;
|
||||
const UNPAID = { state: "pending", lastObservedRemoteState: "UNPAID" };
|
||||
const MINTED = { state: "finalized", lastObservedRemoteState: "ISSUED" };
|
||||
let logged: ReturnType<typeof spyOn>;
|
||||
let warned: ReturnType<typeof spyOn>;
|
||||
|
||||
beforeEach(() => {
|
||||
logged = spyOn(logger, "log").mockImplementation(() => {});
|
||||
warned = spyOn(logger, "warn").mockImplementation(() => {});
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
logged.mockRestore();
|
||||
warned.mockRestore();
|
||||
});
|
||||
|
||||
function pendingMintOp(overrides: Record<string, unknown> = {}) {
|
||||
return {
|
||||
id: "op-1",
|
||||
mintUrl: "https://mint.example.com",
|
||||
quoteId: "quote-1",
|
||||
amount: 21,
|
||||
state: "pending",
|
||||
expiry: UNEXPIRED_S,
|
||||
lastObservedRemoteState: "UNPAID",
|
||||
...overrides,
|
||||
};
|
||||
}
|
||||
|
||||
function fakeSource(
|
||||
ops: Array<Record<string, unknown>>,
|
||||
behavior: {
|
||||
refresh?: (id: string) => Promise<Record<string, unknown>>;
|
||||
get?: (id: string) => Promise<Record<string, unknown> | null>;
|
||||
} = {},
|
||||
) {
|
||||
const refresh = mock(behavior.refresh ?? (async () => UNPAID));
|
||||
const get = mock(behavior.get ?? (async (id: string) => ops.find((op) => op.id === id) ?? null));
|
||||
const failPendingOperation = mock(
|
||||
async (
|
||||
_op: { id: string },
|
||||
_failure: { reason: string; retryable?: boolean; observedAt: number },
|
||||
) => ({}),
|
||||
);
|
||||
const source = {
|
||||
ops: { mint: { listPending: async () => ops, refresh, get } },
|
||||
wallet: {
|
||||
balances: {
|
||||
byMint: async () => ({ "https://mint.example.com": { spendable: 121 } }),
|
||||
},
|
||||
},
|
||||
mintOperationService: { failPendingOperation },
|
||||
} as unknown as PendingMintQuoteSource;
|
||||
return { source, refresh, get, failPendingOperation };
|
||||
}
|
||||
|
||||
const attempted = (refresh: { mock: { calls: unknown[][] } }) =>
|
||||
refresh.mock.calls.map((call) => call[0]);
|
||||
|
||||
it("refreshes every pending quote and logs the credit for the minted one", async () => {
|
||||
const ops = [pendingMintOp({ id: "op-1" }), pendingMintOp({ id: "op-2" })];
|
||||
const { source, refresh } = fakeSource(ops, {
|
||||
refresh: async (id) => (id === "op-1" ? MINTED : UNPAID),
|
||||
});
|
||||
|
||||
const result = await settlePendingMintQuotes(source, NOW_MS);
|
||||
|
||||
expect(result).toEqual({ unreachable: 0 });
|
||||
expect(attempted(refresh)).toEqual(["op-1", "op-2"]);
|
||||
expect(logged).toHaveBeenCalledTimes(1);
|
||||
expect(logged.mock.calls[0]?.[0]).toBe(
|
||||
"Mint quote quote-1 at https://mint.example.com: paid, 21 sat minted, balance now 121 sat",
|
||||
);
|
||||
});
|
||||
|
||||
it("mints an expired quote that was paid before it expired", async () => {
|
||||
const { source, failPendingOperation } = fakeSource(
|
||||
[pendingMintOp({ expiry: EXPIRED_S })],
|
||||
{ refresh: async () => MINTED },
|
||||
);
|
||||
|
||||
await settlePendingMintQuotes(source, NOW_MS);
|
||||
|
||||
expect(logged.mock.calls[0]?.[0]).toContain("21 sat minted");
|
||||
expect(failPendingOperation).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("fails an expired quote locally once its mint confirms it is unpaid", async () => {
|
||||
const { source, failPendingOperation } = fakeSource([
|
||||
pendingMintOp({ id: "expired", expiry: EXPIRED_S }),
|
||||
pendingMintOp({ id: "open" }),
|
||||
]);
|
||||
|
||||
await settlePendingMintQuotes(source, NOW_MS);
|
||||
|
||||
expect(failPendingOperation).toHaveBeenCalledTimes(1);
|
||||
expect(failPendingOperation.mock.calls[0]?.[0]).toEqual({ id: "expired" });
|
||||
});
|
||||
|
||||
it("reports an issued quote whose proofs could not be restored as a warning, not a credit", async () => {
|
||||
const { source } = fakeSource([pendingMintOp()], {
|
||||
refresh: async () => ({ ...MINTED, error: "no proofs could be restored" }),
|
||||
});
|
||||
|
||||
await settlePendingMintQuotes(source, NOW_MS);
|
||||
|
||||
expect(logged).not.toHaveBeenCalled();
|
||||
expect(warned.mock.calls[0]?.[0]).toContain("no proofs could be restored");
|
||||
});
|
||||
|
||||
it("keeps going when a mint is unreachable and counts the failure", async () => {
|
||||
const ops = [pendingMintOp({ id: "op-1" }), pendingMintOp({ id: "op-2" })];
|
||||
const { source, refresh } = fakeSource(ops, {
|
||||
refresh: async (id) => {
|
||||
if (id === "op-1") throw new Error("Network request failed");
|
||||
return MINTED;
|
||||
},
|
||||
});
|
||||
|
||||
const result = await settlePendingMintQuotes(source, NOW_MS);
|
||||
|
||||
expect(result).toEqual({ unreachable: 1 });
|
||||
expect(refresh).toHaveBeenCalledTimes(2);
|
||||
});
|
||||
|
||||
it("warns once for a paid quote whose mint attempt failed and keeps retrying", async () => {
|
||||
const failed = pendingMintOp({ lastObservedRemoteState: "PAID", error: "Mint request failed" });
|
||||
const { source } = fakeSource([pendingMintOp()], {
|
||||
refresh: async () => {
|
||||
throw new Error("Mint request failed");
|
||||
},
|
||||
get: async () => failed,
|
||||
});
|
||||
|
||||
const first = await settlePendingMintQuotes(source, NOW_MS);
|
||||
// Next sweep sees the same persisted error: no second warning.
|
||||
const { source: again } = fakeSource([failed], {
|
||||
refresh: async () => {
|
||||
throw new Error("Mint request failed");
|
||||
},
|
||||
get: async () => failed,
|
||||
});
|
||||
const second = await settlePendingMintQuotes(again, NOW_MS);
|
||||
|
||||
expect(first).toEqual({ unreachable: 0 });
|
||||
expect(second).toEqual({ unreachable: 0 });
|
||||
expect(warned).toHaveBeenCalledTimes(1);
|
||||
expect(warned.mock.calls[0]?.[0]).toContain("paid (21 sat) but proofs not minted yet");
|
||||
});
|
||||
|
||||
it("leaves a quote coco's own watcher is already minting alone", async () => {
|
||||
const { source, get } = fakeSource([pendingMintOp()], {
|
||||
refresh: async () => {
|
||||
throw new OperationInProgressError("op-1");
|
||||
},
|
||||
});
|
||||
|
||||
const result = await settlePendingMintQuotes(source, NOW_MS);
|
||||
|
||||
expect(result).toEqual({ unreachable: 0 });
|
||||
expect(get).not.toHaveBeenCalled();
|
||||
expect(warned).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("survives a reporting failure without an unhandled rejection", async () => {
|
||||
const { source } = fakeSource([pendingMintOp()], {
|
||||
refresh: async () => {
|
||||
throw new Error("Network request failed");
|
||||
},
|
||||
get: async () => {
|
||||
throw new Error("Cannot use a closed database");
|
||||
},
|
||||
});
|
||||
|
||||
const result = await settlePendingMintQuotes(source, NOW_MS);
|
||||
|
||||
expect(result).toEqual({ unreachable: 1 });
|
||||
});
|
||||
|
||||
it("gives a stalled quote one bounded wait and still serves the others", async () => {
|
||||
const ops = [pendingMintOp({ id: "slow" }), pendingMintOp({ id: "paid" })];
|
||||
const { source, refresh } = fakeSource(ops, {
|
||||
refresh: (id) => (id === "slow" ? new Promise(() => {}) : Promise.resolve(MINTED)),
|
||||
});
|
||||
const state: PendingMintSweepState = { outstanding: new Map() };
|
||||
const options = { deadlineMs: 1_000, checkTimeoutMs: 50, state };
|
||||
|
||||
const first = await settlePendingMintQuotes(source, NOW_MS, options);
|
||||
const second = await settlePendingMintQuotes(source, NOW_MS, options);
|
||||
|
||||
expect(first).toEqual({ unreachable: 1 });
|
||||
// Still outstanding, so it is skipped rather than re-sent.
|
||||
expect(second).toEqual({ unreachable: 1 });
|
||||
expect(attempted(refresh)).toEqual(["slow", "paid", "paid"]);
|
||||
});
|
||||
|
||||
it("resumes after the last attempted quote so slow ones cannot starve the rest", async () => {
|
||||
const ops = [
|
||||
pendingMintOp({ id: "s1" }),
|
||||
pendingMintOp({ id: "s2" }),
|
||||
pendingMintOp({ id: "paid" }),
|
||||
];
|
||||
// Slow refreshes settle between sweeps, so nothing is outstanding next time.
|
||||
const { source, refresh } = fakeSource(ops, {
|
||||
refresh: (id) =>
|
||||
id === "paid"
|
||||
? Promise.resolve(MINTED)
|
||||
: new Promise((resolve) => setTimeout(() => resolve(UNPAID), 160)),
|
||||
});
|
||||
const state: PendingMintSweepState = { outstanding: new Map() };
|
||||
const options = { deadlineMs: 200, checkTimeoutMs: 120, state };
|
||||
|
||||
await settlePendingMintQuotes(source, NOW_MS, options);
|
||||
expect(attempted(refresh)).toEqual(["s1", "s2"]);
|
||||
await new Promise((resolve) => setTimeout(resolve, 200));
|
||||
await settlePendingMintQuotes(source, NOW_MS, options);
|
||||
|
||||
expect(attempted(refresh)[2]).toBe("paid");
|
||||
expect(logged.mock.calls[0]?.[0]).toContain("21 sat minted");
|
||||
});
|
||||
});
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
import {
|
||||
Manager,
|
||||
OperationInProgressError,
|
||||
getEncodedToken,
|
||||
normalizeMintUrl,
|
||||
} from "@cashu/coco-core";
|
||||
@@ -146,6 +147,9 @@ const SAFE_COCO_LOG_FIELDS = new Set([
|
||||
"operationId",
|
||||
"quoteId",
|
||||
"state",
|
||||
"code",
|
||||
"detail",
|
||||
"status",
|
||||
"count",
|
||||
"total",
|
||||
"filterCount",
|
||||
@@ -766,6 +770,181 @@ export async function settleExpiredMintQuotes(
|
||||
return settlement;
|
||||
}
|
||||
|
||||
const PENDING_MINT_SWEEP_INTERVAL_MS = 15_000;
|
||||
/** Per-quote wait inside a sweep, so one stalled mint cannot starve the rest. */
|
||||
const PENDING_MINT_CHECK_TIMEOUT_MS = 10_000;
|
||||
const PENDING_MINT_SWEEP_DEADLINE_MS = 30_000;
|
||||
const PENDING_MINT_STOP_DRAIN_MS = 10_000;
|
||||
|
||||
export interface PendingMintQuoteSource {
|
||||
ops: { mint: Pick<Manager["ops"]["mint"], "listPending" | "refresh" | "get"> };
|
||||
wallet: { balances: Pick<Manager["wallet"]["balances"], "byMint"> };
|
||||
mintOperationService: Pick<MintOperationServiceCleanup, "failPendingOperation">;
|
||||
}
|
||||
|
||||
export interface PendingMintSweepState {
|
||||
/** Refreshes that outlived their wait; skipped until they settle. */
|
||||
outstanding: Map<string, Promise<unknown>>;
|
||||
/** Next sweep starts after this operation. */
|
||||
after?: string;
|
||||
}
|
||||
|
||||
export interface PendingMintSweepOptions {
|
||||
deadlineMs?: number;
|
||||
checkTimeoutMs?: number;
|
||||
state?: PendingMintSweepState;
|
||||
}
|
||||
|
||||
type PendingMintOutcome = "unreachable" | "other";
|
||||
type PendingMintOp = Awaited<ReturnType<PendingMintQuoteSource["ops"]["mint"]["listPending"]>>[number];
|
||||
|
||||
/** Startup mint recovery, run while up: coco's live subscriptions can stop. */
|
||||
export async function settlePendingMintQuotes(
|
||||
source: PendingMintQuoteSource,
|
||||
nowMs: number,
|
||||
options: PendingMintSweepOptions = {},
|
||||
): Promise<{ unreachable: number }> {
|
||||
const deadlineMs = options.deadlineMs ?? PENDING_MINT_SWEEP_DEADLINE_MS;
|
||||
const checkTimeoutMs = options.checkTimeoutMs ?? PENDING_MINT_CHECK_TIMEOUT_MS;
|
||||
const state: PendingMintSweepState = options.state ?? { outstanding: new Map() };
|
||||
let unreachable = 0;
|
||||
|
||||
const pending = await source.ops.mint.listPending();
|
||||
const resumeAt = pending.findIndex((op) => op.id === state.after) + 1;
|
||||
const ordered = [...pending.slice(resumeAt), ...pending.slice(0, resumeAt)];
|
||||
const startedAt = Date.now();
|
||||
for (const op of ordered) {
|
||||
const remainingMs = deadlineMs - (Date.now() - startedAt);
|
||||
if (state.outstanding.has(op.id) || remainingMs <= 0) {
|
||||
unreachable++;
|
||||
continue;
|
||||
}
|
||||
state.after = op.id;
|
||||
// Report on the refresh itself so a late result is still logged.
|
||||
const settled = source.ops.mint
|
||||
.refresh(op.id)
|
||||
.then(
|
||||
(result) => reportPendingMintRefresh(source, op, result, nowMs),
|
||||
(error) => reportPendingMintRefreshError(source, op, error),
|
||||
)
|
||||
.finally(() => state.outstanding.delete(op.id));
|
||||
state.outstanding.set(op.id, settled);
|
||||
try {
|
||||
const outcome = await withTimeout(settled, Math.min(checkTimeoutMs, remainingMs));
|
||||
if (outcome === "unreachable") unreachable++;
|
||||
} catch {
|
||||
unreachable++;
|
||||
}
|
||||
}
|
||||
return { unreachable };
|
||||
}
|
||||
|
||||
async function reportPendingMintRefresh(
|
||||
source: PendingMintQuoteSource,
|
||||
op: PendingMintOp,
|
||||
result: Awaited<ReturnType<PendingMintQuoteSource["ops"]["mint"]["refresh"]>>,
|
||||
nowMs: number,
|
||||
): Promise<PendingMintOutcome> {
|
||||
const quote = `Mint quote ${op.quoteId} at ${op.mintUrl}`;
|
||||
if (result.state === "finalized") {
|
||||
if (result.error) {
|
||||
// Already issued, proofs not restored: not a credit.
|
||||
logger.warn(`${quote}: ${result.error}`);
|
||||
return "other";
|
||||
}
|
||||
const balance = (await source.wallet.balances.byMint())[op.mintUrl]?.spendable;
|
||||
logger.log(
|
||||
`${quote}: paid, ${op.amount} sat minted` +
|
||||
(balance === undefined ? "" : `, balance now ${balance} sat`),
|
||||
);
|
||||
return "other";
|
||||
}
|
||||
if (result.state === "failed") {
|
||||
logger.warn(`${quote}: failed at the mint: ${result.error ?? "unknown reason"}`);
|
||||
return "other";
|
||||
}
|
||||
// Paid-before-expiry is only known to the mint, so expired quotes are
|
||||
// checked too; confirmed UNPAID after expiry is failed locally, as at startup.
|
||||
const expired = op.expiry > 0 && op.expiry * 1000 <= nowMs;
|
||||
if (expired && result.state === "pending" && result.lastObservedRemoteState === "UNPAID") {
|
||||
await source.mintOperationService.failPendingOperation(
|
||||
{ id: op.id },
|
||||
{
|
||||
reason: "Expired mint quote confirmed unpaid by mint",
|
||||
retryable: false,
|
||||
observedAt: Date.now(),
|
||||
},
|
||||
);
|
||||
}
|
||||
return "other";
|
||||
}
|
||||
|
||||
async function reportPendingMintRefreshError(
|
||||
source: PendingMintQuoteSource,
|
||||
op: PendingMintOp,
|
||||
error: unknown,
|
||||
): Promise<PendingMintOutcome> {
|
||||
// coco's own watcher is minting this quote right now; let it finish.
|
||||
if (error instanceof OperationInProgressError) return "other";
|
||||
// A failed mint attempt rejects after coco moved the operation back to pending.
|
||||
const current = await source.ops.mint.get(op.id);
|
||||
if (current?.state === "pending" && current.lastObservedRemoteState === "PAID") {
|
||||
if (op.lastObservedRemoteState !== "PAID" || current.error !== op.error) {
|
||||
logger.warn(
|
||||
`Mint quote ${op.quoteId} at ${op.mintUrl}: paid (${op.amount} sat) but proofs not minted yet: ${current.error ?? String(error)}; will retry`,
|
||||
);
|
||||
}
|
||||
return "other";
|
||||
}
|
||||
return "unreachable";
|
||||
}
|
||||
|
||||
/**
|
||||
* Stopping waits for the running sweep, then briefly for stragglers. Anything
|
||||
* still running after that fails against the closed database and is picked
|
||||
* up by startup recovery.
|
||||
*/
|
||||
function startPendingMintSweep(source: PendingMintQuoteSource): () => Promise<void> {
|
||||
let stopped = false;
|
||||
let timer: ReturnType<typeof setTimeout> | undefined;
|
||||
let inFlight: Promise<void> = Promise.resolve();
|
||||
let unreachableBefore = 0;
|
||||
const state: PendingMintSweepState = { outstanding: new Map() };
|
||||
|
||||
const tick = async () => {
|
||||
if (stopped) return;
|
||||
inFlight = settlePendingMintQuotes(source, Date.now(), { state }).then(
|
||||
({ unreachable }) => {
|
||||
// Report a mint becoming unreachable, or reachable again, once.
|
||||
if (unreachable > 0 && unreachableBefore === 0) {
|
||||
logger.warn(`Could not check ${unreachable} pending mint quote(s); will keep retrying`);
|
||||
} else if (unreachable === 0 && unreachableBefore > 0) {
|
||||
logger.log("Pending mint quote checks are reaching the mint again");
|
||||
}
|
||||
unreachableBefore = unreachable;
|
||||
},
|
||||
(error: unknown) => {
|
||||
logger.warn(
|
||||
`Pending mint quote sweep failed: ${error instanceof Error ? error.message : String(error)}`,
|
||||
);
|
||||
},
|
||||
);
|
||||
await inFlight;
|
||||
if (!stopped) timer = setTimeout(tick, PENDING_MINT_SWEEP_INTERVAL_MS);
|
||||
};
|
||||
timer = setTimeout(tick, PENDING_MINT_SWEEP_INTERVAL_MS);
|
||||
|
||||
return async () => {
|
||||
stopped = true;
|
||||
if (timer !== undefined) clearTimeout(timer);
|
||||
await inFlight;
|
||||
await withTimeout(
|
||||
Promise.allSettled(state.outstanding.values()),
|
||||
PENDING_MINT_STOP_DRAIN_MS,
|
||||
).catch(() => {});
|
||||
};
|
||||
}
|
||||
|
||||
interface ReceiveRecoveryInternals {
|
||||
receiveOperationService: {
|
||||
checkProofStatesWithMint(
|
||||
@@ -973,6 +1152,7 @@ export async function createCocoClient(
|
||||
pendingMints: 0,
|
||||
};
|
||||
let recoveryResolve: (() => void) | undefined;
|
||||
let stopPendingMintSweep: (() => Promise<void>) | undefined;
|
||||
const recoveryPromise = new Promise<void>((resolve) => {
|
||||
recoveryResolve = resolve;
|
||||
});
|
||||
@@ -1167,6 +1347,13 @@ export async function createCocoClient(
|
||||
recoveryPhase = "done";
|
||||
recoveryResolve?.();
|
||||
startupProgress("Wallet recovery complete.");
|
||||
stopPendingMintSweep = startPendingMintSweep({
|
||||
ops: coco!.ops,
|
||||
wallet: coco!.wallet,
|
||||
mintOperationService: (
|
||||
coco as unknown as { mintOperationService: MintOperationServiceCleanup }
|
||||
).mintOperationService,
|
||||
});
|
||||
})
|
||||
.catch((error) => {
|
||||
recoveryDone = true;
|
||||
@@ -1489,6 +1676,7 @@ export async function createCocoClient(
|
||||
// Let any in-flight recovery settle before closing the database from
|
||||
// underneath it. The recovery promise resolves on success or failure.
|
||||
await recoveryPromise;
|
||||
await stopPendingMintSweep?.();
|
||||
await coco.dispose();
|
||||
} finally {
|
||||
try {
|
||||
|
||||
Reference in New Issue
Block a user