mirror of
https://github.com/Routstr/routstrd.git
synced 2026-10-05 20:38:22 +00:00
Merge pull request #116 from Routstr/mint-recovery-probe-restore
wallet: probe mint reachability, per-mint recovery gating
This commit is contained in:
@@ -572,6 +572,19 @@ export function createDaemonRequestHandler(deps: {
|
|||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if (req.method === "POST" && url.pathname === "/wallet/recover/operations") {
|
||||||
|
await respond(res, async () => {
|
||||||
|
if (!deps.walletClient.recoverStuckOperations) {
|
||||||
|
throw new CocodHttpError(
|
||||||
|
501,
|
||||||
|
"Stuck operation recovery is not supported by this wallet client.",
|
||||||
|
);
|
||||||
|
}
|
||||||
|
return { output: await deps.walletClient.recoverStuckOperations() };
|
||||||
|
});
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
if (req.method === "POST" && url.pathname === "/wallet/receive/cashu") {
|
if (req.method === "POST" && url.pathname === "/wallet/receive/cashu") {
|
||||||
await respond(res, async () => {
|
await respond(res, async () => {
|
||||||
const body = await readJsonBody(req);
|
const body = await readJsonBody(req);
|
||||||
|
|||||||
@@ -47,3 +47,53 @@ describe("POST /wallet/recover validation", () => {
|
|||||||
expect(recoverMintQuotes).toHaveBeenCalledTimes(1);
|
expect(recoverMintQuotes).toHaveBeenCalledTimes(1);
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
|
async function recoverOperations() {
|
||||||
|
const recoverStuckOperations = mock(async () => ({
|
||||||
|
attempted: 1,
|
||||||
|
busy: 0,
|
||||||
|
skipped: 0,
|
||||||
|
failed: 0,
|
||||||
|
skippedMints: {},
|
||||||
|
}));
|
||||||
|
const handler = createDaemonRequestHandler({ walletClient: { recoverStuckOperations } } as never);
|
||||||
|
const req = new EventEmitter() as any;
|
||||||
|
Object.assign(req, { method: "POST", url: "/wallet/recover/operations", headers: { host: "localhost" } });
|
||||||
|
const res = {
|
||||||
|
status: 0, body: "",
|
||||||
|
writeHead(status: number) { this.status = status; },
|
||||||
|
end(chunk: string) { this.body = chunk; },
|
||||||
|
};
|
||||||
|
setImmediate(() => {
|
||||||
|
req.emit("data", Buffer.from("{}"));
|
||||||
|
req.emit("end");
|
||||||
|
});
|
||||||
|
await handler(req, res as never);
|
||||||
|
return { res, recoverStuckOperations };
|
||||||
|
}
|
||||||
|
|
||||||
|
describe("POST /wallet/recover/operations", () => {
|
||||||
|
it("drives stuck-operation recovery and returns the summary", async () => {
|
||||||
|
const { res, recoverStuckOperations } = await recoverOperations();
|
||||||
|
expect(res.status).toBe(200);
|
||||||
|
expect(recoverStuckOperations).toHaveBeenCalledTimes(1);
|
||||||
|
expect(JSON.parse(res.body).output).toMatchObject({ attempted: 1, busy: 0 });
|
||||||
|
});
|
||||||
|
|
||||||
|
it("returns 501 when the wallet client does not support it", async () => {
|
||||||
|
const handler = createDaemonRequestHandler({ walletClient: {} } as never);
|
||||||
|
const req = new EventEmitter() as any;
|
||||||
|
Object.assign(req, { method: "POST", url: "/wallet/recover/operations", headers: { host: "localhost" } });
|
||||||
|
const res = {
|
||||||
|
status: 0, body: "",
|
||||||
|
writeHead(status: number) { this.status = status; },
|
||||||
|
end(chunk: string) { this.body = chunk; },
|
||||||
|
};
|
||||||
|
setImmediate(() => {
|
||||||
|
req.emit("data", Buffer.from("{}"));
|
||||||
|
req.emit("end");
|
||||||
|
});
|
||||||
|
await handler(req, res as never);
|
||||||
|
expect(res.status).toBe(501);
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|||||||
@@ -14,6 +14,7 @@ import {
|
|||||||
assertLegacyCocodNotRunning,
|
assertLegacyCocodNotRunning,
|
||||||
claimLegacyCocodPidFile,
|
claimLegacyCocodPidFile,
|
||||||
createCocoClient,
|
createCocoClient,
|
||||||
|
createRecoveryGate,
|
||||||
createRunQueue,
|
createRunQueue,
|
||||||
DEFAULT_TRUSTED_MINT_URLS,
|
DEFAULT_TRUSTED_MINT_URLS,
|
||||||
failExpiredMintQuoteIfUnpaid,
|
failExpiredMintQuoteIfUnpaid,
|
||||||
@@ -680,6 +681,29 @@ describe("settleExpiredMintQuotes", () => {
|
|||||||
expect(observePendingOperation).not.toHaveBeenCalled();
|
expect(observePendingOperation).not.toHaveBeenCalled();
|
||||||
expect(failPendingOperation).not.toHaveBeenCalled();
|
expect(failPendingOperation).not.toHaveBeenCalled();
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it("skips quotes at probe-unreachable mints without spending the budget", async () => {
|
||||||
|
const ops = [
|
||||||
|
pendingMintOp({ id: "dead-1", mintUrl: "https://dead.example.com" }),
|
||||||
|
pendingMintOp({ id: "dead-2", mintUrl: "https://dead.example.com/" }),
|
||||||
|
pendingMintOp({ id: "live-1", mintUrl: "https://live.example.com" }),
|
||||||
|
];
|
||||||
|
const { source, observePendingOperation, failPendingOperation } =
|
||||||
|
fakeSource(ops);
|
||||||
|
|
||||||
|
const result = await settleExpiredMintQuotes(
|
||||||
|
source,
|
||||||
|
NOW_MS,
|
||||||
|
undefined,
|
||||||
|
{ unreachableMints: new Set(["https://dead.example.com"]) },
|
||||||
|
);
|
||||||
|
|
||||||
|
expect(result).toEqual({ failed: 1, leftForRecovery: 0, unobserved: 2 });
|
||||||
|
// Only the reachable mint was asked anything.
|
||||||
|
expect(observePendingOperation).toHaveBeenCalledTimes(1);
|
||||||
|
expect(observePendingOperation.mock.calls[0]?.[0]).toBe("live-1");
|
||||||
|
expect(failPendingOperation).toHaveBeenCalledTimes(1);
|
||||||
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
describe("settlePendingMintQuotes", () => {
|
describe("settlePendingMintQuotes", () => {
|
||||||
@@ -1536,7 +1560,7 @@ describe("runMintQuoteRecovery", () => {
|
|||||||
|
|
||||||
it("skips operations whose earlier recovery is still in flight", async () => {
|
it("skips operations whose earlier recovery is still in flight", async () => {
|
||||||
const outstanding = new Map<string, Promise<unknown>>([
|
const outstanding = new Map<string, Promise<unknown>>([
|
||||||
["op-1", new Promise(() => {})],
|
["mint:op-1", new Promise(() => {})],
|
||||||
]);
|
]);
|
||||||
const { source, finalize, observePendingOperation } = fakeSource(
|
const { source, finalize, observePendingOperation } = fakeSource(
|
||||||
[mintOp()],
|
[mintOp()],
|
||||||
@@ -1567,14 +1591,14 @@ describe("runMintQuoteRecovery", () => {
|
|||||||
});
|
});
|
||||||
|
|
||||||
expect(first).toMatchObject({ retryable: 1, recovered: 0 });
|
expect(first).toMatchObject({ retryable: 1, recovered: 0 });
|
||||||
expect(outstanding.has("op-1")).toBe(true);
|
expect(outstanding.has("mint:op-1")).toBe(true);
|
||||||
// The abandoned mint request must not be retried underneath.
|
// The abandoned mint request must not be retried underneath.
|
||||||
expect(second).toMatchObject({ busy: 1, checked: 0 });
|
expect(second).toMatchObject({ busy: 1, checked: 0 });
|
||||||
});
|
});
|
||||||
|
|
||||||
it("does not re-open a failed operation whose recovery is in flight", async () => {
|
it("does not re-open a failed operation whose recovery is in flight", async () => {
|
||||||
const outstanding = new Map<string, Promise<unknown>>([
|
const outstanding = new Map<string, Promise<unknown>>([
|
||||||
["op-1", new Promise(() => {})],
|
["mint:op-1", new Promise(() => {})],
|
||||||
]);
|
]);
|
||||||
const { source, reopenFailedOperation } = fakeSource(
|
const { source, reopenFailedOperation } = fakeSource(
|
||||||
[mintOp({ state: "failed" })],
|
[mintOp({ state: "failed" })],
|
||||||
@@ -1629,8 +1653,78 @@ describe("runMintQuoteRecovery", () => {
|
|||||||
});
|
});
|
||||||
|
|
||||||
expect(first).toMatchObject({ retryable: 1, checked: 1 });
|
expect(first).toMatchObject({ retryable: 1, checked: 1 });
|
||||||
expect(outstanding.has("op-1")).toBe(true);
|
expect(outstanding.has("mint:op-1")).toBe(true);
|
||||||
expect(second).toMatchObject({ busy: 1, checked: 0 });
|
expect(second).toMatchObject({ busy: 1, checked: 0 });
|
||||||
expect(finalize).not.toHaveBeenCalled();
|
expect(finalize).not.toHaveBeenCalled();
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
|
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;
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|||||||
@@ -1,3 +1,4 @@
|
|||||||
|
import { recoveryKey, trackRecovery, drainRecoveryWork, waitForRecoveryWork, createRecoveryDisposer, type RecoveryWork } from "./recovery-work";
|
||||||
import {
|
import {
|
||||||
Manager,
|
Manager,
|
||||||
OperationInProgressError,
|
OperationInProgressError,
|
||||||
@@ -42,6 +43,13 @@ import {
|
|||||||
selectMintQuotesForRecovery,
|
selectMintQuotesForRecovery,
|
||||||
type MintQuoteRecoveryCandidate,
|
type MintQuoteRecoveryCandidate,
|
||||||
} from "./mint-quote-recovery";
|
} from "./mint-quote-recovery";
|
||||||
|
import {
|
||||||
|
collectStuckOperations,
|
||||||
|
probeMintReachability,
|
||||||
|
runTargetedRecovery,
|
||||||
|
type SendRecoveryService,
|
||||||
|
type StuckOperation,
|
||||||
|
} from "./recovery-probe";
|
||||||
import {
|
import {
|
||||||
clearInterruptedReceiveReservations,
|
clearInterruptedReceiveReservations,
|
||||||
deleteReceiveTokenReservation,
|
deleteReceiveTokenReservation,
|
||||||
@@ -724,6 +732,7 @@ const EXPIRED_MINT_OBSERVATION_DEADLINE_MS = 15_000;
|
|||||||
|
|
||||||
/** Rejects when `timeoutMs` elapses before `promise` settles. */
|
/** Rejects when `timeoutMs` elapses before `promise` settles. */
|
||||||
function withTimeout<T>(promise: Promise<T>, timeoutMs: number): Promise<T> {
|
function withTimeout<T>(promise: Promise<T>, timeoutMs: number): Promise<T> {
|
||||||
|
if (timeoutMs === Infinity) return promise;
|
||||||
let timer: ReturnType<typeof setTimeout> | undefined;
|
let timer: ReturnType<typeof setTimeout> | undefined;
|
||||||
const timeout = new Promise<never>((_resolve, reject) => {
|
const timeout = new Promise<never>((_resolve, reject) => {
|
||||||
timer = setTimeout(
|
timer = setTimeout(
|
||||||
@@ -838,6 +847,7 @@ export async function settleExpiredMintQuotes(
|
|||||||
source: ExpiredMintQuoteSource,
|
source: ExpiredMintQuoteSource,
|
||||||
nowMs: number,
|
nowMs: number,
|
||||||
deadlineMs: number = EXPIRED_MINT_OBSERVATION_DEADLINE_MS,
|
deadlineMs: number = EXPIRED_MINT_OBSERVATION_DEADLINE_MS,
|
||||||
|
options: { unreachableMints?: Set<string>; outstanding?: RecoveryWork; shouldStop?: () => boolean } = {},
|
||||||
): Promise<ExpiredMintSettlement> {
|
): Promise<ExpiredMintSettlement> {
|
||||||
const pendingMints = await source.ops.mint.listPending();
|
const pendingMints = await source.ops.mint.listPending();
|
||||||
const selection = selectCleanupOperations({
|
const selection = selectCleanupOperations({
|
||||||
@@ -858,6 +868,25 @@ export async function settleExpiredMintQuotes(
|
|||||||
|
|
||||||
const startedAt = Date.now();
|
const startedAt = Date.now();
|
||||||
for (const op of candidates) {
|
for (const op of candidates) {
|
||||||
|
if (options.shouldStop?.() || options.outstanding?.has(recoveryKey("mint", op.id))) {
|
||||||
|
settlement.unobserved++; continue;
|
||||||
|
}
|
||||||
|
// A mint the startup probe already found unreachable cannot answer an
|
||||||
|
// observation either; skip it without spending the shared wall-clock
|
||||||
|
// budget, leaving the quote pending for a later startup.
|
||||||
|
if (options.unreachableMints) {
|
||||||
|
let mintUrl = op.mintUrl;
|
||||||
|
try {
|
||||||
|
mintUrl = normalizeMintUrl(op.mintUrl);
|
||||||
|
} catch {
|
||||||
|
// Malformed persisted URL: probe keys are raw for those, and the
|
||||||
|
// observation below would fail anyway, landing in `unobserved`.
|
||||||
|
}
|
||||||
|
if (options.unreachableMints.has(mintUrl)) {
|
||||||
|
settlement.unobserved++;
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
}
|
||||||
const remainingMs = deadlineMs - (Date.now() - startedAt);
|
const remainingMs = deadlineMs - (Date.now() - startedAt);
|
||||||
if (remainingMs <= 0) {
|
if (remainingMs <= 0) {
|
||||||
const skipped =
|
const skipped =
|
||||||
@@ -872,11 +901,12 @@ export async function settleExpiredMintQuotes(
|
|||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
|
|
||||||
const check = await failExpiredMintQuoteIfUnpaid(
|
// Track the complete observation/failure chain, not just its bounded wait.
|
||||||
source.mintOperationService,
|
const work = failExpiredMintQuoteIfUnpaid(source.mintOperationService, op.id, Infinity);
|
||||||
op.id,
|
if (options.outstanding) trackRecovery(options.outstanding, recoveryKey("mint", op.id), work);
|
||||||
remainingMs,
|
const check = await waitForRecoveryWork(work, remainingMs).catch(error => ({
|
||||||
);
|
outcome: "unobserved" as const, error, category: undefined,
|
||||||
|
}));
|
||||||
if (check.outcome === "failed") {
|
if (check.outcome === "failed") {
|
||||||
settlement.failed++;
|
settlement.failed++;
|
||||||
} else if (check.outcome === "leftForRecovery") {
|
} else if (check.outcome === "leftForRecovery") {
|
||||||
@@ -937,6 +967,7 @@ export interface MintQuoteRecoverySource {
|
|||||||
}
|
}
|
||||||
|
|
||||||
export interface MintQuoteRecoveryOptions {
|
export interface MintQuoteRecoveryOptions {
|
||||||
|
shouldStop?: () => boolean;
|
||||||
/** Target only these operation ids (may include failed operations). */
|
/** Target only these operation ids (may include failed operations). */
|
||||||
operationIds?: string[];
|
operationIds?: string[];
|
||||||
/** Per-quote budget for observing the mint and finalizing the operation. */
|
/** Per-quote budget for observing the mint and finalizing the operation. */
|
||||||
@@ -948,7 +979,7 @@ export interface MintQuoteRecoveryOptions {
|
|||||||
*/
|
*/
|
||||||
includeFailed?: boolean;
|
includeFailed?: boolean;
|
||||||
/**
|
/**
|
||||||
* In-flight recovery work keyed by operation id, shared across runs.
|
* In-flight recovery work keyed by family:id (mint:<operation id>), shared across runs.
|
||||||
* withTimeout does not cancel the underlying request, so a timed-out quote
|
* withTimeout does not cancel the underlying request, so a timed-out quote
|
||||||
* check or finalize must keep blocking a retry until it actually settles.
|
* check or finalize must keep blocking a retry until it actually settles.
|
||||||
*/
|
*/
|
||||||
@@ -987,9 +1018,9 @@ export interface MintQuoteRecoveryResult {
|
|||||||
* underneath work that outlived its timeout. A rejected task never breaks the
|
* underneath work that outlived its timeout. A rejected task never breaks the
|
||||||
* chain for the next one.
|
* chain for the next one.
|
||||||
*/
|
*/
|
||||||
export function createRunQueue(): <T>(run: () => Promise<T>) => Promise<T> {
|
export function createRunQueue(): (<T>(run: () => Promise<T>) => Promise<T>) & { drain(): Promise<void> } {
|
||||||
let tail: Promise<unknown> = Promise.resolve();
|
let tail: Promise<unknown> = Promise.resolve();
|
||||||
return <T>(run: () => Promise<T>): Promise<T> => {
|
const enqueue = <T>(run: () => Promise<T>): Promise<T> => {
|
||||||
const result = tail.then(run, run);
|
const result = tail.then(run, run);
|
||||||
tail = result.then(
|
tail = result.then(
|
||||||
() => undefined,
|
() => undefined,
|
||||||
@@ -997,6 +1028,7 @@ export function createRunQueue(): <T>(run: () => Promise<T>) => Promise<T> {
|
|||||||
);
|
);
|
||||||
return result;
|
return result;
|
||||||
};
|
};
|
||||||
|
return Object.assign(enqueue, { drain: async () => { await tail; } });
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Per-quote budget for the mint round-trip during explicit recovery. */
|
/** Per-quote budget for the mint round-trip during explicit recovery. */
|
||||||
@@ -1064,21 +1096,8 @@ export async function runMintQuoteRecovery(
|
|||||||
/** coco's fail-fast operation lock rejected the call: another holder exists. */
|
/** coco's fail-fast operation lock rejected the call: another holder exists. */
|
||||||
const isInProgress = (error: unknown) =>
|
const isInProgress = (error: unknown) =>
|
||||||
error instanceof Error && error.name === "OperationInProgressError";
|
error instanceof Error && error.name === "OperationInProgressError";
|
||||||
/**
|
const track = (operationId: string, work: Promise<unknown>) =>
|
||||||
* Register in-flight work for an operation. Entries are cleared only once the
|
trackRecovery(outstanding, recoveryKey("mint", operationId), work);
|
||||||
* work actually settles (withTimeout does not cancel the request behind it),
|
|
||||||
* so a timed-out call keeps blocking a retry. The identity check stops a late
|
|
||||||
* settlement from clearing a newer entry for the same operation.
|
|
||||||
*/
|
|
||||||
const track = (operationId: string, work: Promise<unknown>) => {
|
|
||||||
outstanding.set(operationId, work);
|
|
||||||
const clear = () => {
|
|
||||||
if (outstanding.get(operationId) === work) {
|
|
||||||
outstanding.delete(operationId);
|
|
||||||
}
|
|
||||||
};
|
|
||||||
void work.then(clear, clear);
|
|
||||||
};
|
|
||||||
|
|
||||||
let targets: MintQuoteRecoveryCandidate[];
|
let targets: MintQuoteRecoveryCandidate[];
|
||||||
if (options.operationIds && options.operationIds.length > 0) {
|
if (options.operationIds && options.operationIds.length > 0) {
|
||||||
@@ -1108,8 +1127,9 @@ export async function runMintQuoteRecovery(
|
|||||||
});
|
});
|
||||||
|
|
||||||
for (const op of failed) {
|
for (const op of failed) {
|
||||||
|
if (options.shouldStop?.()) break;
|
||||||
const label = `Mint quote ${op.quoteId ?? op.id} at ${op.mintUrl}`;
|
const label = `Mint quote ${op.quoteId ?? op.id} at ${op.mintUrl}`;
|
||||||
if (outstanding.has(op.id)) {
|
if (outstanding.has(recoveryKey("mint", op.id))) {
|
||||||
result.busy++;
|
result.busy++;
|
||||||
onProgress?.(`${label}: an earlier recovery is still running; skipped`);
|
onProgress?.(`${label}: an earlier recovery is still running; skipped`);
|
||||||
continue;
|
continue;
|
||||||
@@ -1137,13 +1157,16 @@ export async function runMintQuoteRecovery(
|
|||||||
await recoverOne(op);
|
await recoverOne(op);
|
||||||
}
|
}
|
||||||
|
|
||||||
for (const op of pending) await recoverOne(op);
|
for (const op of pending) {
|
||||||
|
if (options.shouldStop?.()) break;
|
||||||
|
await recoverOne(op);
|
||||||
|
}
|
||||||
|
|
||||||
return result;
|
return result;
|
||||||
|
|
||||||
async function recoverOne(op: MintQuoteRecoveryCandidate): Promise<void> {
|
async function recoverOne(op: MintQuoteRecoveryCandidate): Promise<void> {
|
||||||
const label = `Mint quote ${op.quoteId ?? op.id} at ${op.mintUrl}`;
|
const label = `Mint quote ${op.quoteId ?? op.id} at ${op.mintUrl}`;
|
||||||
if (outstanding.has(op.id)) {
|
if (outstanding.has(recoveryKey("mint", op.id))) {
|
||||||
result.busy++;
|
result.busy++;
|
||||||
onProgress?.(`${label}: an earlier recovery is still running; skipped`);
|
onProgress?.(`${label}: an earlier recovery is still running; skipped`);
|
||||||
return;
|
return;
|
||||||
@@ -1298,6 +1321,7 @@ export interface PendingMintSweepOptions {
|
|||||||
deadlineMs?: number;
|
deadlineMs?: number;
|
||||||
checkTimeoutMs?: number;
|
checkTimeoutMs?: number;
|
||||||
state?: PendingMintSweepState;
|
state?: PendingMintSweepState;
|
||||||
|
shouldStop?: () => boolean;
|
||||||
}
|
}
|
||||||
|
|
||||||
type PendingMintOutcome = "unreachable" | "other";
|
type PendingMintOutcome = "unreachable" | "other";
|
||||||
@@ -1319,21 +1343,21 @@ export async function settlePendingMintQuotes(
|
|||||||
const ordered = [...pending.slice(resumeAt), ...pending.slice(0, resumeAt)];
|
const ordered = [...pending.slice(resumeAt), ...pending.slice(0, resumeAt)];
|
||||||
const startedAt = Date.now();
|
const startedAt = Date.now();
|
||||||
for (const op of ordered) {
|
for (const op of ordered) {
|
||||||
|
if (options.shouldStop?.()) break;
|
||||||
const remainingMs = deadlineMs - (Date.now() - startedAt);
|
const remainingMs = deadlineMs - (Date.now() - startedAt);
|
||||||
if (state.outstanding.has(op.id) || remainingMs <= 0) {
|
if (state.outstanding.has(recoveryKey("mint", op.id)) || remainingMs <= 0) {
|
||||||
unreachable++;
|
unreachable++;
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
state.after = op.id;
|
state.after = op.id;
|
||||||
// Report on the refresh itself so a late result is still logged.
|
// Track refresh AND its late-result reporting mutations as one lifetime.
|
||||||
const settled = source.ops.mint
|
const settled = source.ops.mint
|
||||||
.refresh(op.id)
|
.refresh(op.id)
|
||||||
.then(
|
.then(
|
||||||
(result) => reportPendingMintRefresh(source, op, result, nowMs),
|
(result) => reportPendingMintRefresh(source, op, result, nowMs),
|
||||||
(error) => reportPendingMintRefreshError(source, op, error),
|
(error) => reportPendingMintRefreshError(source, op, error),
|
||||||
)
|
);
|
||||||
.finally(() => state.outstanding.delete(op.id));
|
trackRecovery(state.outstanding, recoveryKey("mint", op.id), settled);
|
||||||
state.outstanding.set(op.id, settled);
|
|
||||||
try {
|
try {
|
||||||
const outcome = await withTimeout(settled, Math.min(checkTimeoutMs, remainingMs));
|
const outcome = await withTimeout(settled, Math.min(checkTimeoutMs, remainingMs));
|
||||||
if (outcome === "unreachable") unreachable++;
|
if (outcome === "unreachable") unreachable++;
|
||||||
@@ -1409,16 +1433,16 @@ async function reportPendingMintRefreshError(
|
|||||||
* still running after that fails against the closed database and is picked
|
* still running after that fails against the closed database and is picked
|
||||||
* up by startup recovery.
|
* up by startup recovery.
|
||||||
*/
|
*/
|
||||||
function startPendingMintSweep(source: PendingMintQuoteSource): () => Promise<void> {
|
function startPendingMintSweep(source: PendingMintQuoteSource, outstanding: RecoveryWork): () => Promise<void> {
|
||||||
let stopped = false;
|
let stopped = false;
|
||||||
let timer: ReturnType<typeof setTimeout> | undefined;
|
let timer: ReturnType<typeof setTimeout> | undefined;
|
||||||
let inFlight: Promise<void> = Promise.resolve();
|
let inFlight: Promise<void> = Promise.resolve();
|
||||||
let unreachableBefore = 0;
|
let unreachableBefore = 0;
|
||||||
const state: PendingMintSweepState = { outstanding: new Map() };
|
const state: PendingMintSweepState = { outstanding };
|
||||||
|
|
||||||
const tick = async () => {
|
const tick = async () => {
|
||||||
if (stopped) return;
|
if (stopped) return;
|
||||||
inFlight = settlePendingMintQuotes(source, Date.now(), { state }).then(
|
inFlight = settlePendingMintQuotes(source, Date.now(), { state, shouldStop: () => stopped }).then(
|
||||||
({ unreachable }) => {
|
({ unreachable }) => {
|
||||||
// Report a mint becoming unreachable, or reachable again, once.
|
// Report a mint becoming unreachable, or reachable again, once.
|
||||||
if (unreachable > 0 && unreachableBefore === 0) {
|
if (unreachable > 0 && unreachableBefore === 0) {
|
||||||
@@ -1534,21 +1558,174 @@ 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;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** 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>;
|
||||||
|
} | undefined>;
|
||||||
|
const families = [
|
||||||
|
["send", repo.sendOperationRepository],
|
||||||
|
["melt", repo.meltOperationRepository],
|
||||||
|
["receive", repo.receiveOperationRepository],
|
||||||
|
["mint", repo.mintOperationRepository],
|
||||||
|
] as const;
|
||||||
|
// Fail closed the way reopenFailedMintOperation does: a coco bump that
|
||||||
|
// removes or renames these privates must stop recovery with a clear error
|
||||||
|
// before anything is written, not crash halfway through the loop with the
|
||||||
|
// cleanup half-applied.
|
||||||
|
for (const [kind] of families) {
|
||||||
|
if (typeof services[`${kind}OperationService`]?.recoverInitOperation !== "function") {
|
||||||
|
throw new Error(
|
||||||
|
`coco ${kind}OperationService.recoverInitOperation is unavailable; refusing local recovery cleanup`,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if (typeof services.sendOperationService?.cleanupOrphanedReservations !== "function") {
|
||||||
|
throw new Error(
|
||||||
|
"coco sendOperationService.cleanupOrphanedReservations is unavailable; refusing local recovery cleanup",
|
||||||
|
);
|
||||||
|
}
|
||||||
|
for (const [kind, repository] of families) {
|
||||||
|
for (const op of await repository.getByState("init")) {
|
||||||
|
await services[`${kind}OperationService`]!.recoverInitOperation!(op);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
await services.sendOperationService!.cleanupOrphanedReservations!();
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Gate for value-moving wallet operations while startup recovery runs.
|
||||||
|
*
|
||||||
|
* 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.
|
||||||
|
*/
|
||||||
|
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.
|
||||||
*
|
*
|
||||||
* Expired mint quotes are settled first: quotes their mint confirms as unpaid
|
* Probe first, then settle expired mint quotes: quotes confirmed as unpaid
|
||||||
* are failed locally so `recoverPendingMintOperations()` skips them, while
|
* are failed locally so `recoverPendingMintOperations()` skips them, while
|
||||||
* paid/issued and unreachable-mint quotes stay pending for the sweep.
|
* paid/issued and unreachable-mint quotes stay pending for the sweep.
|
||||||
*/
|
*/
|
||||||
async function runWalletRecovery(
|
export async function runWalletRecovery(
|
||||||
coco: Manager,
|
coco: Manager,
|
||||||
onProgress: (progress: RecoveryPhaseProgress) => void,
|
onProgress: (progress: RecoveryPhaseProgress) => void,
|
||||||
receiveOperationIds?: string[],
|
receiveOperationIds?: string[],
|
||||||
|
onStuckMintsKnown?: (mints: Set<string>) => void,
|
||||||
|
options: { cleanupLocalState?: () => Promise<void>; fetchImpl?: typeof fetch; outstanding?: RecoveryWork; shouldStop?: () => boolean } = {},
|
||||||
): Promise<void> {
|
): Promise<void> {
|
||||||
surfacingRecoveryProgress = true;
|
surfacingRecoveryProgress = true;
|
||||||
let failedMintQuotes = 0;
|
let failedMintQuotes = 0;
|
||||||
try {
|
try {
|
||||||
|
// Probe every mint that has stuck operations once, FIRST, so a dead mint
|
||||||
|
// costs a single short probe instead of taxing settlement's observation
|
||||||
|
// budget plus a network timeout per operation per sweep. Healthy-mint
|
||||||
|
// operations are recovered per op; dead-mint operations stay parked
|
||||||
|
// exactly as coco's own "will retry later" path would leave them.
|
||||||
|
onProgress({ phase: "Probing mints", failedMintQuotes });
|
||||||
|
const stuckOperations = await collectStuckOperations(coco.ops);
|
||||||
|
const unreachableMints = await probeMintReachability(
|
||||||
|
[...new Set(stuckOperations.map((op) => op.mintUrl))],
|
||||||
|
{ fetchImpl: options.fetchImpl },
|
||||||
|
);
|
||||||
|
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;
|
||||||
|
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)));
|
||||||
|
}
|
||||||
|
|
||||||
|
// Settlement runs after the gate opens and only spends its observation
|
||||||
|
// budget on mints the probe found reachable; dead-mint quotes stay
|
||||||
|
// pending untouched. It only reads and locally fails long-expired quotes,
|
||||||
|
// so it cannot conflict with live operations the gate just admitted.
|
||||||
onProgress({ phase: "Settling expired mint quotes", failedMintQuotes });
|
onProgress({ phase: "Settling expired mint quotes", failedMintQuotes });
|
||||||
const settlement = await settleExpiredMintQuotes(
|
const settlement = await settleExpiredMintQuotes(
|
||||||
{
|
{
|
||||||
@@ -1560,6 +1737,8 @@ async function runWalletRecovery(
|
|||||||
).mintOperationService,
|
).mintOperationService,
|
||||||
},
|
},
|
||||||
Date.now(),
|
Date.now(),
|
||||||
|
undefined,
|
||||||
|
{ unreachableMints, outstanding: options.outstanding, shouldStop: options.shouldStop },
|
||||||
);
|
);
|
||||||
failedMintQuotes = settlement.failed;
|
failedMintQuotes = settlement.failed;
|
||||||
if (settlement.leftForRecovery > 0 || settlement.unobserved > 0) {
|
if (settlement.leftForRecovery > 0 || settlement.unobserved > 0) {
|
||||||
@@ -1570,22 +1749,48 @@ async function runWalletRecovery(
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
onProgress({ phase: "Settled expired mint quotes", failedMintQuotes });
|
onProgress({ phase: "Settled expired mint quotes", failedMintQuotes });
|
||||||
|
const targeted = (kinds: Array<StuckOperation["kind"]>) =>
|
||||||
|
runTargetedRecovery(coco.ops, sendRecoveryServiceOf(coco), {
|
||||||
|
kinds,
|
||||||
|
outstanding: options.outstanding,
|
||||||
|
shouldStop: options.shouldStop,
|
||||||
|
stuckOperations,
|
||||||
|
unreachableMints,
|
||||||
|
});
|
||||||
|
|
||||||
|
// Happy path (every mint reachable) keeps coco's global sweeps: they also
|
||||||
|
// 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 });
|
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) {
|
||||||
|
if (options.shouldStop?.()) break;
|
||||||
|
const mintUrl = mintByOperation.get(operationId);
|
||||||
|
if (mintUrl && unreachableMints.has(mintUrl)) continue;
|
||||||
try {
|
try {
|
||||||
await withTimeout(coco.ops.receive.refresh(operationId), 15_000);
|
const work = coco.ops.receive.refresh(operationId);
|
||||||
|
if (options.outstanding) trackRecovery(options.outstanding, recoveryKey("receive", operationId), work);
|
||||||
|
await withTimeout(work, 15_000);
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
logger.warn("Targeted receive recovery did not complete", {
|
logger.warn("Targeted receive recovery did not complete", {
|
||||||
operationId,
|
operationId,
|
||||||
@@ -1593,12 +1798,19 @@ 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 });
|
||||||
|
// A settlement wait may have timed out while an unlocked observation
|
||||||
|
// still runs. Never let a fresh global mint sweep observe it again.
|
||||||
|
if (!degraded && ![...(options.outstanding?.keys() ?? [])].some(key => key.startsWith("mint:"))) {
|
||||||
await coco.recoverPendingMintOperations();
|
await coco.recoverPendingMintOperations();
|
||||||
|
}
|
||||||
|
else await targeted(["mint"]);
|
||||||
|
|
||||||
onProgress({ phase: "done", failedMintQuotes });
|
onProgress({ phase: "done", failedMintQuotes });
|
||||||
} finally {
|
} finally {
|
||||||
@@ -1661,6 +1873,10 @@ export async function createCocoClient(
|
|||||||
const recoveryPromise = new Promise<void>((resolve) => {
|
const recoveryPromise = new Promise<void>((resolve) => {
|
||||||
recoveryResolve = resolve;
|
recoveryResolve = resolve;
|
||||||
});
|
});
|
||||||
|
const recoveryGate = createRecoveryGate();
|
||||||
|
let disposed = false;
|
||||||
|
const enqueueRecovery = createRunQueue();
|
||||||
|
const recoveryOutstanding: RecoveryWork = new Map();
|
||||||
|
|
||||||
try {
|
try {
|
||||||
startupProgress("Opening Cashu wallet database...");
|
startupProgress("Opening Cashu wallet database...");
|
||||||
@@ -1858,11 +2074,14 @@ export async function createCocoClient(
|
|||||||
}
|
}
|
||||||
},
|
},
|
||||||
receiveRecoveryOperationIds,
|
receiveRecoveryOperationIds,
|
||||||
|
(mints) => recoveryGate.publishStuckMints(mints),
|
||||||
|
{ cleanupLocalState: () => cleanupLocalRecoveryState(coco!, repo), outstanding: recoveryOutstanding, shouldStop: () => disposed },
|
||||||
)
|
)
|
||||||
.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.");
|
||||||
stopPendingMintSweep = startPendingMintSweep({
|
stopPendingMintSweep = startPendingMintSweep({
|
||||||
@@ -1871,12 +2090,13 @@ export async function createCocoClient(
|
|||||||
mintOperationService: (
|
mintOperationService: (
|
||||||
coco as unknown as { mintOperationService: MintOperationServiceCleanup }
|
coco as unknown as { mintOperationService: MintOperationServiceCleanup }
|
||||||
).mintOperationService,
|
).mintOperationService,
|
||||||
});
|
}, recoveryOutstanding);
|
||||||
})
|
})
|
||||||
.catch((error) => {
|
.catch((error) => {
|
||||||
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}`);
|
||||||
});
|
});
|
||||||
@@ -1899,23 +2119,33 @@ export async function createCocoClient(
|
|||||||
return api;
|
return api;
|
||||||
};
|
};
|
||||||
|
|
||||||
let disposed = false;
|
const assertOpen = () => { if (disposed) throw new Error("Wallet is shutting down"); };
|
||||||
// Explicit recovery runs are serialized, and finalize work that outlives its
|
|
||||||
// timeout stays in the map so a retry waits for it.
|
|
||||||
const enqueueRecovery = createRunQueue();
|
|
||||||
const recoveryOutstanding = new Map<string, Promise<unknown>>();
|
|
||||||
|
|
||||||
/**
|
// 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 = async (mintUrl?: string): Promise<void> => {
|
||||||
const waitForRecovery = async (): Promise<void> => {
|
assertOpen();
|
||||||
if (!recoveryDone) await recoveryPromise;
|
await recoveryGate.waitForRecovery(mintUrl);
|
||||||
if (recoveryError) {
|
assertOpen();
|
||||||
throw new Error(`Wallet is not ready: ${recoveryError}`);
|
|
||||||
}
|
|
||||||
};
|
};
|
||||||
|
|
||||||
|
const disposeRecovery = createRecoveryDisposer(
|
||||||
|
() => { disposed = true; },
|
||||||
|
async () => {
|
||||||
|
await recoveryPromise;
|
||||||
|
await stopPendingMintSweep?.();
|
||||||
|
await enqueueRecovery.drain();
|
||||||
|
await drainRecoveryWork(recoveryOutstanding);
|
||||||
|
},
|
||||||
|
async () => {
|
||||||
|
await coco.dispose();
|
||||||
|
database.close();
|
||||||
|
releaseLegacyPidClaim();
|
||||||
|
releaseWalletPidClaim();
|
||||||
|
},
|
||||||
|
);
|
||||||
|
|
||||||
return {
|
return {
|
||||||
async ping(): Promise<boolean> {
|
async ping(): Promise<boolean> {
|
||||||
try {
|
try {
|
||||||
@@ -2094,13 +2324,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,
|
||||||
@@ -2126,13 +2356,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,
|
||||||
@@ -2142,13 +2372,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",
|
||||||
@@ -2192,22 +2422,10 @@ export async function createCocoClient(
|
|||||||
},
|
},
|
||||||
|
|
||||||
async dispose(): Promise<void> {
|
async dispose(): Promise<void> {
|
||||||
if (disposed) return;
|
await disposeRecovery().catch(error => {
|
||||||
disposed = true;
|
logger.warn("Wallet shutdown incomplete; database and ownership retained until recovery settles");
|
||||||
try {
|
throw error;
|
||||||
// Let any in-flight recovery settle before closing the database from
|
});
|
||||||
// underneath it. The recovery promise resolves on success or failure.
|
|
||||||
await recoveryPromise;
|
|
||||||
await stopPendingMintSweep?.();
|
|
||||||
await coco.dispose();
|
|
||||||
} finally {
|
|
||||||
try {
|
|
||||||
database.close();
|
|
||||||
} finally {
|
|
||||||
releaseLegacyPidClaim();
|
|
||||||
releaseWalletPidClaim();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
},
|
},
|
||||||
|
|
||||||
async getHistory(offset?: number, limit?: number): Promise<HistoryEntry[]> {
|
async getHistory(offset?: number, limit?: number): Promise<HistoryEntry[]> {
|
||||||
@@ -2402,18 +2620,45 @@ export async function createCocoClient(
|
|||||||
// Serialize explicit recovery: two concurrent requests must not both
|
// Serialize explicit recovery: two concurrent requests must not both
|
||||||
// snapshot the same failed operation, and a retry must not start
|
// snapshot the same failed operation, and a retry must not start
|
||||||
// underneath a finalize that outlived its timeout.
|
// underneath a finalize that outlived its timeout.
|
||||||
return enqueueRecovery(() =>
|
return enqueueRecovery(() => {
|
||||||
runMintQuoteRecovery(
|
assertOpen();
|
||||||
|
return runMintQuoteRecovery(
|
||||||
{
|
{
|
||||||
ops: coco.ops as unknown as MintQuoteRecoverySource["ops"],
|
ops: coco.ops as unknown as MintQuoteRecoverySource["ops"],
|
||||||
mintOperationService: service,
|
mintOperationService: service,
|
||||||
reopenFailedOperation: (operationId) =>
|
reopenFailedOperation: (operationId) =>
|
||||||
reopenFailedMintOperation(service, operationId),
|
reopenFailedMintOperation(service, operationId),
|
||||||
},
|
},
|
||||||
{ ...options, outstanding: recoveryOutstanding },
|
{ ...options, outstanding: recoveryOutstanding, shouldStop: () => disposed },
|
||||||
onProgress,
|
onProgress,
|
||||||
),
|
|
||||||
);
|
);
|
||||||
|
});
|
||||||
|
},
|
||||||
|
|
||||||
|
async recoverStuckOperations() {
|
||||||
|
await waitForRecovery();
|
||||||
|
// Serialized against explicit mint-quote recovery (and itself) through
|
||||||
|
// the same queue and lifetime tracker, so timed-out passes cannot retry the same
|
||||||
|
// operation. Receive stays startup-only: recovering competing receives
|
||||||
|
// safely requires the startup dedup classification (receive-dedup.ts).
|
||||||
|
// Operations a live execute holds come back as busy via coco's
|
||||||
|
// fail-fast operation lock, never driven underneath it.
|
||||||
|
const result = await enqueueRecovery(() => {
|
||||||
|
assertOpen();
|
||||||
|
return runTargetedRecovery(coco!.ops, sendRecoveryServiceOf(coco!), {
|
||||||
|
kinds: ["send", "melt", "mint"],
|
||||||
|
outstanding: recoveryOutstanding,
|
||||||
|
shouldStop: () => disposed,
|
||||||
|
});
|
||||||
|
});
|
||||||
|
return {
|
||||||
|
timedOut: result.timedOut,
|
||||||
|
attempted: result.attempted,
|
||||||
|
busy: result.busy,
|
||||||
|
skipped: result.skipped,
|
||||||
|
failed: result.failed,
|
||||||
|
skippedMints: Object.fromEntries(result.skippedMints),
|
||||||
|
};
|
||||||
},
|
},
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -159,6 +159,22 @@ export interface WalletMintQuoteRecoveryResult {
|
|||||||
errors: Array<{ operationId: string; error: string }>;
|
errors: Array<{ operationId: string; error: string }>;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** Summary of a stuck-operation (send/melt/mint) recovery run. */
|
||||||
|
export interface WalletStuckOperationRecoveryResult {
|
||||||
|
/** Timed-out waits; the underlying operation remains tracked. */
|
||||||
|
timedOut: number;
|
||||||
|
/** Operations for which recovery was attempted (not necessarily completed). */
|
||||||
|
attempted: number;
|
||||||
|
/** Locked operations or unfinished work from another pass; retry later. */
|
||||||
|
busy: number;
|
||||||
|
/** Operations skipped for unreachable mints, shutdown, or pass budget exhaustion. */
|
||||||
|
skipped: number;
|
||||||
|
/** Operations at reachable mints whose recovery still failed. */
|
||||||
|
failed: number;
|
||||||
|
/** Unreachable mint URL -> number of operations skipped there. */
|
||||||
|
skippedMints: Record<string, number>;
|
||||||
|
}
|
||||||
|
|
||||||
export interface CocodClient {
|
export interface CocodClient {
|
||||||
ping(): Promise<boolean>;
|
ping(): Promise<boolean>;
|
||||||
getStatus(): Promise<CocodState>;
|
getStatus(): Promise<CocodState>;
|
||||||
@@ -201,6 +217,12 @@ export interface CocodClient {
|
|||||||
options?: WalletMintQuoteRecoveryOptions,
|
options?: WalletMintQuoteRecoveryOptions,
|
||||||
onProgress?: (message: string) => void,
|
onProgress?: (message: string) => void,
|
||||||
): Promise<WalletMintQuoteRecoveryResult>;
|
): Promise<WalletMintQuoteRecoveryResult>;
|
||||||
|
/**
|
||||||
|
* Recover stuck send/melt/mint operations whose mints answer a
|
||||||
|
* reachability probe. Operations a live execute holds are reported busy,
|
||||||
|
* never driven. Receive stays startup-only (receive dedup classification).
|
||||||
|
*/
|
||||||
|
recoverStuckOperations?(): Promise<WalletStuckOperationRecoveryResult>;
|
||||||
/** Report background wallet recovery progress, when the wallet supports it. */
|
/** Report background wallet recovery progress, when the wallet supports it. */
|
||||||
getRecoveryProgress?(): Promise<WalletRecoveryProgress>;
|
getRecoveryProgress?(): Promise<WalletRecoveryProgress>;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -363,7 +363,7 @@ describe("PAID mint quote recovery with a real Manager and mint", () => {
|
|||||||
outstanding,
|
outstanding,
|
||||||
})) as unknown as Record<string, number>;
|
})) as unknown as Record<string, number>;
|
||||||
expect(first).toMatchObject({ retryable: 1, recovered: 0 });
|
expect(first).toMatchObject({ retryable: 1, recovered: 0 });
|
||||||
expect(outstanding.has(op.id as string)).toBe(true);
|
expect(outstanding.has(`mint:${op.id}`)).toBe(true);
|
||||||
|
|
||||||
const second = (await runMintQuoteRecovery(booted.source() as never, {
|
const second = (await runMintQuoteRecovery(booted.source() as never, {
|
||||||
outstanding,
|
outstanding,
|
||||||
|
|||||||
@@ -0,0 +1,257 @@
|
|||||||
|
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
|
||||||
|
const result = await runTargetedRecovery(coco.ops, internals.sendOperationService, {
|
||||||
|
kinds: ["send"],
|
||||||
|
fetchImpl: (async () => new Response("{}")) as unknown as unknown as typeof fetch,
|
||||||
|
});
|
||||||
|
// On 8005aeb this is "rolled_back" and "in-1" is no longer reserved.
|
||||||
|
expect((await coco.ops.send.get(OP))?.state).toBe("executing");
|
||||||
|
// The live execute holds coco's per-operation lock: busy, never driven.
|
||||||
|
expect(result.attempted).toBe(0);
|
||||||
|
expect(result.busy).toBe(1);
|
||||||
|
|
||||||
|
finishSwap({ send: [{ id: "00aa", amount: 8, secret: "out-1", C: "02bb" }], keep: [] });
|
||||||
|
await live;
|
||||||
|
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("degraded startup probes before settlement and never asks a dead mint", async () => {
|
||||||
|
const db = new Database(":memory:");
|
||||||
|
const repo = new SqliteRepositories({ database: db });
|
||||||
|
await repo.init();
|
||||||
|
const coco = new Manager(repo, async () => new Uint8Array(64));
|
||||||
|
const expiredRow = (id: string, mintUrl: string) => ({
|
||||||
|
id, mintUrl, quoteId: `q-${id}`, state: "pending",
|
||||||
|
createdAt: 1_000, updatedAt: 2_000, method: "bolt11",
|
||||||
|
methodData: { method: "bolt11", data: {} }, amount: 100, unit: "sat",
|
||||||
|
request: "lnbc1example", expiry: 1_000_000, // epoch seconds, long past
|
||||||
|
});
|
||||||
|
await repo.mintOperationRepository.create(expiredRow("dead-quote", "https://dead.example.com") as never);
|
||||||
|
await repo.mintOperationRepository.create(expiredRow("live-quote", "https://live.example.com") as never);
|
||||||
|
const internals = coco as unknown as {
|
||||||
|
mintOperationService: {
|
||||||
|
observePendingOperation(id: string): Promise<{ category: "waiting" }>;
|
||||||
|
failPendingOperation(op: unknown, failure: unknown): Promise<unknown>;
|
||||||
|
};
|
||||||
|
};
|
||||||
|
const events: string[] = [];
|
||||||
|
const observe = spyOn(internals.mintOperationService, "observePendingOperation")
|
||||||
|
.mockImplementation(async (id: string) => {
|
||||||
|
events.push(`observe:${id}`);
|
||||||
|
return { category: "waiting" };
|
||||||
|
});
|
||||||
|
// Fail the row for real (as the production service would) so the targeted
|
||||||
|
// mint pass sees state "failed" and refresh() returns without re-observing.
|
||||||
|
const fail = spyOn(internals.mintOperationService, "failPendingOperation")
|
||||||
|
.mockImplementation(async (op: unknown) => {
|
||||||
|
const id = (op as { id: string }).id;
|
||||||
|
const row = await repo.mintOperationRepository.getById(id);
|
||||||
|
if (row) await repo.mintOperationRepository.update({ ...row, state: "failed" } as never);
|
||||||
|
return {} as never;
|
||||||
|
});
|
||||||
|
const gate = createRecoveryGate();
|
||||||
|
try {
|
||||||
|
const recovery = runWalletRecovery(coco, () => {}, [], (mints) => {
|
||||||
|
events.push("gate");
|
||||||
|
gate.publishStuckMints(mints);
|
||||||
|
}, {
|
||||||
|
fetchImpl: (async (url: unknown) => {
|
||||||
|
if (String(url).includes("dead.example.com")) throw new Error("offline");
|
||||||
|
return new Response("{}");
|
||||||
|
}) as unknown as typeof fetch,
|
||||||
|
cleanupLocalState: async () => { events.push("cleanup"); },
|
||||||
|
}).then(() => gate.complete());
|
||||||
|
await recovery;
|
||||||
|
// The dead mint was probed, never asked to observe its quote; the gate
|
||||||
|
// opened before settlement spent anything on the live mint.
|
||||||
|
expect(events).toEqual(["cleanup", "gate", "observe:live-quote"]);
|
||||||
|
expect(observe).toHaveBeenCalledTimes(1);
|
||||||
|
expect(fail).toHaveBeenCalledTimes(1);
|
||||||
|
} finally {
|
||||||
|
observe.mockRestore();
|
||||||
|
fail.mockRestore();
|
||||||
|
db.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
it("local housekeeping cleans init sends and orphaned reservations without network", async () => {
|
||||||
|
const db = new Database(":memory:");
|
||||||
|
const repo = new SqliteRepositories({ database: db });
|
||||||
|
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();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
it("cleanup checks all private methods before making any local writes", async () => {
|
||||||
|
const db = new Database(":memory:");
|
||||||
|
const repo = new SqliteRepositories({ database: db });
|
||||||
|
await repo.init();
|
||||||
|
const coco = new Manager(repo, async () => new Uint8Array(64));
|
||||||
|
const internals = coco as unknown as { mintOperationService: { recoverInitOperation: unknown } };
|
||||||
|
const original = internals.mintOperationService.recoverInitOperation;
|
||||||
|
try {
|
||||||
|
await repo.sendOperationRepository.create({
|
||||||
|
id: "untouched-init", mintUrl: MINT, amount: 8, state: "init", method: "default",
|
||||||
|
methodData: {}, createdAt: Date.now(), updatedAt: Date.now(),
|
||||||
|
} as never);
|
||||||
|
internals.mintOperationService.recoverInitOperation = undefined;
|
||||||
|
await expect(cleanupLocalRecoveryState(coco, repo)).rejects.toThrow("mintOperationService.recoverInitOperation is unavailable");
|
||||||
|
expect((await repo.sendOperationRepository.getById("untouched-init"))?.state).toBe("init");
|
||||||
|
} finally {
|
||||||
|
internals.mintOperationService.recoverInitOperation = original;
|
||||||
|
db.close();
|
||||||
|
}
|
||||||
|
});
|
||||||
@@ -0,0 +1,448 @@
|
|||||||
|
import { describe, expect, it } from "bun:test";
|
||||||
|
import {
|
||||||
|
collectStuckOperations,
|
||||||
|
probeMintReachability,
|
||||||
|
runTargetedRecovery,
|
||||||
|
type StuckOperationSource,
|
||||||
|
type SendRecoveryService,
|
||||||
|
} from "./recovery-probe";
|
||||||
|
|
||||||
|
import { runMintQuoteRecovery, settlePendingMintQuotes, settleExpiredMintQuotes } from "./coco-client";
|
||||||
|
|
||||||
|
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) => ({
|
||||||
|
diagnostics: { isLocked: () => false },
|
||||||
|
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"),
|
||||||
|
get: async (id: string) =>
|
||||||
|
(full.send.find((o) => o.id === id) ?? null) as never,
|
||||||
|
} as FakeSource["send"],
|
||||||
|
melt: family("melt") as FakeSource["melt"],
|
||||||
|
receive: family("receive") as FakeSource["receive"],
|
||||||
|
mint: family("mint") as FakeSource["mint"],
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
function makeSendService(): SendRecoveryService & {
|
||||||
|
executingRecovered: string[];
|
||||||
|
/** Lock events in order, e.g. "acquire:op-1", "release:op-1". */
|
||||||
|
lockLog: string[];
|
||||||
|
} {
|
||||||
|
const executingRecovered: string[] = [];
|
||||||
|
const lockLog: string[] = [];
|
||||||
|
return {
|
||||||
|
executingRecovered,
|
||||||
|
lockLog,
|
||||||
|
acquireOperationLock: async (id) => {
|
||||||
|
lockLog.push(`acquire:${id}`);
|
||||||
|
return () => {
|
||||||
|
lockLog.push(`release:${id}`);
|
||||||
|
};
|
||||||
|
},
|
||||||
|
recoverExecutingOperation: async (raw) => {
|
||||||
|
const id = (raw as FakeOp).id;
|
||||||
|
executingRecovered.push(id);
|
||||||
|
lockLog.push(`recover:${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({ attempted: 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.attempted).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"),
|
||||||
|
],
|
||||||
|
});
|
||||||
|
const sendService = makeSendService();
|
||||||
|
const result = await runTargetedRecovery(source, sendService, {
|
||||||
|
fetchImpl: liveFetch().fetchImpl,
|
||||||
|
});
|
||||||
|
expect(result.attempted).toBe(2);
|
||||||
|
expect(source.refreshed.send).toEqual(["pending-1"]);
|
||||||
|
expect(sendService.executingRecovered).toEqual(["exec-1"]);
|
||||||
|
// The executing send is driven under coco's per-operation lock.
|
||||||
|
expect(sendService.lockLog).toEqual([
|
||||||
|
"acquire:exec-1",
|
||||||
|
"recover:exec-1",
|
||||||
|
"release:exec-1",
|
||||||
|
]);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("re-reads state under the lock and skips a send that left executing", async () => {
|
||||||
|
const source = makeSource({
|
||||||
|
send: [op("s1", "https://live.example.com", "pending")],
|
||||||
|
});
|
||||||
|
const sendService = makeSendService();
|
||||||
|
// Snapshot taken while the op was still executing; it has since settled.
|
||||||
|
const result = await runTargetedRecovery(source, sendService, {
|
||||||
|
stuckOperations: [
|
||||||
|
{
|
||||||
|
kind: "send",
|
||||||
|
id: "s1",
|
||||||
|
mintUrl: "https://live.example.com",
|
||||||
|
state: "executing",
|
||||||
|
raw: op("s1", "https://live.example.com", "executing"),
|
||||||
|
},
|
||||||
|
],
|
||||||
|
unreachableMints: new Set(),
|
||||||
|
});
|
||||||
|
expect(result.attempted).toBe(1);
|
||||||
|
expect(sendService.executingRecovered).toEqual([]);
|
||||||
|
// The lock is still acquired and released around the re-read.
|
||||||
|
expect(sendService.lockLog).toEqual(["acquire:s1", "release:s1"]);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("counts a send as busy when a live execute wins the lock race", async () => {
|
||||||
|
const source = makeSource({
|
||||||
|
send: [op("s1", "https://live.example.com", "executing")],
|
||||||
|
});
|
||||||
|
const sendService = makeSendService();
|
||||||
|
sendService.acquireOperationLock = async () => {
|
||||||
|
const error = new Error("operation in progress");
|
||||||
|
error.name = "OperationInProgressError";
|
||||||
|
throw error;
|
||||||
|
};
|
||||||
|
const result = await runTargetedRecovery(source, sendService, {
|
||||||
|
fetchImpl: liveFetch().fetchImpl,
|
||||||
|
});
|
||||||
|
expect(result.busy).toBe(1);
|
||||||
|
expect(result.failed).toBe(0);
|
||||||
|
expect(sendService.executingRecovered).toEqual([]);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("counts per-operation failures at reachable mints and continues", async () => {
|
||||||
|
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.attempted).toBe(2);
|
||||||
|
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([]);
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
|
||||||
|
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(result.busy).toBe(1);
|
||||||
|
expect(service.executingRecovered).toEqual([]);
|
||||||
|
expect(service.lockLog).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"]);
|
||||||
|
});
|
||||||
|
|
||||||
|
|
||||||
|
it("shares timed-out mint observations across quote recovery, targeted recovery and the sweep", async () => {
|
||||||
|
const source = makeSource({ mint: [op("q1", "https://live.example.com")] });
|
||||||
|
let finish!: (value: { category: "waiting" }) => void;
|
||||||
|
const observation = new Promise<{ category: "waiting" }>(r => { finish = r; });
|
||||||
|
const outstanding = new Map<string, Promise<unknown>>();
|
||||||
|
const quoteSource = {
|
||||||
|
ops: { mint: {
|
||||||
|
listPending: async () => [{ ...source.stuck.mint[0]!, method: "bolt11" }],
|
||||||
|
get: async () => null,
|
||||||
|
finalize: async () => ({ state: "finalized" }),
|
||||||
|
} },
|
||||||
|
mintOperationService: { observePendingOperation: () => observation },
|
||||||
|
reopenFailedOperation: async () => false,
|
||||||
|
};
|
||||||
|
try {
|
||||||
|
await runMintQuoteRecovery(quoteSource as never, { outstanding, timeoutMs: 5 });
|
||||||
|
expect(outstanding.has("mint:q1")).toBe(true);
|
||||||
|
const result = await runTargetedRecovery(source, makeSendService(), {
|
||||||
|
outstanding, fetchImpl: liveFetch().fetchImpl,
|
||||||
|
});
|
||||||
|
expect(result.busy).toBe(1);
|
||||||
|
await settlePendingMintQuotes({
|
||||||
|
ops: { mint: { ...source.mint, listPending: async () => source.stuck.mint as never } },
|
||||||
|
wallet: { balances: { byMint: async () => ({}) } },
|
||||||
|
mintOperationService: { failPendingOperation: async () => ({}) },
|
||||||
|
} as never, Date.now(), { state: { outstanding } });
|
||||||
|
expect(source.refreshed.mint).toEqual([]);
|
||||||
|
} finally { finish({ category: "waiting" }); await observation; }
|
||||||
|
});
|
||||||
|
|
||||||
|
it("tracks a hung targeted mint so quote recovery skips it while unrelated operations proceed", async () => {
|
||||||
|
const source = makeSource({ mint: [op("q1", "https://live.example.com")], melt: [op("m1", "https://live.example.com")] });
|
||||||
|
let finish!: (value: never) => void;
|
||||||
|
source.mint.refresh = () => new Promise(r => { finish = r; });
|
||||||
|
const outstanding = new Map<string, Promise<unknown>>();
|
||||||
|
try {
|
||||||
|
const result = await runTargetedRecovery(source, makeSendService(), {
|
||||||
|
outstanding, timeoutMs: 5, fetchImpl: liveFetch().fetchImpl,
|
||||||
|
});
|
||||||
|
expect(result.timedOut).toBe(1);
|
||||||
|
expect(source.refreshed.melt).toEqual(["m1"]);
|
||||||
|
const quote = await runMintQuoteRecovery({
|
||||||
|
ops: { mint: { listPending: async () => [{ ...source.stuck.mint[0]!, method: "bolt11" }] } },
|
||||||
|
} as never, { outstanding });
|
||||||
|
expect(quote.busy).toBe(1);
|
||||||
|
} finally { finish({ state: "pending" } as never); }
|
||||||
|
});
|
||||||
|
|
||||||
|
it("retains executing-send lock until a timed-out underlying drive actually finishes", async () => {
|
||||||
|
const source = makeSource({ send: [op("s1", "https://live.example.com", "executing")] });
|
||||||
|
const service = makeSendService();
|
||||||
|
let finish!: () => void;
|
||||||
|
service.recoverExecutingOperation = () => new Promise<void>(r => { finish = r; });
|
||||||
|
const outstanding = new Map<string, Promise<unknown>>();
|
||||||
|
const result = await runTargetedRecovery(source, service, {
|
||||||
|
outstanding, timeoutMs: 5, fetchImpl: liveFetch().fetchImpl,
|
||||||
|
});
|
||||||
|
expect(result.timedOut).toBe(1);
|
||||||
|
expect(service.lockLog).toEqual(["acquire:s1"]);
|
||||||
|
const retry = await runTargetedRecovery(source, service, { outstanding, fetchImpl: liveFetch().fetchImpl });
|
||||||
|
expect(retry.busy).toBe(1);
|
||||||
|
const actual = outstanding.get("send:s1")!;
|
||||||
|
finish();
|
||||||
|
await actual;
|
||||||
|
expect(service.lockLog).toEqual(["acquire:s1", "release:s1"]);
|
||||||
|
expect(outstanding.size).toBe(0);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("refuses missing executing-send internals without driving the operation", async () => {
|
||||||
|
const source = makeSource({ send: [op("s1", "https://live.example.com", "executing")] });
|
||||||
|
const result = await runTargetedRecovery(source, {} as SendRecoveryService, { fetchImpl: liveFetch().fetchImpl });
|
||||||
|
expect(result.failed).toBe(1);
|
||||||
|
});
|
||||||
|
|
||||||
|
|
||||||
|
it("startup settlement keeps its timed-out observation visible to manual recovery", async () => {
|
||||||
|
const source = makeSource({ mint: [op("q1", "https://live.example.com")] });
|
||||||
|
let finish!: (value: { category: "ready" }) => void;
|
||||||
|
const observation = new Promise<{ category: "ready" }>(r => { finish = r; });
|
||||||
|
const outstanding = new Map<string, Promise<unknown>>();
|
||||||
|
const settlement = await settleExpiredMintQuotes({
|
||||||
|
ops: { mint: { listPending: async () => [{ ...source.stuck.mint[0]!, expiry: 1, updatedAt: 1 }] } },
|
||||||
|
mintOperationService: { observePendingOperation: () => observation },
|
||||||
|
} as never, Date.now(), 5, { outstanding });
|
||||||
|
expect(settlement.unobserved).toBe(1);
|
||||||
|
const result = await runTargetedRecovery(source, makeSendService(), { outstanding, fetchImpl: liveFetch().fetchImpl });
|
||||||
|
expect(result.busy).toBe(1);
|
||||||
|
expect(source.refreshed.mint).toEqual([]);
|
||||||
|
const actual = outstanding.get("mint:q1")!;
|
||||||
|
finish({ category: "ready" });
|
||||||
|
await actual;
|
||||||
|
expect(outstanding.size).toBe(0);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("deadline and shutdown leave remaining operations untouched", async () => {
|
||||||
|
const source = makeSource({ mint: [op("q1", "https://live.example.com")] });
|
||||||
|
const result = await runTargetedRecovery(source, makeSendService(), {
|
||||||
|
shouldStop: () => true, fetchImpl: liveFetch().fetchImpl,
|
||||||
|
});
|
||||||
|
expect(result.skipped).toBe(1);
|
||||||
|
expect(result.attempted).toBe(0);
|
||||||
|
expect(source.refreshed.mint).toEqual([]);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("a pass budget bounds total drive waits, not just each operation", async () => {
|
||||||
|
const source = makeSource({ mint: [op("q1", "https://live.example.com"), op("q2", "https://live.example.com")] });
|
||||||
|
let finish!: (value: never) => void;
|
||||||
|
source.mint.refresh = () => new Promise(r => { finish = r; });
|
||||||
|
const outstanding = new Map<string, Promise<unknown>>();
|
||||||
|
const result = await runTargetedRecovery(source, makeSendService(), {
|
||||||
|
outstanding, timeoutMs: 100, deadlineMs: 10, fetchImpl: liveFetch().fetchImpl,
|
||||||
|
});
|
||||||
|
expect(result.timedOut).toBe(1);
|
||||||
|
expect(result.attempted).toBe(1);
|
||||||
|
expect(result.skipped).toBe(1);
|
||||||
|
const actual = outstanding.get("mint:q1")!;
|
||||||
|
finish({ state: "pending" } as never);
|
||||||
|
await actual;
|
||||||
|
});
|
||||||
|
|
||||||
|
it("unfinished work in another operation family does not block the same bare id", async () => {
|
||||||
|
const source = makeSource({ mint: [op("same-id", "https://live.example.com")] });
|
||||||
|
const outstanding = new Map<string, Promise<unknown>>([["send:same-id", new Promise(() => {})]]);
|
||||||
|
const result = await runTargetedRecovery(source, makeSendService(), { outstanding, fetchImpl: liveFetch().fetchImpl });
|
||||||
|
expect(result.busy).toBe(0);
|
||||||
|
expect(source.refreshed.mint).toEqual(["same-id"]);
|
||||||
|
});
|
||||||
@@ -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 an explicit mid-session recovery (see recoverStuckOperations
|
||||||
|
* in coco-client.ts) or a later startup finds the mint reachable.
|
||||||
|
*/
|
||||||
|
import { normalizeMintUrl, type Manager } from "@cashu/coco-core";
|
||||||
|
import { recoveryKey, trackRecovery, waitForRecoveryWork, RecoveryWaitTimeout, type RecoveryWork } from "./recovery-work";
|
||||||
|
import { logger } from "../../utils/logger";
|
||||||
|
|
||||||
|
/** Short probe: a mint that cannot answer /v1/info in 2s slows every op. */
|
||||||
|
export const MINT_PROBE_TIMEOUT_MS = 2_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" | "diagnostics" | "get">;
|
||||||
|
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";
|
||||||
|
|
||||||
|
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 {
|
||||||
|
recoverExecutingOperation(op: unknown): Promise<void>;
|
||||||
|
/**
|
||||||
|
* coco's per-operation lock. It is fail-fast: acquiring an id a live
|
||||||
|
* execute/finalize/recover already holds throws OperationInProgressError
|
||||||
|
* instead of waiting. Holding it across the state re-read and the drive
|
||||||
|
* makes executing-send recovery atomic against a live execute — the same
|
||||||
|
* pattern reopenFailedMintOperation uses for mint operations
|
||||||
|
* (see coco-client.ts).
|
||||||
|
*/
|
||||||
|
acquireOperationLock(operationId: string): Promise<() => void>;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface RecoveryRunResult {
|
||||||
|
/** Operations for which recovery was attempted (not necessarily completed). */
|
||||||
|
attempted: number;
|
||||||
|
/** Timed-out waits, also included in attempted; underlying work remains tracked. */
|
||||||
|
timedOut: number;
|
||||||
|
/** Locked operations or unfinished work from another pass; retry later. */
|
||||||
|
busy: number;
|
||||||
|
/** Operations skipped for unreachable mints, shutdown, or pass budget exhaustion. */
|
||||||
|
skipped: number;
|
||||||
|
/** Attempts that threw a non-busy, non-timeout error. */
|
||||||
|
failed: number;
|
||||||
|
/** Unreachable mint URL -> number of operations skipped there. */
|
||||||
|
skippedMints: Map<string, number>;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface TargetedRecoveryOptions {
|
||||||
|
outstanding?: RecoveryWork;
|
||||||
|
timeoutMs?: number;
|
||||||
|
deadlineMs?: number;
|
||||||
|
shouldStop?: () => boolean;
|
||||||
|
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 {
|
||||||
|
// 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 });
|
||||||
|
}
|
||||||
|
}
|
||||||
|
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(`${normalizeMintUrl(mintUrl)}/v1/info`, {
|
||||||
|
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") {
|
||||||
|
// Fail closed: the lock and the drive are private coco internals
|
||||||
|
// (written against coco-core 1.0.1), so a coco bump that renames them
|
||||||
|
// must fail this operation loudly instead of corrupting it.
|
||||||
|
if (
|
||||||
|
typeof sendService.acquireOperationLock !== "function" ||
|
||||||
|
typeof sendService.recoverExecutingOperation !== "function"
|
||||||
|
) {
|
||||||
|
throw new Error(
|
||||||
|
"coco sendOperationService recovery internals are unavailable; refusing to recover an executing send",
|
||||||
|
);
|
||||||
|
}
|
||||||
|
// recoverExecutingOperation takes no lock and does no state re-read,
|
||||||
|
// so take coco's fail-fast per-operation lock first (throws
|
||||||
|
// OperationInProgressError when a live execute holds the operation —
|
||||||
|
// the driver counts that as busy, not failed) and re-read the state
|
||||||
|
// under it: the snapshot op must still be executing before we drive.
|
||||||
|
const release = await sendService.acquireOperationLock(op.id);
|
||||||
|
try {
|
||||||
|
const latest = await source.send.get(op.id);
|
||||||
|
if (latest?.state === "executing") {
|
||||||
|
await sendService.recoverExecutingOperation(latest);
|
||||||
|
}
|
||||||
|
} finally {
|
||||||
|
release();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
// prepared / rolling_back: the global sweep only warns; nothing to do.
|
||||||
|
return;
|
||||||
|
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 = {
|
||||||
|
attempted: 0,
|
||||||
|
timedOut: 0,
|
||||||
|
busy: 0,
|
||||||
|
skipped: 0,
|
||||||
|
failed: 0,
|
||||||
|
skippedMints: new Map(),
|
||||||
|
};
|
||||||
|
|
||||||
|
const timeoutMs = options.timeoutMs ?? 15_000;
|
||||||
|
const deadlineMs = options.deadlineMs ?? 60_000;
|
||||||
|
if (![timeoutMs, deadlineMs].every(n => Number.isFinite(n) && n > 0)) {
|
||||||
|
throw new Error("Recovery budgets must be positive finite numbers");
|
||||||
|
}
|
||||||
|
const deadline = Date.now() + deadlineMs;
|
||||||
|
const outstanding = options.outstanding ?? new Map();
|
||||||
|
const kinds = options.kinds ?? ["send", "melt", "receive", "mint"];
|
||||||
|
const stuck = (options.stuckOperations ?? (await collectStuckOperations(source))).filter(
|
||||||
|
(op) => kinds.includes(op.kind),
|
||||||
|
);
|
||||||
|
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;
|
||||||
|
if (options.shouldStop?.() || Date.now() >= deadline) { result.skipped++; continue; }
|
||||||
|
const key = recoveryKey(op.kind, op.id);
|
||||||
|
if (outstanding.has(key)) { result.busy++; continue; }
|
||||||
|
// Cheap pre-filter for live operations; the lock inside
|
||||||
|
// recoverStuckOperation is what actually makes the drive atomic.
|
||||||
|
if (source[op.kind].diagnostics.isLocked(op.id)) {
|
||||||
|
result.busy++;
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
if (op.kind === "send" && !["pending", "executing"].includes(op.state)) continue;
|
||||||
|
result.attempted++;
|
||||||
|
try {
|
||||||
|
const work = trackRecovery(outstanding, key, recoverStuckOperation(source, sendService, op));
|
||||||
|
await waitForRecoveryWork(work, Math.min(timeoutMs, Math.max(1, deadline - Date.now())));
|
||||||
|
} catch (error) {
|
||||||
|
if (error instanceof RecoveryWaitTimeout) { result.timedOut++; continue; }
|
||||||
|
// A live execute grabbed the operation between the isLocked pre-filter
|
||||||
|
// and the lock acquisition: busy, not failed — leave it for a later
|
||||||
|
// pass. Name-matched like runMintQuoteRecovery does, because the error
|
||||||
|
// crosses a package boundary.
|
||||||
|
if (error instanceof Error && error.name === "OperationInProgressError") {
|
||||||
|
result.busy++;
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
// Same semantics as coco's tryRecover*: leave the operation for the
|
||||||
|
// next pass. A reachable mint can still reject a specific operation.
|
||||||
|
result.failed++;
|
||||||
|
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;
|
||||||
|
}
|
||||||
@@ -0,0 +1,41 @@
|
|||||||
|
import { expect, it } from "bun:test";
|
||||||
|
import { createRunQueue } from "./coco-client";
|
||||||
|
import { createRecoveryDisposer, drainRecoveryWork, trackRecovery } from "./recovery-work";
|
||||||
|
|
||||||
|
it("incomplete shutdown retains resources until late writes settle; retries join disposal", async () => {
|
||||||
|
const work = new Map<string, Promise<unknown>>();
|
||||||
|
let finish!: () => void;
|
||||||
|
const events: string[] = [];
|
||||||
|
trackRecovery(work, "mint:q", new Promise<void>(r => { finish = r; }).then(() => { events.push("write"); }));
|
||||||
|
const dispose = createRecoveryDisposer(() => {}, () => drainRecoveryWork(work), async () => { events.push("close"); }, 5);
|
||||||
|
await expect(dispose()).rejects.toThrow("Timed out");
|
||||||
|
expect(events).toEqual([]);
|
||||||
|
finish();
|
||||||
|
await dispose();
|
||||||
|
expect(events).toEqual(["write", "close"]);
|
||||||
|
await dispose();
|
||||||
|
expect(events).toEqual(["write", "close"]);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("shutdown drains active queue work and queued callbacks reject before touching DB", async () => {
|
||||||
|
const queue = createRunQueue();
|
||||||
|
let disposed = false;
|
||||||
|
let finish!: () => void;
|
||||||
|
const events: string[] = [];
|
||||||
|
let started!: () => void;
|
||||||
|
const ready = new Promise<void>(r => { started = r; });
|
||||||
|
const active = queue(() => new Promise<void>(r => { finish = r; started(); }).then(() => { events.push("write"); }));
|
||||||
|
await ready;
|
||||||
|
const queued = queue(async () => {
|
||||||
|
if (disposed) throw new Error("Wallet is shutting down");
|
||||||
|
events.push("unexpected");
|
||||||
|
});
|
||||||
|
const rejected = queued.catch(error => error);
|
||||||
|
const dispose = createRecoveryDisposer(() => { disposed = true; }, () => queue.drain(), async () => { events.push("close"); }, 5);
|
||||||
|
await expect(dispose()).rejects.toThrow("Timed out");
|
||||||
|
finish();
|
||||||
|
await active;
|
||||||
|
expect((await rejected).message).toContain("shutting down");
|
||||||
|
await dispose();
|
||||||
|
expect(events).toEqual(["write", "close"]);
|
||||||
|
});
|
||||||
@@ -0,0 +1,45 @@
|
|||||||
|
/** A timed-out wait is not cancellation: retain work until it actually settles. */
|
||||||
|
export type RecoveryWork = Map<string, Promise<unknown>>;
|
||||||
|
export const recoveryKey = (kind: string, id: string): string => `${kind}:${id}`;
|
||||||
|
|
||||||
|
export function trackRecovery<T>(work: RecoveryWork, key: string, promise: Promise<T>): Promise<T> {
|
||||||
|
work.set(key, promise);
|
||||||
|
const clear = () => { if (work.get(key) === promise) work.delete(key); };
|
||||||
|
void promise.then(clear, clear);
|
||||||
|
return promise;
|
||||||
|
}
|
||||||
|
|
||||||
|
export class RecoveryWaitTimeout extends Error {
|
||||||
|
constructor() { super("Timed out waiting for recovery; underlying work is still tracked"); }
|
||||||
|
}
|
||||||
|
|
||||||
|
export async function waitForRecoveryWork<T>(promise: Promise<T>, timeoutMs: number): Promise<T> {
|
||||||
|
let timer: ReturnType<typeof setTimeout> | undefined;
|
||||||
|
try {
|
||||||
|
return await Promise.race([promise, new Promise<never>((_, reject) => {
|
||||||
|
timer = setTimeout(() => reject(new RecoveryWaitTimeout()), timeoutMs);
|
||||||
|
})]);
|
||||||
|
} finally {
|
||||||
|
if (timer !== undefined) clearTimeout(timer);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Drain actual work, including work registered while an earlier task settles. */
|
||||||
|
export async function drainRecoveryWork(work: RecoveryWork): Promise<void> {
|
||||||
|
while (work.size) await Promise.allSettled([...work.values()]);
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Timeout reports incomplete disposal; actual cleanup continues safely. */
|
||||||
|
export function createRecoveryDisposer(
|
||||||
|
quiesce: () => void,
|
||||||
|
settle: () => Promise<void>,
|
||||||
|
close: () => Promise<void>,
|
||||||
|
timeoutMs = 30_000,
|
||||||
|
): () => Promise<void> {
|
||||||
|
let disposal: Promise<void> | undefined;
|
||||||
|
return async () => {
|
||||||
|
quiesce();
|
||||||
|
disposal ??= (async () => { await settle(); await close(); })();
|
||||||
|
await waitForRecoveryWork(disposal, timeoutMs);
|
||||||
|
};
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user