From 217f9d299bf90e9928fe9ea47421079fab469efd Mon Sep 17 00:00:00 2001 From: redshift <213178690+sh1ftred@users.noreply.github.com> Date: Wed, 30 Sep 2026 15:14:54 +0000 Subject: [PATCH] Revert "wallet: probe mint reachability, per-mint recovery gating" --- src/daemon/wallet/coco-client.test.ts | 71 ------ src/daemon/wallet/coco-client.ts | 187 ++------------ src/daemon/wallet/recovery-probe.test.ts | 279 -------------------- src/daemon/wallet/recovery-probe.ts | 312 ----------------------- 4 files changed, 18 insertions(+), 831 deletions(-) delete mode 100644 src/daemon/wallet/recovery-probe.test.ts delete mode 100644 src/daemon/wallet/recovery-probe.ts diff --git a/src/daemon/wallet/coco-client.test.ts b/src/daemon/wallet/coco-client.test.ts index daa8ef1..959e798 100644 --- a/src/daemon/wallet/coco-client.test.ts +++ b/src/daemon/wallet/coco-client.test.ts @@ -14,7 +14,6 @@ import { assertLegacyCocodNotRunning, claimLegacyCocodPidFile, createCocoClient, - createRecoveryGate, DEFAULT_TRUSTED_MINT_URLS, isZombieProcess, settleExpiredMintQuotes, @@ -902,73 +901,3 @@ describe("settlePendingMintQuotes", () => { expect(logged.mock.calls[0]?.[0]).toContain("21 sat minted"); }); }); - -describe("createRecoveryGate", () => { - const HEALTHY = "https://healthy.example.com"; - const STUCK = "https://stuck.example.com"; - - function settled(promise: Promise): Promise { - return Promise.race([ - promise.then(() => true, () => true), - new Promise((resolve) => setTimeout(() => resolve(false), 25)), - ]); - } - - it("lets a mint without stuck operations proceed while another mint recovers", async () => { - const gate = createRecoveryGate(); - gate.publishStuckMints(new Set([STUCK])); - // Recovery is still running (complete() never called), yet the healthy - // mint must not be blocked by the stuck one. - await gate.waitForRecovery(HEALTHY); - }); - - it("holds a mint with stuck operations until recovery completes", async () => { - const gate = createRecoveryGate(); - gate.publishStuckMints(new Set([STUCK])); - const waiting = gate.waitForRecovery(STUCK); - expect(await settled(waiting)).toBe(false); - gate.complete(); - await waiting; - }); - - it("waits for the stuck-mint enumeration before deciding", async () => { - const gate = createRecoveryGate(); - const waiting = gate.waitForRecovery(HEALTHY); - expect(await settled(waiting)).toBe(false); - gate.publishStuckMints(new Set([STUCK])); - await waiting; - }); - - it("holds callers without a target mint until recovery completes", async () => { - const gate = createRecoveryGate(); - gate.publishStuckMints(new Set([STUCK])); - const waiting = gate.waitForRecovery(); - expect(await settled(waiting)).toBe(false); - gate.complete(); - await waiting; - }); - - it("poisons every caller after a recovery failure", async () => { - const gate = createRecoveryGate(); - gate.publishStuckMints(new Set([STUCK])); - gate.fail("disk exploded"); - await expect(gate.waitForRecovery(HEALTHY)).rejects.toThrow( - "Wallet is not ready: disk exploded", - ); - await expect(gate.waitForRecovery(STUCK)).rejects.toThrow( - "Wallet is not ready: disk exploded", - ); - await expect(gate.waitForRecovery()).rejects.toThrow( - "Wallet is not ready: disk exploded", - ); - }); - - it("falls back to the global gate for unparseable mint URLs", async () => { - const gate = createRecoveryGate(); - gate.publishStuckMints(new Set([STUCK])); - const waiting = gate.waitForRecovery("not a url"); - expect(await settled(waiting)).toBe(false); - gate.complete(); - await waiting; - }); -}); diff --git a/src/daemon/wallet/coco-client.ts b/src/daemon/wallet/coco-client.ts index 2768dac..c8077c8 100644 --- a/src/daemon/wallet/coco-client.ts +++ b/src/daemon/wallet/coco-client.ts @@ -37,14 +37,6 @@ import type { WalletRecoveryProgress, } from "./cocod-client"; import { selectCleanupOperations } from "./cleanup"; -import { - collectStuckOperations, - probeMintReachability, - runTargetedRecovery, - startRecoveryRecheck, - type SendRecoveryService, - type StuckOperation, -} from "./recovery-probe"; import { clearInterruptedReceiveReservations, deleteReceiveTokenReservation, @@ -1039,85 +1031,6 @@ interface RecoveryPhaseProgress { failedMintQuotes: number; } -/** - * Coco keeps per-operation recovery private on its services; routstrd already - * reaches into the Manager the same way for `mintOperationService`. Send is - * the only family whose public `refresh()` cannot recover executing ops. - */ -function sendRecoveryServiceOf(coco: Manager): SendRecoveryService { - return (coco as unknown as { sendOperationService: SendRecoveryService }) - .sendOperationService; -} - -/** - * Gate for value-moving wallet operations while startup recovery runs. - * - * The gate is per mint: once recovery has enumerated stuck operations - * (publishStuckMints), callers whose target mint has none proceed immediately - * — a dead or slow mint must not stall spends from a healthy one. Callers - * without a target mint, or whose mint has stuck operations, wait for the - * full sweep. fail() poisons every caller; reads are never gated. - */ -export interface RecoveryGate { - waitForRecovery(mintUrl?: string): Promise; - publishStuckMints(mints: Set): void; - complete(): void; - fail(error: string): void; -} - -export function createRecoveryGate(): RecoveryGate { - let stuckMints: Set | undefined; - let done = false; - let error: string | undefined; - let mintsResolve: (() => void) | undefined; - const mintsPromise = new Promise((resolve) => { - mintsResolve = resolve; - }); - let doneResolve: (() => void) | undefined; - const donePromise = new Promise((resolve) => { - doneResolve = resolve; - }); - - return { - async waitForRecovery(mintUrl?: string): Promise { - if (mintUrl) { - let normalized: string | undefined; - try { - normalized = normalizeMintUrl(mintUrl); - } catch { - // Unparseable URL falls back to the global gate. - normalized = undefined; - } - if (normalized) { - await mintsPromise; - if (!stuckMints?.has(normalized)) { - if (error) throw new Error(`Wallet is not ready: ${error}`); - return; - } - } - } - if (!done) await donePromise; - if (error) throw new Error(`Wallet is not ready: ${error}`); - }, - publishStuckMints(mints: Set): void { - if (stuckMints) return; - stuckMints = mints; - mintsResolve?.(); - }, - complete(): void { - done = true; - mintsResolve?.(); - doneResolve?.(); - }, - fail(message: string): void { - done = true; - error = message; - mintsResolve?.(); - doneResolve?.(); - }, - }; -} - /** * Run the wallet recovery sweeps in order, reporting phase changes. * @@ -1129,7 +1042,6 @@ async function runWalletRecovery( coco: Manager, onProgress: (progress: RecoveryPhaseProgress) => void, receiveOperationIds?: string[], - onStuckMintsKnown?: (mints: Set) => void, ): Promise { surfacingRecoveryProgress = true; let failedMintQuotes = 0; @@ -1156,60 +1068,19 @@ async function runWalletRecovery( } onProgress({ phase: "Settled expired mint quotes", failedMintQuotes }); - // Probe every mint that has stuck operations once, up front, so a dead - // mint costs a single short probe instead of a network timeout per - // operation per sweep. Healthy-mint operations are recovered per op; - // dead-mint operations stay parked exactly as coco's own "will retry - // later" path would leave them. - onProgress({ phase: "Probing mints", failedMintQuotes }); - const stuckOperations = await collectStuckOperations(coco.ops); - // Value-moving operations gate per mint on this set (see waitForRecovery); - // publish it as soon as it is known so healthy mints unblock immediately. - onStuckMintsKnown?.(new Set(stuckOperations.map((op) => op.mintUrl))); - const unreachableMints = await probeMintReachability( - [...new Set(stuckOperations.map((op) => op.mintUrl))], - ); - for (const mintUrl of unreachableMints) { - const count = stuckOperations.filter((op) => op.mintUrl === mintUrl).length; - startupProgress( - `Skipping recovery for unreachable mint: ${mintUrl} (${count} op${count === 1 ? "" : "s"})`, - ); - } - const degraded = unreachableMints.size > 0; - const targeted = (kinds: Array) => - runTargetedRecovery(coco.ops, sendRecoveryServiceOf(coco), { - kinds, - stuckOperations, - unreachableMints, - }); - - // Happy path (every mint reachable) keeps coco's global sweeps: they also - // clean up init operations and orphaned proof reservations, which the - // per-op driver cannot enumerate. The targeted driver only runs when a - // dead mint would otherwise tax every stuck op with a network timeout. onProgress({ phase: "Send recovery", failedMintQuotes }); - if (!degraded) await coco.ops.send.recovery.run(); - else await targeted(["send"]); + await coco.ops.send.recovery.run(); onProgress({ phase: "Melt recovery", failedMintQuotes }); - if (!degraded) await coco.ops.melt.recovery.run(); - else await targeted(["melt"]); + await coco.ops.melt.recovery.run(); onProgress({ phase: "Receive recovery", failedMintQuotes }); if (receiveOperationIds) { // The pre-check already classified every executing receive by unique // input set. Recover only the conclusive retained operations; unresolved // groups stay untouched instead of falling back to Coco 1's expensive - // per-row sweep on this startup. Operations at mints the probe found - // unreachable are skipped rather than costing their 15s timeout each. - const mintByOperation = new Map( - stuckOperations - .filter((op) => op.kind === "receive") - .map((op) => [op.id, op.mintUrl]), - ); + // per-row sweep on this startup. for (const operationId of receiveOperationIds) { - const mintUrl = mintByOperation.get(operationId); - if (mintUrl && unreachableMints.has(mintUrl)) continue; try { await withTimeout(coco.ops.receive.refresh(operationId), 15_000); } catch (error) { @@ -1219,15 +1090,12 @@ async function runWalletRecovery( }); } } - } else if (!degraded) { - await coco.ops.receive.recovery.run(); } else { - await targeted(["receive"]); + await coco.ops.receive.recovery.run(); } onProgress({ phase: "Mint recovery", failedMintQuotes }); - if (!degraded) await coco.recoverPendingMintOperations(); - else await targeted(["mint"]); + await coco.recoverPendingMintOperations(); onProgress({ phase: "done", failedMintQuotes }); } finally { @@ -1287,11 +1155,9 @@ export async function createCocoClient( }; let recoveryResolve: (() => void) | undefined; let stopPendingMintSweep: (() => Promise) | undefined; - let stopRecoveryRecheck: (() => Promise) | undefined; const recoveryPromise = new Promise((resolve) => { recoveryResolve = resolve; }); - const recoveryGate = createRecoveryGate(); try { startupProgress("Opening Cashu wallet database..."); @@ -1489,33 +1355,13 @@ export async function createCocoClient( } }, receiveRecoveryOperationIds, - (mints) => recoveryGate.publishStuckMints(mints), ) .then(async () => { await syncReceiveReservations(); recoveryDone = true; recoveryPhase = "done"; - recoveryGate.complete(); recoveryResolve?.(); startupProgress("Wallet recovery complete."); - // Re-check stuck operations periodically: a mint that comes back has - // its parked operations recovered without a daemon restart, and - // operations stuck mid-session are reconciled too. Receive stays - // startup-only because recovering competing receives safely requires - // the startup dedup classification (see receive-dedup.ts). - stopRecoveryRecheck = startRecoveryRecheck( - coco!.ops, - sendRecoveryServiceOf(coco!), - { - kinds: ["send", "melt", "mint"], - onMintDown: (mintUrl, opCount) => - logger.warn( - `Mint ${mintUrl} is unreachable; ${opCount} stuck operation(s) will keep retrying`, - ), - onMintBack: (mintUrl) => - logger.log(`Mint ${mintUrl} is reachable again; resumed recovering its operations`), - }, - ); stopPendingMintSweep = startPendingMintSweep({ ops: coco!.ops, wallet: coco!.wallet, @@ -1528,7 +1374,6 @@ export async function createCocoClient( recoveryDone = true; recoveryPhase = "error"; recoveryError = error instanceof Error ? error.message : String(error); - recoveryGate.fail(recoveryError); recoveryResolve?.(); startupProgress(`Wallet recovery failed: ${recoveryError}`); }); @@ -1553,11 +1398,16 @@ export async function createCocoClient( let disposed = false; - // Block a value-moving operation until background recovery has settled for - // its target mint (see createRecoveryGate). Reads stay ungated so the - // daemon can report balances/status immediately. - const waitForRecovery = (mintUrl?: string): Promise => - recoveryGate.waitForRecovery(mintUrl); + /** + * Block a value-moving operation until background recovery has settled. + * Reads stay ungated so the daemon can report balances/status immediately. + */ + const waitForRecovery = async (): Promise => { + if (!recoveryDone) await recoveryPromise; + if (recoveryError) { + throw new Error(`Wallet is not ready: ${recoveryError}`); + } + }; return { async ping(): Promise { @@ -1737,13 +1587,13 @@ export async function createCocoClient( }, async receiveBolt11(amount: number, mintUrl?: string) { + await waitForRecovery(); const targetMint = mintUrl ? normalizeMintUrl(mintUrl) : walletConfig.defaultMintUrl; if (!targetMint) { throw new Error("No trusted mint available for Lightning invoice"); } - await waitForRecovery(targetMint); const op = await coco.ops.mint.prepare({ mintUrl: targetMint, amount, @@ -1769,13 +1619,13 @@ export async function createCocoClient( }, async sendCashu(amount: number, mintUrl?: string): Promise { + await waitForRecovery(); const targetMint = mintUrl ? normalizeMintUrl(mintUrl) : walletConfig.defaultMintUrl; if (!targetMint) { throw new Error("No trusted mint available for sending"); } - await waitForRecovery(targetMint); const prepared = await coco.ops.send.prepare({ mintUrl: targetMint, amount, @@ -1785,13 +1635,13 @@ export async function createCocoClient( }, async sendBolt11(invoice: string, mintUrl?: string): Promise { + await waitForRecovery(); const targetMint = mintUrl ? normalizeMintUrl(mintUrl) : walletConfig.defaultMintUrl; if (!targetMint) { throw new Error("No trusted mint available for Lightning payment"); } - await waitForRecovery(targetMint); const prepared = await coco.ops.melt.prepare({ mintUrl: targetMint, method: "bolt11", @@ -1837,7 +1687,6 @@ export async function createCocoClient( async dispose(): Promise { if (disposed) return; disposed = true; - await stopRecoveryRecheck?.(); try { // Let any in-flight recovery settle before closing the database from // underneath it. The recovery promise resolves on success or failure. diff --git a/src/daemon/wallet/recovery-probe.test.ts b/src/daemon/wallet/recovery-probe.test.ts deleted file mode 100644 index 2468cb9..0000000 --- a/src/daemon/wallet/recovery-probe.test.ts +++ /dev/null @@ -1,279 +0,0 @@ -import { describe, expect, it, mock } from "bun:test"; -import { - collectStuckOperations, - probeMintReachability, - runTargetedRecovery, - startRecoveryRecheck, - type StuckOperationSource, - type SendRecoveryService, -} from "./recovery-probe"; - -interface FakeOp { - id: string; - mintUrl: string; - state: string; -} - -function op(id: string, mintUrl: string, state = "pending"): FakeOp { - return { id, mintUrl, state }; -} - -interface FakeSource extends StuckOperationSource { - stuck: Record<"send" | "melt" | "receive" | "mint", FakeOp[]>; - refreshed: Record<"send" | "melt" | "receive" | "mint", string[]>; - failOnRefresh?: Set; -} - -function makeSource( - stuck: Partial>, -): FakeSource { - const refreshed: FakeSource["refreshed"] = { - send: [], - melt: [], - receive: [], - mint: [], - }; - const full = { - send: stuck.send ?? [], - melt: stuck.melt ?? [], - receive: stuck.receive ?? [], - mint: stuck.mint ?? [], - }; - const family = (kind: keyof typeof full) => ({ - listInFlight: async () => full[kind] as never, - refresh: async (id: string) => { - refreshed[kind].push(id); - if (full[kind].some((o) => o.id === id && o.id.startsWith("boom"))) { - throw new Error("mint rejected the operation"); - } - return full[kind].find((o) => o.id === id) as never; - }, - }); - return { - stuck: full, - refreshed, - send: family("send") as FakeSource["send"], - melt: family("melt") as FakeSource["melt"], - receive: family("receive") as FakeSource["receive"], - mint: family("mint") as FakeSource["mint"], - }; -} - -function makeSendService(): SendRecoveryService & { - initRecovered: string[]; - executingRecovered: string[]; -} { - const initRecovered: string[] = []; - const executingRecovered: string[] = []; - return { - initRecovered, - executingRecovered, - tryRecoverInitOperation: async (raw) => { - initRecovered.push((raw as FakeOp).id); - }, - tryRecoverExecutingOperation: async (raw) => { - executingRecovered.push((raw as FakeOp).id); - }, - }; -} - -/** fetch stub: mints whose URL contains "dead" hang/fail; others answer. */ -function makeFetch(deadPredicate: (url: string) => boolean) { - const calls: string[] = []; - const fetchImpl = (async (url: string | URL | Request) => { - const href = String(url); - calls.push(href); - if (deadPredicate(href)) throw new Error("connect ECONNREFUSED"); - return new Response("{}", { status: 200 }); - }) as unknown as typeof fetch; - return { calls, fetchImpl }; -} - -const isDeadUrl = (url: string) => url.includes("dead"); - -const liveFetch = () => makeFetch(() => false); - -describe("collectStuckOperations", () => { - it("aggregates all four operation families with normalized mint URLs", async () => { - const source = makeSource({ - send: [op("s1", "https://mint.example.com/")], - melt: [op("m1", "https://mint.example.com")], - receive: [op("r1", "https://other.example.com", "executing")], - mint: [op("q1", "https://mint.example.com")], - }); - const stuck = await collectStuckOperations(source); - expect(stuck).toHaveLength(4); - expect(stuck.map((s) => s.kind).sort()).toEqual([ - "melt", - "mint", - "receive", - "send", - ]); - // Trailing slash normalized so the same mint dedupes to one probe target. - const mintUrls = new Set(stuck.map((s) => s.mintUrl)); - expect(mintUrls.size).toBe(2); - }); - - it("returns empty when nothing is stuck", async () => { - const stuck = await collectStuckOperations(makeSource({})); - expect(stuck).toEqual([]); - }); -}); - -describe("probeMintReachability", () => { - it("marks only mints whose fetch fails as unreachable", async () => { - const deadFetch = makeFetch(isDeadUrl); - const unreachable = await probeMintReachability( - ["https://live.example.com", "https://dead.example.com"], - { fetchImpl: deadFetch.fetchImpl }, - ); - expect([...unreachable]).toEqual(["https://dead.example.com"]); - expect(deadFetch.calls).toHaveLength(2); - expect(deadFetch.calls.every((c) => c.endsWith("/v1/info"))).toBe(true); - }); - - it("treats HTTP error responses as reachable", async () => { - const fetchImpl = (async () => - new Response("oops", { status: 500 })) as unknown as typeof fetch; - const unreachable = await probeMintReachability(["https://live.example.com"], { - fetchImpl, - }); - expect(unreachable.size).toBe(0); - }); -}); - -describe("runTargetedRecovery", () => { - it("does not probe when nothing is stuck", async () => { - const deadFetch = makeFetch(isDeadUrl); - const result = await runTargetedRecovery(makeSource({}), makeSendService(), { - fetchImpl: deadFetch.fetchImpl, - }); - expect(result).toMatchObject({ recovered: 0, skipped: 0, failed: 0 }); - expect(deadFetch.calls).toHaveLength(0); - }); - - it("skips operations at unreachable mints and recovers the rest", async () => { - const source = makeSource({ - send: [op("s1", "https://dead.example.com"), op("s2", "https://live.example.com")], - melt: [op("m1", "https://dead.example.com"), op("m2", "https://live.example.com")], - }); - const deadFetch = makeFetch(isDeadUrl); - const skippedMints: Array<[string, number]> = []; - const result = await runTargetedRecovery(source, makeSendService(), { - fetchImpl: deadFetch.fetchImpl, - onSkippedMint: (mintUrl, count) => skippedMints.push([mintUrl, count]), - }); - expect(result.recovered).toBe(2); - expect(result.skipped).toBe(2); - expect(result.failed).toBe(0); - expect(skippedMints).toEqual([["https://dead.example.com", 2]]); - expect(result.skippedMints.get("https://dead.example.com")).toBe(2); - // Only live-mint operations were driven. - expect(source.refreshed.send).toEqual(["s2"]); - expect(source.refreshed.melt).toEqual(["m2"]); - }); - - it("recovers pending sends via refresh and executing sends via the service", async () => { - const source = makeSource({ - send: [ - op("pending-1", "https://live.example.com", "pending"), - op("exec-1", "https://live.example.com", "executing"), - op("init-1", "https://live.example.com", "init"), - ], - }); - const sendService = makeSendService(); - const result = await runTargetedRecovery(source, sendService, { - fetchImpl: liveFetch().fetchImpl, - }); - expect(result.recovered).toBe(3); - expect(source.refreshed.send).toEqual(["pending-1"]); - expect(sendService.executingRecovered).toEqual(["exec-1"]); - expect(sendService.initRecovered).toEqual(["init-1"]); - }); - - it("counts per-operation failures at reachable mints and continues", async () => { - const source = makeSource({ - melt: [op("boom-1", "https://live.example.com"), op("m2", "https://live.example.com")], - }); - const result = await runTargetedRecovery(source, makeSendService(), { - fetchImpl: liveFetch().fetchImpl, - }); - expect(result.recovered).toBe(1); - expect(result.failed).toBe(1); - expect(source.refreshed.melt).toEqual(["boom-1", "m2"]); - }); - - it("respects the kinds filter and a pre-computed unreachable set", async () => { - const deadFetch = makeFetch(isDeadUrl); - const source = makeSource({ - send: [op("s1", "https://dead.example.com")], - mint: [op("q1", "https://dead.example.com")], - }); - const result = await runTargetedRecovery(source, makeSendService(), { - kinds: ["mint"], - unreachableMints: new Set(["https://dead.example.com"]), - fetchImpl: deadFetch.fetchImpl, - }); - // No probe ran (pre-computed set) and the send op was not even counted. - expect(deadFetch.calls).toHaveLength(0); - expect(result.skipped).toBe(1); - expect(source.refreshed.send).toEqual([]); - }); -}); - -describe("startRecoveryRecheck", () => { - it("recovers stuck operations on the interval and stops cleanly", async () => { - const source = makeSource({ - send: [op("s1", "https://live.example.com")], - }); - const stop = startRecoveryRecheck(source, makeSendService(), { - intervalMs: 20, - fetchImpl: liveFetch().fetchImpl, - }); - await new Promise((resolve) => setTimeout(resolve, 60)); - await stop(); - expect(source.refreshed.send.length).toBeGreaterThanOrEqual(1); - const after = source.refreshed.send.length; - await new Promise((resolve) => setTimeout(resolve, 60)); - expect(source.refreshed.send.length).toBe(after); - }); - - it("is idle without network traffic when nothing is stuck", async () => { - const { calls, fetchImpl } = makeFetch(() => false); - const stop = startRecoveryRecheck(makeSource({}), makeSendService(), { - intervalMs: 20, - fetchImpl, - }); - await new Promise((resolve) => setTimeout(resolve, 60)); - await stop(); - expect(calls).toHaveLength(0); - }); - - it("reports mint down/up transitions once instead of repeating", async () => { - let dead = true; - const { fetchImpl } = makeFetch(() => dead); - const source = makeSource({ - melt: [op("m1", "https://flaky.example.com")], - }); - const down: string[] = []; - const back: string[] = []; - const stop = startRecoveryRecheck(source, makeSendService(), { - intervalMs: 20, - fetchImpl, - onMintDown: (mintUrl) => down.push(mintUrl), - onMintBack: (mintUrl) => back.push(mintUrl), - }); - await new Promise((resolve) => setTimeout(resolve, 90)); - // Dead across several ticks: reported once, never re-probed per op. - expect(down).toEqual(["https://flaky.example.com"]); - expect(back).toEqual([]); - expect(source.refreshed.melt).toEqual([]); - - dead = false; - await new Promise((resolve) => setTimeout(resolve, 90)); - await stop(); - expect(down).toEqual(["https://flaky.example.com"]); - expect(back).toEqual(["https://flaky.example.com"]); - expect(source.refreshed.melt.length).toBeGreaterThanOrEqual(1); - }); -}); diff --git a/src/daemon/wallet/recovery-probe.ts b/src/daemon/wallet/recovery-probe.ts deleted file mode 100644 index 952b3ab..0000000 --- a/src/daemon/wallet/recovery-probe.ts +++ /dev/null @@ -1,312 +0,0 @@ -/** - * Mint reachability probing and targeted (per-operation) wallet recovery. - * - * Coco's global recovery sweeps walk every non-terminal operation one at a - * time; each operation at an unreachable mint costs a full network timeout. - * This module probes every mint that has stuck operations once, in parallel, - * and drives recovery per operation only for mints that answer — so one dead - * mint costs a single short probe instead of N sequential timeouts, and its - * operations stay parked (exactly as coco's "will retry later" path leaves - * them) until the mint comes back. - */ -import { normalizeMintUrl, type Manager } from "@cashu/coco-core"; -import { logger } from "../../utils/logger"; - -/** Short probe: a mint that cannot answer /v1/info in 2s slows every op. */ -export const MINT_PROBE_TIMEOUT_MS = 2_000; -/** How often stuck operations are re-checked (and dead mints re-probed). */ -export const RECOVERY_RECHECK_INTERVAL_MS = 300_000; - -type OpsApi = Manager["ops"]; - -/** Structural subset of the ops APIs used to enumerate and recover operations. */ -export interface StuckOperationSource { - send: Pick; - melt: Pick; - receive: Pick; - mint: Pick; -} - -export type StuckOperationKind = "send" | "melt" | "receive" | "mint"; - -export interface StuckOperation { - kind: StuckOperationKind; - id: string; - /** Normalized mint URL. */ - mintUrl: string; - state: string; - /** The operation object as returned by the API (needed by service-level recovery). */ - raw: unknown; -} - -/** - * The per-operation send recovery entry points coco keeps private. Send is the - * only operation family whose public `refresh()` does not cover `executing` - * operations; routstrd already reaches into coco internals the same way for - * `mintOperationService` (see coco-client.ts). - */ -export interface SendRecoveryService { - tryRecoverInitOperation(op: unknown): Promise; - tryRecoverExecutingOperation(op: unknown): Promise; -} - -export interface RecoveryRunResult { - /** Operations at reachable mints whose recovery completed. */ - recovered: number; - /** Operations skipped because their mint did not answer the probe. */ - skipped: number; - /** Operations at reachable mints whose recovery still failed. */ - failed: number; - /** Unreachable mint URL -> number of operations skipped there. */ - skippedMints: Map; -} - -export interface TargetedRecoveryOptions { - probeTimeoutMs?: number; - fetchImpl?: typeof fetch; - /** Operation families to recover. Defaults to all four. */ - kinds?: StuckOperationKind[]; - /** Pre-collected operations (e.g. from startup gating); default: enumerate now. */ - stuckOperations?: StuckOperation[]; - /** Pre-probed unreachable mints; default: probe now. */ - unreachableMints?: Set; - /** Called once per unreachable mint with the number of skipped operations. */ - onSkippedMint?: (mintUrl: string, opCount: number) => void; -} - -function asStuckOperations( - kind: StuckOperationKind, - ops: Array<{ id: string; mintUrl: string; state: string }>, -): StuckOperation[] { - const stuck: StuckOperation[] = []; - for (const op of ops) { - try { - stuck.push({ - kind, - id: op.id, - mintUrl: normalizeMintUrl(op.mintUrl), - state: op.state, - raw: op, - }); - } catch { - // Unparseable mint URL: keep the operation recoverable by treating it as - // reachable (probe only covers successfully normalized URLs). - stuck.push({ kind, id: op.id, mintUrl: op.mintUrl, state: op.state, raw: op }); - } - } - return stuck; -} - -/** - * Enumerate every non-terminal operation across all four operation families. - * Returns [] quickly when nothing is stuck, which is the common case. - */ -export async function collectStuckOperations( - source: StuckOperationSource, -): Promise { - const [sends, melts, receives, mints] = await Promise.all([ - source.send.listInFlight(), - source.melt.listInFlight(), - source.receive.listInFlight(), - source.mint.listInFlight(), - ]); - return [ - ...asStuckOperations("send", sends), - ...asStuckOperations("melt", melts), - ...asStuckOperations("receive", receives), - ...asStuckOperations("mint", mints), - ]; -} - -/** - * Probe each mint once, in parallel, and return the set of mint URLs that did - * not answer `GET /v1/info` in time. A mint that answers with an HTTP error is - * still "reachable" — its operations will fail with a real mint error instead - * of a network timeout, which is the information recovery needs. - */ -export async function probeMintReachability( - mintUrls: string[], - options: { timeoutMs?: number; fetchImpl?: typeof fetch } = {}, -): Promise> { - const fetcher = options.fetchImpl ?? fetch; - const timeoutMs = options.timeoutMs ?? MINT_PROBE_TIMEOUT_MS; - const unreachable = new Set(); - - await Promise.all( - mintUrls.map(async (mintUrl) => { - try { - await fetcher(new URL("/v1/info", mintUrl).toString(), { - signal: AbortSignal.timeout(timeoutMs), - }); - } catch (error) { - logger.debug("Mint did not answer recovery probe", { - mintUrl, - error: error instanceof Error ? error.message : String(error), - }); - unreachable.add(mintUrl); - } - }), - ); - - return unreachable; -} - -async function recoverStuckOperation( - source: StuckOperationSource, - sendService: SendRecoveryService, - op: StuckOperation, -): Promise { - switch (op.kind) { - case "send": - if (op.state === "pending") { - // Public API: actively re-checks the proofs with the mint. - await source.send.refresh(op.id); - } else if (op.state === "executing") { - // No public per-op path exists for executing sends (coco keeps it - // private); tryRecover* swallows per-op errors and leaves the - // operation for the next pass, matching the global sweep's behavior. - await sendService.tryRecoverExecutingOperation(op.raw); - } else if (op.state === "init") { - await sendService.tryRecoverInitOperation(op.raw); - } - // prepared / rolling_back: the global sweep only warns; nothing to do. - return; - case "melt": - // refresh() covers both pending and executing melt operations. - await source.melt.refresh(op.id); - return; - case "receive": - // refresh() actively recovers executing receive operations. - await source.receive.refresh(op.id); - return; - case "mint": - // refresh() covers both pending and executing mint operations. - await source.mint.refresh(op.id); - return; - } -} - -/** - * Recover every stuck operation whose mint answers a reachability probe, - * skipping operations at unreachable mints. Operations are recovered - * sequentially per mint (matching the global sweep's ordering guarantees); - * skipped operations are left untouched for a later pass, which is exactly - * what coco's own "Could not reach mint for recovery, will retry later" path - * does with them. - */ -export async function runTargetedRecovery( - source: StuckOperationSource, - sendService: SendRecoveryService, - options: TargetedRecoveryOptions = {}, -): Promise { - const result: RecoveryRunResult = { - recovered: 0, - skipped: 0, - failed: 0, - skippedMints: new Map(), - }; - - const kinds = options.kinds ?? ["send", "melt", "receive", "mint"]; - const stuck = (options.stuckOperations ?? (await collectStuckOperations(source))).filter( - (op) => kinds.includes(op.kind), - ); - if (stuck.length === 0) return result; - - const unreachable = - options.unreachableMints ?? - (await probeMintReachability([...new Set(stuck.map((op) => op.mintUrl))], { - timeoutMs: options.probeTimeoutMs, - fetchImpl: options.fetchImpl, - })); - - for (const mintUrl of unreachable) { - const count = stuck.filter((op) => op.mintUrl === mintUrl).length; - result.skippedMints.set(mintUrl, count); - result.skipped += count; - options.onSkippedMint?.(mintUrl, count); - } - - for (const op of stuck) { - if (unreachable.has(op.mintUrl)) continue; - try { - await recoverStuckOperation(source, sendService, op); - result.recovered++; - } catch (error) { - // Same semantics as coco's tryRecover*: leave the operation for the - // next pass. A reachable mint can still reject a specific operation. - result.failed++; - logger.warn("Targeted operation recovery did not complete", { - kind: op.kind, - operationId: op.id, - mintUrl: op.mintUrl, - error: error instanceof Error ? error.message : String(error), - }); - } - } - - return result; -} - -export interface RecoveryRecheckOptions extends TargetedRecoveryOptions { - intervalMs?: number; - /** Called when a mint transitions unreachable -> reachable with recovered op count. */ - onMintBack?: (mintUrl: string) => void; - /** Called when a mint transitions reachable -> unreachable. */ - onMintDown?: (mintUrl: string, opCount: number) => void; -} - -/** - * Periodically re-run targeted recovery so a mint that comes back online has - * its parked operations recovered without a daemon restart, and operations - * that get stuck mid-session are reconciled too. Idle ticks (no stuck - * operations) cost one DB query and no network traffic. Logging is - * transition-based: a mint that stays dead produces no repeated output. - * - * Returns a stop function that waits for any in-flight tick. - */ -export function startRecoveryRecheck( - source: StuckOperationSource, - sendService: SendRecoveryService, - options: RecoveryRecheckOptions = {}, -): () => Promise { - const intervalMs = options.intervalMs ?? RECOVERY_RECHECK_INTERVAL_MS; - let stopped = false; - let timer: ReturnType | undefined; - let inFlight: Promise = Promise.resolve(); - const knownDead = new Set(); - - const tick = async () => { - if (stopped) return; - inFlight = (async () => { - const result = await runTargetedRecovery(source, sendService, { - ...options, - onSkippedMint: (mintUrl, opCount) => { - if (!knownDead.has(mintUrl)) { - knownDead.add(mintUrl); - options.onMintDown?.(mintUrl, opCount); - } - options.onSkippedMint?.(mintUrl, opCount); - }, - }); - for (const mintUrl of [...knownDead]) { - if (!result.skippedMints.has(mintUrl)) { - knownDead.delete(mintUrl); - options.onMintBack?.(mintUrl); - } - } - })().catch((error: unknown) => { - logger.warn( - `Stuck-operation recheck failed: ${error instanceof Error ? error.message : String(error)}`, - ); - }); - await inFlight; - if (!stopped) timer = setTimeout(tick, intervalMs); - }; - timer = setTimeout(tick, intervalMs); - - return async () => { - stopped = true; - if (timer) clearTimeout(timer); - await inFlight; - }; -}