From 6ff7fe2de37fb02ccfb39b00cd8676caf93fe6e1 Mon Sep 17 00:00:00 2001 From: redshift <213178690+1ftredsh@users.noreply.github.com> Date: Fri, 2 Oct 2026 14:45:21 +0800 Subject: [PATCH] wallet: gate recovery per mint only in degraded mode, guard live ops Review fixes for the mint reachability probe (#116): - Open the per-mint recovery gate only on the degraded path, after the probe and local housekeeping finish. On the happy path every value-moving caller waits for the global sweeps as before, so coco's send/melt recovery never runs beside a live execute. - Skip operations whose per-operation lock is held (diagnostics.isLocked) in the targeted driver, and drive sends only in pending/executing states, so recovery can never race a swap that execute() still holds. Drive sends through recoverExecutingOperation instead of the error-swallowing tryRecover* wrappers. - Run coco's local-only crash cleanup (init operations, orphaned proof reservations) on the degraded path before the gate opens, so a mint that never returns cannot leave that housekeeping undone forever. - Probe with a path-preserving /v1/info join so subpath mints (e.g. https://host/Bitcoin) are checked at their real endpoint, and classify malformed persisted URLs as unreachable. - Drop the 5-minute stuck-operation recheck: it introduced the live-send race, an unbounded shutdown wait, and duplicate polling next to the pending-mint sweep. Parked operations are recovered on the next startup instead. - Count recovery attempts rather than completions, since the driver cannot observe tryRecover*-style swallowed errors. Adds recovery-integration.test.ts with real coco Manager + SQLite regression tests: live execute vs targeted recovery, gate ordering on both paths, and local housekeeping without network. --- src/daemon/wallet/coco-client.ts | 75 ++++---- .../wallet/recovery-integration.test.ts | 174 ++++++++++++++++++ src/daemon/wallet/recovery-probe.test.ts | 110 +++++------ src/daemon/wallet/recovery-probe.ts | 103 ++--------- 4 files changed, 276 insertions(+), 186 deletions(-) create mode 100644 src/daemon/wallet/recovery-integration.test.ts diff --git a/src/daemon/wallet/coco-client.ts b/src/daemon/wallet/coco-client.ts index 60c0f2e..8fae8c9 100644 --- a/src/daemon/wallet/coco-client.ts +++ b/src/daemon/wallet/coco-client.ts @@ -46,7 +46,6 @@ import { collectStuckOperations, probeMintReachability, runTargetedRecovery, - startRecoveryRecheck, type SendRecoveryService, type StuckOperation, } from "./recovery-probe"; @@ -1552,12 +1551,38 @@ function sendRecoveryServiceOf(coco: Manager): SendRecoveryService { .sendOperationService; } +/** Local-only crash cleanup, run before the degraded gate opens. */ +export async function cleanupLocalRecoveryState( + coco: Manager, + repo: SqliteRepositories, +): Promise { + // Coco 1.0.1 implements these as local repository/proof operations only. + // Keep this version-sensitive bridge together with the send recovery bridge. + const services = coco as unknown as Record; + cleanupOrphanedReservations(): Promise; + }>; + const families = [ + ["send", repo.sendOperationRepository], + ["melt", repo.meltOperationRepository], + ["receive", repo.receiveOperationRepository], + ["mint", repo.mintOperationRepository], + ] as const; + for (const [kind, repository] of families) { + for (const op of await repository.getByState("init")) { + await services[`${kind}OperationService`]!.recoverInitOperation(op); + } + } + await services.sendOperationService!.cleanupOrphanedReservations(); +} + /** * 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 + * On degraded startup only, publishStuckMints opens the gate for callers + * whose target mint has no stuck operations after probing and local cleanup. + * A dead mint must not stall spends from a healthy one. On the happy path + * all callers wait until the global sweeps finish. Callers * without a target mint, or whose mint has stuck operations, wait for the * full sweep. fail() poisons every caller; reads are never gated. */ @@ -1628,11 +1653,12 @@ export function createRecoveryGate(): RecoveryGate { * are failed locally so `recoverPendingMintOperations()` skips them, while * paid/issued and unreachable-mint quotes stay pending for the sweep. */ -async function runWalletRecovery( +export async function runWalletRecovery( coco: Manager, onProgress: (progress: RecoveryPhaseProgress) => void, receiveOperationIds?: string[], onStuckMintsKnown?: (mints: Set) => void, + options: { cleanupLocalState?: () => Promise; fetchImpl?: typeof fetch } = {}, ): Promise { surfacingRecoveryProgress = true; let failedMintQuotes = 0; @@ -1666,11 +1692,9 @@ async function runWalletRecovery( // 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))], + { fetchImpl: options.fetchImpl }, ); for (const mintUrl of unreachableMints) { const count = stuckOperations.filter((op) => op.mintUrl === mintUrl).length; @@ -1679,6 +1703,13 @@ async function runWalletRecovery( ); } const degraded = unreachableMints.size > 0; + if (degraded) { + // Local-only housekeeping must finish before any new operation is allowed. + await options.cleanupLocalState?.(); + // Global sweeps enumerate fresh state and are unsafe beside live sends. + // Only the snapshot-based degraded path may open the per-mint gate. + onStuckMintsKnown?.(new Set(stuckOperations.map((op) => op.mintUrl))); + } const targeted = (kinds: Array) => runTargetedRecovery(coco.ops, sendRecoveryServiceOf(coco), { kinds, @@ -1687,9 +1718,10 @@ async function runWalletRecovery( }); // 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. + // clean up init operations and orphaned proof reservations. Degraded + // startup runs that local housekeeping before opening its gate, and only + // drives the previously collected snapshot 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"]); @@ -1790,7 +1822,6 @@ export async function createCocoClient( }; let recoveryResolve: (() => void) | undefined; let stopPendingMintSweep: (() => Promise) | undefined; - let stopRecoveryRecheck: (() => Promise) | undefined; const recoveryPromise = new Promise((resolve) => { recoveryResolve = resolve; }); @@ -1993,6 +2024,7 @@ export async function createCocoClient( }, receiveRecoveryOperationIds, (mints) => recoveryGate.publishStuckMints(mints), + { cleanupLocalState: () => cleanupLocalRecoveryState(coco!, repo) }, ) .then(async () => { await syncReceiveReservations(); @@ -2001,24 +2033,6 @@ export async function createCocoClient( 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, @@ -2344,7 +2358,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-integration.test.ts b/src/daemon/wallet/recovery-integration.test.ts new file mode 100644 index 0000000..59aa30e --- /dev/null +++ b/src/daemon/wallet/recovery-integration.test.ts @@ -0,0 +1,174 @@ +import { expect, it, spyOn } from "bun:test"; +import { Manager } from "@cashu/coco-core"; +import { SqliteRepositories } from "@cashu/coco-sqlite-bun"; +import { Database } from "bun:sqlite"; +import { runTargetedRecovery, type SendRecoveryService } from "./recovery-probe"; + +import { + cleanupLocalRecoveryState, + createRecoveryGate, + runWalletRecovery, +} from "./coco-client"; + +const MINT = "https://mint.example.com"; +const OP = "op-live"; + +it("targeted recovery leaves a send that execute() holds alone", async () => { + const repo = new SqliteRepositories({ database: new Database(":memory:") }); + await repo.init(); + const coco = new Manager(repo, async () => new Uint8Array(64)); + const internals = coco as unknown as { + sendOperationService: SendRecoveryService; + walletService: { getWalletWithActiveKeysetId: (m: string) => Promise }; + }; + + let swapStarted!: () => void; + const started = new Promise((r) => (swapStarted = r)); + let finishSwap!: (v: { send: unknown[]; keep: unknown[] }) => void; + const swap = new Promise<{ send: unknown[]; keep: unknown[] }>((r) => (finishSwap = r)); + internals.walletService.getWalletWithActiveKeysetId = async () => ({ + wallet: { + unit: "sat", + send: async () => (swapStarted(), swap), + checkProofsStates: async () => [{ state: "UNSPENT" }], + getFeesForProofs: () => 0, + }, + }); + + await repo.proofRepository.saveProofs(MINT, [ + { id: "00aa", amount: 8, secret: "in-1", C: "02aa", mintUrl: MINT, state: "ready" } as never, + ]); + await repo.proofRepository.reserveProofs(MINT, ["in-1"], OP); + await repo.sendOperationRepository.create({ + id: OP, mintUrl: MINT, amount: 8, state: "prepared", method: "default", methodData: {}, + createdAt: Date.now(), updatedAt: Date.now(), needsSwap: true, fee: 0, inputAmount: 8, + inputProofSecrets: ["in-1"], + outputData: { + keep: [], + send: [{ + blindedMessage: { amount: 8, id: "00aa", B_: "02" + "11".repeat(32) }, + blindingFactor: "01", + secret: Buffer.from("out-1").toString("hex"), + }], + }, + } as never); + + const live = coco.ops.send.execute(OP); + await started; // swap is at the mint + await runTargetedRecovery(coco.ops, internals.sendOperationService, { + kinds: ["send"], + fetchImpl: (async () => new Response("{}")) as unknown as unknown as typeof fetch, + }); + // On 8005aeb this is "rolled_back" and "in-1" is no longer reserved. + expect((await coco.ops.send.get(OP))?.state).toBe("executing"); + + finishSwap({ send: [{ id: "00aa", amount: 8, secret: "out-1", C: "02bb" }], keep: [] }); + await live; + expect((await coco.ops.send.get(OP))?.state).toBe("pending"); +}); + + +it("keeps healthy-mint callers gated while happy-path global recovery runs", async () => { + const db = new Database(":memory:"); + const repo = new SqliteRepositories({ database: db }); + await repo.init(); + const coco = new Manager(repo, async () => new Uint8Array(64)); + const gate = createRecoveryGate(); + let enter!: () => void; + const entered = new Promise((resolve) => { enter = resolve; }); + let release!: () => void; + const barrier = new Promise((resolve) => { release = resolve; }); + const spies = [ + spyOn(coco.ops.send.recovery, "run").mockImplementation(async () => { enter(); await barrier; }), + spyOn(coco.ops.melt.recovery, "run").mockResolvedValue(undefined), + spyOn(coco.ops.receive.recovery, "run").mockResolvedValue(undefined), + spyOn(coco, "recoverPendingMintOperations").mockResolvedValue(undefined), + ]; + try { + const recovery = runWalletRecovery(coco, () => {}, undefined, + (mints) => gate.publishStuckMints(mints), + { fetchImpl: (async () => new Response("{}")) as unknown as typeof fetch }, + ).then(() => gate.complete()); + await entered; + let released = false; + const waiting = gate.waitForRecovery(MINT).then(() => { released = true; }); + await new Promise((resolve) => setTimeout(resolve, 10)); + expect(released).toBe(false); + release(); + await recovery; + await waiting; + expect(released).toBe(true); + } finally { + release(); + for (const spy of spies) spy.mockRestore(); + db.close(); + } +}); + +it("degraded startup finishes local housekeeping before opening the per-mint gate", async () => { + const db = new Database(":memory:"); + const repo = new SqliteRepositories({ database: db }); + await repo.init(); + const coco = new Manager(repo, async () => new Uint8Array(64)); + await repo.sendOperationRepository.create({ + id: "stuck", mintUrl: MINT, amount: 8, state: "rolling_back", + method: "default", methodData: {}, createdAt: Date.now(), updatedAt: Date.now(), + } as never); + const gate = createRecoveryGate(); + const events: string[] = []; + const sweep = spyOn(coco.ops.send.recovery, "run"); + try { + const recovery = runWalletRecovery(coco, () => {}, [], (mints) => { + events.push("gate"); + gate.publishStuckMints(mints); + }, { + fetchImpl: (async () => { throw new Error("offline"); }) as unknown as unknown as typeof fetch, + cleanupLocalState: async () => { events.push("cleanup"); }, + }).then(() => gate.complete()); + await gate.waitForRecovery("https://healthy.example.com"); + expect(events).toEqual(["cleanup", "gate"]); + await recovery; + expect(sweep).not.toHaveBeenCalled(); + expect((await coco.ops.send.get("stuck"))?.state).toBe("rolling_back"); + } finally { + sweep.mockRestore(); + db.close(); + } +}); + +it("local housekeeping cleans init sends and orphaned reservations without network", async () => { + const db = new Database(":memory:"); + const repo = new SqliteRepositories({ database: db }); + await repo.init(); + const coco = new Manager(repo, async () => new Uint8Array(64)); + try { + await repo.sendOperationRepository.create({ + id: "init-send", mintUrl: MINT, amount: 8, state: "init", method: "default", + methodData: {}, createdAt: Date.now(), updatedAt: Date.now(), + } as never); + for (const [id, repository] of [ + ["init-melt", repo.meltOperationRepository], + ["init-receive", repo.receiveOperationRepository], + ["init-mint", repo.mintOperationRepository], + ] as const) { + await repository.create({ + id, mintUrl: MINT, amount: 8, state: "init", method: "bolt11", + methodData: {}, inputProofs: [], createdAt: Date.now(), updatedAt: Date.now(), + } as never); + } + await repo.proofRepository.saveProofs(MINT, [ + { id: "00aa", amount: 8, secret: "init-input", C: "02aa", mintUrl: MINT, state: "ready" }, + { id: "00aa", amount: 8, secret: "orphan-input", C: "02aa", mintUrl: MINT, state: "ready" }, + ] as never); + await repo.proofRepository.reserveProofs(MINT, ["init-input"], "init-send"); + await repo.proofRepository.reserveProofs(MINT, ["orphan-input"], "missing-send"); + await cleanupLocalRecoveryState(coco, repo); + expect(await coco.ops.send.get("init-send")).toBeNull(); + expect(await repo.proofRepository.getReservedProofs()).toEqual([]); + expect(await repo.meltOperationRepository.getById("init-melt")).toBeNull(); + expect(await repo.receiveOperationRepository.getById("init-receive")).toBeNull(); + expect(await repo.mintOperationRepository.getById("init-mint")).toBeNull(); + } finally { + db.close(); + } +}); diff --git a/src/daemon/wallet/recovery-probe.test.ts b/src/daemon/wallet/recovery-probe.test.ts index 2468cb9..0be1a06 100644 --- a/src/daemon/wallet/recovery-probe.test.ts +++ b/src/daemon/wallet/recovery-probe.test.ts @@ -1,9 +1,8 @@ -import { describe, expect, it, mock } from "bun:test"; +import { describe, expect, it } from "bun:test"; import { collectStuckOperations, probeMintReachability, runTargetedRecovery, - startRecoveryRecheck, type StuckOperationSource, type SendRecoveryService, } from "./recovery-probe"; @@ -40,6 +39,7 @@ function makeSource( mint: stuck.mint ?? [], }; const family = (kind: keyof typeof full) => ({ + diagnostics: { isLocked: () => false }, listInFlight: async () => full[kind] as never, refresh: async (id: string) => { refreshed[kind].push(id); @@ -60,18 +60,12 @@ function makeSource( } 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) => { + recoverExecutingOperation: async (raw) => { executingRecovered.push((raw as FakeOp).id); }, }; @@ -148,7 +142,7 @@ describe("runTargetedRecovery", () => { const result = await runTargetedRecovery(makeSource({}), makeSendService(), { fetchImpl: deadFetch.fetchImpl, }); - expect(result).toMatchObject({ recovered: 0, skipped: 0, failed: 0 }); + expect(result).toMatchObject({ attempted: 0, skipped: 0, failed: 0 }); expect(deadFetch.calls).toHaveLength(0); }); @@ -163,7 +157,7 @@ describe("runTargetedRecovery", () => { fetchImpl: deadFetch.fetchImpl, onSkippedMint: (mintUrl, count) => skippedMints.push([mintUrl, count]), }); - expect(result.recovered).toBe(2); + expect(result.attempted).toBe(2); expect(result.skipped).toBe(2); expect(result.failed).toBe(0); expect(skippedMints).toEqual([["https://dead.example.com", 2]]); @@ -178,17 +172,15 @@ describe("runTargetedRecovery", () => { 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(result.attempted).toBe(2); 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 () => { @@ -198,7 +190,7 @@ describe("runTargetedRecovery", () => { const result = await runTargetedRecovery(source, makeSendService(), { fetchImpl: liveFetch().fetchImpl, }); - expect(result.recovered).toBe(1); + expect(result.attempted).toBe(2); expect(result.failed).toBe(1); expect(source.refreshed.melt).toEqual(["boom-1", "m2"]); }); @@ -221,59 +213,39 @@ describe("runTargetedRecovery", () => { }); }); -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); - }); +it("preserves subpath mint URLs when probing", async () => { + const { fetchImpl, calls } = liveFetch(); + await probeMintReachability(["https://mint.example.com/Bitcoin/"], { fetchImpl }); + expect(calls).toEqual(["https://mint.example.com/Bitcoin/v1/info"]); +}); + +it("does not drive locked operations or count rolling-back sends as attempts", async () => { + const source = makeSource({ send: [ + op("live", "https://mint.example.com", "executing"), + op("rollback", "https://mint.example.com", "rolling_back"), + ] }); + source.send.diagnostics.isLocked = (id) => id === "live"; + const service = makeSendService(); + const result = await runTargetedRecovery(source, service, { fetchImpl: liveFetch().fetchImpl }); + expect(result.attempted).toBe(0); + expect(service.executingRecovered).toEqual([]); +}); + +it("classifies malformed persisted URLs as unreachable without fetching", async () => { + const { fetchImpl, calls } = liveFetch(); + const unreachable = await probeMintReachability(["not-a-url"], { fetchImpl }); + expect([...unreachable]).toEqual(["not-a-url"]); + expect(calls).toEqual([]); +}); + +it("bounds a hanging probe with an abort signal", async () => { + const fetchImpl = (async (_url: unknown, options: RequestInit) => + new Promise((_resolve, reject) => { + options.signal!.addEventListener("abort", () => reject(options.signal!.reason), { once: true }); + })) as unknown as typeof fetch; + const unreachable = await probeMintReachability(["https://slow.example.com"], { + fetchImpl, timeoutMs: 10, + }); + expect([...unreachable]).toEqual(["https://slow.example.com"]); }); diff --git a/src/daemon/wallet/recovery-probe.ts b/src/daemon/wallet/recovery-probe.ts index 952b3ab..9dce35f 100644 --- a/src/daemon/wallet/recovery-probe.ts +++ b/src/daemon/wallet/recovery-probe.ts @@ -7,24 +7,22 @@ * 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. + * them) until a later startup finds the mint reachable. */ 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; + send: Pick; + melt: Pick; + receive: Pick; + mint: Pick; } export type StuckOperationKind = "send" | "melt" | "receive" | "mint"; @@ -46,13 +44,12 @@ export interface StuckOperation { * `mintOperationService` (see coco-client.ts). */ export interface SendRecoveryService { - tryRecoverInitOperation(op: unknown): Promise; - tryRecoverExecutingOperation(op: unknown): Promise; + recoverExecutingOperation(op: unknown): Promise; } export interface RecoveryRunResult { - /** Operations at reachable mints whose recovery completed. */ - recovered: number; + /** Operations for which recovery was attempted (not necessarily completed). */ + attempted: number; /** Operations skipped because their mint did not answer the probe. */ skipped: number; /** Operations at reachable mints whose recovery still failed. */ @@ -89,8 +86,8 @@ function asStuckOperations( raw: op, }); } catch { - // Unparseable mint URL: keep the operation recoverable by treating it as - // reachable (probe only covers successfully normalized URLs). + // Preserve malformed persisted URLs; the probe will classify them as + // unreachable rather than attempting recovery against an invalid URL. stuck.push({ kind, id: op.id, mintUrl: op.mintUrl, state: op.state, raw: op }); } } @@ -135,7 +132,7 @@ export async function probeMintReachability( await Promise.all( mintUrls.map(async (mintUrl) => { try { - await fetcher(new URL("/v1/info", mintUrl).toString(), { + await fetcher(`${normalizeMintUrl(mintUrl)}/v1/info`, { signal: AbortSignal.timeout(timeoutMs), }); } catch (error) { @@ -162,12 +159,8 @@ async function recoverStuckOperation( // 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); + // Startup snapshot only; skip live operations in the driver below. + await sendService.recoverExecutingOperation(op.raw); } // prepared / rolling_back: the global sweep only warns; nothing to do. return; @@ -200,7 +193,7 @@ export async function runTargetedRecovery( options: TargetedRecoveryOptions = {}, ): Promise { const result: RecoveryRunResult = { - recovered: 0, + attempted: 0, skipped: 0, failed: 0, skippedMints: new Map(), @@ -228,9 +221,11 @@ export async function runTargetedRecovery( for (const op of stuck) { if (unreachable.has(op.mintUrl)) continue; + if (source[op.kind].diagnostics.isLocked(op.id)) continue; + if (op.kind === "send" && !["pending", "executing"].includes(op.state)) continue; + result.attempted++; 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. @@ -246,67 +241,3 @@ export async function runTargetedRecovery( 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; - }; -}