Revert "wallet: probe mint reachability, per-mint recovery gating"

This commit is contained in:
redshift
2026-09-30 15:14:54 +00:00
committed by GitHub
parent b437cd5eb3
commit 217f9d299b
4 changed files with 18 additions and 831 deletions
-71
View File
@@ -14,7 +14,6 @@ import {
assertLegacyCocodNotRunning,
claimLegacyCocodPidFile,
createCocoClient,
createRecoveryGate,
DEFAULT_TRUSTED_MINT_URLS,
isZombieProcess,
settleExpiredMintQuotes,
@@ -902,73 +901,3 @@ describe("settlePendingMintQuotes", () => {
expect(logged.mock.calls[0]?.[0]).toContain("21 sat minted");
});
});
describe("createRecoveryGate", () => {
const HEALTHY = "https://healthy.example.com";
const STUCK = "https://stuck.example.com";
function settled(promise: Promise<unknown>): Promise<boolean> {
return Promise.race([
promise.then(() => true, () => true),
new Promise<boolean>((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;
});
});
+18 -169
View File
@@ -37,14 +37,6 @@ import type {
WalletRecoveryProgress,
} from "./cocod-client";
import { selectCleanupOperations } from "./cleanup";
import {
collectStuckOperations,
probeMintReachability,
runTargetedRecovery,
startRecoveryRecheck,
type SendRecoveryService,
type StuckOperation,
} from "./recovery-probe";
import {
clearInterruptedReceiveReservations,
deleteReceiveTokenReservation,
@@ -1039,85 +1031,6 @@ interface RecoveryPhaseProgress {
failedMintQuotes: number;
}
/**
* Coco keeps per-operation recovery private on its services; routstrd already
* reaches into the Manager the same way for `mintOperationService`. Send is
* the only family whose public `refresh()` cannot recover executing ops.
*/
function sendRecoveryServiceOf(coco: Manager): SendRecoveryService {
return (coco as unknown as { sendOperationService: SendRecoveryService })
.sendOperationService;
}
/**
* Gate for value-moving wallet operations while startup recovery runs.
*
* The gate is per mint: once recovery has enumerated stuck operations
* (publishStuckMints), callers whose target mint has none proceed immediately
* — a dead or slow mint must not stall spends from a healthy one. Callers
* without a target mint, or whose mint has stuck operations, wait for the
* full sweep. fail() poisons every caller; reads are never gated.
*/
export interface RecoveryGate {
waitForRecovery(mintUrl?: string): Promise<void>;
publishStuckMints(mints: Set<string>): void;
complete(): void;
fail(error: string): void;
}
export function createRecoveryGate(): RecoveryGate {
let stuckMints: Set<string> | undefined;
let done = false;
let error: string | undefined;
let mintsResolve: (() => void) | undefined;
const mintsPromise = new Promise<void>((resolve) => {
mintsResolve = resolve;
});
let doneResolve: (() => void) | undefined;
const donePromise = new Promise<void>((resolve) => {
doneResolve = resolve;
});
return {
async waitForRecovery(mintUrl?: string): Promise<void> {
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<string>): void {
if (stuckMints) return;
stuckMints = mints;
mintsResolve?.();
},
complete(): void {
done = true;
mintsResolve?.();
doneResolve?.();
},
fail(message: string): void {
done = true;
error = message;
mintsResolve?.();
doneResolve?.();
},
};
}
/**
* Run the wallet recovery sweeps in order, reporting phase changes.
*
@@ -1129,7 +1042,6 @@ async function runWalletRecovery(
coco: Manager,
onProgress: (progress: RecoveryPhaseProgress) => void,
receiveOperationIds?: string[],
onStuckMintsKnown?: (mints: Set<string>) => void,
): Promise<void> {
surfacingRecoveryProgress = true;
let failedMintQuotes = 0;
@@ -1156,60 +1068,19 @@ async function runWalletRecovery(
}
onProgress({ phase: "Settled expired mint quotes", failedMintQuotes });
// Probe every mint that has stuck operations once, up front, so a dead
// mint costs a single short probe instead of a network timeout per
// operation per sweep. Healthy-mint operations are recovered per op;
// dead-mint operations stay parked exactly as coco's own "will retry
// later" path would leave them.
onProgress({ phase: "Probing mints", failedMintQuotes });
const stuckOperations = await collectStuckOperations(coco.ops);
// Value-moving operations gate per mint on this set (see waitForRecovery);
// publish it as soon as it is known so healthy mints unblock immediately.
onStuckMintsKnown?.(new Set(stuckOperations.map((op) => op.mintUrl)));
const unreachableMints = await probeMintReachability(
[...new Set(stuckOperations.map((op) => op.mintUrl))],
);
for (const mintUrl of unreachableMints) {
const count = stuckOperations.filter((op) => op.mintUrl === mintUrl).length;
startupProgress(
`Skipping recovery for unreachable mint: ${mintUrl} (${count} op${count === 1 ? "" : "s"})`,
);
}
const degraded = unreachableMints.size > 0;
const targeted = (kinds: Array<StuckOperation["kind"]>) =>
runTargetedRecovery(coco.ops, sendRecoveryServiceOf(coco), {
kinds,
stuckOperations,
unreachableMints,
});
// Happy path (every mint reachable) keeps coco's global sweeps: they also
// clean up init operations and orphaned proof reservations, which the
// per-op driver cannot enumerate. The targeted driver only runs when a
// dead mint would otherwise tax every stuck op with a network timeout.
onProgress({ phase: "Send recovery", failedMintQuotes });
if (!degraded) await coco.ops.send.recovery.run();
else await targeted(["send"]);
await coco.ops.send.recovery.run();
onProgress({ phase: "Melt recovery", failedMintQuotes });
if (!degraded) await coco.ops.melt.recovery.run();
else await targeted(["melt"]);
await coco.ops.melt.recovery.run();
onProgress({ phase: "Receive recovery", failedMintQuotes });
if (receiveOperationIds) {
// The pre-check already classified every executing receive by unique
// input set. Recover only the conclusive retained operations; unresolved
// groups stay untouched instead of falling back to Coco 1's expensive
// per-row sweep on this startup. Operations at mints the probe found
// unreachable are skipped rather than costing their 15s timeout each.
const mintByOperation = new Map(
stuckOperations
.filter((op) => op.kind === "receive")
.map((op) => [op.id, op.mintUrl]),
);
// per-row sweep on this startup.
for (const operationId of receiveOperationIds) {
const mintUrl = mintByOperation.get(operationId);
if (mintUrl && unreachableMints.has(mintUrl)) continue;
try {
await withTimeout(coco.ops.receive.refresh(operationId), 15_000);
} catch (error) {
@@ -1219,15 +1090,12 @@ async function runWalletRecovery(
});
}
}
} else if (!degraded) {
await coco.ops.receive.recovery.run();
} else {
await targeted(["receive"]);
await coco.ops.receive.recovery.run();
}
onProgress({ phase: "Mint recovery", failedMintQuotes });
if (!degraded) await coco.recoverPendingMintOperations();
else await targeted(["mint"]);
await coco.recoverPendingMintOperations();
onProgress({ phase: "done", failedMintQuotes });
} finally {
@@ -1287,11 +1155,9 @@ export async function createCocoClient(
};
let recoveryResolve: (() => void) | undefined;
let stopPendingMintSweep: (() => Promise<void>) | undefined;
let stopRecoveryRecheck: (() => Promise<void>) | undefined;
const recoveryPromise = new Promise<void>((resolve) => {
recoveryResolve = resolve;
});
const recoveryGate = createRecoveryGate();
try {
startupProgress("Opening Cashu wallet database...");
@@ -1489,33 +1355,13 @@ export async function createCocoClient(
}
},
receiveRecoveryOperationIds,
(mints) => recoveryGate.publishStuckMints(mints),
)
.then(async () => {
await syncReceiveReservations();
recoveryDone = true;
recoveryPhase = "done";
recoveryGate.complete();
recoveryResolve?.();
startupProgress("Wallet recovery complete.");
// Re-check stuck operations periodically: a mint that comes back has
// its parked operations recovered without a daemon restart, and
// operations stuck mid-session are reconciled too. Receive stays
// startup-only because recovering competing receives safely requires
// the startup dedup classification (see receive-dedup.ts).
stopRecoveryRecheck = startRecoveryRecheck(
coco!.ops,
sendRecoveryServiceOf(coco!),
{
kinds: ["send", "melt", "mint"],
onMintDown: (mintUrl, opCount) =>
logger.warn(
`Mint ${mintUrl} is unreachable; ${opCount} stuck operation(s) will keep retrying`,
),
onMintBack: (mintUrl) =>
logger.log(`Mint ${mintUrl} is reachable again; resumed recovering its operations`),
},
);
stopPendingMintSweep = startPendingMintSweep({
ops: coco!.ops,
wallet: coco!.wallet,
@@ -1528,7 +1374,6 @@ export async function createCocoClient(
recoveryDone = true;
recoveryPhase = "error";
recoveryError = error instanceof Error ? error.message : String(error);
recoveryGate.fail(recoveryError);
recoveryResolve?.();
startupProgress(`Wallet recovery failed: ${recoveryError}`);
});
@@ -1553,11 +1398,16 @@ export async function createCocoClient(
let disposed = false;
// Block a value-moving operation until background recovery has settled for
// its target mint (see createRecoveryGate). Reads stay ungated so the
// daemon can report balances/status immediately.
const waitForRecovery = (mintUrl?: string): Promise<void> =>
recoveryGate.waitForRecovery(mintUrl);
/**
* Block a value-moving operation until background recovery has settled.
* Reads stay ungated so the daemon can report balances/status immediately.
*/
const waitForRecovery = async (): Promise<void> => {
if (!recoveryDone) await recoveryPromise;
if (recoveryError) {
throw new Error(`Wallet is not ready: ${recoveryError}`);
}
};
return {
async ping(): Promise<boolean> {
@@ -1737,13 +1587,13 @@ export async function createCocoClient(
},
async receiveBolt11(amount: number, mintUrl?: string) {
await waitForRecovery();
const targetMint = mintUrl
? normalizeMintUrl(mintUrl)
: walletConfig.defaultMintUrl;
if (!targetMint) {
throw new Error("No trusted mint available for Lightning invoice");
}
await waitForRecovery(targetMint);
const op = await coco.ops.mint.prepare({
mintUrl: targetMint,
amount,
@@ -1769,13 +1619,13 @@ export async function createCocoClient(
},
async sendCashu(amount: number, mintUrl?: string): Promise<string> {
await waitForRecovery();
const targetMint = mintUrl
? normalizeMintUrl(mintUrl)
: walletConfig.defaultMintUrl;
if (!targetMint) {
throw new Error("No trusted mint available for sending");
}
await waitForRecovery(targetMint);
const prepared = await coco.ops.send.prepare({
mintUrl: targetMint,
amount,
@@ -1785,13 +1635,13 @@ export async function createCocoClient(
},
async sendBolt11(invoice: string, mintUrl?: string): Promise<string> {
await waitForRecovery();
const targetMint = mintUrl
? normalizeMintUrl(mintUrl)
: walletConfig.defaultMintUrl;
if (!targetMint) {
throw new Error("No trusted mint available for Lightning payment");
}
await waitForRecovery(targetMint);
const prepared = await coco.ops.melt.prepare({
mintUrl: targetMint,
method: "bolt11",
@@ -1837,7 +1687,6 @@ export async function createCocoClient(
async dispose(): Promise<void> {
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.
-279
View File
@@ -1,279 +0,0 @@
import { describe, expect, it, mock } from "bun:test";
import {
collectStuckOperations,
probeMintReachability,
runTargetedRecovery,
startRecoveryRecheck,
type StuckOperationSource,
type SendRecoveryService,
} from "./recovery-probe";
interface FakeOp {
id: string;
mintUrl: string;
state: string;
}
function op(id: string, mintUrl: string, state = "pending"): FakeOp {
return { id, mintUrl, state };
}
interface FakeSource extends StuckOperationSource {
stuck: Record<"send" | "melt" | "receive" | "mint", FakeOp[]>;
refreshed: Record<"send" | "melt" | "receive" | "mint", string[]>;
failOnRefresh?: Set<string>;
}
function makeSource(
stuck: Partial<Record<"send" | "melt" | "receive" | "mint", FakeOp[]>>,
): 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);
});
});
-312
View File
@@ -1,312 +0,0 @@
/**
* Mint reachability probing and targeted (per-operation) wallet recovery.
*
* Coco's global recovery sweeps walk every non-terminal operation one at a
* time; each operation at an unreachable mint costs a full network timeout.
* This module probes every mint that has stuck operations once, in parallel,
* and drives recovery per operation only for mints that answer — so one dead
* mint costs a single short probe instead of N sequential timeouts, and its
* operations stay parked (exactly as coco's "will retry later" path leaves
* them) until the mint comes back.
*/
import { normalizeMintUrl, type Manager } from "@cashu/coco-core";
import { logger } from "../../utils/logger";
/** Short probe: a mint that cannot answer /v1/info in 2s slows every op. */
export const MINT_PROBE_TIMEOUT_MS = 2_000;
/** How often stuck operations are re-checked (and dead mints re-probed). */
export const RECOVERY_RECHECK_INTERVAL_MS = 300_000;
type OpsApi = Manager["ops"];
/** Structural subset of the ops APIs used to enumerate and recover operations. */
export interface StuckOperationSource {
send: Pick<OpsApi["send"], "listInFlight" | "refresh">;
melt: Pick<OpsApi["melt"], "listInFlight" | "refresh">;
receive: Pick<OpsApi["receive"], "listInFlight" | "refresh">;
mint: Pick<OpsApi["mint"], "listInFlight" | "refresh">;
}
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<void>;
tryRecoverExecutingOperation(op: unknown): Promise<void>;
}
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<string, number>;
}
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<string>;
/** 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<StuckOperation[]> {
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<Set<string>> {
const fetcher = options.fetchImpl ?? fetch;
const timeoutMs = options.timeoutMs ?? MINT_PROBE_TIMEOUT_MS;
const unreachable = new Set<string>();
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<void> {
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<RecoveryRunResult> {
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<void> {
const intervalMs = options.intervalMs ?? RECOVERY_RECHECK_INTERVAL_MS;
let stopped = false;
let timer: ReturnType<typeof setTimeout> | undefined;
let inFlight: Promise<void> = Promise.resolve();
const knownDead = new Set<string>();
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;
};
}