diff --git a/src/daemon/wallet/coco-client.test.ts b/src/daemon/wallet/coco-client.test.ts index 5997421..2a44170 100644 --- a/src/daemon/wallet/coco-client.test.ts +++ b/src/daemon/wallet/coco-client.test.ts @@ -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[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; + let warned: ReturnType; + + beforeEach(() => { + logged = spyOn(logger, "log").mockImplementation(() => {}); + warned = spyOn(logger, "warn").mockImplementation(() => {}); + }); + + afterEach(() => { + logged.mockRestore(); + warned.mockRestore(); + }); + + function pendingMintOp(overrides: Record = {}) { + 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>, + behavior: { + refresh?: (id: string) => Promise>; + get?: (id: string) => Promise | 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"); + }); +}); diff --git a/src/daemon/wallet/coco-client.ts b/src/daemon/wallet/coco-client.ts index b8e91fc..53c0141 100644 --- a/src/daemon/wallet/coco-client.ts +++ b/src/daemon/wallet/coco-client.ts @@ -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 }; + wallet: { balances: Pick }; + mintOperationService: Pick; +} + +export interface PendingMintSweepState { + /** Refreshes that outlived their wait; skipped until they settle. */ + outstanding: Map>; + /** Next sweep starts after this operation. */ + after?: string; +} + +export interface PendingMintSweepOptions { + deadlineMs?: number; + checkTimeoutMs?: number; + state?: PendingMintSweepState; +} + +type PendingMintOutcome = "unreachable" | "other"; +type PendingMintOp = Awaited>[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>, + nowMs: number, +): Promise { + 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 { + // 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 { + let stopped = false; + let timer: ReturnType | undefined; + let inFlight: Promise = 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) | undefined; const recoveryPromise = new Promise((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 {