diff --git a/src/daemon/wallet/coco-client.test.ts b/src/daemon/wallet/coco-client.test.ts index 5c11927..ed175f3 100644 --- a/src/daemon/wallet/coco-client.test.ts +++ b/src/daemon/wallet/coco-client.test.ts @@ -14,6 +14,7 @@ import { assertLegacyCocodNotRunning, claimLegacyCocodPidFile, createCocoClient, + createRecoveryGate, DEFAULT_TRUSTED_MINT_URLS, isZombieProcess, settleExpiredMintQuotes, @@ -901,3 +902,73 @@ 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 c8077c8..2768dac 100644 --- a/src/daemon/wallet/coco-client.ts +++ b/src/daemon/wallet/coco-client.ts @@ -37,6 +37,14 @@ 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, @@ -1031,6 +1039,85 @@ 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. * @@ -1042,6 +1129,7 @@ async function runWalletRecovery( coco: Manager, onProgress: (progress: RecoveryPhaseProgress) => void, receiveOperationIds?: string[], + onStuckMintsKnown?: (mints: Set) => void, ): Promise { surfacingRecoveryProgress = true; let failedMintQuotes = 0; @@ -1068,19 +1156,60 @@ 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 }); - await coco.ops.send.recovery.run(); + if (!degraded) await coco.ops.send.recovery.run(); + else await targeted(["send"]); onProgress({ phase: "Melt recovery", failedMintQuotes }); - await coco.ops.melt.recovery.run(); + if (!degraded) await coco.ops.melt.recovery.run(); + else await targeted(["melt"]); 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. + // 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]), + ); 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) { @@ -1090,12 +1219,15 @@ async function runWalletRecovery( }); } } - } else { + } else if (!degraded) { await coco.ops.receive.recovery.run(); + } else { + await targeted(["receive"]); } onProgress({ phase: "Mint recovery", failedMintQuotes }); - await coco.recoverPendingMintOperations(); + if (!degraded) await coco.recoverPendingMintOperations(); + else await targeted(["mint"]); onProgress({ phase: "done", failedMintQuotes }); } finally { @@ -1155,9 +1287,11 @@ 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..."); @@ -1355,13 +1489,33 @@ 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, @@ -1374,6 +1528,7 @@ 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}`); }); @@ -1398,16 +1553,11 @@ export async function createCocoClient( let disposed = false; - /** - * 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}`); - } - }; + // 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); return { async ping(): Promise { @@ -1587,13 +1737,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, @@ -1619,13 +1769,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, @@ -1635,13 +1785,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", @@ -1687,6 +1837,7 @@ 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 new file mode 100644 index 0000000..2468cb9 --- /dev/null +++ b/src/daemon/wallet/recovery-probe.test.ts @@ -0,0 +1,279 @@ +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 new file mode 100644 index 0000000..952b3ab --- /dev/null +++ b/src/daemon/wallet/recovery-probe.ts @@ -0,0 +1,312 @@ +/** + * 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; + }; +}