mirror of
https://github.com/Routstr/routstrd.git
synced 2026-10-05 12:28:23 +00:00
wallet: probe mint reachability and gate recovery per mint
Startup recovery swept every non-terminal operation through coco's global recovery.run(), so each op at an unreachable mint cost a full network timeout (one 'Send recovery: mint unreachable' line per op per startup), and waitForRecovery gated ALL value-moving operations on the global sweep — a dead mint blocked sends from healthy mints too. - recovery-probe.ts: enumerate stuck ops across send/melt/receive/mint, probe each affected mint once in parallel (GET /v1/info, 2s), and drive recovery per op for reachable mints only. Dead-mint ops stay parked, exactly as coco's own 'will retry later' path leaves them. - runWalletRecovery keeps coco's global sweeps when every mint answers (they also clean up init ops and orphaned reservations); the targeted driver only runs when a mint is dead. Executing sends use coco's private tryRecover* via cast, matching the mintOperationService precedent. - createRecoveryGate: per-mint waitForRecovery. Mints with no stuck ops are never gated; sendCashu/sendBolt11/receiveBolt11 resolve their target mint before gating. Global recovery errors still poison all callers. - startRecoveryRecheck: re-run targeted recovery every 5 minutes so a mint that comes back is recovered without a restart and mid-session stuck ops get reconciled. Transition-based logging only (down once, back once); idle ticks cost one DB query and no network. Receive ops stay startup-only because recovering competing receives needs the startup dedup classification. Known limitation: in degraded mode, init-op and orphaned-reservation housekeeping is deferred to a startup where all stuck mints answer.
This commit is contained in:
@@ -14,6 +14,7 @@ import {
|
|||||||
assertLegacyCocodNotRunning,
|
assertLegacyCocodNotRunning,
|
||||||
claimLegacyCocodPidFile,
|
claimLegacyCocodPidFile,
|
||||||
createCocoClient,
|
createCocoClient,
|
||||||
|
createRecoveryGate,
|
||||||
DEFAULT_TRUSTED_MINT_URLS,
|
DEFAULT_TRUSTED_MINT_URLS,
|
||||||
isZombieProcess,
|
isZombieProcess,
|
||||||
settleExpiredMintQuotes,
|
settleExpiredMintQuotes,
|
||||||
@@ -901,3 +902,73 @@ describe("settlePendingMintQuotes", () => {
|
|||||||
expect(logged.mock.calls[0]?.[0]).toContain("21 sat minted");
|
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;
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|||||||
@@ -37,6 +37,14 @@ import type {
|
|||||||
WalletRecoveryProgress,
|
WalletRecoveryProgress,
|
||||||
} from "./cocod-client";
|
} from "./cocod-client";
|
||||||
import { selectCleanupOperations } from "./cleanup";
|
import { selectCleanupOperations } from "./cleanup";
|
||||||
|
import {
|
||||||
|
collectStuckOperations,
|
||||||
|
probeMintReachability,
|
||||||
|
runTargetedRecovery,
|
||||||
|
startRecoveryRecheck,
|
||||||
|
type SendRecoveryService,
|
||||||
|
type StuckOperation,
|
||||||
|
} from "./recovery-probe";
|
||||||
import {
|
import {
|
||||||
clearInterruptedReceiveReservations,
|
clearInterruptedReceiveReservations,
|
||||||
deleteReceiveTokenReservation,
|
deleteReceiveTokenReservation,
|
||||||
@@ -1031,6 +1039,85 @@ interface RecoveryPhaseProgress {
|
|||||||
failedMintQuotes: number;
|
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.
|
* Run the wallet recovery sweeps in order, reporting phase changes.
|
||||||
*
|
*
|
||||||
@@ -1042,6 +1129,7 @@ async function runWalletRecovery(
|
|||||||
coco: Manager,
|
coco: Manager,
|
||||||
onProgress: (progress: RecoveryPhaseProgress) => void,
|
onProgress: (progress: RecoveryPhaseProgress) => void,
|
||||||
receiveOperationIds?: string[],
|
receiveOperationIds?: string[],
|
||||||
|
onStuckMintsKnown?: (mints: Set<string>) => void,
|
||||||
): Promise<void> {
|
): Promise<void> {
|
||||||
surfacingRecoveryProgress = true;
|
surfacingRecoveryProgress = true;
|
||||||
let failedMintQuotes = 0;
|
let failedMintQuotes = 0;
|
||||||
@@ -1068,19 +1156,60 @@ async function runWalletRecovery(
|
|||||||
}
|
}
|
||||||
onProgress({ phase: "Settled expired mint quotes", failedMintQuotes });
|
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 });
|
onProgress({ phase: "Send recovery", failedMintQuotes });
|
||||||
await coco.ops.send.recovery.run();
|
if (!degraded) await coco.ops.send.recovery.run();
|
||||||
|
else await targeted(["send"]);
|
||||||
|
|
||||||
onProgress({ phase: "Melt recovery", failedMintQuotes });
|
onProgress({ phase: "Melt recovery", failedMintQuotes });
|
||||||
await coco.ops.melt.recovery.run();
|
if (!degraded) await coco.ops.melt.recovery.run();
|
||||||
|
else await targeted(["melt"]);
|
||||||
|
|
||||||
onProgress({ phase: "Receive recovery", failedMintQuotes });
|
onProgress({ phase: "Receive recovery", failedMintQuotes });
|
||||||
if (receiveOperationIds) {
|
if (receiveOperationIds) {
|
||||||
// The pre-check already classified every executing receive by unique
|
// The pre-check already classified every executing receive by unique
|
||||||
// input set. Recover only the conclusive retained operations; unresolved
|
// input set. Recover only the conclusive retained operations; unresolved
|
||||||
// groups stay untouched instead of falling back to Coco 1's expensive
|
// groups stay untouched instead of falling back to Coco 1's expensive
|
||||||
// per-row sweep on this startup.
|
// per-row sweep on this startup. Operations at mints the probe found
|
||||||
|
// unreachable are skipped rather than costing their 15s timeout each.
|
||||||
|
const mintByOperation = new Map(
|
||||||
|
stuckOperations
|
||||||
|
.filter((op) => op.kind === "receive")
|
||||||
|
.map((op) => [op.id, op.mintUrl]),
|
||||||
|
);
|
||||||
for (const operationId of receiveOperationIds) {
|
for (const operationId of receiveOperationIds) {
|
||||||
|
const mintUrl = mintByOperation.get(operationId);
|
||||||
|
if (mintUrl && unreachableMints.has(mintUrl)) continue;
|
||||||
try {
|
try {
|
||||||
await withTimeout(coco.ops.receive.refresh(operationId), 15_000);
|
await withTimeout(coco.ops.receive.refresh(operationId), 15_000);
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
@@ -1090,12 +1219,15 @@ async function runWalletRecovery(
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
} else {
|
} else if (!degraded) {
|
||||||
await coco.ops.receive.recovery.run();
|
await coco.ops.receive.recovery.run();
|
||||||
|
} else {
|
||||||
|
await targeted(["receive"]);
|
||||||
}
|
}
|
||||||
|
|
||||||
onProgress({ phase: "Mint recovery", failedMintQuotes });
|
onProgress({ phase: "Mint recovery", failedMintQuotes });
|
||||||
await coco.recoverPendingMintOperations();
|
if (!degraded) await coco.recoverPendingMintOperations();
|
||||||
|
else await targeted(["mint"]);
|
||||||
|
|
||||||
onProgress({ phase: "done", failedMintQuotes });
|
onProgress({ phase: "done", failedMintQuotes });
|
||||||
} finally {
|
} finally {
|
||||||
@@ -1155,9 +1287,11 @@ export async function createCocoClient(
|
|||||||
};
|
};
|
||||||
let recoveryResolve: (() => void) | undefined;
|
let recoveryResolve: (() => void) | undefined;
|
||||||
let stopPendingMintSweep: (() => Promise<void>) | undefined;
|
let stopPendingMintSweep: (() => Promise<void>) | undefined;
|
||||||
|
let stopRecoveryRecheck: (() => Promise<void>) | undefined;
|
||||||
const recoveryPromise = new Promise<void>((resolve) => {
|
const recoveryPromise = new Promise<void>((resolve) => {
|
||||||
recoveryResolve = resolve;
|
recoveryResolve = resolve;
|
||||||
});
|
});
|
||||||
|
const recoveryGate = createRecoveryGate();
|
||||||
|
|
||||||
try {
|
try {
|
||||||
startupProgress("Opening Cashu wallet database...");
|
startupProgress("Opening Cashu wallet database...");
|
||||||
@@ -1355,13 +1489,33 @@ export async function createCocoClient(
|
|||||||
}
|
}
|
||||||
},
|
},
|
||||||
receiveRecoveryOperationIds,
|
receiveRecoveryOperationIds,
|
||||||
|
(mints) => recoveryGate.publishStuckMints(mints),
|
||||||
)
|
)
|
||||||
.then(async () => {
|
.then(async () => {
|
||||||
await syncReceiveReservations();
|
await syncReceiveReservations();
|
||||||
recoveryDone = true;
|
recoveryDone = true;
|
||||||
recoveryPhase = "done";
|
recoveryPhase = "done";
|
||||||
|
recoveryGate.complete();
|
||||||
recoveryResolve?.();
|
recoveryResolve?.();
|
||||||
startupProgress("Wallet recovery complete.");
|
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({
|
stopPendingMintSweep = startPendingMintSweep({
|
||||||
ops: coco!.ops,
|
ops: coco!.ops,
|
||||||
wallet: coco!.wallet,
|
wallet: coco!.wallet,
|
||||||
@@ -1374,6 +1528,7 @@ export async function createCocoClient(
|
|||||||
recoveryDone = true;
|
recoveryDone = true;
|
||||||
recoveryPhase = "error";
|
recoveryPhase = "error";
|
||||||
recoveryError = error instanceof Error ? error.message : String(error);
|
recoveryError = error instanceof Error ? error.message : String(error);
|
||||||
|
recoveryGate.fail(recoveryError);
|
||||||
recoveryResolve?.();
|
recoveryResolve?.();
|
||||||
startupProgress(`Wallet recovery failed: ${recoveryError}`);
|
startupProgress(`Wallet recovery failed: ${recoveryError}`);
|
||||||
});
|
});
|
||||||
@@ -1398,16 +1553,11 @@ export async function createCocoClient(
|
|||||||
|
|
||||||
let disposed = false;
|
let disposed = false;
|
||||||
|
|
||||||
/**
|
// Block a value-moving operation until background recovery has settled for
|
||||||
* Block a value-moving operation until background recovery has settled.
|
// its target mint (see createRecoveryGate). Reads stay ungated so the
|
||||||
* Reads stay ungated so the daemon can report balances/status immediately.
|
// daemon can report balances/status immediately.
|
||||||
*/
|
const waitForRecovery = (mintUrl?: string): Promise<void> =>
|
||||||
const waitForRecovery = async (): Promise<void> => {
|
recoveryGate.waitForRecovery(mintUrl);
|
||||||
if (!recoveryDone) await recoveryPromise;
|
|
||||||
if (recoveryError) {
|
|
||||||
throw new Error(`Wallet is not ready: ${recoveryError}`);
|
|
||||||
}
|
|
||||||
};
|
|
||||||
|
|
||||||
return {
|
return {
|
||||||
async ping(): Promise<boolean> {
|
async ping(): Promise<boolean> {
|
||||||
@@ -1587,13 +1737,13 @@ export async function createCocoClient(
|
|||||||
},
|
},
|
||||||
|
|
||||||
async receiveBolt11(amount: number, mintUrl?: string) {
|
async receiveBolt11(amount: number, mintUrl?: string) {
|
||||||
await waitForRecovery();
|
|
||||||
const targetMint = mintUrl
|
const targetMint = mintUrl
|
||||||
? normalizeMintUrl(mintUrl)
|
? normalizeMintUrl(mintUrl)
|
||||||
: walletConfig.defaultMintUrl;
|
: walletConfig.defaultMintUrl;
|
||||||
if (!targetMint) {
|
if (!targetMint) {
|
||||||
throw new Error("No trusted mint available for Lightning invoice");
|
throw new Error("No trusted mint available for Lightning invoice");
|
||||||
}
|
}
|
||||||
|
await waitForRecovery(targetMint);
|
||||||
const op = await coco.ops.mint.prepare({
|
const op = await coco.ops.mint.prepare({
|
||||||
mintUrl: targetMint,
|
mintUrl: targetMint,
|
||||||
amount,
|
amount,
|
||||||
@@ -1619,13 +1769,13 @@ export async function createCocoClient(
|
|||||||
},
|
},
|
||||||
|
|
||||||
async sendCashu(amount: number, mintUrl?: string): Promise<string> {
|
async sendCashu(amount: number, mintUrl?: string): Promise<string> {
|
||||||
await waitForRecovery();
|
|
||||||
const targetMint = mintUrl
|
const targetMint = mintUrl
|
||||||
? normalizeMintUrl(mintUrl)
|
? normalizeMintUrl(mintUrl)
|
||||||
: walletConfig.defaultMintUrl;
|
: walletConfig.defaultMintUrl;
|
||||||
if (!targetMint) {
|
if (!targetMint) {
|
||||||
throw new Error("No trusted mint available for sending");
|
throw new Error("No trusted mint available for sending");
|
||||||
}
|
}
|
||||||
|
await waitForRecovery(targetMint);
|
||||||
const prepared = await coco.ops.send.prepare({
|
const prepared = await coco.ops.send.prepare({
|
||||||
mintUrl: targetMint,
|
mintUrl: targetMint,
|
||||||
amount,
|
amount,
|
||||||
@@ -1635,13 +1785,13 @@ export async function createCocoClient(
|
|||||||
},
|
},
|
||||||
|
|
||||||
async sendBolt11(invoice: string, mintUrl?: string): Promise<string> {
|
async sendBolt11(invoice: string, mintUrl?: string): Promise<string> {
|
||||||
await waitForRecovery();
|
|
||||||
const targetMint = mintUrl
|
const targetMint = mintUrl
|
||||||
? normalizeMintUrl(mintUrl)
|
? normalizeMintUrl(mintUrl)
|
||||||
: walletConfig.defaultMintUrl;
|
: walletConfig.defaultMintUrl;
|
||||||
if (!targetMint) {
|
if (!targetMint) {
|
||||||
throw new Error("No trusted mint available for Lightning payment");
|
throw new Error("No trusted mint available for Lightning payment");
|
||||||
}
|
}
|
||||||
|
await waitForRecovery(targetMint);
|
||||||
const prepared = await coco.ops.melt.prepare({
|
const prepared = await coco.ops.melt.prepare({
|
||||||
mintUrl: targetMint,
|
mintUrl: targetMint,
|
||||||
method: "bolt11",
|
method: "bolt11",
|
||||||
@@ -1687,6 +1837,7 @@ export async function createCocoClient(
|
|||||||
async dispose(): Promise<void> {
|
async dispose(): Promise<void> {
|
||||||
if (disposed) return;
|
if (disposed) return;
|
||||||
disposed = true;
|
disposed = true;
|
||||||
|
await stopRecoveryRecheck?.();
|
||||||
try {
|
try {
|
||||||
// Let any in-flight recovery settle before closing the database from
|
// Let any in-flight recovery settle before closing the database from
|
||||||
// underneath it. The recovery promise resolves on success or failure.
|
// underneath it. The recovery promise resolves on success or failure.
|
||||||
|
|||||||
@@ -0,0 +1,279 @@
|
|||||||
|
import { describe, expect, it, mock } from "bun:test";
|
||||||
|
import {
|
||||||
|
collectStuckOperations,
|
||||||
|
probeMintReachability,
|
||||||
|
runTargetedRecovery,
|
||||||
|
startRecoveryRecheck,
|
||||||
|
type StuckOperationSource,
|
||||||
|
type SendRecoveryService,
|
||||||
|
} from "./recovery-probe";
|
||||||
|
|
||||||
|
interface FakeOp {
|
||||||
|
id: string;
|
||||||
|
mintUrl: string;
|
||||||
|
state: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
function op(id: string, mintUrl: string, state = "pending"): FakeOp {
|
||||||
|
return { id, mintUrl, state };
|
||||||
|
}
|
||||||
|
|
||||||
|
interface FakeSource extends StuckOperationSource {
|
||||||
|
stuck: Record<"send" | "melt" | "receive" | "mint", FakeOp[]>;
|
||||||
|
refreshed: Record<"send" | "melt" | "receive" | "mint", string[]>;
|
||||||
|
failOnRefresh?: Set<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);
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -0,0 +1,312 @@
|
|||||||
|
/**
|
||||||
|
* Mint reachability probing and targeted (per-operation) wallet recovery.
|
||||||
|
*
|
||||||
|
* Coco's global recovery sweeps walk every non-terminal operation one at a
|
||||||
|
* time; each operation at an unreachable mint costs a full network timeout.
|
||||||
|
* This module probes every mint that has stuck operations once, in parallel,
|
||||||
|
* and drives recovery per operation only for mints that answer — so one dead
|
||||||
|
* mint costs a single short probe instead of N sequential timeouts, and its
|
||||||
|
* operations stay parked (exactly as coco's "will retry later" path leaves
|
||||||
|
* them) until the mint comes back.
|
||||||
|
*/
|
||||||
|
import { normalizeMintUrl, type Manager } from "@cashu/coco-core";
|
||||||
|
import { logger } from "../../utils/logger";
|
||||||
|
|
||||||
|
/** Short probe: a mint that cannot answer /v1/info in 2s slows every op. */
|
||||||
|
export const MINT_PROBE_TIMEOUT_MS = 2_000;
|
||||||
|
/** How often stuck operations are re-checked (and dead mints re-probed). */
|
||||||
|
export const RECOVERY_RECHECK_INTERVAL_MS = 300_000;
|
||||||
|
|
||||||
|
type OpsApi = Manager["ops"];
|
||||||
|
|
||||||
|
/** Structural subset of the ops APIs used to enumerate and recover operations. */
|
||||||
|
export interface StuckOperationSource {
|
||||||
|
send: Pick<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;
|
||||||
|
};
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user