Merge branch 'main' into fix/nwc-status-hang

Resolve the single conflict in coco-client.ts: keep main's Infinity
fast-path in the local withTimeout (recovery passes Infinity for an
unbounded remaining budget, and setTimeout(Infinity) would overflow and
fire immediately), while still delegating the bounded path to the shared
utils/with-timeout helper.
This commit is contained in:
redshift
2026-10-02 16:50:31 +08:00
24 changed files with 4903 additions and 112 deletions
+2
View File
@@ -121,6 +121,8 @@ Show wallet transaction history.
|--------|---------|-------------|
| `-n, --limit <number>` | 50 | Number of entries to show |
| `--offset <number>` | 0 | Number of entries to skip |
| `-t, --type <type...>` | all | Filter by transaction type (`send`, `receive`, `mint`, `melt`). Repeatable or comma-separated |
| `-i, --id <id>` | | Show a single transaction by its ID (prints one summary line; add `--verbose` for full details) |
| `-v, --verbose` | false | Show full details including encoded Cashu tokens |
| `--json` | false | Output raw JSON with token objects (no encoding) |
+71
View File
@@ -0,0 +1,71 @@
# Mint quote recovery: scope and troubleshooting
`routstrd wallet recover` explicitly retries mint operations through coco using
**their existing stored outputs**. It can restore signatures when a quote is
already issued, and reopen failed operations when explicitly requested:
```sh
routstrd history --json
routstrd wallet recover --op <operationId> --include-failed
```
Failed operations require explicit IDs over both HTTP and the CLI. A successful
re-run on an already finalized operation is a no-op. Requests that exceed their
wait budget are not cancelled; explicit retries skip the operation while the
underlying work is outstanding.
## What this fixes—and what it does not
Coco already checks pending quotes on startup and the daemon periodically
refreshes them. A quote paid while the daemon was offline does not, by itself,
require a new issuance implementation.
This change makes normal cleanup confirm UNPAID with the mint before failing an
expired quote. PAID, ISSUED and unverified quotes remain pending. It also gives
operators a recovery path for operations previously marked failed.
It does **not** replace rejected outputs with fresh outputs on an active keyset.
An inactive-keyset rejection can therefore remain retryable with zero recovery.
Recovery reports coco's persisted mint error when available, rather than only a
generic “remains pending” error.
Do not infer that the production incidents were caused by keyset retirement.
Before claiming those incidents are fixed, collect:
- The affected operation IDs, quote IDs, state and persisted `error`.
- A fresh remote quote state and, where provided, paid/issued amounts.
- The keyset IDs in the stored outputs and the mint's current keyset metadata.
- A reproduction showing existing recovery fails and the proposed fix succeeds.
Inspect persisted operation data through a read-only database copy; do not edit
rows or run recovery scripts concurrently with a daemon against the same wallet.
Never share the mnemonic, output secrets, or full wallet database in a PR.
A future fresh-output path must preserve original outputs for uncertain issuance
and NUT-09 restore, allocate fresh deterministic counters safely, and coordinate
with coco's watcher/processor. It needs its own integration tests before handling
real funds.
## Cleanup preview and force
`wallet cleanup --dry-run` is local-only: it reports `mintQuoteCandidates`, not
confirmed failures. `failedMintQuotes` and `leftForRecovery` are zero because no
mint check or cleanup transition was performed. Send/melt counts remain planned
cleanup counts in dry-run mode.
`--force` deliberately bypasses mint confirmation and can strand paid sats in a
failed operation. Prefer normal cleanup. Forced operations can be retried with
`--op <operationId> --include-failed`, but recovery still depends on the mint
accepting their stored outputs or restoring their signatures.
## Integration and release notes
The reopen helper uses private coco-core 1.0.1 methods. Retain real-Manager and
HTTP fake-mint coverage, use frozen dependency installs, and re-run integration
tests on coco upgrades. A controlled low-value live-mint smoke test remains
recommended before release.
PR #118 removes `cocod-client.ts`. When integrating that change, move recovery
and cleanup contracts into its replacement `wallet-client.ts`, rename HTTP error
references accordingly, and make recovery mandatory for the in-process client.
This follow-up does not pull in #118's unrelated removal.
+191 -7
View File
@@ -58,6 +58,10 @@ import {
} from "./daemon/wallet/paths";
import * as QRCode from "qrcode";
import { normalizeNostrPubkey, npubFromPubkey, npubFromSecretKey } from "./utils/nip98";
import {
HISTORY_ENTRY_TYPES,
isHistoryEntryType,
} from "./utils/history";
import { generateSecretKey, nip19 } from "nostr-tools";
import { generateMnemonic } from "@scure/bip39";
import { wordlist } from "@scure/bip39/wordlists/english.js";
@@ -1789,17 +1793,50 @@ program
.description("Show wallet transaction history")
.option("-n, --limit <number>", "Number of entries to show", "50")
.option("--offset <number>", "Number of entries to skip", "0")
.option("-v, --verbose", "Show full details including encoded Cashu tokens")
.option(
"-t, --type <type...>",
"Filter by transaction type (send, receive, mint, melt). Repeatable or comma-separated",
)
.option("-i, --id <id>", "Show a single transaction by its ID")
.option("-v, --verbose", "Show full details instead of the summary line")
.option("--json", "Output raw JSON with token objects (no encoding)")
.action(async (options: { limit: string; offset: string; verbose: boolean; json: boolean }) => {
.action(async (options: {
limit: string;
offset: string;
type?: string[];
id?: string;
verbose: boolean;
json: boolean;
}) => {
await ensureDaemonRunning();
const limit = Math.min(parseInt(options.limit, 10) || 50, 1000);
const offset = parseInt(options.offset, 10) || 0;
const result = await callDaemon(
`/wallet/history?offset=${offset}&limit=${limit}`,
// Support both repeated flags (-t send -t mint) and comma-separated
// values (-t send,mint).
const requestedTypes = (options.type ?? [])
.flatMap((value) => value.split(","))
.map((value) => value.trim().toLowerCase())
.filter((value) => value.length > 0);
const unknownTypes = requestedTypes.filter(
(value) => !isHistoryEntryType(value),
);
if (unknownTypes.length > 0) {
console.error(
`Unknown transaction type: ${unknownTypes.join(", ")}. Valid types: ${HISTORY_ENTRY_TYPES.join(", ")}.`,
);
process.exit(1);
}
const query = new URLSearchParams({
offset: String(offset),
limit: String(limit),
});
if (requestedTypes.length > 0) query.set("type", requestedTypes.join(","));
if (options.id) query.set("id", options.id);
const result = await callDaemon(`/wallet/history?${query.toString()}`);
if (result.error) {
console.log(result.error);
@@ -1812,7 +1849,15 @@ program
const entries = data?.entries || [];
if (entries.length === 0) {
if (options.id) {
console.log(`No transaction found with ID ${options.id}.`);
} else if (requestedTypes.length > 0) {
console.log(
`No ${requestedTypes.join(", ")} transactions found.`,
);
} else {
console.log("No transaction history yet.");
}
return;
}
@@ -1869,10 +1914,13 @@ program
const pad = (s: string, w: number) => s.padEnd(w);
const sep = Object.values(widths).map((w) => "-".repeat(w)).join(" | ");
// A single transaction lookup prints just its summary line.
if (!options.id) {
console.log(
`${pad(idCol, widths.id)} | ${pad(timeCol, widths.time)} | ${pad(typeCol, widths.type)} | ${pad(mintCol, widths.mint)} | ${pad(amtCol, widths.amount)}`,
);
console.log(sep);
}
for (const row of rows) {
console.log(
@@ -2014,16 +2062,22 @@ walletCmd
.option("--mint-url <url>", "Only clean up operations for this mint URL")
.option(
"--min-age <hours>",
"Minimum age for reclaiming sends/cancelling melts, in hours (default: 168, one week; expired mint quotes are always failed)",
"Minimum age for reclaiming sends/cancelling melts, in hours (default: 168, one week; expired mint quotes are checked with their mint)",
"168",
)
.option("--dry-run", "Report what would be cleaned without applying changes", false)
.option(
"--force",
"Fail expired mint quotes without confirming UNPAID with the mint (may strand paid quotes)",
false,
)
.option("-y, --yes", "Skip confirmation prompt", false)
.action(
async (options: {
mintUrl?: string;
minAge: string;
dryRun: boolean;
force: boolean;
yes: boolean;
}) => {
const minAgeHours = Number.parseFloat(options.minAge);
@@ -2039,7 +2093,9 @@ walletCmd
});
const answer = await new Promise<string>((resolve) => {
rl.question(
"This will fail expired mint quotes, reclaim old pending sends, and cancel prepared melts. Continue? [y/N] ",
options.force
? "WARNING: --force fails expired mint quotes WITHOUT checking the mint and may strand paid sats. It also reclaims old pending sends and cancels prepared melts. Continue? [y/N] "
: "This will fail expired mint quotes confirmed unpaid, reclaim old pending sends, and cancel prepared melts. Continue? [y/N] ",
(value: string) => {
rl.close();
resolve(value.trim().toLowerCase());
@@ -2061,6 +2117,7 @@ walletCmd
mintUrl: options.mintUrl,
minAgeMs: Math.round(minAgeHours * 60 * 60 * 1000),
dryRun: options.dryRun === true,
force: options.force === true,
},
});
@@ -2073,6 +2130,8 @@ walletCmd
| {
dryRun?: boolean;
failedMintQuotes?: number;
mintQuoteCandidates?: number;
leftForRecovery?: number;
reclaimedSends?: number;
cancelledMelts?: number;
skipped?: number;
@@ -2083,8 +2142,15 @@ walletCmd
if (output) {
const prefix = output.dryRun ? "Would clean up:" : "Cleaned up:";
console.log(prefix);
if (output.dryRun) {
console.log(
` Expired mint quotes failed: ${output.failedMintQuotes ?? 0}`,
` Expired mint quote candidates (not checked with mint): ${output.mintQuoteCandidates ?? 0}`,
);
} else {
console.log(` Expired mint quotes failed: ${output.failedMintQuotes ?? 0}`);
}
console.log(
` Expired quotes kept for recovery (paid/issued/unverified): ${output.leftForRecovery ?? 0}`,
);
console.log(` Pending sends reclaimed: ${output.reclaimedSends ?? 0}`);
console.log(
@@ -2115,6 +2181,124 @@ walletCmd
},
);
walletCmd
.command("recover")
.description(
"Retry mint quotes using their stored outputs (does not replace rejected outputs)",
)
.option(
"--op <id>",
"Recover this operation id (repeatable; find IDs with routstrd history --json)",
(value: string, previous: string[]) => [...previous, value],
[] as string[],
)
.option(
"--include-failed",
"Also re-open operations coco already gave up on (requires --op)",
false,
)
.option("-y, --yes", "Skip confirmation prompt", false)
.action(
async (options: {
op: string[];
includeFailed: boolean;
yes: boolean;
}) => {
const operationIds = options.op ?? [];
if (options.includeFailed && operationIds.length === 0) {
console.error(
"--include-failed can only target operations named with --op",
);
process.exit(1);
}
if (!options.yes) {
const rl = require("readline").createInterface({
input: process.stdin,
output: process.stdout,
});
const prompt =
operationIds.length > 0
? `Recover ${operationIds.length} mint quote operation(s)? [y/N] `
: "Check every pending mint quote with its mint and claim any paid sats? [y/N] ";
const answer = await new Promise<string>((resolve) => {
rl.question(prompt, (value: string) => {
rl.close();
resolve(value.trim().toLowerCase());
});
});
if (answer !== "y" && answer !== "yes") {
console.log("Aborted.");
return;
}
}
try {
await ensureDaemonRunning();
const result = await callDaemon("/wallet/recover", {
method: "POST",
body: {
operationIds: operationIds.length > 0 ? operationIds : undefined,
includeFailed: options.includeFailed === true,
},
});
if (result.error) {
console.log(result.error);
process.exit(1);
}
const output = result.output as
| {
checked?: number;
recovered?: number;
waiting?: number;
terminal?: number;
reopened?: number;
retryable?: number;
busy?: number;
errors?: Array<{ operationId: string; error: string }>;
}
| undefined;
if (output) {
console.log("Mint quote recovery:");
console.log(` Checked with mint: ${output.checked ?? 0}`);
console.log(
` Recovered (paid sats claimed): ${output.recovered ?? 0}`,
);
console.log(` Still unpaid: ${output.waiting ?? 0}`);
console.log(` No longer issuable: ${output.terminal ?? 0}`);
console.log(
` Re-opened failed operations: ${output.reopened ?? 0}`,
);
console.log(` Left for a later run: ${output.retryable ?? 0}`);
console.log(
` Skipped (recovery still running): ${output.busy ?? 0}`,
);
if (output.errors && output.errors.length > 0) {
console.log("\nErrors:");
for (const e of output.errors) {
console.log(` - ${e.operationId}: ${e.error}`);
}
}
}
} catch (error) {
const message = (error as Error).message;
if (
message?.includes("fetch failed") ||
message?.includes("Connection refused")
) {
console.error("Daemon is not running");
process.exit(1);
}
console.error(message);
process.exit(1);
}
},
);
const walletReceiveCmd = walletCmd
.command("receive")
.description("Wallet receive operations");
+165
View File
@@ -0,0 +1,165 @@
import { describe, expect, it } from "bun:test";
import { EventEmitter } from "events";
import type { HistoryEntry } from "@cashu/coco-core";
import {
createDaemonRequestHandler,
getHistoryByTypes,
parseHistoryTypes,
} from "./index";
const MINT_URL = "https://mint.example/";
function makeEntry(
id: string,
type: HistoryEntry["type"],
createdAt: number,
): HistoryEntry {
return {
id,
type,
createdAt,
mintUrl: MINT_URL,
unit: "sat",
amount: 21,
} as HistoryEntry;
}
const ENTRIES: HistoryEntry[] = [
makeEntry("send-1", "send", 4),
makeEntry("receive-1", "receive", 3),
makeEntry("mint-1", "mint", 2),
makeEntry("melt-1", "melt", 1),
makeEntry("send-2", "send", 0),
];
function makeWalletClient(entries: HistoryEntry[] = ENTRIES) {
return {
getHistory: async (offset = 0, limit = 50) =>
entries.slice(offset, offset + limit),
getHistoryEntryById: async (id: string) =>
entries.find((entry) => entry.id === id) ?? null,
};
}
function makeReq(method: string, path: string) {
const req = new EventEmitter() as any;
req.method = method;
req.url = path;
req.headers = { host: "localhost" };
return req;
}
function makeRes() {
const res: any = {
status: 0,
body: "",
writeHead(status: number) {
res.status = status;
return res;
},
end(chunk?: string) {
if (chunk) res.body += chunk;
return res;
},
json() {
return JSON.parse(res.body);
},
};
return res;
}
async function callHistory(path: string, entries?: HistoryEntry[]) {
const handler = createDaemonRequestHandler({
walletClient: makeWalletClient(entries),
} as any);
const res = makeRes();
await handler(makeReq("GET", path), res);
return res;
}
describe("parseHistoryTypes", () => {
it("splits, trims, and lowercases comma-separated values", () => {
expect(parseHistoryTypes(" Send , MINT ")).toEqual(["send", "mint"]);
});
it("returns an empty list for missing or blank input", () => {
expect(parseHistoryTypes(null)).toEqual([]);
expect(parseHistoryTypes("")).toEqual([]);
expect(parseHistoryTypes(" , ")).toEqual([]);
});
});
describe("getHistoryByTypes", () => {
it("filters by type and applies offset/limit after filtering", async () => {
const client = makeWalletClient();
expect(await getHistoryByTypes(client, ["send"], 0, 10)).toEqual([
ENTRIES[0]!,
ENTRIES[4]!,
]);
expect(await getHistoryByTypes(client, ["send"], 1, 1)).toEqual([
ENTRIES[4]!,
]);
});
it("scans past the page size boundary to find matches", async () => {
const filler = Array.from({ length: 250 }, (_, index) =>
makeEntry(`receive-${index}`, "receive", index),
);
const target = makeEntry("mint-late", "mint", -1);
const client = makeWalletClient([...filler, target]);
expect(await getHistoryByTypes(client, ["mint"], 0, 10)).toEqual([target]);
});
it("returns no entries when nothing matches", async () => {
const client = makeWalletClient([makeEntry("mint-1", "mint", 1)]);
expect(await getHistoryByTypes(client, ["send"], 0, 10)).toEqual([]);
});
});
describe("GET /wallet/history", () => {
it("returns all entries without a filter", async () => {
const res = await callHistory("/wallet/history");
expect(res.status).toBe(200);
const body = res.json();
expect(body.output.entries.map((e: HistoryEntry) => e.id)).toEqual([
"send-1",
"receive-1",
"mint-1",
"melt-1",
"send-2",
]);
});
it("filters entries by a comma-separated type list", async () => {
const res = await callHistory("/wallet/history?type=send,melt");
const body = res.json();
expect(body.output.entries.map((e: HistoryEntry) => e.id)).toEqual([
"send-1",
"melt-1",
"send-2",
]);
});
it("honors offset/limit for filtered results", async () => {
const res = await callHistory("/wallet/history?type=send&offset=1&limit=1");
const body = res.json();
expect(body.output.entries.map((e: HistoryEntry) => e.id)).toEqual([
"send-2",
]);
});
it("looks up a single entry by id", async () => {
const res = await callHistory("/wallet/history?id=mint-1");
const body = res.json();
expect(body.output.entries.map((e: HistoryEntry) => e.id)).toEqual([
"mint-1",
]);
});
it("returns an empty list when the id is unknown", async () => {
const res = await callHistory("/wallet/history?id=does-not-exist");
const body = res.json();
expect(body.output.entries).toEqual([]);
});
});
+123 -1
View File
@@ -283,6 +283,21 @@ function optionalStringField(
return typeof value === "string" && value.trim() ? value.trim() : undefined;
}
function optionalStringArrayField(
body: Record<string, unknown>,
field: string,
): string[] | undefined {
const value = body[field];
if (value === undefined) return undefined;
if (
!Array.isArray(value) ||
value.some((item) => typeof item !== "string" || !item.trim())
) {
throw new CocodHttpError(400, `'${field}' must be an array of non-empty strings.`);
}
return value.map((item: string) => item.trim());
}
function getCurrentMode(deps: DaemonDeps): ClientMode {
const stateMode = deps.store.getState()?.mode;
return stateMode || deps.mode || "apikeys";
@@ -375,6 +390,46 @@ function makeSdkLogger(...parts: string[]): SdkLogger {
};
}
/**
* Parse a `type` query parameter into a list of normalized transaction types.
* Accepts comma-separated values and is case-insensitive.
*/
export function parseHistoryTypes(raw: string | null | undefined): string[] {
if (!raw) return [];
return raw
.split(",")
.map((value) => value.trim().toLowerCase())
.filter((value) => value.length > 0);
}
/** Page size used when scanning history to apply a type filter. */
const HISTORY_TYPE_SCAN_PAGE_SIZE = 200;
/**
* Return history entries matching `types`, with offset/limit applied after
* filtering. The coco history repository only paginates by raw position, so we
* scan pages until enough matches accumulate (or history is exhausted).
*/
export async function getHistoryByTypes(
client: Pick<CocodClient, "getHistory">,
types: string[],
offset: number,
limit: number,
): Promise<HistoryEntry[]> {
const wanted = new Set(types);
const matched: HistoryEntry[] = [];
let scanOffset = 0;
while (matched.length < offset + limit) {
const page = await client.getHistory(scanOffset, HISTORY_TYPE_SCAN_PAGE_SIZE);
for (const entry of page) {
if (wanted.has(entry.type)) matched.push(entry);
}
if (page.length < HISTORY_TYPE_SCAN_PAGE_SIZE) break;
scanOffset += HISTORY_TYPE_SCAN_PAGE_SIZE;
}
return matched.slice(offset, offset + limit);
}
export function createDaemonRequestHandler(deps: {
provider: string | null;
server: { close(cb?: () => void): void };
@@ -473,12 +528,63 @@ export function createDaemonRequestHandler(deps: {
? body.minAgeMs
: undefined,
dryRun: body.dryRun === true,
force: body.force === true,
});
return { output: result };
});
return;
}
if (req.method === "POST" && url.pathname === "/wallet/recover") {
await respond(res, async () => {
if (!deps.walletClient.recoverMintQuotes) {
throw new CocodHttpError(
501,
"Mint quote recovery is not supported by this wallet client.",
);
}
const body = await readJsonBody(req);
const operationIds = optionalStringArrayField(body, "operationIds");
if (body.includeFailed === true && !operationIds?.length) {
throw new CocodHttpError(
400,
"'includeFailed' requires non-empty 'operationIds'.",
);
}
if (
body.timeoutMs !== undefined &&
(typeof body.timeoutMs !== "number" ||
!Number.isFinite(body.timeoutMs) || body.timeoutMs <= 0)
) {
throw new CocodHttpError(
400,
"'timeoutMs' must be a positive finite number.",
);
}
const result = await deps.walletClient.recoverMintQuotes({
operationIds,
includeFailed: body.includeFailed === true,
timeoutMs: body.timeoutMs as number | undefined,
});
return { output: result };
});
return;
}
if (req.method === "POST" && url.pathname === "/wallet/recover/operations") {
await respond(res, async () => {
if (!deps.walletClient.recoverStuckOperations) {
throw new CocodHttpError(
501,
"Stuck operation recovery is not supported by this wallet client.",
);
}
return { output: await deps.walletClient.recoverStuckOperations() };
});
return;
}
if (req.method === "POST" && url.pathname === "/wallet/receive/cashu") {
await respond(res, async () => {
const body = await readJsonBody(req);
@@ -596,9 +702,25 @@ export function createDaemonRequestHandler(deps: {
await respond(res, async () => {
const offsetParam = url.searchParams.get("offset");
const limitParam = url.searchParams.get("limit");
const idParam = url.searchParams.get("id")?.trim() || "";
const types = parseHistoryTypes(url.searchParams.get("type"));
const offset = offsetParam ? parseInt(offsetParam, 10) || 0 : 0;
const limit = limitParam ? parseInt(limitParam, 10) || 50 : 50;
const entries = await deps.walletClient.getHistory(offset, limit);
let entries: HistoryEntry[];
if (idParam) {
const entry = await deps.walletClient.getHistoryEntryById(idParam);
entries = entry ? [entry] : [];
} else if (types.length > 0) {
entries = await getHistoryByTypes(
deps.walletClient,
types,
offset,
limit,
);
} else {
entries = await deps.walletClient.getHistory(offset, limit);
}
const encoded = entries.map((entry: HistoryEntry) => {
const base = { ...entry } as Record<string, unknown>;
+99
View File
@@ -0,0 +1,99 @@
import { describe, expect, it, mock } from "bun:test";
import { EventEmitter } from "node:events";
import { createDaemonRequestHandler } from "./index";
async function recover(body: unknown) {
const recoverMintQuotes = mock(async (_options: unknown) => ({ recovered: 0 }));
const handler = createDaemonRequestHandler({ walletClient: { recoverMintQuotes } } as never);
const req = new EventEmitter() as any;
Object.assign(req, { method: "POST", url: "/wallet/recover", 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(JSON.stringify(body)));
req.emit("end");
});
await handler(req, res as never);
return { res, recoverMintQuotes };
}
describe("POST /wallet/recover validation", () => {
it.each([{}, { operationIds: [] }])("rejects includeFailed without explicit IDs: %j", async (body) => {
const { res, recoverMintQuotes } = await recover({ ...body, includeFailed: true });
expect(res.status).toBe(400);
expect(recoverMintQuotes).not.toHaveBeenCalled();
});
it.each([[""], [" "], [42]])("rejects invalid operation IDs: %j", async (operationIds) => {
const { res, recoverMintQuotes } = await recover({ operationIds });
expect(res.status).toBe(400);
expect(recoverMintQuotes).not.toHaveBeenCalled();
});
it.each([0, -1, "1000"])("rejects invalid timeout %j", async (timeoutMs) => {
const { res, recoverMintQuotes } = await recover({ timeoutMs });
expect(res.status).toBe(400);
expect(recoverMintQuotes).not.toHaveBeenCalled();
});
it("passes normalized explicit IDs and a positive timeout", async () => {
const { res, recoverMintQuotes } = await recover({ operationIds: [" op-1 "], includeFailed: true, timeoutMs: 1000 });
expect(res.status).toBe(200);
expect(recoverMintQuotes).toHaveBeenCalledWith({ operationIds: ["op-1"], includeFailed: true, timeoutMs: 1000 });
});
it("still permits checking pending quotes without IDs", async () => {
const { res, recoverMintQuotes } = await recover({});
expect(res.status).toBe(200);
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);
});
});
+12 -1
View File
@@ -1,5 +1,5 @@
import { describe, expect, it } from "bun:test";
import { selectCleanupOperations } from "./cleanup";
import { selectCleanupOperations, summarizeMintCleanup } from "./cleanup";
const NOW_MS = 1_800_000_000_000;
const DAY_MS = 24 * 60 * 60 * 1000;
@@ -166,3 +166,14 @@ describe("selectCleanupOperations", () => {
expect(result.meltsToCancel).toEqual([]);
});
});
describe("mint cleanup reporting", () => {
it("reports dry-run candidates, not confirmed failures", () => {
expect(summarizeMintCleanup({ dryRun: true, candidates: 3, failed: 0, leftForRecovery: 0 }))
.toEqual({ mintQuoteCandidates: 3, failedMintQuotes: 0, leftForRecovery: 0 });
});
it("reports only actual failures in a real run", () => {
expect(summarizeMintCleanup({ dryRun: false, candidates: 3, failed: 1, leftForRecovery: 2 }))
.toEqual({ mintQuoteCandidates: 3, failedMintQuotes: 1, leftForRecovery: 2 });
});
});
+16 -2
View File
@@ -64,8 +64,8 @@ export interface CleanupSelection<
* can have happened before expiry while the daemon was down, leaving no
* local observation. Callers that fail quotes automatically at startup must
* therefore confirm UNPAID with the mint first (see
* settleExpiredMintQuotes in coco-client.ts); only the explicit,
* user-invoked cleanup command may fail candidates purely locally.
* failExpiredMintQuoteIfUnpaid in coco-client.ts). Explicit cleanup follows
* the same rule unless the operator opts into unsafe `--force` behaviour.
* - Pending sends are reclaimed (rolled back) only when they are older than
* `minAgeMs`, so we never roll back a token that a receiver might still
* legitimately claim.
@@ -102,3 +102,17 @@ export function selectCleanupOperations<
return { mintsToFail, sendsToReclaim, meltsToCancel };
}
/** Keep a local-only dry-run preview distinct from mint-confirmed outcomes. */
export function summarizeMintCleanup(input: {
dryRun: boolean;
candidates: number;
failed: number;
leftForRecovery: number;
}) {
return {
mintQuoteCandidates: input.candidates,
failedMintQuotes: input.dryRun ? 0 : input.failed,
leftForRecovery: input.dryRun ? 0 : input.leftForRecovery,
};
}
+827
View File
@@ -14,12 +14,18 @@ import {
assertLegacyCocodNotRunning,
claimLegacyCocodPidFile,
createCocoClient,
createRecoveryGate,
createRunQueue,
DEFAULT_TRUSTED_MINT_URLS,
failExpiredMintQuoteIfUnpaid,
isZombieProcess,
reopenFailedMintOperation,
runMintQuoteRecovery,
settleExpiredMintQuotes,
settlePendingMintQuotes,
stopLegacyCocod,
type ExpiredMintQuoteSource,
type MintQuoteRecoverySource,
type PendingMintQuoteSource,
type PendingMintSweepState,
} from "./coco-client";
@@ -675,6 +681,29 @@ describe("settleExpiredMintQuotes", () => {
expect(observePendingOperation).not.toHaveBeenCalled();
expect(failPendingOperation).not.toHaveBeenCalled();
});
it("skips quotes at probe-unreachable mints without spending the budget", async () => {
const ops = [
pendingMintOp({ id: "dead-1", mintUrl: "https://dead.example.com" }),
pendingMintOp({ id: "dead-2", mintUrl: "https://dead.example.com/" }),
pendingMintOp({ id: "live-1", mintUrl: "https://live.example.com" }),
];
const { source, observePendingOperation, failPendingOperation } =
fakeSource(ops);
const result = await settleExpiredMintQuotes(
source,
NOW_MS,
undefined,
{ unreachableMints: new Set(["https://dead.example.com"]) },
);
expect(result).toEqual({ failed: 1, leftForRecovery: 0, unobserved: 2 });
// Only the reachable mint was asked anything.
expect(observePendingOperation).toHaveBeenCalledTimes(1);
expect(observePendingOperation.mock.calls[0]?.[0]).toBe("live-1");
expect(failPendingOperation).toHaveBeenCalledTimes(1);
});
});
describe("settlePendingMintQuotes", () => {
@@ -901,3 +930,801 @@ describe("settlePendingMintQuotes", () => {
expect(logged.mock.calls[0]?.[0]).toContain("21 sat minted");
});
});
describe("createRunQueue", () => {
it("runs tasks strictly one after another", async () => {
const enqueue = createRunQueue();
const order: string[] = [];
let active = 0;
let maxActive = 0;
const task = (name: string, delay: number) => async () => {
active++;
maxActive = Math.max(maxActive, active);
order.push(`${name}:start`);
await new Promise((resolve) => setTimeout(resolve, delay));
order.push(`${name}:end`);
active--;
return name;
};
const results = await Promise.all([
enqueue(task("a", 20)),
enqueue(task("b", 1)),
enqueue(task("c", 1)),
]);
expect(results).toEqual(["a", "b", "c"]);
expect(maxActive).toBe(1);
expect(order).toEqual([
"a:start",
"a:end",
"b:start",
"b:end",
"c:start",
"c:end",
]);
});
it("keeps the chain alive after a rejected task", async () => {
const enqueue = createRunQueue();
const failed = enqueue(async () => {
throw new Error("boom");
});
const next = enqueue(async () => "ok");
await expect(failed).rejects.toThrow("boom");
expect(await next).toBe("ok");
});
});
describe("failExpiredMintQuoteIfUnpaid", () => {
function fakeMintService(observe: (id: string) => Promise<{ category: "waiting" | "ready" | "completed" | "terminal" }>) {
const failPendingOperation = mock(
async (
_op: { id: string },
_failure: { reason: string; retryable?: boolean; observedAt: number },
) => ({}),
);
return {
mintService: {
observePendingOperation: mock(observe),
failPendingOperation,
},
failPendingOperation,
};
}
it("fails a quote its mint confirms unpaid", async () => {
const { mintService, failPendingOperation } = fakeMintService(async () => ({
category: "waiting",
}));
const result = await failExpiredMintQuoteIfUnpaid(mintService, "op-1", 1000);
expect(result.outcome).toBe("failed");
expect(failPendingOperation).toHaveBeenCalledTimes(1);
expect(failPendingOperation.mock.calls[0]?.[1]?.reason).toContain(
"confirmed unpaid by mint",
);
});
it.each(["ready", "completed", "terminal"] as const)(
"leaves a quote observed as %s for recovery",
async (category) => {
const { mintService, failPendingOperation } = fakeMintService(
async () => ({ category }),
);
const result = await failExpiredMintQuoteIfUnpaid(mintService, "op-1", 1000);
expect(result).toEqual({ outcome: "leftForRecovery", category });
expect(failPendingOperation).not.toHaveBeenCalled();
},
);
it("leaves a quote pending when the mint cannot be reached", async () => {
const { mintService, failPendingOperation } = fakeMintService(async () => {
throw new Error("Network request failed");
});
const result = await failExpiredMintQuoteIfUnpaid(mintService, "op-1", 1000);
expect(result.outcome).toBe("unobserved");
expect(failPendingOperation).not.toHaveBeenCalled();
});
it("gives up waiting on a hung mint without failing the quote", async () => {
const { mintService, failPendingOperation } = fakeMintService(
() => new Promise(() => {}),
);
const result = await failExpiredMintQuoteIfUnpaid(mintService, "op-1", 20);
expect(result.outcome).toBe("unobserved");
expect(failPendingOperation).not.toHaveBeenCalled();
});
});
describe("reopenFailedMintOperation", () => {
/**
* Mirrors coco's OperationIdLock, which is fail-fast: acquiring an id that is
* already locked throws OperationInProgressError instead of waiting.
*/
function makeLock() {
let held = false;
return {
get held() {
return held;
},
async acquire() {
if (held) {
const error = new Error("Operation op-1 is already in progress");
error.name = "OperationInProgressError";
throw error;
}
held = true;
return () => {
held = false;
};
},
};
}
function fakeService(
current: Record<string, unknown> | null,
hooks: {
lock?: ReturnType<typeof makeLock>;
onWrite?: (lock: ReturnType<typeof makeLock>) => void;
} = {},
) {
const lock = hooks.lock ?? makeLock();
const transitionToPending = mock(
async (_op: Record<string, unknown>, _error?: string) => {
hooks.onWrite?.(lock);
return {};
},
);
return {
service: {
acquireOperationLock: mock(async (_id: string) => lock.acquire()),
getOperation: mock(async (_id: string) => current),
transitionToPending,
},
transitionToPending,
lock,
};
}
it("re-opens a failed operation with the full persisted row", async () => {
// coco spreads whatever it is handed and the sqlite repository rewrites
// every column, so a partial object would erase the stored outputs.
const row = {
id: "op-1",
state: "failed",
mintUrl: "https://mint.example.com",
quoteId: "quote-1",
method: "bolt11",
amount: 210_000,
unit: "sat",
request: "lnbc...",
expiry: 1_800_000_000,
outputDataJson: "[{\"secret\":\"abc\"}]",
terminalFailure: { reason: "expired" },
};
const { service, transitionToPending } = fakeService(row);
const reopened = await reopenFailedMintOperation(service, "op-1");
expect(reopened).toBe(true);
expect(transitionToPending).toHaveBeenCalledTimes(1);
const passed = transitionToPending.mock.calls[0]?.[0];
expect(passed).toMatchObject({
id: "op-1",
quoteId: "quote-1",
amount: 210_000,
unit: "sat",
outputDataJson: "[{\"secret\":\"abc\"}]",
});
// The stale terminal marker must not survive the re-open.
expect(passed?.terminalFailure).toBeUndefined();
});
it("does nothing when the operation is no longer failed", async () => {
const { service, transitionToPending, lock } = fakeService({
id: "op-1",
state: "finalized",
});
expect(await reopenFailedMintOperation(service, "op-1")).toBe(false);
expect(transitionToPending).not.toHaveBeenCalled();
// The lock must be released even on the no-op path.
expect(lock.held).toBe(false);
});
it("throws when the operation is missing", async () => {
const { service, lock } = fakeService(null);
await expect(reopenFailedMintOperation(service, "op-1")).rejects.toThrow(
"not found",
);
expect(lock.held).toBe(false);
});
it("fails closed when coco no longer exposes the operation lock", async () => {
const transitionToPending = mock(
async (_op: Record<string, unknown>, _error?: string) => ({}),
);
const service = {
getOperation: mock(async () => ({ id: "op-1", state: "failed" })),
transitionToPending,
} as unknown as Parameters<typeof reopenFailedMintOperation>[0];
await expect(
reopenFailedMintOperation(service, "op-1"),
).rejects.toThrow("acquireOperationLock");
expect(transitionToPending).not.toHaveBeenCalled();
});
it("holds the operation lock across read-check-write", async () => {
const lock = makeLock();
const order: string[] = [];
const service = {
acquireOperationLock: mock(async (_id: string) => {
order.push("lock");
const release = await lock.acquire();
return () => {
order.push("unlock");
release();
};
}),
getOperation: mock(async (_id: string) => {
expect(lock.held).toBe(true);
order.push("read");
return { id: "op-1", state: "failed", quoteId: "quote-1" };
}),
transitionToPending: mock(async () => {
expect(lock.held).toBe(true);
order.push("write");
return {};
}),
};
const reopened = await reopenFailedMintOperation(service, "op-1");
expect(reopened).toBe(true);
expect(order).toEqual(["lock", "read", "write", "unlock"]);
expect(lock.held).toBe(false);
});
it("refuses to re-open while the operation lock is held elsewhere", async () => {
const lock = makeLock();
const release = await lock.acquire();
const { service, transitionToPending } = fakeService(
{ id: "op-1", state: "failed" },
{ lock },
);
await expect(reopenFailedMintOperation(service, "op-1")).rejects.toThrow(
/in progress/,
);
expect(transitionToPending).not.toHaveBeenCalled();
release();
});
});
describe("runMintQuoteRecovery", () => {
function mintOp(overrides: Record<string, unknown> = {}) {
return {
id: "op-1",
mintUrl: "https://mint.example.com",
quoteId: "quote-1",
state: "pending",
amount: 210_000,
expiry: 0,
...overrides,
};
}
function fakeSource(
ops: Array<Record<string, unknown>>,
behavior: {
observe?: (id: string) => Promise<{
category: "waiting" | "ready" | "completed" | "terminal";
}>;
finalize?: (id: string) => Promise<unknown>;
reopen?: (id: string) => Promise<boolean>;
} = {},
) {
const finalize = mock(
behavior.finalize ??
(async (_id: string) => ({ state: "finalized" })),
);
const observePendingOperation = mock(
behavior.observe ?? (async () => ({ category: "waiting" as const })),
);
const reopenFailedOperation = mock(
behavior.reopen ?? (async (_id: string) => true),
);
const byId = new Map(ops.map((op) => [op.id as string, op]));
const source = {
ops: {
mint: {
listPending: async () =>
ops.filter(
(op) => op.state === "pending" || op.state === "executing",
),
get: async (id: string) => byId.get(id) ?? null,
finalize,
},
},
mintOperationService: { observePendingOperation },
reopenFailedOperation,
} as unknown as MintQuoteRecoverySource;
return { source, finalize, observePendingOperation, reopenFailedOperation };
}
it("mints the stored outputs for a quote the mint reports PAID", async () => {
const { source, finalize } = fakeSource([mintOp()], {
observe: async () => ({ category: "ready" }),
});
const result = await runMintQuoteRecovery(source);
expect(result).toMatchObject({ checked: 1, recovered: 1, waiting: 0 });
expect(finalize).toHaveBeenCalledTimes(1);
expect(finalize.mock.calls[0]?.[0]).toBe("op-1");
});
it("restores proofs for a quote already issued at the mint", async () => {
const { source, finalize } = fakeSource([mintOp()], {
observe: async () => ({ category: "completed" }),
});
const result = await runMintQuoteRecovery(source);
expect(result).toMatchObject({ recovered: 1 });
expect(finalize).toHaveBeenCalledTimes(1);
});
it("leaves an unpaid quote pending", async () => {
const { source, finalize } = fakeSource([mintOp()], {
observe: async () => ({ category: "waiting" }),
});
const result = await runMintQuoteRecovery(source);
expect(result).toMatchObject({ checked: 1, recovered: 0, waiting: 1 });
expect(finalize).not.toHaveBeenCalled();
});
it("reports a quote the mint can no longer issue", async () => {
const { source, finalize } = fakeSource([mintOp()], {
observe: async () => ({ category: "terminal" }),
});
const result = await runMintQuoteRecovery(source);
expect(result).toMatchObject({ terminal: 1, recovered: 0 });
expect(finalize).not.toHaveBeenCalled();
});
it("retries later when the mint is unreachable", async () => {
const { source, finalize } = fakeSource([mintOp()], {
observe: async () => {
throw new Error("fetch failed");
},
});
const result = await runMintQuoteRecovery(source);
expect(result).toMatchObject({ retryable: 1, recovered: 0 });
expect(result.errors).toHaveLength(1);
expect(finalize).not.toHaveBeenCalled();
});
it("does not count a failed finalize as a recovery", async () => {
// coco returns a terminal operation instead of throwing when the mint
// refuses, so a fulfilled finalize is not evidence that sats were claimed.
const { source, finalize } = fakeSource([mintOp()], {
observe: async () => ({ category: "ready" }),
finalize: async () => ({
state: "failed",
error: "Recovered: quote quote-1 expired while executing mint",
}),
});
const result = await runMintQuoteRecovery(source);
expect(result).toMatchObject({ recovered: 0, terminal: 1 });
expect(result.errors[0]?.error).toContain("expired");
expect(finalize).toHaveBeenCalledTimes(1);
});
it("does not count a finalized-with-error operation as a recovery", async () => {
const { source } = fakeSource([mintOp()], {
observe: async () => ({ category: "completed" }),
finalize: async () => ({
state: "finalized",
error: "Recovered issued quote quote-1 but no proofs could be restored",
}),
});
const result = await runMintQuoteRecovery(source);
expect(result).toMatchObject({ recovered: 0, terminal: 1 });
expect(result.errors).toHaveLength(1);
});
it("falls back to the original failure when diagnostic lookup fails", async () => {
const { source } = fakeSource([mintOp()], {
observe: async () => ({ category: "ready" }),
finalize: async () => { throw new Error("original failure"); },
});
source.ops.mint.get = async () => { throw new Error("lookup failed"); };
const result = await runMintQuoteRecovery(source);
expect(result).toMatchObject({ retryable: 1, recovered: 0 });
expect(result.errors).toEqual([{ operationId: "op-1", error: "original failure" }]);
});
it("bounds a hung diagnostic lookup after finalize fails", async () => {
const { source } = fakeSource([mintOp()], {
observe: async () => ({ category: "ready" }),
finalize: async () => { throw new Error("original failure"); },
});
source.ops.mint.get = () => new Promise(() => {});
const result = await runMintQuoteRecovery(source, { timeoutMs: 20 });
expect(result).toMatchObject({ retryable: 1, recovered: 0 });
expect(result.errors).toEqual([{ operationId: "op-1", error: "original failure" }]);
});
it("surfaces the persisted mint error even when the mint budget is nearly spent", async () => {
const { source } = fakeSource([mintOp()], {
observe: async () => ({ category: "ready" }),
finalize: async () => {
// Consume almost the whole per-op mint budget before throwing: exactly
// the slow-mint case where a diagnostic bound to remaining() would
// starve and silently fall back to the generic message.
await new Promise((resolve) => setTimeout(resolve, 25));
throw new Error("remains pending");
},
});
source.ops.mint.get = async () => {
await new Promise((resolve) => setTimeout(resolve, 50));
return {
...mintOp(),
state: "pending",
error: "keyset id inactive.",
};
};
const result = await runMintQuoteRecovery(source, { timeoutMs: 30 });
expect(result).toMatchObject({ retryable: 1, recovered: 0 });
expect(result.errors).toEqual([
{ operationId: "op-1", error: "keyset id inactive." },
]);
});
it("bounds finalize so one hung mint cannot block recovery", async () => {
const { source } = fakeSource([mintOp()], {
observe: async () => ({ category: "ready" }),
finalize: () => new Promise(() => {}),
});
const started = Date.now();
const result = await runMintQuoteRecovery(source, { timeoutMs: 20 });
expect(Date.now() - started).toBeLessThan(5_000);
expect(result).toMatchObject({ retryable: 1, recovered: 0 });
});
it("recovers an interrupted mint without re-checking the quote", async () => {
const { source, finalize, observePendingOperation } = fakeSource([
mintOp({ state: "executing" }),
]);
const result = await runMintQuoteRecovery(source);
expect(result).toMatchObject({ recovered: 1 });
expect(finalize).toHaveBeenCalledTimes(1);
expect(observePendingOperation).not.toHaveBeenCalled();
});
it("skips failed operations unless the caller opts in", async () => {
const { source, reopenFailedOperation } = fakeSource([
mintOp({ state: "failed", lastObservedRemoteState: "PAID" }),
]);
const result = await runMintQuoteRecovery(source);
expect(result).toMatchObject({ checked: 0, recovered: 0, reopened: 0 });
expect(reopenFailedOperation).not.toHaveBeenCalled();
});
it("requires explicit IDs when including failed operations", async () => {
const { source, reopenFailedOperation, observePendingOperation } = fakeSource([]);
await expect(runMintQuoteRecovery(source, { includeFailed: true })).rejects.toThrow(
"includeFailed requires explicit operationIds",
);
await expect(runMintQuoteRecovery(source, { includeFailed: true, operationIds: [] })).rejects.toThrow(
"includeFailed requires explicit operationIds",
);
expect(reopenFailedOperation).not.toHaveBeenCalled();
expect(observePendingOperation).not.toHaveBeenCalled();
});
it.each([0, -1, NaN, Infinity])("rejects invalid recovery timeout %s", async (timeoutMs) => {
const { source, observePendingOperation } = fakeSource([]);
await expect(runMintQuoteRecovery(source, { timeoutMs })).rejects.toThrow(
"timeoutMs must be a positive finite number",
);
expect(observePendingOperation).not.toHaveBeenCalled();
});
it("re-opens a named failed operation, then mints it", async () => {
const { source, reopenFailedOperation, finalize } = fakeSource(
[mintOp({ state: "failed", lastObservedRemoteState: "PAID" })],
{ observe: async () => ({ category: "ready" }) },
);
const result = await runMintQuoteRecovery(source, {
operationIds: ["op-1"],
includeFailed: true,
});
expect(result).toMatchObject({ reopened: 1, recovered: 1 });
expect(reopenFailedOperation).toHaveBeenCalledWith("op-1");
expect(finalize).toHaveBeenCalledTimes(1);
});
it("re-opens a named failed operation even without a PAID observation", async () => {
// The old local-fail bug left quotes with a stale or missing observation,
// which is exactly when an operator needs to retry them.
const { source, reopenFailedOperation } = fakeSource(
[mintOp({ state: "failed" })],
{ observe: async () => ({ category: "ready" }) },
);
const result = await runMintQuoteRecovery(source, {
operationIds: ["op-1"],
includeFailed: true,
});
expect(result).toMatchObject({ reopened: 1, recovered: 1 });
expect(reopenFailedOperation).toHaveBeenCalledTimes(1);
});
it("skips an operation that is no longer failed when re-opened", async () => {
const { source, finalize } = fakeSource(
[mintOp({ state: "failed" })],
{ reopen: async () => false },
);
const result = await runMintQuoteRecovery(source, {
operationIds: ["op-1"],
includeFailed: true,
});
expect(result).toMatchObject({ reopened: 0, checked: 0, recovered: 0 });
expect(finalize).not.toHaveBeenCalled();
});
it("reports an unknown operation id instead of throwing", async () => {
const { source } = fakeSource([]);
const result = await runMintQuoteRecovery(source, {
operationIds: ["missing"],
});
expect(result.checked).toBe(0);
expect(result.errors).toEqual([
{ operationId: "missing", error: "operation not found" },
]);
});
it("deduplicates repeated operation ids", async () => {
const { source, finalize } = fakeSource([mintOp()], {
observe: async () => ({ category: "ready" }),
});
const result = await runMintQuoteRecovery(source, {
operationIds: ["op-1", "op-1"],
});
expect(result.checked).toBe(1);
expect(finalize).toHaveBeenCalledTimes(1);
});
it("ignores finalized operations even when targeted", async () => {
const { source, finalize } = fakeSource([mintOp({ state: "finalized" })]);
const result = await runMintQuoteRecovery(source, {
operationIds: ["op-1"],
});
expect(result.checked).toBe(0);
expect(finalize).not.toHaveBeenCalled();
});
it("leaves a non-terminal finalize result for a later run, not terminal", async () => {
const { source } = fakeSource([mintOp()], {
observe: async () => ({ category: "ready" }),
finalize: async () => ({ state: "pending" }),
});
const result = await runMintQuoteRecovery(source);
expect(result).toMatchObject({ recovered: 0, terminal: 0, retryable: 1 });
expect(result.errors[0]?.error).toContain("will retry");
});
it("skips operations whose earlier recovery is still in flight", async () => {
const outstanding = new Map<string, Promise<unknown>>([
["mint:op-1", new Promise(() => {})],
]);
const { source, finalize, observePendingOperation } = fakeSource(
[mintOp()],
{ observe: async () => ({ category: "ready" }) },
);
const result = await runMintQuoteRecovery(source, { outstanding });
expect(result).toMatchObject({ busy: 1, checked: 0, recovered: 0 });
expect(observePendingOperation).not.toHaveBeenCalled();
expect(finalize).not.toHaveBeenCalled();
});
it("keeps a timed-out finalize registered so a retry waits", async () => {
const outstanding = new Map<string, Promise<unknown>>();
const { source } = fakeSource([mintOp()], {
observe: async () => ({ category: "ready" }),
finalize: () => new Promise(() => {}),
});
const first = await runMintQuoteRecovery(source, {
timeoutMs: 20,
outstanding,
});
const second = await runMintQuoteRecovery(source, {
timeoutMs: 20,
outstanding,
});
expect(first).toMatchObject({ retryable: 1, recovered: 0 });
expect(outstanding.has("mint:op-1")).toBe(true);
// The abandoned mint request must not be retried underneath.
expect(second).toMatchObject({ busy: 1, checked: 0 });
});
it("does not re-open a failed operation whose recovery is in flight", async () => {
const outstanding = new Map<string, Promise<unknown>>([
["mint:op-1", new Promise(() => {})],
]);
const { source, reopenFailedOperation } = fakeSource(
[mintOp({ state: "failed" })],
{},
);
const result = await runMintQuoteRecovery(source, {
operationIds: ["op-1"],
includeFailed: true,
outstanding,
});
expect(result).toMatchObject({ busy: 1, reopened: 0 });
expect(reopenFailedOperation).not.toHaveBeenCalled();
});
it("counts an in-progress operation as busy rather than retryable", async () => {
const { source } = fakeSource(
[mintOp({ state: "failed" })],
{
reopen: async () => {
const error = new Error("Operation op-1 is already in progress");
error.name = "OperationInProgressError";
throw error;
},
},
);
const result = await runMintQuoteRecovery(source, {
operationIds: ["op-1"],
includeFailed: true,
});
expect(result).toMatchObject({ busy: 1, retryable: 0, reopened: 0 });
});
it("keeps a timed-out quote check registered so a retry waits", async () => {
// observePendingOperation is not read-only, so a hung check must not be
// retried underneath: it could persist a stale observation later.
const outstanding = new Map<string, Promise<unknown>>();
const { source, finalize } = fakeSource([mintOp()], {
observe: () => new Promise(() => {}),
});
const first = await runMintQuoteRecovery(source, {
timeoutMs: 20,
outstanding,
});
const second = await runMintQuoteRecovery(source, {
timeoutMs: 20,
outstanding,
});
expect(first).toMatchObject({ retryable: 1, checked: 1 });
expect(outstanding.has("mint:op-1")).toBe(true);
expect(second).toMatchObject({ busy: 1, checked: 0 });
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;
});
});
File diff suppressed because it is too large Load Diff
+75 -1
View File
@@ -83,13 +83,23 @@ export interface WalletCleanupOptions {
minAgeMs?: number;
/** Report what would be cleaned without applying changes. */
dryRun?: boolean;
/**
* Fail expired mint quotes without confirming UNPAID with the mint. Only for
* operators who accept the risk of stranding a quote that was paid before
* its invoice expired; recovery is the safe default.
*/
force?: boolean;
}
/** Summary of a wallet cleanup run. */
export interface WalletCleanupResult {
dryRun: boolean;
/** Number of expired pending mint quotes marked as failed. */
/** Expired quotes selected for checking; dry runs do not contact the mint. */
mintQuoteCandidates: number;
/** Number actually marked failed (always zero in a dry run). */
failedMintQuotes: number;
/** Expired quotes kept pending because they are paid/issued or unverified. */
leftForRecovery: number;
/** Number of stale pending send operations reclaimed. */
reclaimedSends: number;
/** Number of stale prepared melt operations cancelled. */
@@ -120,6 +130,51 @@ export interface MintQuoteStatus {
error?: string;
}
/** Options for explicit PAID mint-quote recovery. */
export interface WalletMintQuoteRecoveryOptions {
/** Target only these operation ids (may include failed operations). */
operationIds?: string[];
/** Re-open failed operations instead of skipping them. */
includeFailed?: boolean;
/** Per-quote mint timeout in milliseconds. */
timeoutMs?: number;
}
/** Summary of a PAID mint-quote recovery run. */
export interface WalletMintQuoteRecoveryResult {
/** Operations whose quote state was checked with the mint. */
checked: number;
/** Operations whose paid sats were minted or restored. */
recovered: number;
/** Quotes the mint still reports UNPAID; left pending. */
waiting: number;
/** Quotes the mint can no longer issue. */
terminal: number;
/** Failed operations moved back to pending before checking. */
reopened: number;
/** Operations left to a later run (mint unreachable, budget spent, non-terminal). */
retryable: number;
/** Operations skipped because an earlier recovery of them is still running. */
busy: number;
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 {
ping(): Promise<boolean>;
getStatus(): Promise<CocodState>;
@@ -142,6 +197,8 @@ export interface CocodClient {
/** Release resources held by in-process wallet implementations. */
dispose?(): Promise<void>;
getHistory(offset?: number, limit?: number): Promise<HistoryEntry[]>;
/** Look up a single transaction by its history entry ID. */
getHistoryEntryById(id: string): Promise<HistoryEntry | null>;
/** NPC (npubx.cash) Lightning address for this wallet. */
getNpcAddress(): Promise<NpcAddress>;
/** Claim an NPC username; pass confirm=true to pay the claim fee from the wallet. */
@@ -152,6 +209,20 @@ export interface CocodClient {
cleanupStuckOperations?(
options?: WalletCleanupOptions,
): Promise<WalletCleanupResult>;
/**
* Re-issue PAID mint quotes whose sats were never claimed, optionally
* targeting specific operations (including ones coco already failed).
*/
recoverMintQuotes?(
options?: WalletMintQuoteRecoveryOptions,
onProgress?: (message: string) => void,
): 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. */
getRecoveryProgress?(): Promise<WalletRecoveryProgress>;
}
@@ -465,6 +536,9 @@ export function createCocodClient(
async getHistory(_offset?: number, _limit?: number): Promise<HistoryEntry[]> {
return [];
},
async getHistoryEntryById(_id: string): Promise<HistoryEntry | null> {
return null;
},
async getNpcAddress(): Promise<NpcAddress> {
const address = await callDaemon<string>("/npc/address");
if (typeof address !== "string" || !address.trim()) {
@@ -0,0 +1,163 @@
/**
* Re-opening a failed mint operation is the one genuinely destructive step of
* PAID-quote recovery, because coco's private `transitionToPending` spreads
* whatever it is handed and `SqliteMintOperationRepository.update` rewrites
* every column. A partial object such as `{ id }` is therefore rejected by the
* NOT NULL schema, and on a more permissive adapter would overwrite `quoteId`,
* `amount`, `request`, `lastObservedRemoteState` and `outputDataJson` with NULL,
* destroying the material needed to claim the paid sats.
*
* These tests drive the production `reopenFailedMintOperation` helper against a
* real coco sqlite repository through adapter-backed getOperation and
* transitionToPending implementations, so a regression at the helper/service
* boundary is caught rather than a mock standing in for it.
*/
import { afterEach, describe, expect, it } from "bun:test";
import { Database } from "bun:sqlite";
import { mkdtempSync, rmSync } from "node:fs";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { SqliteRepositories } from "@cashu/coco-sqlite-bun";
import { reopenFailedMintOperation } from "./coco-client";
const OUTPUT_DATA = [
{
blindedMessage: { amount: "210000", id: "00deadbeef", B_: "02deadbeef" },
blindingFactor: "1234567890",
secret: "aabbccdd",
},
];
function failedRow(state = "failed") {
return {
id: "op-1",
mintUrl: "https://mint.example.com",
quoteId: "quote-1",
state,
createdAt: 1_000,
updatedAt: 2_000,
error: state === "failed" ? "expired" : undefined,
method: "bolt11",
methodData: { method: "bolt11", data: {} },
amount: 210_000,
unit: "sat",
request: "lnbc1example",
expiry: 1_800_000_000,
pubkey: undefined,
lastObservedRemoteState: "PAID",
lastObservedRemoteStateAt: 3_000,
terminalFailure:
state === "failed" ? { reason: "expired", observedAt: 3_000 } : undefined,
outputData: OUTPUT_DATA,
};
}
describe("reopenFailedMintOperation against real coco sqlite", () => {
let dir: string | undefined;
let database: Database | undefined;
afterEach(() => {
database?.close();
database = undefined;
if (dir) rmSync(dir, { recursive: true, force: true });
dir = undefined;
});
async function repos() {
dir = mkdtempSync(join(tmpdir(), "routstrd-reopen-"));
database = new Database(join(dir, "coco.db"));
const repositories = new SqliteRepositories({ database });
await repositories.init();
return repositories;
}
function readRow(repositories: SqliteRepositories, id: string) {
return repositories.mintOperationRepository.getById(id) as unknown as Promise<
Record<string, unknown> | null
>;
}
/**
* Adapter-backed stand-in for the private coco service methods the helper
* uses: getOperation reads and transitionToPending mirrors coco's
* spread-and-update implementation.
*/
function serviceOver(repositories: SqliteRepositories) {
return {
acquireOperationLock: async (_id: string) => () => {},
getOperation: (id: string) => readRow(repositories, id),
transitionToPending: async (
op: Record<string, unknown>,
error?: string,
) => {
await repositories.mintOperationRepository.update({
...op,
state: "pending",
error,
} as never);
},
};
}
it("re-opens a failed row without losing quote metadata or stored outputs", async () => {
const repositories = await repos();
await repositories.mintOperationRepository.create(
failedRow() as never,
);
const reopened = await reopenFailedMintOperation(
serviceOver(repositories),
"op-1",
);
expect(reopened).toBe(true);
const row = await readRow(repositories, "op-1");
expect(row?.state).toBe("pending");
expect(row?.quoteId).toBe("quote-1");
expect(row?.amount).toBe(210_000);
expect(row?.unit).toBe("sat");
expect(row?.request).toBe("lnbc1example");
expect(row?.lastObservedRemoteState).toBe("PAID");
expect(row?.outputData).toEqual(OUTPUT_DATA);
expect(row?.terminalFailure ?? undefined).toBeUndefined();
});
it("is a no-op when the operation is no longer failed", async () => {
const repositories = await repos();
await repositories.mintOperationRepository.create(
failedRow("finalized") as never,
);
const reopened = await reopenFailedMintOperation(
serviceOver(repositories),
"op-1",
);
expect(reopened).toBe(false);
const row = await readRow(repositories, "op-1");
expect(row?.state).toBe("finalized");
expect(row?.outputData).toEqual(OUTPUT_DATA);
});
it("rejects a partial row, which is why the helper reloads in full", async () => {
// Locks in the reason for reloading. If a future coco version accepts
// partial updates this fails, and the helper can be simplified rather than
// silently losing paid sats.
const repositories = await repos();
await repositories.mintOperationRepository.create(
failedRow() as never,
);
await expect(
repositories.mintOperationRepository.update({
id: "op-1",
state: "pending",
updatedAt: Date.now(),
} as never),
).rejects.toThrow();
const unchanged = await readRow(repositories, "op-1");
expect(unchanged?.state).toBe("failed");
expect(unchanged?.outputData).toEqual(OUTPUT_DATA);
});
});
@@ -0,0 +1,386 @@
/**
* End-to-end PAID mint-quote recovery against a real coco Manager, real sqlite
* and a real in-process mint that produces genuine blind signatures.
*
* Nothing here mocks the wallet: a quote is created through coco, the mint is
* told what to report, and the production `runMintQuoteRecovery` drives the
* outcome. These are the release-gating scenarios for the feature, and they are
* the only tests that exercise issuance and NUT-09 restore over HTTP.
*/
import { afterEach, describe, expect, it } from "bun:test";
import { Database } from "bun:sqlite";
import { mkdtempSync, rmSync } from "node:fs";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { Manager } from "@cashu/coco-core";
import { SqliteRepositories } from "@cashu/coco-sqlite-bun";
import { QUOTE_EXPIRED, FakeMint } from "./testing/fake-mint";
import { reopenFailedMintOperation, runMintQuoteRecovery } from "./coco-client";
type AnyRecord = Record<string, unknown>;
interface Booted {
manager: Manager;
repositories: SqliteRepositories;
mint: FakeMint;
/** Build the recovery source the production function expects. */
source: () => AnyRecord;
spendable: () => Promise<number>;
close: () => Promise<void>;
}
async function boot(options: { quoteExpiry?: number | null } = {}) {
const mint = new FakeMint();
mint.quoteExpiry = options.quoteExpiry ?? null;
mint.start();
const dir = mkdtempSync(join(tmpdir(), "routstrd-fakemint-"));
const database = new Database(join(dir, "coco.db"));
const repositories = new SqliteRepositories({ database });
await repositories.init();
const manager = new Manager(repositories, async () => new Uint8Array(64).fill(7));
await manager.mint.addMint(mint.url, { trusted: true });
const service = (manager as unknown as { mintOperationService: AnyRecord })
.mintOperationService;
const booted: Booted = {
manager,
repositories,
mint,
spendable: async () => {
const balances = (await manager.wallet.balances.byMint()) as Record<
string,
{ spendable: number }
>;
return balances[mint.url]?.spendable ?? 0;
},
source: () => ({
ops: {
mint: {
listPending: () => manager.ops.mint.listPending(),
get: (id: string) => manager.ops.mint.get(id),
finalize: (id: string) => manager.ops.mint.finalize(id),
},
},
mintOperationService: service,
reopenFailedOperation: (id: string) =>
reopenFailedMintOperation(service as never, id),
}),
close: async () => {
await manager.dispose().catch(() => undefined);
database.close();
mint.stop();
rmSync(dir, { recursive: true, force: true });
},
};
return booted;
}
async function prepareQuote(booted: Booted, amount: number) {
const op = (await booted.manager.ops.mint.prepare({
mintUrl: booted.mint.url,
amount,
method: "bolt11",
} as never)) as unknown as AnyRecord;
return op;
}
function outputsOf(op: AnyRecord) {
// coco stores mint outputs as { keep, send }, like the on-disk output JSON.
const outputData = (op.outputData as AnyRecord).keep as Array<{
blindedMessage: { amount: unknown; id: string; B_: string };
}>;
return outputData.map((output) => ({
amount: Number(String(output.blindedMessage.amount)),
id: output.blindedMessage.id,
B_: output.blindedMessage.B_,
}));
}
let booted: Booted | undefined;
afterEach(async () => {
await booted?.close();
booted = undefined;
});
describe("PAID mint quote recovery with a real Manager and mint", () => {
it("issues an expired-but-PAID quote with the operation's own outputs, once", async () => {
// The motivating case: the invoice expired, but the mint says PAID and has
// issued nothing.
booted = await boot({ quoteExpiry: -60 });
const op = await prepareQuote(booted, 210_000);
const expectedOutputs = outputsOf(op).map((o) => o.B_);
booted.mint.markPaid(op.quoteId as string);
const result = (await runMintQuoteRecovery(
booted.source() as never,
)) as unknown as Record<string, number>;
expect(result).toMatchObject({ checked: 1, recovered: 1, terminal: 0 });
expect(await booted.spendable()).toBe(210_000);
// Issuance used exactly the blinded outputs stored on the operation.
expect(booted.mint.requests).toHaveLength(1);
expect(booted.mint.requests[0]?.outputs.map((o) => o.B_).sort()).toEqual(
[...expectedOutputs].sort(),
);
// A second run must not mint again or double-credit.
const again = (await runMintQuoteRecovery(
booted.source() as never,
)) as unknown as Record<string, number>;
expect(again).toMatchObject({ checked: 0, recovered: 0 });
expect(booted.mint.requests).toHaveLength(1);
expect(await booted.spendable()).toBe(210_000);
});
it("existing coco recovery already issues expired paid pending quotes", async () => {
booted = await boot({ quoteExpiry: -60 });
const op = await prepareQuote(booted, 100);
booted.mint.markPaid(op.quoteId as string);
await booted.manager.recoverPendingMintOperations();
expect(await booted.spendable()).toBe(100);
expect(booted.mint.getQuote(op.quoteId as string)?.state).toBe("ISSUED");
});
it("keeps rejected stored outputs and reports the actionable mint error", async () => {
booted = await boot({ quoteExpiry: -60 });
const op = await prepareQuote(booted, 100);
const outputs = outputsOf(op);
booted.mint.markPaid(op.quoteId as string);
// Model the mint refusing the stored outputs, not invoice expiry. This is
// not evidence that the production quotes used an inactive keyset.
booted.mint.mintError = { code: 12001, detail: "keyset id inactive." };
await booted.manager.recoverPendingMintOperations();
expect(await booted.spendable()).toBe(0);
const result = await runMintQuoteRecovery(booted.source() as never, {
operationIds: [op.id as string],
});
expect(result).toMatchObject({ recovered: 0, retryable: 1 });
expect(result.errors.some((entry) => entry.error.includes("keyset id inactive"))).toBe(true);
expect(await booted.spendable()).toBe(0);
expect(booted.mint.getQuote(op.quoteId as string)?.state).toBe("PAID");
expect(outputsOf(await booted.manager.ops.mint.get(op.id as string) as unknown as AnyRecord)).toEqual(outputs);
for (const request of booted.mint.requests) expect(request.outputs).toEqual(outputs);
});
it("restores proofs for a quote already issued at the mint", async () => {
booted = await boot({ quoteExpiry: null });
const op = await prepareQuote(booted, 210_000);
// Another wallet issued it: the signatures exist at the mint for the very
// outputs this operation stored.
booted.mint.signFor(op.quoteId as string, outputsOf(op));
const result = (await runMintQuoteRecovery(
booted.source() as never,
)) as unknown as Record<string, number>;
expect(result).toMatchObject({ recovered: 1, terminal: 0, retryable: 0 });
expect(await booted.spendable()).toBe(210_000);
});
it("reports terminal without credit when the mint refuses issuance", async () => {
booted = await boot({ quoteExpiry: -60 });
const op = await prepareQuote(booted, 21_000);
booted.mint.markPaid(op.quoteId as string);
booted.mint.mintError = { code: QUOTE_EXPIRED, detail: "quote expired" };
const result = (await runMintQuoteRecovery(
booted.source() as never,
)) as unknown as Record<string, unknown>;
expect(result).toMatchObject({ recovered: 0, terminal: 1 });
expect((result.errors as unknown[]).length).toBeGreaterThan(0);
expect(await booted.spendable()).toBe(0);
});
it("reports terminal without credit when an issued quote cannot be restored", async () => {
booted = await boot({ quoteExpiry: null });
const op = await prepareQuote(booted, 21_000);
// Issued at the mint, but the signatures for our outputs are gone.
booted.mint.markIssued(op.quoteId as string);
const result = (await runMintQuoteRecovery(
booted.source() as never,
)) as unknown as Record<string, number>;
expect(result).toMatchObject({ recovered: 0, terminal: 1 });
expect(await booted.spendable()).toBe(0);
});
it("does not credit a quote when the mint returns null NUT-09 signatures", async () => {
// KNOWN INTEROP GAP, not a supported path. NUT-09 permits `null` in the
// positional `signatures` array for outputs the mint never signed, but
// cashu-ts 3.7.1 - which coco depends on - dereferences every entry while
// normalising amounts, so the wallet throws instead of skipping the null.
// Recovery therefore surfaces an error and credits nothing, leaving the
// operation for a later run.
//
// This test pins the current behaviour so the gap cannot quietly disappear.
// Revisit (and change this expectation) once coco's cashu-ts parses
// positional nulls, and file/track it upstream in the meantime.
booted = await boot({ quoteExpiry: null });
const op = await prepareQuote(booted, 21_000);
// Issued at the mint with nothing signed for this operation's outputs.
booted.mint.markIssued(op.quoteId as string);
booted.mint.restoreIncludesNulls = true;
const result = (await runMintQuoteRecovery(
booted.source() as never,
)) as unknown as Record<string, number>;
expect(result).toMatchObject({ recovered: 0, terminal: 0, retryable: 1 });
expect(await booted.spendable()).toBe(0);
});
it("recovers a failed operation whose stale history says UNPAID", async () => {
// The composed path the feature exists for: a real prepared operation,
// failed locally, with a stale UNPAID observation from the old local-fail
// behaviour, recovered by explicit id and issued with its own outputs.
booted = await boot({ quoteExpiry: -60 });
const op = await prepareQuote(booted, 21_000);
const storedOutputs = outputsOf(op).map((o) => o.B_);
booted.mint.markPaid(op.quoteId as string);
const row = (await booted.repositories.mintOperationRepository.getById(
op.id as string,
)) as unknown as AnyRecord;
expect(row.outputData).toBeDefined();
await booted.repositories.mintOperationRepository.update({
...row,
state: "failed",
lastObservedRemoteState: "UNPAID",
error: "Expired unpaid mint quote cleaned up by routstrd",
terminalFailure: { reason: "expired", observedAt: Date.now() },
} as never);
// A purely local decision would skip this row; the mint has the last word.
const result = (await runMintQuoteRecovery(booted.source() as never, {
operationIds: [op.id as string],
includeFailed: true,
})) as unknown as Record<string, number>;
expect(result).toMatchObject({ reopened: 1, recovered: 1, terminal: 0 });
expect(await booted.spendable()).toBe(21_000);
expect(booted.mint.requests).toHaveLength(1);
expect(booted.mint.requests[0]?.outputs.map((o) => o.B_).sort()).toEqual(
[...storedOutputs].sort(),
);
});
it("counts a locked operation as busy instead of minting underneath it", async () => {
booted = await boot({ quoteExpiry: -60 });
const op = await prepareQuote(booted, 21_000);
booted.mint.markPaid(op.quoteId as string);
const service = (
booted.manager as unknown as {
mintOperationService: { acquireOperationLock(id: string): Promise<() => void> };
}
).mintOperationService;
// Hold the operation, as a processor or another recovery would.
const release = await service.acquireOperationLock(op.id as string);
const blocked = (await runMintQuoteRecovery(
booted.source() as never,
)) as unknown as Record<string, number>;
expect(blocked).toMatchObject({ recovered: 0 });
expect((blocked.busy ?? 0) + (blocked.retryable ?? 0)).toBeGreaterThan(0);
expect(booted.mint.requests).toHaveLength(0);
expect(await booted.spendable()).toBe(0);
release();
const after = (await runMintQuoteRecovery(
booted.source() as never,
)) as unknown as Record<string, number>;
expect(after).toMatchObject({ recovered: 1 });
expect(await booted.spendable()).toBe(21_000);
expect(booted.mint.requests).toHaveLength(1);
});
it("does not mint a second time while issuance is in flight", async () => {
booted = await boot({ quoteExpiry: -60 });
const op = await prepareQuote(booted, 21_000);
booted.mint.markPaid(op.quoteId as string);
// Hold the mint's response so the first recovery is visibly in flight.
let releaseGate!: () => void;
booted.mint.gate = new Promise<void>((resolve) => {
releaseGate = resolve;
});
const first = runMintQuoteRecovery(
booted.source() as never,
) as unknown as Promise<Record<string, number>>;
for (let i = 0; i < 400 && booted.mint.requests.length === 0; i++) {
await new Promise((resolve) => setTimeout(resolve, 5));
}
expect(booted.mint.requests).toHaveLength(1);
// A concurrent run must not issue again while that request is outstanding.
await runMintQuoteRecovery(booted.source() as never);
expect(booted.mint.requests).toHaveLength(1);
releaseGate();
expect(await first).toMatchObject({ recovered: 1 });
expect(await booted.spendable()).toBe(21_000);
expect(booted.mint.requests).toHaveLength(1);
});
it("coexists with coco's own mint operation watcher and processor", async () => {
booted = await boot({ quoteExpiry: -60 });
// The real background machinery coco uses to settle pending mint quotes.
await booted.manager.enableMintOperationWatcher();
await booted.manager.enableMintOperationProcessor();
const op = await prepareQuote(booted, 21_000);
booted.mint.markPaid(op.quoteId as string);
// Either our recovery or coco's processor may win; both are safe.
await runMintQuoteRecovery(booted.source() as never);
expect(await booted.spendable()).toBe(21_000);
expect(booted.mint.requests.length).toBeLessThanOrEqual(1);
// Give the processor time to act and confirm nothing is credited twice.
await new Promise((resolve) => setTimeout(resolve, 250));
expect(await booted.spendable()).toBe(21_000);
expect(booted.mint.requests.length).toBeLessThanOrEqual(1);
});
it("tracks a hung quote check so a later run waits instead of re-reading", async () => {
booted = await boot({ quoteExpiry: -60 });
const op = await prepareQuote(booted, 21_000);
booted.mint.markPaid(op.quoteId as string);
let releaseObserve!: () => void;
booted.mint.observeGate = new Promise<void>((resolve) => {
releaseObserve = resolve;
});
const outstanding = new Map<string, Promise<unknown>>();
const first = (await runMintQuoteRecovery(booted.source() as never, {
timeoutMs: 30,
outstanding,
})) as unknown as Record<string, number>;
expect(first).toMatchObject({ retryable: 1, recovered: 0 });
expect(outstanding.has(`mint:${op.id}`)).toBe(true);
const second = (await runMintQuoteRecovery(booted.source() as never, {
outstanding,
})) as unknown as Record<string, number>;
expect(second).toMatchObject({ checked: 0, busy: 1 });
// Once the held check settles the tracking drains, and recovery proceeds.
releaseObserve();
for (let i = 0; i < 400 && outstanding.size > 0; i++) {
await new Promise((resolve) => setTimeout(resolve, 5));
}
expect(outstanding.size).toBe(0);
const third = (await runMintQuoteRecovery(
booted.source() as never,
)) as unknown as Record<string, number>;
expect(third).toMatchObject({ recovered: 1 });
expect(await booted.spendable()).toBe(21_000);
});
});
@@ -0,0 +1,159 @@
/**
* Real-Manager integration for re-opening a failed mint operation.
*
* The unit tests exercise `reopenFailedMintOperation` against a hand-written
* service double, which cannot catch a change in coco's own
* `MintOperationService.transitionToPending` semantics or in its per-operation
* lock. This drives the production helper against an actual coco `Manager`
* backed by sqlite, so the private-service boundary that the helper depends on
* is exercised for real: full-row preservation, the shared operation lock, and
* the `mint-op:pending` event.
*
* No mint or network access is involved; nothing here enables the mint watcher
* or processor, which the fake-mint integration covers.
*/
import { afterEach, describe, expect, it } from "bun:test";
import { Database } from "bun:sqlite";
import { mkdtempSync, rmSync } from "node:fs";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { Manager } from "@cashu/coco-core";
import { SqliteRepositories } from "@cashu/coco-sqlite-bun";
import { reopenFailedMintOperation } from "./coco-client";
const OUTPUT_DATA = [
{
blindedMessage: { amount: "210000", id: "00deadbeef", B_: "02deadbeef" },
blindingFactor: "1234567890",
secret: "aabbccdd",
},
];
function failedRow() {
return {
id: "op-1",
mintUrl: "https://mint.invalid",
quoteId: "quote-1",
state: "failed",
createdAt: 1_000,
updatedAt: 2_000,
error: "expired",
method: "bolt11",
methodData: { method: "bolt11", data: {} },
amount: 210_000,
unit: "sat",
request: "lnbc1example",
expiry: 1_800_000_000,
pubkey: undefined,
lastObservedRemoteState: "PAID",
lastObservedRemoteStateAt: 3_000,
terminalFailure: { reason: "expired", observedAt: 3_000 },
outputData: OUTPUT_DATA,
};
}
describe("reopenFailedMintOperation with a real coco Manager", () => {
let dir: string | undefined;
let database: Database | undefined;
let manager: Manager | undefined;
let repositories: SqliteRepositories | undefined;
afterEach(async () => {
await manager?.dispose().catch(() => undefined);
manager = undefined;
database?.close();
database = undefined;
repositories = undefined;
if (dir) rmSync(dir, { recursive: true, force: true });
dir = undefined;
});
async function boot() {
dir = mkdtempSync(join(tmpdir(), "routstrd-manager-"));
database = new Database(join(dir, "coco.db"));
repositories = new SqliteRepositories({ database });
await repositories.init();
manager = new Manager(repositories, async () => new Uint8Array(64).fill(7));
await repositories.mintOperationRepository.create(failedRow() as never);
const service = (
manager as unknown as {
mintOperationService: {
acquireOperationLock(id: string): Promise<() => void>;
getOperation(id: string): Promise<Record<string, unknown> | null>;
transitionToPending(
op: Record<string, unknown>,
error?: string,
): Promise<unknown>;
};
}
).mintOperationService;
const eventBus = (
manager as unknown as {
eventBus: { on(event: string, handler: () => void): () => void };
}
).eventBus;
return { service, eventBus };
}
function readRow(id: string) {
return repositories!.mintOperationRepository.getById(id) as unknown as Promise<
Record<string, unknown> | null
>;
}
it("re-opens through the real service and emits mint-op:pending", async () => {
const { service, eventBus } = await boot();
const events: string[] = [];
const off = eventBus.on("mint-op:pending", () => events.push("pending"));
const reopened = await reopenFailedMintOperation(service, "op-1");
off();
expect(reopened).toBe(true);
const row = await readRow("op-1");
expect(row?.state).toBe("pending");
expect(row?.quoteId).toBe("quote-1");
expect(row?.amount).toBe(210_000);
expect(row?.lastObservedRemoteState).toBe("PAID");
expect(row?.outputData).toEqual(OUTPUT_DATA);
expect(row?.terminalFailure ?? undefined).toBeUndefined();
expect(events).toEqual(["pending"]);
});
it("refuses to re-open while coco's operation lock is held", async () => {
const { service } = await boot();
// coco's OperationIdLock is fail-fast: a holder blocks the re-open by
// making it throw, and nothing is written.
const release = await service.acquireOperationLock("op-1");
await expect(reopenFailedMintOperation(service, "op-1")).rejects.toThrow(
/already in progress/,
);
const blocked = await readRow("op-1");
expect(blocked?.state).toBe("failed");
expect(blocked?.outputData).toEqual(OUTPUT_DATA);
release();
expect(await reopenFailedMintOperation(service, "op-1")).toBe(true);
const reopened = await readRow("op-1");
expect(reopened?.state).toBe("pending");
expect(reopened?.outputData).toEqual(OUTPUT_DATA);
});
it("leaves an operation a processor finalized alone", async () => {
const { service } = await boot();
// Emulate a processor winning the race while holding the same lock.
const release = await service.acquireOperationLock("op-1");
const row = await readRow("op-1");
await repositories!.mintOperationRepository.update({
...row,
state: "finalized",
} as never);
release();
expect(await reopenFailedMintOperation(service, "op-1")).toBe(false);
const after = await readRow("op-1");
expect(after?.state).toBe("finalized");
expect(after?.outputData).toEqual(OUTPUT_DATA);
});
});
@@ -0,0 +1,97 @@
import { describe, expect, it } from "bun:test";
import {
classifyMintQuoteObservation,
selectMintQuotesForRecovery,
type MintQuoteRecoveryCandidate,
} from "./mint-quote-recovery";
function mint(
overrides: Partial<MintQuoteRecoveryCandidate> = {},
): MintQuoteRecoveryCandidate {
return {
id: "mint-1",
mintUrl: "https://mint.example",
quoteId: "quote-1",
state: "pending",
amount: 1000,
expiry: 0,
...overrides,
};
}
describe("classifyMintQuoteObservation", () => {
it("finalizes a paid-but-unissued quote by minting its stored outputs", () => {
expect(classifyMintQuoteObservation("ready")).toEqual({
action: "finalize",
observedRemoteState: "PAID",
});
});
it("finalizes an already-issued quote by restoring its proofs", () => {
expect(classifyMintQuoteObservation("completed")).toEqual({
action: "finalize",
observedRemoteState: "ISSUED",
});
});
it("leaves an unpaid quote alone", () => {
expect(classifyMintQuoteObservation("waiting")).toEqual({
action: "waiting",
});
});
it("reports a quote the mint can no longer issue", () => {
expect(classifyMintQuoteObservation("terminal")).toEqual({
action: "terminal",
});
});
});
describe("selectMintQuotesForRecovery", () => {
it("selects every pending quote because only the mint knows if it was paid", () => {
const result = selectMintQuotesForRecovery({
mints: [
mint({ id: "expired", expiry: 1 }),
mint({ id: "fresh" }),
mint({ id: "observed-unpaid", lastObservedRemoteState: "UNPAID" }),
mint({ id: "observed-paid", lastObservedRemoteState: "PAID" }),
],
});
expect(result.pending.map((op) => op.id)).toEqual([
"expired",
"fresh",
"observed-unpaid",
"observed-paid",
]);
});
it("recovers executing operations left behind by a crash mid-mint", () => {
const result = selectMintQuotesForRecovery({
mints: [mint({ id: "executing", state: "executing" })],
});
expect(result.pending.map((op) => op.id)).toEqual(["executing"]);
});
it("ignores finalized operations", () => {
const result = selectMintQuotesForRecovery({
mints: [mint({ id: "done", state: "finalized" })],
});
expect(result.pending).toEqual([]);
expect(result.failed).toEqual([]);
});
it("does not re-open failed operations by default", () => {
const result = selectMintQuotesForRecovery({
mints: [mint({ id: "given-up", state: "failed" })],
});
expect(result.failed).toEqual([]);
});
it("returns failed operations when the caller opts in", () => {
const result = selectMintQuotesForRecovery({
mints: [mint({ id: "given-up", state: "failed" })],
includeFailed: true,
});
expect(result.failed.map((op) => op.id)).toEqual(["given-up"]);
});
});
+119
View File
@@ -0,0 +1,119 @@
/**
* Pure helpers for PAID mint-quote recovery.
*
* A mint quote can be PAID at the mint while its local operation is still
* `pending` (the Lightning payment landed before expiry while the daemon was
* down, so no local observation was ever recorded) or even terminally
* `failed` (coco gives up when the mint refuses to sign, for example after the
* invoice expiry). Claimability still depends on the mint accepting issuance.
* This feature retries the stored outputs or restores their signatures; it
* does not regenerate outputs rejected by the mint (for example an inactive
* keyset). coco already reconciles pending paid quotes at startup and in the
* periodic sweep. The new capability is operator-targeted recovery, including
* explicitly reopening failed operations, alongside safer cleanup.
*
* Recovery asks the mint what it thinks, then retries issuance or restore. These helpers decide *what* to do from a remote observation; the
* actual state transitions are applied by the in-process coco wallet client
* so coco-core's operation services emit their normal events and release
* proof reservations. Keeping the decisions here makes them unit testable
* without a wallet database or network access.
*/
/** Subset of coco's mint operation rows that recovery needs. */
export interface MintQuoteRecoveryCandidate {
id: string;
mintUrl: string;
quoteId?: string;
state: string;
/** Quote amount in sats. */
amount: number;
/** Quote expiry in epoch seconds. `0` means unknown/not applicable. */
expiry: number;
/** Last quote state observed from the mint (UNPAID, PAID, ISSUED). */
lastObservedRemoteState?: string;
error?: string;
}
/** Coco's classification of a fresh remote quote check. */
export type PendingMintCheckCategory =
| "waiting"
| "ready"
| "completed"
| "terminal";
/** What recovery should do with a quote after checking it with the mint. */
export type MintQuoteRecoveryDecision =
| { action: "finalize"; observedRemoteState: "PAID" | "ISSUED" }
| { action: "waiting" }
| { action: "terminal" };
/**
* Map a remote quote check onto a recovery action.
*
* - `ready` means the mint reports the quote PAID but never issued: submit the
* operation's stored outputs to claim the sats.
* - `completed` means the mint already issued it: recover the signatures
* (NUT-09) instead of minting again.
* - `waiting` means the mint still reports the quote UNPAID: nothing is
* claimable, so leave the operation alone.
* - `terminal` means the quote can no longer be issued (for example the mint
* refused an expired quote). coco persists that verdict as a failed
* operation, so recovery must report it rather than treat it as progress.
*/
export function classifyMintQuoteObservation(
category: PendingMintCheckCategory,
): MintQuoteRecoveryDecision {
switch (category) {
case "ready":
return { action: "finalize", observedRemoteState: "PAID" };
case "completed":
return { action: "finalize", observedRemoteState: "ISSUED" };
case "waiting":
return { action: "waiting" };
case "terminal":
return { action: "terminal" };
}
}
export interface MintQuoteRecoverySelectionOptions<
T extends MintQuoteRecoveryCandidate,
> {
mints: T[];
/**
* Also consider terminally failed operations. Off by default: re-opening a
* failed operation is a mutation, so only an explicit user-invoked recovery
* may do it. Startup recovery must never resurrect quotes on its own.
*/
includeFailed?: boolean;
}
export interface MintQuoteRecoverySelection<
T extends MintQuoteRecoveryCandidate,
> {
/** Pending (or executing) operations that need a fresh mint observation. */
pending: T[];
/** Failed operations the caller may re-open and retry. */
failed: T[];
}
/**
* Split operations into those recovery should check and those that were
* already given up on.
*
* Every `pending` operation is selected: only the mint knows whether an
* expired quote was paid before the local invoice ran out. `executing`
* operations are recovered too, since a crash mid-mint leaves outputs that
* may already be signed.
*/
export function selectMintQuotesForRecovery<
T extends MintQuoteRecoveryCandidate,
>(options: MintQuoteRecoverySelectionOptions<T>): MintQuoteRecoverySelection<T> {
const { mints, includeFailed = false } = options;
const pending = mints.filter(
(op) => op.state === "pending" || op.state === "executing",
);
const failed = includeFailed
? mints.filter((op) => op.state === "failed")
: [];
return { pending, failed };
}
@@ -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();
}
});
+448
View File
@@ -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"]);
});
+312
View File
@@ -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;
}
+41
View File
@@ -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"]);
});
+45
View File
@@ -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);
};
}
+347
View File
@@ -0,0 +1,347 @@
/**
* In-process Cashu mint for integration tests.
*
* Implements just enough of NUT-01/02/04/06/07/09 to drive a real coco
* `Manager` against real HTTP: keyset publication, bolt11 mint quotes, minting
* blinded outputs with a real secp256k1 blind signature, NUT-09 restore and
* NUT-07 proof states. No Lightning, no network, no NPC.
*
* The signing keys and signatures are genuine (`@cashu/cashu-ts` mint-side
* helpers), so a wallet that receives these signatures can unblind and verify
* them exactly as with a production mint.
*/
import {
createBlindSignature,
createNewMintKeys,
pointFromHex,
} from "@cashu/cashu-ts";
/** NUT error codes used by the scenarios. */
export const QUOTE_EXPIRED = 20007;
export const ALREADY_ISSUED = 20002;
export type FakeQuoteState = "UNPAID" | "PAID" | "ISSUED";
export interface FakeMintQuote {
quote: string;
request: string;
amount: number;
unit: string;
state: FakeQuoteState;
/** Epoch seconds, or null for a quote that never expires. */
expiry: number | null;
amountPaid: number;
amountIssued: number;
pubkey: null;
}
export interface FakeMintRequest {
quote: string;
outputs: Array<{ amount: number; id: string; B_: string }>;
}
interface StoredSignature {
amount: number;
id: string;
C_: string;
}
const toHex = (bytes: Uint8Array) => Buffer.from(bytes).toString("hex");
export class FakeMint {
readonly keysetId: string;
readonly keysByAmount: Record<string, string>;
readonly requests: FakeMintRequest[] = [];
/** Every output the mint has ever signed, keyed by B_. */
readonly signed = new Map<string, StoredSignature>();
/** When set, POST /v1/mint/bolt11 fails with this NUT error. */
mintError: { code: number; detail: string } | null = null;
/** When set, quote creation returns this expiry (epoch seconds). */
quoteExpiry: number | null = 3_600;
/**
* Awaited before responding to a mint request, so a test can hold minting
* open and interleave another recovery attempt.
*/
gate: Promise<void> | null = null;
/** Awaited before answering a quote-state check, to hold observe open. */
observeGate: Promise<void> | null = null;
/**
* When true, POST /v1/restore answers with NUT-09's spec-legal positional
* arrays, including `null` for outputs the mint never signed.
*
* This is a known interop gap, not a supported path: cashu-ts 3.7.1 (which
* coco depends on) dereferences every entry of `signatures` while normalising
* amounts, so a `null` makes the wallet throw instead of skipping it. The
* switch exists so tests keep that behaviour visible; see
* mint-quote-recovery.fake-mint.test.ts.
*/
restoreIncludesNulls = false;
private readonly quotes = new Map<string, FakeMintQuote>();
private counter = 0;
private server?: ReturnType<typeof Bun.serve>;
constructor() {
const pair = createNewMintKeys(20, new Uint8Array(32).fill(9));
this.keysetId = pair.keysetId;
this.keysByAmount = Object.fromEntries(
Object.entries(pair.pubKeys).map(([amount, key]) => [
amount,
typeof key === "string" ? key : toHex(key),
]),
);
this.privKeys = Object.fromEntries(
Object.entries(pair.privKeys) as Array<[string, Uint8Array]>,
);
}
private readonly privKeys: Record<string, Uint8Array>;
get url(): string {
if (!this.server) throw new Error("fake mint not started");
return `http://127.0.0.1:${this.server.port}`;
}
start(): void {
this.server = Bun.serve({
hostname: "127.0.0.1",
port: 0,
fetch: (request) => this.handle(request),
});
}
stop(): void {
this.server?.stop(true);
this.server = undefined;
}
/** Test control: the mint sees the invoice as paid but has issued nothing. */
markPaid(quoteId: string): void {
const quote = this.quotes.get(quoteId);
if (!quote) throw new Error(`unknown fake quote ${quoteId}`);
quote.state = "PAID";
quote.amountPaid = quote.amount;
}
/**
* Test control: mark the quote issued without going through the mint
* endpoint, simulating a wallet that lost the signatures.
*/
markIssued(quoteId: string): void {
const quote = this.quotes.get(quoteId);
if (!quote) throw new Error(`unknown fake quote ${quoteId}`);
quote.state = "ISSUED";
quote.amountPaid = quote.amount;
quote.amountIssued = quote.amount;
}
/** Test control: sign outputs directly, as if another wallet had issued them. */
signFor(quoteId: string, outputs: FakeMintRequest["outputs"]): void {
for (const output of outputs) {
this.signOutput(output);
}
this.markIssued(quoteId);
}
getQuote(quoteId: string): FakeMintQuote | undefined {
return this.quotes.get(quoteId);
}
private quoteBody(quote: FakeMintQuote) {
return {
quote: quote.quote,
request: quote.request,
amount: quote.amount,
unit: quote.unit,
state: quote.state,
expiry: quote.expiry,
amount_paid: quote.amountPaid,
amount_issued: quote.amountIssued,
pubkey: quote.pubkey,
};
}
private signOutput(output: {
amount: number;
id: string;
B_: string;
}): StoredSignature {
const privKey = this.privKeys[String(output.amount)];
if (!privKey) {
throw new Error(`fake mint has no key for amount ${output.amount}`);
}
const signature = createBlindSignature(
pointFromHex(output.B_),
privKey,
this.keysetId,
);
const stored: StoredSignature = {
amount: output.amount,
id: this.keysetId,
C_: toHex(signature.C_.toBytes(false)),
};
this.signed.set(output.B_, stored);
return stored;
}
private json(body: unknown, status = 200): Response {
return new Response(JSON.stringify(body), {
status,
headers: { "content-type": "application/json" },
});
}
private error(code: number, detail: string, status = 400): Response {
return this.json({ code, detail }, status);
}
private async handle(request: Request): Promise<Response> {
const { pathname } = new URL(request.url);
const body = async () => {
try {
return (await request.json()) as Record<string, unknown>;
} catch {
return {};
}
};
if (pathname === "/v1/info") {
return this.json({
name: "fake-mint",
version: "0.0.1",
nuts: {
4: {
methods: [
{
method: "bolt11",
unit: "sat",
min_amount: 1,
max_amount: 1_000_000,
},
],
},
5: {
methods: [
{
method: "bolt11",
unit: "sat",
min_amount: 1,
max_amount: 1_000_000,
},
],
},
7: { supported: true },
9: { supported: true },
},
});
}
if (pathname === "/v1/keys" || pathname.startsWith("/v1/keys/")) {
return this.json({
keysets: [
{ id: this.keysetId, unit: "sat", keys: this.keysByAmount },
],
});
}
if (pathname === "/v1/keysets") {
return this.json({
keysets: [{ id: this.keysetId, unit: "sat", active: true }],
});
}
if (pathname === "/v1/mint/quote/bolt11" && request.method === "POST") {
const input = await body();
const amount = Number(input.amount);
const unit = typeof input.unit === "string" ? input.unit : "sat";
const quote: FakeMintQuote = {
quote: `fake-quote-${++this.counter}`,
request: `lnbcfake${this.counter}`,
amount,
unit,
state: "UNPAID",
expiry:
this.quoteExpiry === null
? null
: Math.floor(Date.now() / 1000) + this.quoteExpiry,
amountPaid: 0,
amountIssued: 0,
pubkey: null,
};
this.quotes.set(quote.quote, quote);
return this.json(this.quoteBody(quote));
}
const quoteMatch = pathname.match(/^\/v1\/mint\/quote\/bolt11\/(.+)$/);
if (quoteMatch?.[1] && request.method === "GET") {
if (this.observeGate) await this.observeGate;
const quote = this.quotes.get(decodeURIComponent(quoteMatch[1]));
if (!quote) return this.error(50000, "Unknown quote");
return this.json(this.quoteBody(quote));
}
if (pathname === "/v1/mint/bolt11" && request.method === "POST") {
const input = await body();
const quoteId = String(input.quote);
const outputs = (input.outputs ?? []) as FakeMintRequest["outputs"];
this.requests.push({ quote: quoteId, outputs });
if (this.gate) await this.gate;
if (this.mintError) {
return this.error(this.mintError.code, this.mintError.detail);
}
const quote = this.quotes.get(quoteId);
if (!quote) return this.error(50000, "Unknown quote");
if (quote.state === "ISSUED" || quote.amountIssued > 0) {
return this.error(ALREADY_ISSUED, "Quote already issued");
}
if (quote.state !== "PAID") {
return this.error(20001, "Quote is not paid");
}
const total = outputs.reduce((sum, o) => sum + Number(o.amount), 0);
if (total > quote.amountPaid - quote.amountIssued) {
return this.error(10002, "Outputs exceed the paid amount");
}
const signatures = outputs.map((output) => ({
amount: output.amount,
id: this.keysetId,
C_: this.signOutput(output).C_,
}));
quote.state = "ISSUED";
quote.amountIssued += total;
return this.json({ signatures });
}
if (pathname === "/v1/restore" && request.method === "POST") {
const input = await body();
const outputs = (input.outputs ?? []) as FakeMintRequest["outputs"];
if (this.restoreIncludesNulls) {
// Spec-legal NUT-09 shape, including nulls. Kept behind a switch
// because it is currently unusable with coco's cashu-ts version.
return this.json({
outputs,
signatures: outputs.map(
(output) => this.signed.get(output.B_) ?? null,
),
});
}
// Return only the signed outputs; coco matches them by B_ and treats the
// rest as "nothing to restore".
const signed = outputs.filter((output) => this.signed.has(output.B_));
return this.json({
outputs: signed,
signatures: signed.map((output) => this.signed.get(output.B_)),
});
}
if (pathname === "/v1/checkstate" && request.method === "POST") {
const input = await body();
const ys = (input.Ys ?? []) as string[];
return this.json({
states: ys.map((Y) => ({ Y, state: "UNSPENT", witness: null })),
});
}
return this.error(404, `fake mint has no route for ${pathname}`, 404);
}
}
+20
View File
@@ -0,0 +1,20 @@
import { describe, expect, it } from "bun:test";
import { HISTORY_ENTRY_TYPES, isHistoryEntryType } from "./history";
describe("history entry types", () => {
it("lists the four supported transaction types", () => {
expect([...HISTORY_ENTRY_TYPES]).toEqual([
"mint",
"melt",
"send",
"receive",
]);
});
it("recognizes known types and rejects unknown ones", () => {
expect(isHistoryEntryType("send")).toBe(true);
expect(isHistoryEntryType("melt")).toBe(true);
expect(isHistoryEntryType("SEND")).toBe(false);
expect(isHistoryEntryType("refund")).toBe(false);
});
});
+9
View File
@@ -0,0 +1,9 @@
/** Transaction types understood by the wallet history type filter. */
export const HISTORY_ENTRY_TYPES = ["mint", "melt", "send", "receive"] as const;
export type HistoryEntryType = (typeof HISTORY_ENTRY_TYPES)[number];
/** True when `value` is a recognized history transaction type. */
export function isHistoryEntryType(value: string): value is HistoryEntryType {
return (HISTORY_ENTRY_TYPES as readonly string[]).includes(value);
}