diff --git a/src/daemon/http/index.ts b/src/daemon/http/index.ts index 04d15e9..31cae7c 100644 --- a/src/daemon/http/index.ts +++ b/src/daemon/http/index.ts @@ -572,6 +572,19 @@ export function createDaemonRequestHandler(deps: { return; } + if (req.method === "POST" && url.pathname === "/wallet/recover/operations") { + await respond(res, async () => { + if (!deps.walletClient.recoverStuckOperations) { + throw new CocodHttpError( + 501, + "Stuck operation recovery is not supported by this wallet client.", + ); + } + return { output: await deps.walletClient.recoverStuckOperations() }; + }); + return; + } + if (req.method === "POST" && url.pathname === "/wallet/receive/cashu") { await respond(res, async () => { const body = await readJsonBody(req); diff --git a/src/daemon/http/wallet-recovery.test.ts b/src/daemon/http/wallet-recovery.test.ts index b49a984..4c61737 100644 --- a/src/daemon/http/wallet-recovery.test.ts +++ b/src/daemon/http/wallet-recovery.test.ts @@ -47,3 +47,53 @@ describe("POST /wallet/recover validation", () => { expect(recoverMintQuotes).toHaveBeenCalledTimes(1); }); }); + +async function recoverOperations() { + const recoverStuckOperations = mock(async () => ({ + attempted: 1, + busy: 0, + skipped: 0, + failed: 0, + skippedMints: {}, + })); + const handler = createDaemonRequestHandler({ walletClient: { recoverStuckOperations } } as never); + const req = new EventEmitter() as any; + Object.assign(req, { method: "POST", url: "/wallet/recover/operations", headers: { host: "localhost" } }); + const res = { + status: 0, body: "", + writeHead(status: number) { this.status = status; }, + end(chunk: string) { this.body = chunk; }, + }; + setImmediate(() => { + req.emit("data", Buffer.from("{}")); + req.emit("end"); + }); + await handler(req, res as never); + return { res, recoverStuckOperations }; +} + +describe("POST /wallet/recover/operations", () => { + it("drives stuck-operation recovery and returns the summary", async () => { + const { res, recoverStuckOperations } = await recoverOperations(); + expect(res.status).toBe(200); + expect(recoverStuckOperations).toHaveBeenCalledTimes(1); + expect(JSON.parse(res.body).output).toMatchObject({ attempted: 1, busy: 0 }); + }); + + it("returns 501 when the wallet client does not support it", async () => { + const handler = createDaemonRequestHandler({ walletClient: {} } as never); + const req = new EventEmitter() as any; + Object.assign(req, { method: "POST", url: "/wallet/recover/operations", headers: { host: "localhost" } }); + const res = { + status: 0, body: "", + writeHead(status: number) { this.status = status; }, + end(chunk: string) { this.body = chunk; }, + }; + setImmediate(() => { + req.emit("data", Buffer.from("{}")); + req.emit("end"); + }); + await handler(req, res as never); + expect(res.status).toBe(501); + }); +}); diff --git a/src/daemon/wallet/coco-client.test.ts b/src/daemon/wallet/coco-client.test.ts index 5da56a9..b020237 100644 --- a/src/daemon/wallet/coco-client.test.ts +++ b/src/daemon/wallet/coco-client.test.ts @@ -681,6 +681,29 @@ describe("settleExpiredMintQuotes", () => { expect(observePendingOperation).not.toHaveBeenCalled(); expect(failPendingOperation).not.toHaveBeenCalled(); }); + + it("skips quotes at probe-unreachable mints without spending the budget", async () => { + const ops = [ + pendingMintOp({ id: "dead-1", mintUrl: "https://dead.example.com" }), + pendingMintOp({ id: "dead-2", mintUrl: "https://dead.example.com/" }), + pendingMintOp({ id: "live-1", mintUrl: "https://live.example.com" }), + ]; + const { source, observePendingOperation, failPendingOperation } = + fakeSource(ops); + + const result = await settleExpiredMintQuotes( + source, + NOW_MS, + undefined, + { unreachableMints: new Set(["https://dead.example.com"]) }, + ); + + expect(result).toEqual({ failed: 1, leftForRecovery: 0, unobserved: 2 }); + // Only the reachable mint was asked anything. + expect(observePendingOperation).toHaveBeenCalledTimes(1); + expect(observePendingOperation.mock.calls[0]?.[0]).toBe("live-1"); + expect(failPendingOperation).toHaveBeenCalledTimes(1); + }); }); describe("settlePendingMintQuotes", () => { @@ -1537,7 +1560,7 @@ describe("runMintQuoteRecovery", () => { it("skips operations whose earlier recovery is still in flight", async () => { const outstanding = new Map>([ - ["op-1", new Promise(() => {})], + ["mint:op-1", new Promise(() => {})], ]); const { source, finalize, observePendingOperation } = fakeSource( [mintOp()], @@ -1568,14 +1591,14 @@ describe("runMintQuoteRecovery", () => { }); expect(first).toMatchObject({ retryable: 1, recovered: 0 }); - expect(outstanding.has("op-1")).toBe(true); + expect(outstanding.has("mint:op-1")).toBe(true); // The abandoned mint request must not be retried underneath. expect(second).toMatchObject({ busy: 1, checked: 0 }); }); it("does not re-open a failed operation whose recovery is in flight", async () => { const outstanding = new Map>([ - ["op-1", new Promise(() => {})], + ["mint:op-1", new Promise(() => {})], ]); const { source, reopenFailedOperation } = fakeSource( [mintOp({ state: "failed" })], @@ -1630,7 +1653,7 @@ describe("runMintQuoteRecovery", () => { }); expect(first).toMatchObject({ retryable: 1, checked: 1 }); - expect(outstanding.has("op-1")).toBe(true); + expect(outstanding.has("mint:op-1")).toBe(true); expect(second).toMatchObject({ busy: 1, checked: 0 }); expect(finalize).not.toHaveBeenCalled(); }); diff --git a/src/daemon/wallet/coco-client.ts b/src/daemon/wallet/coco-client.ts index 8fae8c9..bbd1f84 100644 --- a/src/daemon/wallet/coco-client.ts +++ b/src/daemon/wallet/coco-client.ts @@ -1,3 +1,4 @@ +import { recoveryKey, trackRecovery, drainRecoveryWork, waitForRecoveryWork, createRecoveryDisposer, type RecoveryWork } from "./recovery-work"; import { Manager, OperationInProgressError, @@ -731,6 +732,7 @@ const EXPIRED_MINT_OBSERVATION_DEADLINE_MS = 15_000; /** Rejects when `timeoutMs` elapses before `promise` settles. */ function withTimeout(promise: Promise, timeoutMs: number): Promise { + if (timeoutMs === Infinity) return promise; let timer: ReturnType | undefined; const timeout = new Promise((_resolve, reject) => { timer = setTimeout( @@ -845,6 +847,7 @@ export async function settleExpiredMintQuotes( source: ExpiredMintQuoteSource, nowMs: number, deadlineMs: number = EXPIRED_MINT_OBSERVATION_DEADLINE_MS, + options: { unreachableMints?: Set; outstanding?: RecoveryWork; shouldStop?: () => boolean } = {}, ): Promise { const pendingMints = await source.ops.mint.listPending(); const selection = selectCleanupOperations({ @@ -865,6 +868,25 @@ export async function settleExpiredMintQuotes( const startedAt = Date.now(); for (const op of candidates) { + if (options.shouldStop?.() || options.outstanding?.has(recoveryKey("mint", op.id))) { + settlement.unobserved++; continue; + } + // A mint the startup probe already found unreachable cannot answer an + // observation either; skip it without spending the shared wall-clock + // budget, leaving the quote pending for a later startup. + if (options.unreachableMints) { + let mintUrl = op.mintUrl; + try { + mintUrl = normalizeMintUrl(op.mintUrl); + } catch { + // Malformed persisted URL: probe keys are raw for those, and the + // observation below would fail anyway, landing in `unobserved`. + } + if (options.unreachableMints.has(mintUrl)) { + settlement.unobserved++; + continue; + } + } const remainingMs = deadlineMs - (Date.now() - startedAt); if (remainingMs <= 0) { const skipped = @@ -879,11 +901,12 @@ export async function settleExpiredMintQuotes( break; } - const check = await failExpiredMintQuoteIfUnpaid( - source.mintOperationService, - op.id, - remainingMs, - ); + // Track the complete observation/failure chain, not just its bounded wait. + const work = failExpiredMintQuoteIfUnpaid(source.mintOperationService, op.id, Infinity); + if (options.outstanding) trackRecovery(options.outstanding, recoveryKey("mint", op.id), work); + const check = await waitForRecoveryWork(work, remainingMs).catch(error => ({ + outcome: "unobserved" as const, error, category: undefined, + })); if (check.outcome === "failed") { settlement.failed++; } else if (check.outcome === "leftForRecovery") { @@ -944,6 +967,7 @@ export interface MintQuoteRecoverySource { } export interface MintQuoteRecoveryOptions { + shouldStop?: () => boolean; /** Target only these operation ids (may include failed operations). */ operationIds?: string[]; /** Per-quote budget for observing the mint and finalizing the operation. */ @@ -955,7 +979,7 @@ export interface MintQuoteRecoveryOptions { */ includeFailed?: boolean; /** - * In-flight recovery work keyed by operation id, shared across runs. + * In-flight recovery work keyed by family:id (mint:), shared across runs. * withTimeout does not cancel the underlying request, so a timed-out quote * check or finalize must keep blocking a retry until it actually settles. */ @@ -994,9 +1018,9 @@ export interface MintQuoteRecoveryResult { * underneath work that outlived its timeout. A rejected task never breaks the * chain for the next one. */ -export function createRunQueue(): (run: () => Promise) => Promise { +export function createRunQueue(): ((run: () => Promise) => Promise) & { drain(): Promise } { let tail: Promise = Promise.resolve(); - return (run: () => Promise): Promise => { + const enqueue = (run: () => Promise): Promise => { const result = tail.then(run, run); tail = result.then( () => undefined, @@ -1004,6 +1028,7 @@ export function createRunQueue(): (run: () => Promise) => Promise { ); return result; }; + return Object.assign(enqueue, { drain: async () => { await tail; } }); } /** Per-quote budget for the mint round-trip during explicit recovery. */ @@ -1071,21 +1096,8 @@ export async function runMintQuoteRecovery( /** coco's fail-fast operation lock rejected the call: another holder exists. */ const isInProgress = (error: unknown) => error instanceof Error && error.name === "OperationInProgressError"; - /** - * Register in-flight work for an operation. Entries are cleared only once the - * work actually settles (withTimeout does not cancel the request behind it), - * so a timed-out call keeps blocking a retry. The identity check stops a late - * settlement from clearing a newer entry for the same operation. - */ - const track = (operationId: string, work: Promise) => { - outstanding.set(operationId, work); - const clear = () => { - if (outstanding.get(operationId) === work) { - outstanding.delete(operationId); - } - }; - void work.then(clear, clear); - }; + const track = (operationId: string, work: Promise) => + trackRecovery(outstanding, recoveryKey("mint", operationId), work); let targets: MintQuoteRecoveryCandidate[]; if (options.operationIds && options.operationIds.length > 0) { @@ -1115,8 +1127,9 @@ export async function runMintQuoteRecovery( }); for (const op of failed) { + if (options.shouldStop?.()) break; const label = `Mint quote ${op.quoteId ?? op.id} at ${op.mintUrl}`; - if (outstanding.has(op.id)) { + if (outstanding.has(recoveryKey("mint", op.id))) { result.busy++; onProgress?.(`${label}: an earlier recovery is still running; skipped`); continue; @@ -1144,13 +1157,16 @@ export async function runMintQuoteRecovery( await recoverOne(op); } - for (const op of pending) await recoverOne(op); + for (const op of pending) { + if (options.shouldStop?.()) break; + await recoverOne(op); + } return result; async function recoverOne(op: MintQuoteRecoveryCandidate): Promise { const label = `Mint quote ${op.quoteId ?? op.id} at ${op.mintUrl}`; - if (outstanding.has(op.id)) { + if (outstanding.has(recoveryKey("mint", op.id))) { result.busy++; onProgress?.(`${label}: an earlier recovery is still running; skipped`); return; @@ -1305,6 +1321,7 @@ export interface PendingMintSweepOptions { deadlineMs?: number; checkTimeoutMs?: number; state?: PendingMintSweepState; + shouldStop?: () => boolean; } type PendingMintOutcome = "unreachable" | "other"; @@ -1326,21 +1343,21 @@ export async function settlePendingMintQuotes( const ordered = [...pending.slice(resumeAt), ...pending.slice(0, resumeAt)]; const startedAt = Date.now(); for (const op of ordered) { + if (options.shouldStop?.()) break; const remainingMs = deadlineMs - (Date.now() - startedAt); - if (state.outstanding.has(op.id) || remainingMs <= 0) { + if (state.outstanding.has(recoveryKey("mint", op.id)) || remainingMs <= 0) { unreachable++; continue; } state.after = op.id; - // Report on the refresh itself so a late result is still logged. + // Track refresh AND its late-result reporting mutations as one lifetime. const settled = source.ops.mint .refresh(op.id) .then( (result) => reportPendingMintRefresh(source, op, result, nowMs), (error) => reportPendingMintRefreshError(source, op, error), - ) - .finally(() => state.outstanding.delete(op.id)); - state.outstanding.set(op.id, settled); + ); + trackRecovery(state.outstanding, recoveryKey("mint", op.id), settled); try { const outcome = await withTimeout(settled, Math.min(checkTimeoutMs, remainingMs)); if (outcome === "unreachable") unreachable++; @@ -1416,16 +1433,16 @@ async function reportPendingMintRefreshError( * still running after that fails against the closed database and is picked * up by startup recovery. */ -function startPendingMintSweep(source: PendingMintQuoteSource): () => Promise { +function startPendingMintSweep(source: PendingMintQuoteSource, outstanding: RecoveryWork): () => Promise { let stopped = false; let timer: ReturnType | undefined; let inFlight: Promise = Promise.resolve(); let unreachableBefore = 0; - const state: PendingMintSweepState = { outstanding: new Map() }; + const state: PendingMintSweepState = { outstanding }; const tick = async () => { if (stopped) return; - inFlight = settlePendingMintQuotes(source, Date.now(), { state }).then( + inFlight = settlePendingMintQuotes(source, Date.now(), { state, shouldStop: () => stopped }).then( ({ unreachable }) => { // Report a mint becoming unreachable, or reachable again, once. if (unreachable > 0 && unreachableBefore === 0) { @@ -1559,21 +1576,37 @@ export async function cleanupLocalRecoveryState( // 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; - }>; + recoverInitOperation?(op: unknown): Promise; + cleanupOrphanedReservations?(): Promise; + } | undefined>; 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); + // Fail closed the way reopenFailedMintOperation does: a coco bump that + // removes or renames these privates must stop recovery with a clear error + // before anything is written, not crash halfway through the loop with the + // cleanup half-applied. + for (const [kind] of families) { + if (typeof services[`${kind}OperationService`]?.recoverInitOperation !== "function") { + throw new Error( + `coco ${kind}OperationService.recoverInitOperation is unavailable; refusing local recovery cleanup`, + ); } } - await services.sendOperationService!.cleanupOrphanedReservations(); + if (typeof services.sendOperationService?.cleanupOrphanedReservations !== "function") { + throw new Error( + "coco sendOperationService.cleanupOrphanedReservations is unavailable; refusing local recovery cleanup", + ); + } + for (const [kind, repository] of families) { + for (const op of await repository.getByState("init")) { + await services[`${kind}OperationService`]!.recoverInitOperation!(op); + } + } + await services.sendOperationService!.cleanupOrphanedReservations!(); } /** @@ -1649,7 +1682,7 @@ export function createRecoveryGate(): RecoveryGate { /** * Run the wallet recovery sweeps in order, reporting phase changes. * - * Expired mint quotes are settled first: quotes their mint confirms as unpaid + * Probe first, then settle expired mint quotes: quotes confirmed as unpaid * are failed locally so `recoverPendingMintOperations()` skips them, while * paid/issued and unreachable-mint quotes stay pending for the sweep. */ @@ -1658,38 +1691,16 @@ export async function runWalletRecovery( onProgress: (progress: RecoveryPhaseProgress) => void, receiveOperationIds?: string[], onStuckMintsKnown?: (mints: Set) => void, - options: { cleanupLocalState?: () => Promise; fetchImpl?: typeof fetch } = {}, + options: { cleanupLocalState?: () => Promise; fetchImpl?: typeof fetch; outstanding?: RecoveryWork; shouldStop?: () => boolean } = {}, ): Promise { surfacingRecoveryProgress = true; let failedMintQuotes = 0; try { - onProgress({ phase: "Settling expired mint quotes", failedMintQuotes }); - const settlement = await settleExpiredMintQuotes( - { - ops: coco.ops, - mintOperationService: ( - coco as unknown as { - mintOperationService: MintOperationServiceCleanup; - } - ).mintOperationService, - }, - Date.now(), - ); - failedMintQuotes = settlement.failed; - if (settlement.leftForRecovery > 0 || settlement.unobserved > 0) { - startupProgress( - `Expired mint quotes: ${settlement.failed} failed locally, ` + - `${settlement.leftForRecovery} paid/issued (kept for recovery), ` + - `${settlement.unobserved} unverifiable (kept pending).`, - ); - } - 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. + // Probe every mint that has stuck operations once, FIRST, so a dead mint + // costs a single short probe instead of taxing settlement's observation + // budget plus 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); const unreachableMints = await probeMintReachability( @@ -1710,9 +1721,39 @@ export async function runWalletRecovery( // Only the snapshot-based degraded path may open the per-mint gate. onStuckMintsKnown?.(new Set(stuckOperations.map((op) => op.mintUrl))); } + + // Settlement runs after the gate opens and only spends its observation + // budget on mints the probe found reachable; dead-mint quotes stay + // pending untouched. It only reads and locally fails long-expired quotes, + // so it cannot conflict with live operations the gate just admitted. + onProgress({ phase: "Settling expired mint quotes", failedMintQuotes }); + const settlement = await settleExpiredMintQuotes( + { + ops: coco.ops, + mintOperationService: ( + coco as unknown as { + mintOperationService: MintOperationServiceCleanup; + } + ).mintOperationService, + }, + Date.now(), + undefined, + { unreachableMints, outstanding: options.outstanding, shouldStop: options.shouldStop }, + ); + failedMintQuotes = settlement.failed; + if (settlement.leftForRecovery > 0 || settlement.unobserved > 0) { + startupProgress( + `Expired mint quotes: ${settlement.failed} failed locally, ` + + `${settlement.leftForRecovery} paid/issued (kept for recovery), ` + + `${settlement.unobserved} unverifiable (kept pending).`, + ); + } + onProgress({ phase: "Settled expired mint quotes", failedMintQuotes }); const targeted = (kinds: Array) => runTargetedRecovery(coco.ops, sendRecoveryServiceOf(coco), { kinds, + outstanding: options.outstanding, + shouldStop: options.shouldStop, stuckOperations, unreachableMints, }); @@ -1743,10 +1784,13 @@ export async function runWalletRecovery( .map((op) => [op.id, op.mintUrl]), ); for (const operationId of receiveOperationIds) { + if (options.shouldStop?.()) break; const mintUrl = mintByOperation.get(operationId); if (mintUrl && unreachableMints.has(mintUrl)) continue; try { - await withTimeout(coco.ops.receive.refresh(operationId), 15_000); + const work = coco.ops.receive.refresh(operationId); + if (options.outstanding) trackRecovery(options.outstanding, recoveryKey("receive", operationId), work); + await withTimeout(work, 15_000); } catch (error) { logger.warn("Targeted receive recovery did not complete", { operationId, @@ -1761,7 +1805,11 @@ export async function runWalletRecovery( } onProgress({ phase: "Mint recovery", failedMintQuotes }); - if (!degraded) await coco.recoverPendingMintOperations(); + // A settlement wait may have timed out while an unlocked observation + // still runs. Never let a fresh global mint sweep observe it again. + if (!degraded && ![...(options.outstanding?.keys() ?? [])].some(key => key.startsWith("mint:"))) { + await coco.recoverPendingMintOperations(); + } else await targeted(["mint"]); onProgress({ phase: "done", failedMintQuotes }); @@ -1826,6 +1874,9 @@ export async function createCocoClient( recoveryResolve = resolve; }); const recoveryGate = createRecoveryGate(); + let disposed = false; + const enqueueRecovery = createRunQueue(); + const recoveryOutstanding: RecoveryWork = new Map(); try { startupProgress("Opening Cashu wallet database..."); @@ -2024,7 +2075,7 @@ export async function createCocoClient( }, receiveRecoveryOperationIds, (mints) => recoveryGate.publishStuckMints(mints), - { cleanupLocalState: () => cleanupLocalRecoveryState(coco!, repo) }, + { cleanupLocalState: () => cleanupLocalRecoveryState(coco!, repo), outstanding: recoveryOutstanding, shouldStop: () => disposed }, ) .then(async () => { await syncReceiveReservations(); @@ -2039,7 +2090,7 @@ export async function createCocoClient( mintOperationService: ( coco as unknown as { mintOperationService: MintOperationServiceCleanup } ).mintOperationService, - }); + }, recoveryOutstanding); }) .catch((error) => { recoveryDone = true; @@ -2068,17 +2119,32 @@ export async function createCocoClient( return api; }; - let disposed = false; - // Explicit recovery runs are serialized, and finalize work that outlives its - // timeout stays in the map so a retry waits for it. - const enqueueRecovery = createRunQueue(); - const recoveryOutstanding = new Map>(); + const assertOpen = () => { if (disposed) throw new Error("Wallet is shutting down"); }; // 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); + const waitForRecovery = async (mintUrl?: string): Promise => { + assertOpen(); + await recoveryGate.waitForRecovery(mintUrl); + assertOpen(); + }; + + const disposeRecovery = createRecoveryDisposer( + () => { disposed = true; }, + async () => { + await recoveryPromise; + await stopPendingMintSweep?.(); + await enqueueRecovery.drain(); + await drainRecoveryWork(recoveryOutstanding); + }, + async () => { + await coco.dispose(); + database.close(); + releaseLegacyPidClaim(); + releaseWalletPidClaim(); + }, + ); return { async ping(): Promise { @@ -2356,22 +2422,10 @@ export async function createCocoClient( }, async dispose(): Promise { - if (disposed) return; - disposed = true; - try { - // Let any in-flight recovery settle before closing the database from - // underneath it. The recovery promise resolves on success or failure. - await recoveryPromise; - await stopPendingMintSweep?.(); - await coco.dispose(); - } finally { - try { - database.close(); - } finally { - releaseLegacyPidClaim(); - releaseWalletPidClaim(); - } - } + await disposeRecovery().catch(error => { + logger.warn("Wallet shutdown incomplete; database and ownership retained until recovery settles"); + throw error; + }); }, async getHistory(offset?: number, limit?: number): Promise { @@ -2566,18 +2620,45 @@ export async function createCocoClient( // Serialize explicit recovery: two concurrent requests must not both // snapshot the same failed operation, and a retry must not start // underneath a finalize that outlived its timeout. - return enqueueRecovery(() => - runMintQuoteRecovery( + return enqueueRecovery(() => { + assertOpen(); + return runMintQuoteRecovery( { ops: coco.ops as unknown as MintQuoteRecoverySource["ops"], mintOperationService: service, reopenFailedOperation: (operationId) => reopenFailedMintOperation(service, operationId), }, - { ...options, outstanding: recoveryOutstanding }, + { ...options, outstanding: recoveryOutstanding, shouldStop: () => disposed }, onProgress, - ), - ); + ); + }); + }, + + async recoverStuckOperations() { + await waitForRecovery(); + // Serialized against explicit mint-quote recovery (and itself) through + // the same queue and lifetime tracker, so timed-out passes cannot retry the same + // operation. Receive stays startup-only: recovering competing receives + // safely requires the startup dedup classification (receive-dedup.ts). + // Operations a live execute holds come back as busy via coco's + // fail-fast operation lock, never driven underneath it. + const result = await enqueueRecovery(() => { + assertOpen(); + return runTargetedRecovery(coco!.ops, sendRecoveryServiceOf(coco!), { + kinds: ["send", "melt", "mint"], + outstanding: recoveryOutstanding, + shouldStop: () => disposed, + }); + }); + return { + timedOut: result.timedOut, + attempted: result.attempted, + busy: result.busy, + skipped: result.skipped, + failed: result.failed, + skippedMints: Object.fromEntries(result.skippedMints), + }; }, }; } diff --git a/src/daemon/wallet/cocod-client.ts b/src/daemon/wallet/cocod-client.ts index 2fab2c6..93d4d49 100644 --- a/src/daemon/wallet/cocod-client.ts +++ b/src/daemon/wallet/cocod-client.ts @@ -159,6 +159,22 @@ export interface WalletMintQuoteRecoveryResult { errors: Array<{ operationId: string; error: string }>; } +/** Summary of a stuck-operation (send/melt/mint) recovery run. */ +export interface WalletStuckOperationRecoveryResult { + /** Timed-out waits; the underlying operation remains tracked. */ + timedOut: number; + /** Operations for which recovery was attempted (not necessarily completed). */ + attempted: number; + /** Locked operations or unfinished work from another pass; retry later. */ + busy: number; + /** Operations skipped for unreachable mints, shutdown, or pass budget exhaustion. */ + skipped: number; + /** Operations at reachable mints whose recovery still failed. */ + failed: number; + /** Unreachable mint URL -> number of operations skipped there. */ + skippedMints: Record; +} + export interface CocodClient { ping(): Promise; getStatus(): Promise; @@ -201,6 +217,12 @@ export interface CocodClient { options?: WalletMintQuoteRecoveryOptions, onProgress?: (message: string) => void, ): Promise; + /** + * Recover stuck send/melt/mint operations whose mints answer a + * reachability probe. Operations a live execute holds are reported busy, + * never driven. Receive stays startup-only (receive dedup classification). + */ + recoverStuckOperations?(): Promise; /** Report background wallet recovery progress, when the wallet supports it. */ getRecoveryProgress?(): Promise; } diff --git a/src/daemon/wallet/mint-quote-recovery.fake-mint.test.ts b/src/daemon/wallet/mint-quote-recovery.fake-mint.test.ts index 2554665..69b0c06 100644 --- a/src/daemon/wallet/mint-quote-recovery.fake-mint.test.ts +++ b/src/daemon/wallet/mint-quote-recovery.fake-mint.test.ts @@ -363,7 +363,7 @@ describe("PAID mint quote recovery with a real Manager and mint", () => { outstanding, })) as unknown as Record; expect(first).toMatchObject({ retryable: 1, recovered: 0 }); - expect(outstanding.has(op.id as string)).toBe(true); + expect(outstanding.has(`mint:${op.id}`)).toBe(true); const second = (await runMintQuoteRecovery(booted.source() as never, { outstanding, diff --git a/src/daemon/wallet/recovery-integration.test.ts b/src/daemon/wallet/recovery-integration.test.ts index 59aa30e..f279f93 100644 --- a/src/daemon/wallet/recovery-integration.test.ts +++ b/src/daemon/wallet/recovery-integration.test.ts @@ -55,12 +55,15 @@ it("targeted recovery leaves a send that execute() holds alone", async () => { const live = coco.ops.send.execute(OP); await started; // swap is at the mint - await runTargetedRecovery(coco.ops, internals.sendOperationService, { + const result = 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"); + // The live execute holds coco's per-operation lock: busy, never driven. + expect(result.attempted).toBe(0); + expect(result.busy).toBe(1); finishSwap({ send: [{ id: "00aa", amount: 8, secret: "out-1", C: "02bb" }], keep: [] }); await live; @@ -136,6 +139,65 @@ it("degraded startup finishes local housekeeping before opening the per-mint gat } }); +it("degraded startup probes before settlement and never asks a dead mint", 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 expiredRow = (id: string, mintUrl: string) => ({ + id, mintUrl, quoteId: `q-${id}`, state: "pending", + createdAt: 1_000, updatedAt: 2_000, method: "bolt11", + methodData: { method: "bolt11", data: {} }, amount: 100, unit: "sat", + request: "lnbc1example", expiry: 1_000_000, // epoch seconds, long past + }); + await repo.mintOperationRepository.create(expiredRow("dead-quote", "https://dead.example.com") as never); + await repo.mintOperationRepository.create(expiredRow("live-quote", "https://live.example.com") as never); + const internals = coco as unknown as { + mintOperationService: { + observePendingOperation(id: string): Promise<{ category: "waiting" }>; + failPendingOperation(op: unknown, failure: unknown): Promise; + }; + }; + const events: string[] = []; + const observe = spyOn(internals.mintOperationService, "observePendingOperation") + .mockImplementation(async (id: string) => { + events.push(`observe:${id}`); + return { category: "waiting" }; + }); + // Fail the row for real (as the production service would) so the targeted + // mint pass sees state "failed" and refresh() returns without re-observing. + const fail = spyOn(internals.mintOperationService, "failPendingOperation") + .mockImplementation(async (op: unknown) => { + const id = (op as { id: string }).id; + const row = await repo.mintOperationRepository.getById(id); + if (row) await repo.mintOperationRepository.update({ ...row, state: "failed" } as never); + return {} as never; + }); + const gate = createRecoveryGate(); + try { + const recovery = runWalletRecovery(coco, () => {}, [], (mints) => { + events.push("gate"); + gate.publishStuckMints(mints); + }, { + fetchImpl: (async (url: unknown) => { + if (String(url).includes("dead.example.com")) throw new Error("offline"); + return new Response("{}"); + }) as unknown as typeof fetch, + cleanupLocalState: async () => { events.push("cleanup"); }, + }).then(() => gate.complete()); + await recovery; + // The dead mint was probed, never asked to observe its quote; the gate + // opened before settlement spent anything on the live mint. + expect(events).toEqual(["cleanup", "gate", "observe:live-quote"]); + expect(observe).toHaveBeenCalledTimes(1); + expect(fail).toHaveBeenCalledTimes(1); + } finally { + observe.mockRestore(); + fail.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 }); @@ -172,3 +234,24 @@ it("local housekeeping cleans init sends and orphaned reservations without netwo db.close(); } }); + +it("cleanup checks all private methods before making any local writes", 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 internals = coco as unknown as { mintOperationService: { recoverInitOperation: unknown } }; + const original = internals.mintOperationService.recoverInitOperation; + try { + await repo.sendOperationRepository.create({ + id: "untouched-init", mintUrl: MINT, amount: 8, state: "init", method: "default", + methodData: {}, createdAt: Date.now(), updatedAt: Date.now(), + } as never); + internals.mintOperationService.recoverInitOperation = undefined; + await expect(cleanupLocalRecoveryState(coco, repo)).rejects.toThrow("mintOperationService.recoverInitOperation is unavailable"); + expect((await repo.sendOperationRepository.getById("untouched-init"))?.state).toBe("init"); + } finally { + internals.mintOperationService.recoverInitOperation = original; + db.close(); + } +}); diff --git a/src/daemon/wallet/recovery-probe.test.ts b/src/daemon/wallet/recovery-probe.test.ts index 0be1a06..18b7ca8 100644 --- a/src/daemon/wallet/recovery-probe.test.ts +++ b/src/daemon/wallet/recovery-probe.test.ts @@ -7,6 +7,8 @@ import { type SendRecoveryService, } from "./recovery-probe"; +import { runMintQuoteRecovery, settlePendingMintQuotes, settleExpiredMintQuotes } from "./coco-client"; + interface FakeOp { id: string; mintUrl: string; @@ -52,7 +54,11 @@ function makeSource( return { stuck: full, refreshed, - send: family("send") as FakeSource["send"], + send: { + ...family("send"), + get: async (id: string) => + (full.send.find((o) => o.id === id) ?? null) as never, + } as FakeSource["send"], melt: family("melt") as FakeSource["melt"], receive: family("receive") as FakeSource["receive"], mint: family("mint") as FakeSource["mint"], @@ -61,12 +67,24 @@ function makeSource( function makeSendService(): SendRecoveryService & { executingRecovered: string[]; + /** Lock events in order, e.g. "acquire:op-1", "release:op-1". */ + lockLog: string[]; } { const executingRecovered: string[] = []; + const lockLog: string[] = []; return { executingRecovered, + lockLog, + acquireOperationLock: async (id) => { + lockLog.push(`acquire:${id}`); + return () => { + lockLog.push(`release:${id}`); + }; + }, recoverExecutingOperation: async (raw) => { - executingRecovered.push((raw as FakeOp).id); + const id = (raw as FakeOp).id; + executingRecovered.push(id); + lockLog.push(`recover:${id}`); }, }; } @@ -181,6 +199,54 @@ describe("runTargetedRecovery", () => { expect(result.attempted).toBe(2); expect(source.refreshed.send).toEqual(["pending-1"]); expect(sendService.executingRecovered).toEqual(["exec-1"]); + // The executing send is driven under coco's per-operation lock. + expect(sendService.lockLog).toEqual([ + "acquire:exec-1", + "recover:exec-1", + "release:exec-1", + ]); + }); + + it("re-reads state under the lock and skips a send that left executing", async () => { + const source = makeSource({ + send: [op("s1", "https://live.example.com", "pending")], + }); + const sendService = makeSendService(); + // Snapshot taken while the op was still executing; it has since settled. + const result = await runTargetedRecovery(source, sendService, { + stuckOperations: [ + { + kind: "send", + id: "s1", + mintUrl: "https://live.example.com", + state: "executing", + raw: op("s1", "https://live.example.com", "executing"), + }, + ], + unreachableMints: new Set(), + }); + expect(result.attempted).toBe(1); + expect(sendService.executingRecovered).toEqual([]); + // The lock is still acquired and released around the re-read. + expect(sendService.lockLog).toEqual(["acquire:s1", "release:s1"]); + }); + + it("counts a send as busy when a live execute wins the lock race", async () => { + const source = makeSource({ + send: [op("s1", "https://live.example.com", "executing")], + }); + const sendService = makeSendService(); + sendService.acquireOperationLock = async () => { + const error = new Error("operation in progress"); + error.name = "OperationInProgressError"; + throw error; + }; + const result = await runTargetedRecovery(source, sendService, { + fetchImpl: liveFetch().fetchImpl, + }); + expect(result.busy).toBe(1); + expect(result.failed).toBe(0); + expect(sendService.executingRecovered).toEqual([]); }); it("counts per-operation failures at reachable mints and continues", async () => { @@ -229,7 +295,9 @@ it("does not drive locked operations or count rolling-back sends as attempts", a const service = makeSendService(); const result = await runTargetedRecovery(source, service, { fetchImpl: liveFetch().fetchImpl }); expect(result.attempted).toBe(0); + expect(result.busy).toBe(1); expect(service.executingRecovered).toEqual([]); + expect(service.lockLog).toEqual([]); }); it("classifies malformed persisted URLs as unreachable without fetching", async () => { @@ -249,3 +317,132 @@ it("bounds a hanging probe with an abort signal", async () => { }); expect([...unreachable]).toEqual(["https://slow.example.com"]); }); + + +it("shares timed-out mint observations across quote recovery, targeted recovery and the sweep", async () => { + const source = makeSource({ mint: [op("q1", "https://live.example.com")] }); + let finish!: (value: { category: "waiting" }) => void; + const observation = new Promise<{ category: "waiting" }>(r => { finish = r; }); + const outstanding = new Map>(); + const quoteSource = { + ops: { mint: { + listPending: async () => [{ ...source.stuck.mint[0]!, method: "bolt11" }], + get: async () => null, + finalize: async () => ({ state: "finalized" }), + } }, + mintOperationService: { observePendingOperation: () => observation }, + reopenFailedOperation: async () => false, + }; + try { + await runMintQuoteRecovery(quoteSource as never, { outstanding, timeoutMs: 5 }); + expect(outstanding.has("mint:q1")).toBe(true); + const result = await runTargetedRecovery(source, makeSendService(), { + outstanding, fetchImpl: liveFetch().fetchImpl, + }); + expect(result.busy).toBe(1); + await settlePendingMintQuotes({ + ops: { mint: { ...source.mint, listPending: async () => source.stuck.mint as never } }, + wallet: { balances: { byMint: async () => ({}) } }, + mintOperationService: { failPendingOperation: async () => ({}) }, + } as never, Date.now(), { state: { outstanding } }); + expect(source.refreshed.mint).toEqual([]); + } finally { finish({ category: "waiting" }); await observation; } +}); + +it("tracks a hung targeted mint so quote recovery skips it while unrelated operations proceed", async () => { + const source = makeSource({ mint: [op("q1", "https://live.example.com")], melt: [op("m1", "https://live.example.com")] }); + let finish!: (value: never) => void; + source.mint.refresh = () => new Promise(r => { finish = r; }); + const outstanding = new Map>(); + try { + const result = await runTargetedRecovery(source, makeSendService(), { + outstanding, timeoutMs: 5, fetchImpl: liveFetch().fetchImpl, + }); + expect(result.timedOut).toBe(1); + expect(source.refreshed.melt).toEqual(["m1"]); + const quote = await runMintQuoteRecovery({ + ops: { mint: { listPending: async () => [{ ...source.stuck.mint[0]!, method: "bolt11" }] } }, + } as never, { outstanding }); + expect(quote.busy).toBe(1); + } finally { finish({ state: "pending" } as never); } +}); + +it("retains executing-send lock until a timed-out underlying drive actually finishes", async () => { + const source = makeSource({ send: [op("s1", "https://live.example.com", "executing")] }); + const service = makeSendService(); + let finish!: () => void; + service.recoverExecutingOperation = () => new Promise(r => { finish = r; }); + const outstanding = new Map>(); + const result = await runTargetedRecovery(source, service, { + outstanding, timeoutMs: 5, fetchImpl: liveFetch().fetchImpl, + }); + expect(result.timedOut).toBe(1); + expect(service.lockLog).toEqual(["acquire:s1"]); + const retry = await runTargetedRecovery(source, service, { outstanding, fetchImpl: liveFetch().fetchImpl }); + expect(retry.busy).toBe(1); + const actual = outstanding.get("send:s1")!; + finish(); + await actual; + expect(service.lockLog).toEqual(["acquire:s1", "release:s1"]); + expect(outstanding.size).toBe(0); +}); + +it("refuses missing executing-send internals without driving the operation", async () => { + const source = makeSource({ send: [op("s1", "https://live.example.com", "executing")] }); + const result = await runTargetedRecovery(source, {} as SendRecoveryService, { fetchImpl: liveFetch().fetchImpl }); + expect(result.failed).toBe(1); +}); + + +it("startup settlement keeps its timed-out observation visible to manual recovery", async () => { + const source = makeSource({ mint: [op("q1", "https://live.example.com")] }); + let finish!: (value: { category: "ready" }) => void; + const observation = new Promise<{ category: "ready" }>(r => { finish = r; }); + const outstanding = new Map>(); + const settlement = await settleExpiredMintQuotes({ + ops: { mint: { listPending: async () => [{ ...source.stuck.mint[0]!, expiry: 1, updatedAt: 1 }] } }, + mintOperationService: { observePendingOperation: () => observation }, + } as never, Date.now(), 5, { outstanding }); + expect(settlement.unobserved).toBe(1); + const result = await runTargetedRecovery(source, makeSendService(), { outstanding, fetchImpl: liveFetch().fetchImpl }); + expect(result.busy).toBe(1); + expect(source.refreshed.mint).toEqual([]); + const actual = outstanding.get("mint:q1")!; + finish({ category: "ready" }); + await actual; + expect(outstanding.size).toBe(0); +}); + +it("deadline and shutdown leave remaining operations untouched", async () => { + const source = makeSource({ mint: [op("q1", "https://live.example.com")] }); + const result = await runTargetedRecovery(source, makeSendService(), { + shouldStop: () => true, fetchImpl: liveFetch().fetchImpl, + }); + expect(result.skipped).toBe(1); + expect(result.attempted).toBe(0); + expect(source.refreshed.mint).toEqual([]); +}); + +it("a pass budget bounds total drive waits, not just each operation", async () => { + const source = makeSource({ mint: [op("q1", "https://live.example.com"), op("q2", "https://live.example.com")] }); + let finish!: (value: never) => void; + source.mint.refresh = () => new Promise(r => { finish = r; }); + const outstanding = new Map>(); + const result = await runTargetedRecovery(source, makeSendService(), { + outstanding, timeoutMs: 100, deadlineMs: 10, fetchImpl: liveFetch().fetchImpl, + }); + expect(result.timedOut).toBe(1); + expect(result.attempted).toBe(1); + expect(result.skipped).toBe(1); + const actual = outstanding.get("mint:q1")!; + finish({ state: "pending" } as never); + await actual; +}); + +it("unfinished work in another operation family does not block the same bare id", async () => { + const source = makeSource({ mint: [op("same-id", "https://live.example.com")] }); + const outstanding = new Map>([["send:same-id", new Promise(() => {})]]); + const result = await runTargetedRecovery(source, makeSendService(), { outstanding, fetchImpl: liveFetch().fetchImpl }); + expect(result.busy).toBe(0); + expect(source.refreshed.mint).toEqual(["same-id"]); +}); diff --git a/src/daemon/wallet/recovery-probe.ts b/src/daemon/wallet/recovery-probe.ts index 9dce35f..3dddfaa 100644 --- a/src/daemon/wallet/recovery-probe.ts +++ b/src/daemon/wallet/recovery-probe.ts @@ -7,9 +7,11 @@ * 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 a later startup finds the mint reachable. + * them) until an explicit mid-session recovery (see recoverStuckOperations + * in coco-client.ts) or a later startup finds the mint reachable. */ import { normalizeMintUrl, type Manager } from "@cashu/coco-core"; +import { recoveryKey, trackRecovery, waitForRecoveryWork, RecoveryWaitTimeout, type RecoveryWork } from "./recovery-work"; import { logger } from "../../utils/logger"; /** Short probe: a mint that cannot answer /v1/info in 2s slows every op. */ @@ -19,7 +21,7 @@ type OpsApi = Manager["ops"]; /** Structural subset of the ops APIs used to enumerate and recover operations. */ export interface StuckOperationSource { - send: Pick; + send: Pick; melt: Pick; receive: Pick; mint: Pick; @@ -45,20 +47,37 @@ export interface StuckOperation { */ export interface SendRecoveryService { recoverExecutingOperation(op: unknown): Promise; + /** + * coco's per-operation lock. It is fail-fast: acquiring an id a live + * execute/finalize/recover already holds throws OperationInProgressError + * instead of waiting. Holding it across the state re-read and the drive + * makes executing-send recovery atomic against a live execute — the same + * pattern reopenFailedMintOperation uses for mint operations + * (see coco-client.ts). + */ + acquireOperationLock(operationId: string): Promise<() => void>; } export interface RecoveryRunResult { /** Operations for which recovery was attempted (not necessarily completed). */ attempted: number; - /** Operations skipped because their mint did not answer the probe. */ + /** Timed-out waits, also included in attempted; underlying work remains tracked. */ + timedOut: number; + /** Locked operations or unfinished work from another pass; retry later. */ + busy: number; + /** Operations skipped for unreachable mints, shutdown, or pass budget exhaustion. */ skipped: number; - /** Operations at reachable mints whose recovery still failed. */ + /** Attempts that threw a non-busy, non-timeout error. */ failed: number; /** Unreachable mint URL -> number of operations skipped there. */ skippedMints: Map; } export interface TargetedRecoveryOptions { + outstanding?: RecoveryWork; + timeoutMs?: number; + deadlineMs?: number; + shouldStop?: () => boolean; probeTimeoutMs?: number; fetchImpl?: typeof fetch; /** Operation families to recover. Defaults to all four. */ @@ -159,8 +178,31 @@ async function recoverStuckOperation( // Public API: actively re-checks the proofs with the mint. await source.send.refresh(op.id); } else if (op.state === "executing") { - // Startup snapshot only; skip live operations in the driver below. - await sendService.recoverExecutingOperation(op.raw); + // Fail closed: the lock and the drive are private coco internals + // (written against coco-core 1.0.1), so a coco bump that renames them + // must fail this operation loudly instead of corrupting it. + if ( + typeof sendService.acquireOperationLock !== "function" || + typeof sendService.recoverExecutingOperation !== "function" + ) { + throw new Error( + "coco sendOperationService recovery internals are unavailable; refusing to recover an executing send", + ); + } + // recoverExecutingOperation takes no lock and does no state re-read, + // so take coco's fail-fast per-operation lock first (throws + // OperationInProgressError when a live execute holds the operation — + // the driver counts that as busy, not failed) and re-read the state + // under it: the snapshot op must still be executing before we drive. + const release = await sendService.acquireOperationLock(op.id); + try { + const latest = await source.send.get(op.id); + if (latest?.state === "executing") { + await sendService.recoverExecutingOperation(latest); + } + } finally { + release(); + } } // prepared / rolling_back: the global sweep only warns; nothing to do. return; @@ -194,11 +236,20 @@ export async function runTargetedRecovery( ): Promise { const result: RecoveryRunResult = { attempted: 0, + timedOut: 0, + busy: 0, skipped: 0, failed: 0, skippedMints: new Map(), }; + const timeoutMs = options.timeoutMs ?? 15_000; + const deadlineMs = options.deadlineMs ?? 60_000; + if (![timeoutMs, deadlineMs].every(n => Number.isFinite(n) && n > 0)) { + throw new Error("Recovery budgets must be positive finite numbers"); + } + const deadline = Date.now() + deadlineMs; + const outstanding = options.outstanding ?? new Map(); const kinds = options.kinds ?? ["send", "melt", "receive", "mint"]; const stuck = (options.stuckOperations ?? (await collectStuckOperations(source))).filter( (op) => kinds.includes(op.kind), @@ -221,12 +272,30 @@ 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 (options.shouldStop?.() || Date.now() >= deadline) { result.skipped++; continue; } + const key = recoveryKey(op.kind, op.id); + if (outstanding.has(key)) { result.busy++; continue; } + // Cheap pre-filter for live operations; the lock inside + // recoverStuckOperation is what actually makes the drive atomic. + if (source[op.kind].diagnostics.isLocked(op.id)) { + result.busy++; + continue; + } if (op.kind === "send" && !["pending", "executing"].includes(op.state)) continue; result.attempted++; try { - await recoverStuckOperation(source, sendService, op); + const work = trackRecovery(outstanding, key, recoverStuckOperation(source, sendService, op)); + await waitForRecoveryWork(work, Math.min(timeoutMs, Math.max(1, deadline - Date.now()))); } catch (error) { + if (error instanceof RecoveryWaitTimeout) { result.timedOut++; continue; } + // A live execute grabbed the operation between the isLocked pre-filter + // and the lock acquisition: busy, not failed — leave it for a later + // pass. Name-matched like runMintQuoteRecovery does, because the error + // crosses a package boundary. + if (error instanceof Error && error.name === "OperationInProgressError") { + result.busy++; + continue; + } // Same semantics as coco's tryRecover*: leave the operation for the // next pass. A reachable mint can still reject a specific operation. result.failed++; diff --git a/src/daemon/wallet/recovery-work.test.ts b/src/daemon/wallet/recovery-work.test.ts new file mode 100644 index 0000000..06bfb51 --- /dev/null +++ b/src/daemon/wallet/recovery-work.test.ts @@ -0,0 +1,41 @@ +import { expect, it } from "bun:test"; +import { createRunQueue } from "./coco-client"; +import { createRecoveryDisposer, drainRecoveryWork, trackRecovery } from "./recovery-work"; + +it("incomplete shutdown retains resources until late writes settle; retries join disposal", async () => { + const work = new Map>(); + let finish!: () => void; + const events: string[] = []; + trackRecovery(work, "mint:q", new Promise(r => { finish = r; }).then(() => { events.push("write"); })); + const dispose = createRecoveryDisposer(() => {}, () => drainRecoveryWork(work), async () => { events.push("close"); }, 5); + await expect(dispose()).rejects.toThrow("Timed out"); + expect(events).toEqual([]); + finish(); + await dispose(); + expect(events).toEqual(["write", "close"]); + await dispose(); + expect(events).toEqual(["write", "close"]); +}); + +it("shutdown drains active queue work and queued callbacks reject before touching DB", async () => { + const queue = createRunQueue(); + let disposed = false; + let finish!: () => void; + const events: string[] = []; + let started!: () => void; + const ready = new Promise(r => { started = r; }); + const active = queue(() => new Promise(r => { finish = r; started(); }).then(() => { events.push("write"); })); + await ready; + const queued = queue(async () => { + if (disposed) throw new Error("Wallet is shutting down"); + events.push("unexpected"); + }); + const rejected = queued.catch(error => error); + const dispose = createRecoveryDisposer(() => { disposed = true; }, () => queue.drain(), async () => { events.push("close"); }, 5); + await expect(dispose()).rejects.toThrow("Timed out"); + finish(); + await active; + expect((await rejected).message).toContain("shutting down"); + await dispose(); + expect(events).toEqual(["write", "close"]); +}); diff --git a/src/daemon/wallet/recovery-work.ts b/src/daemon/wallet/recovery-work.ts new file mode 100644 index 0000000..a776f74 --- /dev/null +++ b/src/daemon/wallet/recovery-work.ts @@ -0,0 +1,45 @@ +/** A timed-out wait is not cancellation: retain work until it actually settles. */ +export type RecoveryWork = Map>; +export const recoveryKey = (kind: string, id: string): string => `${kind}:${id}`; + +export function trackRecovery(work: RecoveryWork, key: string, promise: Promise): Promise { + work.set(key, promise); + const clear = () => { if (work.get(key) === promise) work.delete(key); }; + void promise.then(clear, clear); + return promise; +} + +export class RecoveryWaitTimeout extends Error { + constructor() { super("Timed out waiting for recovery; underlying work is still tracked"); } +} + +export async function waitForRecoveryWork(promise: Promise, timeoutMs: number): Promise { + let timer: ReturnType | undefined; + try { + return await Promise.race([promise, new Promise((_, reject) => { + timer = setTimeout(() => reject(new RecoveryWaitTimeout()), timeoutMs); + })]); + } finally { + if (timer !== undefined) clearTimeout(timer); + } +} + +/** Drain actual work, including work registered while an earlier task settles. */ +export async function drainRecoveryWork(work: RecoveryWork): Promise { + while (work.size) await Promise.allSettled([...work.values()]); +} + +/** Timeout reports incomplete disposal; actual cleanup continues safely. */ +export function createRecoveryDisposer( + quiesce: () => void, + settle: () => Promise, + close: () => Promise, + timeoutMs = 30_000, +): () => Promise { + let disposal: Promise | undefined; + return async () => { + quiesce(); + disposal ??= (async () => { await settle(); await close(); })(); + await waitForRecoveryWork(disposal, timeoutMs); + }; +}