wallet: gate recovery per mint only in degraded mode, guard live ops

Review fixes for the mint reachability probe (#116):

- Open the per-mint recovery gate only on the degraded path, after the
  probe and local housekeeping finish. On the happy path every
  value-moving caller waits for the global sweeps as before, so coco's
  send/melt recovery never runs beside a live execute.
- Skip operations whose per-operation lock is held (diagnostics.isLocked)
  in the targeted driver, and drive sends only in pending/executing
  states, so recovery can never race a swap that execute() still holds.
  Drive sends through recoverExecutingOperation instead of the
  error-swallowing tryRecover* wrappers.
- Run coco's local-only crash cleanup (init operations, orphaned proof
  reservations) on the degraded path before the gate opens, so a mint
  that never returns cannot leave that housekeeping undone forever.
- Probe with a path-preserving /v1/info join so subpath mints (e.g.
  https://host/Bitcoin) are checked at their real endpoint, and classify
  malformed persisted URLs as unreachable.
- Drop the 5-minute stuck-operation recheck: it introduced the live-send
  race, an unbounded shutdown wait, and duplicate polling next to the
  pending-mint sweep. Parked operations are recovered on the next
  startup instead.
- Count recovery attempts rather than completions, since the driver
  cannot observe tryRecover*-style swallowed errors.

Adds recovery-integration.test.ts with real coco Manager + SQLite
regression tests: live execute vs targeted recovery, gate ordering on
both paths, and local housekeeping without network.
This commit is contained in:
redshift
2026-10-02 14:50:42 +08:00
parent b97a4e3798
commit 6ff7fe2de3
4 changed files with 276 additions and 186 deletions
+44 -31
View File
@@ -46,7 +46,6 @@ import {
collectStuckOperations,
probeMintReachability,
runTargetedRecovery,
startRecoveryRecheck,
type SendRecoveryService,
type StuckOperation,
} from "./recovery-probe";
@@ -1552,12 +1551,38 @@ function sendRecoveryServiceOf(coco: Manager): SendRecoveryService {
.sendOperationService;
}
/** Local-only crash cleanup, run before the degraded gate opens. */
export async function cleanupLocalRecoveryState(
coco: Manager,
repo: SqliteRepositories,
): Promise<void> {
// 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<string, {
recoverInitOperation(op: unknown): Promise<void>;
cleanupOrphanedReservations(): Promise<number>;
}>;
const families = [
["send", repo.sendOperationRepository],
["melt", repo.meltOperationRepository],
["receive", repo.receiveOperationRepository],
["mint", repo.mintOperationRepository],
] as const;
for (const [kind, repository] of families) {
for (const op of await repository.getByState("init")) {
await services[`${kind}OperationService`]!.recoverInitOperation(op);
}
}
await services.sendOperationService!.cleanupOrphanedReservations();
}
/**
* Gate for value-moving wallet operations while startup recovery runs.
*
* The gate is per mint: once recovery has enumerated stuck operations
* (publishStuckMints), callers whose target mint has none proceed immediately
* — a dead or slow mint must not stall spends from a healthy one. Callers
* On degraded startup only, publishStuckMints opens the gate for callers
* whose target mint has no stuck operations after probing and local cleanup.
* A dead mint must not stall spends from a healthy one. On the happy path
* all callers wait until the global sweeps finish. Callers
* without a target mint, or whose mint has stuck operations, wait for the
* full sweep. fail() poisons every caller; reads are never gated.
*/
@@ -1628,11 +1653,12 @@ export function createRecoveryGate(): RecoveryGate {
* are failed locally so `recoverPendingMintOperations()` skips them, while
* paid/issued and unreachable-mint quotes stay pending for the sweep.
*/
async function runWalletRecovery(
export async function runWalletRecovery(
coco: Manager,
onProgress: (progress: RecoveryPhaseProgress) => void,
receiveOperationIds?: string[],
onStuckMintsKnown?: (mints: Set<string>) => void,
options: { cleanupLocalState?: () => Promise<void>; fetchImpl?: typeof fetch } = {},
): Promise<void> {
surfacingRecoveryProgress = true;
let failedMintQuotes = 0;
@@ -1666,11 +1692,9 @@ async function runWalletRecovery(
// later" path would leave them.
onProgress({ phase: "Probing mints", failedMintQuotes });
const stuckOperations = await collectStuckOperations(coco.ops);
// Value-moving operations gate per mint on this set (see waitForRecovery);
// publish it as soon as it is known so healthy mints unblock immediately.
onStuckMintsKnown?.(new Set(stuckOperations.map((op) => op.mintUrl)));
const unreachableMints = await probeMintReachability(
[...new Set(stuckOperations.map((op) => op.mintUrl))],
{ fetchImpl: options.fetchImpl },
);
for (const mintUrl of unreachableMints) {
const count = stuckOperations.filter((op) => op.mintUrl === mintUrl).length;
@@ -1679,6 +1703,13 @@ async function runWalletRecovery(
);
}
const degraded = unreachableMints.size > 0;
if (degraded) {
// Local-only housekeeping must finish before any new operation is allowed.
await options.cleanupLocalState?.();
// Global sweeps enumerate fresh state and are unsafe beside live sends.
// Only the snapshot-based degraded path may open the per-mint gate.
onStuckMintsKnown?.(new Set(stuckOperations.map((op) => op.mintUrl)));
}
const targeted = (kinds: Array<StuckOperation["kind"]>) =>
runTargetedRecovery(coco.ops, sendRecoveryServiceOf(coco), {
kinds,
@@ -1687,9 +1718,10 @@ async function runWalletRecovery(
});
// Happy path (every mint reachable) keeps coco's global sweeps: they also
// clean up init operations and orphaned proof reservations, which the
// per-op driver cannot enumerate. The targeted driver only runs when a
// dead mint would otherwise tax every stuck op with a network timeout.
// clean up init operations and orphaned proof reservations. Degraded
// startup runs that local housekeeping before opening its gate, and only
// drives the previously collected snapshot when a dead mint would
// otherwise tax every stuck op with a network timeout.
onProgress({ phase: "Send recovery", failedMintQuotes });
if (!degraded) await coco.ops.send.recovery.run();
else await targeted(["send"]);
@@ -1790,7 +1822,6 @@ export async function createCocoClient(
};
let recoveryResolve: (() => void) | undefined;
let stopPendingMintSweep: (() => Promise<void>) | undefined;
let stopRecoveryRecheck: (() => Promise<void>) | undefined;
const recoveryPromise = new Promise<void>((resolve) => {
recoveryResolve = resolve;
});
@@ -1993,6 +2024,7 @@ export async function createCocoClient(
},
receiveRecoveryOperationIds,
(mints) => recoveryGate.publishStuckMints(mints),
{ cleanupLocalState: () => cleanupLocalRecoveryState(coco!, repo) },
)
.then(async () => {
await syncReceiveReservations();
@@ -2001,24 +2033,6 @@ export async function createCocoClient(
recoveryGate.complete();
recoveryResolve?.();
startupProgress("Wallet recovery complete.");
// Re-check stuck operations periodically: a mint that comes back has
// its parked operations recovered without a daemon restart, and
// operations stuck mid-session are reconciled too. Receive stays
// startup-only because recovering competing receives safely requires
// the startup dedup classification (see receive-dedup.ts).
stopRecoveryRecheck = startRecoveryRecheck(
coco!.ops,
sendRecoveryServiceOf(coco!),
{
kinds: ["send", "melt", "mint"],
onMintDown: (mintUrl, opCount) =>
logger.warn(
`Mint ${mintUrl} is unreachable; ${opCount} stuck operation(s) will keep retrying`,
),
onMintBack: (mintUrl) =>
logger.log(`Mint ${mintUrl} is reachable again; resumed recovering its operations`),
},
);
stopPendingMintSweep = startPendingMintSweep({
ops: coco!.ops,
wallet: coco!.wallet,
@@ -2344,7 +2358,6 @@ export async function createCocoClient(
async dispose(): Promise<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.
@@ -0,0 +1,174 @@
import { expect, it, spyOn } from "bun:test";
import { Manager } from "@cashu/coco-core";
import { SqliteRepositories } from "@cashu/coco-sqlite-bun";
import { Database } from "bun:sqlite";
import { runTargetedRecovery, type SendRecoveryService } from "./recovery-probe";
import {
cleanupLocalRecoveryState,
createRecoveryGate,
runWalletRecovery,
} from "./coco-client";
const MINT = "https://mint.example.com";
const OP = "op-live";
it("targeted recovery leaves a send that execute() holds alone", async () => {
const repo = new SqliteRepositories({ database: new Database(":memory:") });
await repo.init();
const coco = new Manager(repo, async () => new Uint8Array(64));
const internals = coco as unknown as {
sendOperationService: SendRecoveryService;
walletService: { getWalletWithActiveKeysetId: (m: string) => Promise<unknown> };
};
let swapStarted!: () => void;
const started = new Promise<void>((r) => (swapStarted = r));
let finishSwap!: (v: { send: unknown[]; keep: unknown[] }) => void;
const swap = new Promise<{ send: unknown[]; keep: unknown[] }>((r) => (finishSwap = r));
internals.walletService.getWalletWithActiveKeysetId = async () => ({
wallet: {
unit: "sat",
send: async () => (swapStarted(), swap),
checkProofsStates: async () => [{ state: "UNSPENT" }],
getFeesForProofs: () => 0,
},
});
await repo.proofRepository.saveProofs(MINT, [
{ id: "00aa", amount: 8, secret: "in-1", C: "02aa", mintUrl: MINT, state: "ready" } as never,
]);
await repo.proofRepository.reserveProofs(MINT, ["in-1"], OP);
await repo.sendOperationRepository.create({
id: OP, mintUrl: MINT, amount: 8, state: "prepared", method: "default", methodData: {},
createdAt: Date.now(), updatedAt: Date.now(), needsSwap: true, fee: 0, inputAmount: 8,
inputProofSecrets: ["in-1"],
outputData: {
keep: [],
send: [{
blindedMessage: { amount: 8, id: "00aa", B_: "02" + "11".repeat(32) },
blindingFactor: "01",
secret: Buffer.from("out-1").toString("hex"),
}],
},
} as never);
const live = coco.ops.send.execute(OP);
await started; // swap is at the mint
await runTargetedRecovery(coco.ops, internals.sendOperationService, {
kinds: ["send"],
fetchImpl: (async () => new Response("{}")) as unknown as unknown as typeof fetch,
});
// On 8005aeb this is "rolled_back" and "in-1" is no longer reserved.
expect((await coco.ops.send.get(OP))?.state).toBe("executing");
finishSwap({ send: [{ id: "00aa", amount: 8, secret: "out-1", C: "02bb" }], keep: [] });
await live;
expect((await coco.ops.send.get(OP))?.state).toBe("pending");
});
it("keeps healthy-mint callers gated while happy-path global recovery runs", async () => {
const db = new Database(":memory:");
const repo = new SqliteRepositories({ database: db });
await repo.init();
const coco = new Manager(repo, async () => new Uint8Array(64));
const gate = createRecoveryGate();
let enter!: () => void;
const entered = new Promise<void>((resolve) => { enter = resolve; });
let release!: () => void;
const barrier = new Promise<void>((resolve) => { release = resolve; });
const spies = [
spyOn(coco.ops.send.recovery, "run").mockImplementation(async () => { enter(); await barrier; }),
spyOn(coco.ops.melt.recovery, "run").mockResolvedValue(undefined),
spyOn(coco.ops.receive.recovery, "run").mockResolvedValue(undefined),
spyOn(coco, "recoverPendingMintOperations").mockResolvedValue(undefined),
];
try {
const recovery = runWalletRecovery(coco, () => {}, undefined,
(mints) => gate.publishStuckMints(mints),
{ fetchImpl: (async () => new Response("{}")) as unknown as typeof fetch },
).then(() => gate.complete());
await entered;
let released = false;
const waiting = gate.waitForRecovery(MINT).then(() => { released = true; });
await new Promise((resolve) => setTimeout(resolve, 10));
expect(released).toBe(false);
release();
await recovery;
await waiting;
expect(released).toBe(true);
} finally {
release();
for (const spy of spies) spy.mockRestore();
db.close();
}
});
it("degraded startup finishes local housekeeping before opening the per-mint gate", async () => {
const db = new Database(":memory:");
const repo = new SqliteRepositories({ database: db });
await repo.init();
const coco = new Manager(repo, async () => new Uint8Array(64));
await repo.sendOperationRepository.create({
id: "stuck", mintUrl: MINT, amount: 8, state: "rolling_back",
method: "default", methodData: {}, createdAt: Date.now(), updatedAt: Date.now(),
} as never);
const gate = createRecoveryGate();
const events: string[] = [];
const sweep = spyOn(coco.ops.send.recovery, "run");
try {
const recovery = runWalletRecovery(coco, () => {}, [], (mints) => {
events.push("gate");
gate.publishStuckMints(mints);
}, {
fetchImpl: (async () => { throw new Error("offline"); }) as unknown as unknown as typeof fetch,
cleanupLocalState: async () => { events.push("cleanup"); },
}).then(() => gate.complete());
await gate.waitForRecovery("https://healthy.example.com");
expect(events).toEqual(["cleanup", "gate"]);
await recovery;
expect(sweep).not.toHaveBeenCalled();
expect((await coco.ops.send.get("stuck"))?.state).toBe("rolling_back");
} finally {
sweep.mockRestore();
db.close();
}
});
it("local housekeeping cleans init sends and orphaned reservations without network", async () => {
const db = new Database(":memory:");
const repo = new SqliteRepositories({ database: db });
await repo.init();
const coco = new Manager(repo, async () => new Uint8Array(64));
try {
await repo.sendOperationRepository.create({
id: "init-send", mintUrl: MINT, amount: 8, state: "init", method: "default",
methodData: {}, createdAt: Date.now(), updatedAt: Date.now(),
} as never);
for (const [id, repository] of [
["init-melt", repo.meltOperationRepository],
["init-receive", repo.receiveOperationRepository],
["init-mint", repo.mintOperationRepository],
] as const) {
await repository.create({
id, mintUrl: MINT, amount: 8, state: "init", method: "bolt11",
methodData: {}, inputProofs: [], createdAt: Date.now(), updatedAt: Date.now(),
} as never);
}
await repo.proofRepository.saveProofs(MINT, [
{ id: "00aa", amount: 8, secret: "init-input", C: "02aa", mintUrl: MINT, state: "ready" },
{ id: "00aa", amount: 8, secret: "orphan-input", C: "02aa", mintUrl: MINT, state: "ready" },
] as never);
await repo.proofRepository.reserveProofs(MINT, ["init-input"], "init-send");
await repo.proofRepository.reserveProofs(MINT, ["orphan-input"], "missing-send");
await cleanupLocalRecoveryState(coco, repo);
expect(await coco.ops.send.get("init-send")).toBeNull();
expect(await repo.proofRepository.getReservedProofs()).toEqual([]);
expect(await repo.meltOperationRepository.getById("init-melt")).toBeNull();
expect(await repo.receiveOperationRepository.getById("init-receive")).toBeNull();
expect(await repo.mintOperationRepository.getById("init-mint")).toBeNull();
} finally {
db.close();
}
});
+41 -69
View File
@@ -1,9 +1,8 @@
import { describe, expect, it, mock } from "bun:test";
import { describe, expect, it } from "bun:test";
import {
collectStuckOperations,
probeMintReachability,
runTargetedRecovery,
startRecoveryRecheck,
type StuckOperationSource,
type SendRecoveryService,
} from "./recovery-probe";
@@ -40,6 +39,7 @@ function makeSource(
mint: stuck.mint ?? [],
};
const family = (kind: keyof typeof full) => ({
diagnostics: { isLocked: () => false },
listInFlight: async () => full[kind] as never,
refresh: async (id: string) => {
refreshed[kind].push(id);
@@ -60,18 +60,12 @@ function makeSource(
}
function makeSendService(): SendRecoveryService & {
initRecovered: string[];
executingRecovered: string[];
} {
const initRecovered: string[] = [];
const executingRecovered: string[] = [];
return {
initRecovered,
executingRecovered,
tryRecoverInitOperation: async (raw) => {
initRecovered.push((raw as FakeOp).id);
},
tryRecoverExecutingOperation: async (raw) => {
recoverExecutingOperation: async (raw) => {
executingRecovered.push((raw as FakeOp).id);
},
};
@@ -148,7 +142,7 @@ describe("runTargetedRecovery", () => {
const result = await runTargetedRecovery(makeSource({}), makeSendService(), {
fetchImpl: deadFetch.fetchImpl,
});
expect(result).toMatchObject({ recovered: 0, skipped: 0, failed: 0 });
expect(result).toMatchObject({ attempted: 0, skipped: 0, failed: 0 });
expect(deadFetch.calls).toHaveLength(0);
});
@@ -163,7 +157,7 @@ describe("runTargetedRecovery", () => {
fetchImpl: deadFetch.fetchImpl,
onSkippedMint: (mintUrl, count) => skippedMints.push([mintUrl, count]),
});
expect(result.recovered).toBe(2);
expect(result.attempted).toBe(2);
expect(result.skipped).toBe(2);
expect(result.failed).toBe(0);
expect(skippedMints).toEqual([["https://dead.example.com", 2]]);
@@ -178,17 +172,15 @@ describe("runTargetedRecovery", () => {
send: [
op("pending-1", "https://live.example.com", "pending"),
op("exec-1", "https://live.example.com", "executing"),
op("init-1", "https://live.example.com", "init"),
],
});
const sendService = makeSendService();
const result = await runTargetedRecovery(source, sendService, {
fetchImpl: liveFetch().fetchImpl,
});
expect(result.recovered).toBe(3);
expect(result.attempted).toBe(2);
expect(source.refreshed.send).toEqual(["pending-1"]);
expect(sendService.executingRecovered).toEqual(["exec-1"]);
expect(sendService.initRecovered).toEqual(["init-1"]);
});
it("counts per-operation failures at reachable mints and continues", async () => {
@@ -198,7 +190,7 @@ describe("runTargetedRecovery", () => {
const result = await runTargetedRecovery(source, makeSendService(), {
fetchImpl: liveFetch().fetchImpl,
});
expect(result.recovered).toBe(1);
expect(result.attempted).toBe(2);
expect(result.failed).toBe(1);
expect(source.refreshed.melt).toEqual(["boom-1", "m2"]);
});
@@ -221,59 +213,39 @@ describe("runTargetedRecovery", () => {
});
});
describe("startRecoveryRecheck", () => {
it("recovers stuck operations on the interval and stops cleanly", async () => {
const source = makeSource({
send: [op("s1", "https://live.example.com")],
});
const stop = startRecoveryRecheck(source, makeSendService(), {
intervalMs: 20,
fetchImpl: liveFetch().fetchImpl,
});
await new Promise((resolve) => setTimeout(resolve, 60));
await stop();
expect(source.refreshed.send.length).toBeGreaterThanOrEqual(1);
const after = source.refreshed.send.length;
await new Promise((resolve) => setTimeout(resolve, 60));
expect(source.refreshed.send.length).toBe(after);
});
it("is idle without network traffic when nothing is stuck", async () => {
const { calls, fetchImpl } = makeFetch(() => false);
const stop = startRecoveryRecheck(makeSource({}), makeSendService(), {
intervalMs: 20,
fetchImpl,
});
await new Promise((resolve) => setTimeout(resolve, 60));
await stop();
expect(calls).toHaveLength(0);
});
it("reports mint down/up transitions once instead of repeating", async () => {
let dead = true;
const { fetchImpl } = makeFetch(() => dead);
const source = makeSource({
melt: [op("m1", "https://flaky.example.com")],
});
const down: string[] = [];
const back: string[] = [];
const stop = startRecoveryRecheck(source, makeSendService(), {
intervalMs: 20,
fetchImpl,
onMintDown: (mintUrl) => down.push(mintUrl),
onMintBack: (mintUrl) => back.push(mintUrl),
});
await new Promise((resolve) => setTimeout(resolve, 90));
// Dead across several ticks: reported once, never re-probed per op.
expect(down).toEqual(["https://flaky.example.com"]);
expect(back).toEqual([]);
expect(source.refreshed.melt).toEqual([]);
dead = false;
await new Promise((resolve) => setTimeout(resolve, 90));
await stop();
expect(down).toEqual(["https://flaky.example.com"]);
expect(back).toEqual(["https://flaky.example.com"]);
expect(source.refreshed.melt.length).toBeGreaterThanOrEqual(1);
});
it("preserves subpath mint URLs when probing", async () => {
const { fetchImpl, calls } = liveFetch();
await probeMintReachability(["https://mint.example.com/Bitcoin/"], { fetchImpl });
expect(calls).toEqual(["https://mint.example.com/Bitcoin/v1/info"]);
});
it("does not drive locked operations or count rolling-back sends as attempts", async () => {
const source = makeSource({ send: [
op("live", "https://mint.example.com", "executing"),
op("rollback", "https://mint.example.com", "rolling_back"),
] });
source.send.diagnostics.isLocked = (id) => id === "live";
const service = makeSendService();
const result = await runTargetedRecovery(source, service, { fetchImpl: liveFetch().fetchImpl });
expect(result.attempted).toBe(0);
expect(service.executingRecovered).toEqual([]);
});
it("classifies malformed persisted URLs as unreachable without fetching", async () => {
const { fetchImpl, calls } = liveFetch();
const unreachable = await probeMintReachability(["not-a-url"], { fetchImpl });
expect([...unreachable]).toEqual(["not-a-url"]);
expect(calls).toEqual([]);
});
it("bounds a hanging probe with an abort signal", async () => {
const fetchImpl = (async (_url: unknown, options: RequestInit) =>
new Promise<Response>((_resolve, reject) => {
options.signal!.addEventListener("abort", () => reject(options.signal!.reason), { once: true });
})) as unknown as typeof fetch;
const unreachable = await probeMintReachability(["https://slow.example.com"], {
fetchImpl, timeoutMs: 10,
});
expect([...unreachable]).toEqual(["https://slow.example.com"]);
});
+17 -86
View File
@@ -7,24 +7,22 @@
* and drives recovery per operation only for mints that answer — so one dead
* mint costs a single short probe instead of N sequential timeouts, and its
* operations stay parked (exactly as coco's "will retry later" path leaves
* them) until the mint comes back.
* them) until a later startup finds the mint reachable.
*/
import { normalizeMintUrl, type Manager } from "@cashu/coco-core";
import { logger } from "../../utils/logger";
/** Short probe: a mint that cannot answer /v1/info in 2s slows every op. */
export const MINT_PROBE_TIMEOUT_MS = 2_000;
/** How often stuck operations are re-checked (and dead mints re-probed). */
export const RECOVERY_RECHECK_INTERVAL_MS = 300_000;
type OpsApi = Manager["ops"];
/** Structural subset of the ops APIs used to enumerate and recover operations. */
export interface StuckOperationSource {
send: Pick<OpsApi["send"], "listInFlight" | "refresh">;
melt: Pick<OpsApi["melt"], "listInFlight" | "refresh">;
receive: Pick<OpsApi["receive"], "listInFlight" | "refresh">;
mint: Pick<OpsApi["mint"], "listInFlight" | "refresh">;
send: Pick<OpsApi["send"], "listInFlight" | "refresh" | "diagnostics">;
melt: Pick<OpsApi["melt"], "listInFlight" | "refresh" | "diagnostics">;
receive: Pick<OpsApi["receive"], "listInFlight" | "refresh" | "diagnostics">;
mint: Pick<OpsApi["mint"], "listInFlight" | "refresh" | "diagnostics">;
}
export type StuckOperationKind = "send" | "melt" | "receive" | "mint";
@@ -46,13 +44,12 @@ export interface StuckOperation {
* `mintOperationService` (see coco-client.ts).
*/
export interface SendRecoveryService {
tryRecoverInitOperation(op: unknown): Promise<void>;
tryRecoverExecutingOperation(op: unknown): Promise<void>;
recoverExecutingOperation(op: unknown): Promise<void>;
}
export interface RecoveryRunResult {
/** Operations at reachable mints whose recovery completed. */
recovered: number;
/** Operations for which recovery was attempted (not necessarily completed). */
attempted: number;
/** Operations skipped because their mint did not answer the probe. */
skipped: number;
/** Operations at reachable mints whose recovery still failed. */
@@ -89,8 +86,8 @@ function asStuckOperations(
raw: op,
});
} catch {
// Unparseable mint URL: keep the operation recoverable by treating it as
// reachable (probe only covers successfully normalized URLs).
// Preserve malformed persisted URLs; the probe will classify them as
// unreachable rather than attempting recovery against an invalid URL.
stuck.push({ kind, id: op.id, mintUrl: op.mintUrl, state: op.state, raw: op });
}
}
@@ -135,7 +132,7 @@ export async function probeMintReachability(
await Promise.all(
mintUrls.map(async (mintUrl) => {
try {
await fetcher(new URL("/v1/info", mintUrl).toString(), {
await fetcher(`${normalizeMintUrl(mintUrl)}/v1/info`, {
signal: AbortSignal.timeout(timeoutMs),
});
} catch (error) {
@@ -162,12 +159,8 @@ async function recoverStuckOperation(
// Public API: actively re-checks the proofs with the mint.
await source.send.refresh(op.id);
} else if (op.state === "executing") {
// No public per-op path exists for executing sends (coco keeps it
// private); tryRecover* swallows per-op errors and leaves the
// operation for the next pass, matching the global sweep's behavior.
await sendService.tryRecoverExecutingOperation(op.raw);
} else if (op.state === "init") {
await sendService.tryRecoverInitOperation(op.raw);
// Startup snapshot only; skip live operations in the driver below.
await sendService.recoverExecutingOperation(op.raw);
}
// prepared / rolling_back: the global sweep only warns; nothing to do.
return;
@@ -200,7 +193,7 @@ export async function runTargetedRecovery(
options: TargetedRecoveryOptions = {},
): Promise<RecoveryRunResult> {
const result: RecoveryRunResult = {
recovered: 0,
attempted: 0,
skipped: 0,
failed: 0,
skippedMints: new Map(),
@@ -228,9 +221,11 @@ export async function runTargetedRecovery(
for (const op of stuck) {
if (unreachable.has(op.mintUrl)) continue;
if (source[op.kind].diagnostics.isLocked(op.id)) continue;
if (op.kind === "send" && !["pending", "executing"].includes(op.state)) continue;
result.attempted++;
try {
await recoverStuckOperation(source, sendService, op);
result.recovered++;
} catch (error) {
// Same semantics as coco's tryRecover*: leave the operation for the
// next pass. A reachable mint can still reject a specific operation.
@@ -246,67 +241,3 @@ export async function runTargetedRecovery(
return result;
}
export interface RecoveryRecheckOptions extends TargetedRecoveryOptions {
intervalMs?: number;
/** Called when a mint transitions unreachable -> reachable with recovered op count. */
onMintBack?: (mintUrl: string) => void;
/** Called when a mint transitions reachable -> unreachable. */
onMintDown?: (mintUrl: string, opCount: number) => void;
}
/**
* Periodically re-run targeted recovery so a mint that comes back online has
* its parked operations recovered without a daemon restart, and operations
* that get stuck mid-session are reconciled too. Idle ticks (no stuck
* operations) cost one DB query and no network traffic. Logging is
* transition-based: a mint that stays dead produces no repeated output.
*
* Returns a stop function that waits for any in-flight tick.
*/
export function startRecoveryRecheck(
source: StuckOperationSource,
sendService: SendRecoveryService,
options: RecoveryRecheckOptions = {},
): () => Promise<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;
};
}