mirror of
https://github.com/Routstr/routstrd.git
synced 2026-10-05 12:28:23 +00:00
fix(wallet): deduplicate Cashu receives before recovery
Prevent repeated submissions of the same encoded Cashu token from creating an unbounded number of Coco receive operations. Routstrd now hashes each exact token with SHA-256 and atomically reserves that hash in persistent wallet metadata before preparing a Coco receive operation. Persist the token hash, Coco operation ID, processing state, success state, and failure details so duplicate protection survives daemon restarts. Reconcile interrupted reservations against Coco operations at startup, cancel prepared operations that never reached the mint, preserve unresolved executing operations, and convert finalized operations into idempotent successful responses. When a retry references an unresolved operation, refresh that existing operation instead of creating another one. If Coco finalized an operation while returning an earlier request error, report success. If a rolled-back operation has a finalized sibling with the same proofs, report the token as already received instead of recording a false failure. Add a pre-recovery reconciliation pass for legacy executing receive operations. Group rows by normalized mint, unit, and complete sorted proof material so thousands of duplicate rows are handled as a small number of unique input sets. For each unique group, check proof state only once. If inputs are unspent, retain one valid canonical operation and retire redundant copies. If inputs are spent, probe all stored deterministic outputs in bounded restore batches, retain operations whose outputs are recoverable, and roll back only non-owning duplicates. If restore conclusively returns no matching outputs, mark the group rolled back using Coco 2's terminal recovery reason. Leave groups untouched when proof-state responses are incomplete or mixed, output data is malformed, the mint is unavailable, or any persistence step fails. Add per-request timeouts, a global 45-second startup budget, and per-mint fail-fast behavior so an offline mint cannot recreate the multi-hour startup incident. Create and record a standalone SQLite backup before the first receive cleanup, before Coco watchers and processors are enabled. Continue normal send, melt, and mint recovery while limiting receive recovery to retained operations, avoiding Coco 1.0.1's expensive per-row sweep over unresolved duplicates. Add focused tests for exact-token reservation, persistent state, interrupted reservation cleanup, grouping, unspent canonical selection, spent empty restore, restored output ownership, incomplete and mixed proof states, and a production-shaped case containing 1,924 operations collapsed into 47 mint checks. Also add the detailed Coco 2.0.0 migration plan covering staged database migration, receive preflight, rollback, compatibility, and eventual upgrade work. Validation: - bun run lint - bun run build - 40 focused wallet tests pass - full suite reaches 96 passing tests; one pre-existing Hermes config test remains unrelated and failing
This commit is contained in:
File diff suppressed because it is too large
Load Diff
@@ -5,6 +5,7 @@ import {
|
||||
} from "@cashu/coco-core";
|
||||
import type {
|
||||
HistoryEntry,
|
||||
ReceiveOperation,
|
||||
Logger as CocoLogger,
|
||||
Plugin as CocoPlugin,
|
||||
} from "@cashu/coco-core";
|
||||
@@ -35,6 +36,20 @@ import type {
|
||||
WalletRecoveryProgress,
|
||||
} from "./cocod-client";
|
||||
import { selectCleanupOperations } from "./cleanup";
|
||||
import {
|
||||
clearInterruptedReceiveReservations,
|
||||
deleteReceiveTokenReservation,
|
||||
getReceiveReconcileBackup,
|
||||
initReceiveDedupSchema,
|
||||
listProcessingReceiveTokens,
|
||||
receiveInputFingerprint,
|
||||
reconcileExecutingReceives,
|
||||
releaseReceiveToken,
|
||||
reserveReceiveToken,
|
||||
setReceiveReconcileBackup,
|
||||
updateReceiveToken,
|
||||
type ReceiveReconcileSource,
|
||||
} from "./receive-dedup";
|
||||
import { cocoLogger, logger } from "../../utils/logger";
|
||||
import {
|
||||
legacyCocodPidPath,
|
||||
@@ -582,17 +597,19 @@ export interface CreateCocoClientOptions {
|
||||
* receive/mint recovery passes, so the daemon can serve wallet reads while
|
||||
* recovery proceeds in the background.
|
||||
*/
|
||||
async function buildCocoManager(
|
||||
function constructCocoManager(
|
||||
repo: SqliteRepositories,
|
||||
seed: Uint8Array,
|
||||
): Promise<Manager> {
|
||||
const coco = new Manager(repo, async () => seed, createCocoLogger());
|
||||
): Manager {
|
||||
return new Manager(repo, async () => seed, createCocoLogger());
|
||||
}
|
||||
|
||||
async function enableCocoManager(coco: Manager): Promise<void> {
|
||||
await coco.initPlugins();
|
||||
await coco.reconcileLegacyMintQuotes();
|
||||
await coco.enableMintOperationWatcher();
|
||||
await coco.enableProofStateWatcher();
|
||||
await coco.enableMintOperationProcessor();
|
||||
return coco;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -748,6 +765,85 @@ export async function settleExpiredMintQuotes(
|
||||
return settlement;
|
||||
}
|
||||
|
||||
interface ReceiveRecoveryInternals {
|
||||
receiveOperationService: {
|
||||
checkProofStatesWithMint(
|
||||
mintUrl: string,
|
||||
proofs: Array<{ secret: string }>,
|
||||
): Promise<Array<{ state: string }>>;
|
||||
hasSavedOutputs(operation: unknown): Promise<boolean>;
|
||||
markAsRolledBack(operation: unknown, error: string): Promise<unknown>;
|
||||
};
|
||||
mintAdapter: {
|
||||
getCashuMint(mintUrl: string): {
|
||||
restore(input: { outputs: Array<{ amount: number; id: string; B_: string }> }): Promise<{
|
||||
outputs: Array<{ B_: string }>;
|
||||
}>;
|
||||
};
|
||||
};
|
||||
}
|
||||
|
||||
async function reconcileDuplicateReceiveOperations(
|
||||
coco: Manager,
|
||||
repo: SqliteRepositories,
|
||||
): Promise<Awaited<ReturnType<typeof reconcileExecutingReceives>>> {
|
||||
const internals = coco as unknown as ReceiveRecoveryInternals;
|
||||
const unavailableMints = new Set<string>();
|
||||
const deadline = Date.now() + 45_000;
|
||||
const ensureMintBudget = (mintUrl: string): void => {
|
||||
if (Date.now() >= deadline) throw new Error("Receive cleanup time budget exhausted");
|
||||
if (unavailableMints.has(mintUrl)) throw new Error("Mint already failed receive cleanup");
|
||||
};
|
||||
const markMintFailure = (mintUrl: string, error: unknown): never => {
|
||||
unavailableMints.add(mintUrl);
|
||||
throw error;
|
||||
};
|
||||
const source: ReceiveReconcileSource = {
|
||||
listExecuting: () => repo.receiveOperationRepository.getByState("executing"),
|
||||
checkProofStates: async (operation) => {
|
||||
ensureMintBudget(operation.mintUrl);
|
||||
try {
|
||||
return await withTimeout(
|
||||
internals.receiveOperationService.checkProofStatesWithMint(
|
||||
operation.mintUrl,
|
||||
operation.inputProofs,
|
||||
),
|
||||
Math.min(15_000, Math.max(1, deadline - Date.now())),
|
||||
);
|
||||
} catch (error) {
|
||||
return markMintFailure(operation.mintUrl, error);
|
||||
}
|
||||
},
|
||||
restoreOutputs: async (mintUrl, outputs) => {
|
||||
ensureMintBudget(mintUrl);
|
||||
// Probe deterministic outputs in bounded batches. The Cashu restore
|
||||
// response echoes only blinded messages with stored signatures, which
|
||||
// identifies the owning operation. Coco later performs the real proof
|
||||
// recovery for the retained operation.
|
||||
const restored: Array<{ B_: string }> = [];
|
||||
const mint = internals.mintAdapter.getCashuMint(mintUrl);
|
||||
for (let index = 0; index < outputs.length; index += 300) {
|
||||
try {
|
||||
const response = await withTimeout(
|
||||
mint.restore({ outputs: outputs.slice(index, index + 300) }),
|
||||
Math.min(15_000, Math.max(1, deadline - Date.now())),
|
||||
);
|
||||
restored.push(...response.outputs);
|
||||
} catch (error) {
|
||||
return markMintFailure(mintUrl, error);
|
||||
}
|
||||
}
|
||||
return restored;
|
||||
},
|
||||
hasSavedOutputs: (operation) =>
|
||||
internals.receiveOperationService.hasSavedOutputs(operation),
|
||||
rollBack: async (operation, reason) => {
|
||||
await internals.receiveOperationService.markAsRolledBack(operation, reason);
|
||||
},
|
||||
};
|
||||
return reconcileExecutingReceives(source);
|
||||
}
|
||||
|
||||
interface RecoveryPhaseProgress {
|
||||
phase: string;
|
||||
failedMintQuotes: number;
|
||||
@@ -763,6 +859,7 @@ interface RecoveryPhaseProgress {
|
||||
async function runWalletRecovery(
|
||||
coco: Manager,
|
||||
onProgress: (progress: RecoveryPhaseProgress) => void,
|
||||
receiveOperationIds?: string[],
|
||||
): Promise<void> {
|
||||
surfacingRecoveryProgress = true;
|
||||
let failedMintQuotes = 0;
|
||||
@@ -796,7 +893,24 @@ async function runWalletRecovery(
|
||||
await coco.ops.melt.recovery.run();
|
||||
|
||||
onProgress({ phase: "Receive recovery", failedMintQuotes });
|
||||
if (receiveOperationIds) {
|
||||
// The pre-check already classified every executing receive by unique
|
||||
// input set. Recover only the conclusive retained operations; unresolved
|
||||
// groups stay untouched instead of falling back to Coco 1's expensive
|
||||
// per-row sweep on this startup.
|
||||
for (const operationId of receiveOperationIds) {
|
||||
try {
|
||||
await withTimeout(coco.ops.receive.refresh(operationId), 15_000);
|
||||
} catch (error) {
|
||||
logger.warn("Targeted receive recovery did not complete", {
|
||||
operationId,
|
||||
error: error instanceof Error ? error.message : String(error),
|
||||
});
|
||||
}
|
||||
}
|
||||
} else {
|
||||
await coco.ops.receive.recovery.run();
|
||||
}
|
||||
|
||||
onProgress({ phase: "Mint recovery", failedMintQuotes });
|
||||
await coco.recoverPendingMintOperations();
|
||||
@@ -843,6 +957,9 @@ export async function createCocoClient(
|
||||
|
||||
let database: Database | undefined;
|
||||
let coco: Manager | undefined;
|
||||
let findFinalizedReceiveSibling: (
|
||||
operation: ReceiveOperation | null,
|
||||
) => Promise<string | null> = async () => null;
|
||||
let walletConfig = loadConfig(configFile);
|
||||
|
||||
let recoveryPhase = "queued";
|
||||
@@ -869,6 +986,13 @@ export async function createCocoClient(
|
||||
database = new Database(dbPath);
|
||||
const repo = new SqliteRepositories({ database });
|
||||
await repo.init();
|
||||
initReceiveDedupSchema(database);
|
||||
const interruptedReservations = clearInterruptedReceiveReservations(database);
|
||||
if (interruptedReservations > 0) {
|
||||
logger.warn("Cleared interrupted receive reservations with no Coco operation", {
|
||||
count: interruptedReservations,
|
||||
});
|
||||
}
|
||||
|
||||
const [pendingSends, inflightProofs, pendingMints] = await Promise.all([
|
||||
repo.sendOperationRepository.getPending(),
|
||||
@@ -892,7 +1016,95 @@ export async function createCocoClient(
|
||||
startupProgress("Initializing Cashu wallet...");
|
||||
}
|
||||
|
||||
coco = await buildCocoManager(repo, seed);
|
||||
// Construct Coco so the pre-recovery checker can reuse its mint adapter,
|
||||
// but do not enable watchers/processors until the backup and cleanup finish.
|
||||
coco = constructCocoManager(repo, seed);
|
||||
|
||||
const executingReceives = await repo.receiveOperationRepository.getByState("executing");
|
||||
let receiveRecoveryOperationIds: string[] | undefined;
|
||||
if (executingReceives.length > 0) {
|
||||
startupProgress(
|
||||
`Checking ${executingReceives.length} unfinished Cashu receive operation(s) for duplicates...`,
|
||||
);
|
||||
const recordedBackup = getReceiveReconcileBackup(database);
|
||||
if (!recordedBackup || !existsSync(recordedBackup)) {
|
||||
const backupPath = `${dbPath}.pre-receive-reconcile-${Date.now()}`;
|
||||
database.exec(`VACUUM INTO '${backupPath.replaceAll("'", "''")}'`);
|
||||
setReceiveReconcileBackup(database, backupPath);
|
||||
startupProgress(`Created wallet backup before receive cleanup: ${backupPath}`);
|
||||
}
|
||||
const receiveReconcile = await reconcileDuplicateReceiveOperations(coco, repo);
|
||||
receiveRecoveryOperationIds = receiveReconcile.recoveryOperationIds;
|
||||
startupProgress(
|
||||
`Receive cleanup: ${receiveReconcile.executing} operation(s), ` +
|
||||
`${receiveReconcile.uniqueGroups} unique input set(s), ` +
|
||||
`${receiveReconcile.rolledBack} stale duplicate(s) retired, ` +
|
||||
`${receiveReconcile.unresolved} unresolved.`,
|
||||
);
|
||||
}
|
||||
|
||||
await enableCocoManager(coco);
|
||||
const openDatabase = database;
|
||||
|
||||
findFinalizedReceiveSibling = async (
|
||||
operation: ReceiveOperation | null,
|
||||
): Promise<string | null> => {
|
||||
if (!operation) return null;
|
||||
const fingerprint = receiveInputFingerprint(operation);
|
||||
const siblings = await repo.receiveOperationRepository.getByMintUrl(operation.mintUrl);
|
||||
return (
|
||||
siblings.find(
|
||||
(candidate) =>
|
||||
candidate.state === "finalized" &&
|
||||
receiveInputFingerprint(candidate) === fingerprint,
|
||||
)?.id ?? null
|
||||
);
|
||||
};
|
||||
|
||||
const syncReceiveReservations = async (): Promise<void> => {
|
||||
for (const reservation of listProcessingReceiveTokens(openDatabase)) {
|
||||
if (!reservation.operationId) continue;
|
||||
try {
|
||||
const operation = await coco!.ops.receive.get(reservation.operationId);
|
||||
if (operation?.state === "finalized") {
|
||||
updateReceiveToken(openDatabase, reservation.tokenHash, {
|
||||
state: "succeeded",
|
||||
operationId: operation.id,
|
||||
});
|
||||
} else if (operation?.state === "rolled_back") {
|
||||
const finalizedSibling = await findFinalizedReceiveSibling(operation);
|
||||
updateReceiveToken(openDatabase, reservation.tokenHash, finalizedSibling
|
||||
? {
|
||||
state: "succeeded",
|
||||
operationId: finalizedSibling,
|
||||
}
|
||||
: {
|
||||
state: "failed",
|
||||
operationId: operation.id,
|
||||
error: operation.error || "Token receive was rolled back",
|
||||
});
|
||||
} else if (operation?.state === "prepared") {
|
||||
// A crash before execute had no mint side effect. Cancel the stale
|
||||
// prepared operation and permit a fresh exact-token attempt.
|
||||
await coco!.ops.receive.cancel(
|
||||
operation.id,
|
||||
"Cancelled interrupted receive before execution",
|
||||
);
|
||||
deleteReceiveTokenReservation(openDatabase, reservation.tokenHash);
|
||||
} else if (!operation) {
|
||||
deleteReceiveTokenReservation(openDatabase, reservation.tokenHash);
|
||||
}
|
||||
} catch (error) {
|
||||
// One damaged/stale reservation must never block daemon startup or
|
||||
// make every wallet write fail after recovery.
|
||||
logger.warn("Could not reconcile receive token reservation", {
|
||||
operationId: reservation.operationId,
|
||||
error: error instanceof Error ? error.message : String(error),
|
||||
});
|
||||
}
|
||||
}
|
||||
};
|
||||
await syncReceiveReservations();
|
||||
|
||||
const trustedMints = await coco.mint.getAllTrustedMints();
|
||||
const configuredDefault = walletConfig.defaultMintUrl;
|
||||
@@ -937,14 +1149,19 @@ export async function createCocoClient(
|
||||
|
||||
// Recovery runs in the background so the daemon can serve wallet reads
|
||||
// immediately. Value-moving operations await the same promise below.
|
||||
runWalletRecovery(coco, (progress) => {
|
||||
runWalletRecovery(
|
||||
coco,
|
||||
(progress) => {
|
||||
recoveryPhase = progress.phase;
|
||||
recoveryFailedMintQuotes = progress.failedMintQuotes;
|
||||
if (progress.phase !== "done") {
|
||||
startupProgress(`Wallet recovery: ${progress.phase}...`);
|
||||
}
|
||||
})
|
||||
.then(() => {
|
||||
},
|
||||
receiveRecoveryOperationIds,
|
||||
)
|
||||
.then(async () => {
|
||||
await syncReceiveReservations();
|
||||
recoveryDone = true;
|
||||
recoveryPhase = "done";
|
||||
recoveryResolve?.();
|
||||
@@ -1039,9 +1256,131 @@ export async function createCocoClient(
|
||||
},
|
||||
|
||||
async receiveCashu(token: string): Promise<string> {
|
||||
await waitForRecovery();
|
||||
await coco.wallet.receive(token);
|
||||
const reservation = reserveReceiveToken(database, token);
|
||||
if (!reservation.acquired) {
|
||||
if (reservation.existing?.state === "succeeded") {
|
||||
return "Token already received successfully";
|
||||
}
|
||||
if (
|
||||
reservation.existing?.state === "processing" &&
|
||||
reservation.existing.operationId
|
||||
) {
|
||||
// A prior request may have stopped after the mint call became
|
||||
// uncertain. Re-drive that one Coco operation instead of creating a
|
||||
// duplicate. Refresh is idempotent and uses its stored output data.
|
||||
try {
|
||||
const operation = await withTimeout(
|
||||
coco.ops.receive.refresh(reservation.existing.operationId),
|
||||
15_000,
|
||||
);
|
||||
if (operation.state === "finalized") {
|
||||
updateReceiveToken(database, reservation.tokenHash, {
|
||||
state: "succeeded",
|
||||
operationId: operation.id,
|
||||
});
|
||||
return "Token received successfully";
|
||||
}
|
||||
if (operation.state === "rolled_back") {
|
||||
const finalizedSibling = await findFinalizedReceiveSibling(operation);
|
||||
if (finalizedSibling) {
|
||||
updateReceiveToken(database, reservation.tokenHash, {
|
||||
state: "succeeded",
|
||||
operationId: finalizedSibling,
|
||||
});
|
||||
return "Token already received successfully";
|
||||
}
|
||||
updateReceiveToken(database, reservation.tokenHash, {
|
||||
state: "failed",
|
||||
operationId: operation.id,
|
||||
error: operation.error || "Token receive was rolled back",
|
||||
});
|
||||
throw new Error(operation.error || "Token receive was rolled back");
|
||||
}
|
||||
} catch (error) {
|
||||
const latest = await coco.ops.receive.get(reservation.existing.operationId);
|
||||
if (latest?.state === "finalized") {
|
||||
updateReceiveToken(database, reservation.tokenHash, {
|
||||
state: "succeeded",
|
||||
operationId: latest.id,
|
||||
});
|
||||
return "Token received successfully";
|
||||
}
|
||||
throw error;
|
||||
}
|
||||
throw new Error("Token receive is still unresolved");
|
||||
}
|
||||
if (reservation.existing?.state === "processing") {
|
||||
throw new Error("Token receive is already in progress");
|
||||
}
|
||||
throw new Error(reservation.existing?.error || "Token receive previously failed");
|
||||
}
|
||||
|
||||
let preparedOperationId: string | undefined;
|
||||
try {
|
||||
await waitForRecovery();
|
||||
const prepared = await coco.ops.receive.prepare({ token });
|
||||
preparedOperationId = prepared.id;
|
||||
updateReceiveToken(database, reservation.tokenHash, {
|
||||
state: "processing",
|
||||
operationId: prepared.id,
|
||||
});
|
||||
await coco.ops.receive.execute(prepared.id);
|
||||
updateReceiveToken(database, reservation.tokenHash, {
|
||||
state: "succeeded",
|
||||
operationId: prepared.id,
|
||||
});
|
||||
return "Token received successfully";
|
||||
} catch (error) {
|
||||
const message = error instanceof Error ? error.message : String(error);
|
||||
if (preparedOperationId) {
|
||||
let latest: Awaited<ReturnType<typeof coco.ops.receive.get>> = null;
|
||||
try {
|
||||
latest = await coco.ops.receive.get(preparedOperationId);
|
||||
} catch (lookupError) {
|
||||
logger.warn("Could not inspect failed receive operation", {
|
||||
operationId: preparedOperationId,
|
||||
error:
|
||||
lookupError instanceof Error ? lookupError.message : String(lookupError),
|
||||
});
|
||||
}
|
||||
if (latest?.state === "finalized") {
|
||||
updateReceiveToken(database, reservation.tokenHash, {
|
||||
state: "succeeded",
|
||||
operationId: preparedOperationId,
|
||||
});
|
||||
return "Token received successfully";
|
||||
}
|
||||
if (latest?.state === "executing") {
|
||||
updateReceiveToken(database, reservation.tokenHash, {
|
||||
state: "processing",
|
||||
operationId: preparedOperationId,
|
||||
error: message,
|
||||
});
|
||||
} else if (latest?.state === "rolled_back") {
|
||||
const finalizedSibling = await findFinalizedReceiveSibling(latest);
|
||||
if (finalizedSibling) {
|
||||
updateReceiveToken(database, reservation.tokenHash, {
|
||||
state: "succeeded",
|
||||
operationId: finalizedSibling,
|
||||
});
|
||||
return "Token already received successfully";
|
||||
}
|
||||
updateReceiveToken(database, reservation.tokenHash, {
|
||||
state: "failed",
|
||||
operationId: preparedOperationId,
|
||||
error: latest.error || message,
|
||||
});
|
||||
} else {
|
||||
// A prepared or missing operation had no known mint side effect.
|
||||
deleteReceiveTokenReservation(database, reservation.tokenHash);
|
||||
}
|
||||
} else {
|
||||
// Decode/validation failed before coco created an operation. Do not
|
||||
// permanently reserve malformed input or transient mint-fetch errors.
|
||||
releaseReceiveToken(database, reservation.tokenHash);
|
||||
}
|
||||
throw error;
|
||||
}
|
||||
},
|
||||
|
||||
async receiveBolt11(amount: number, mintUrl?: string): Promise<string> {
|
||||
|
||||
@@ -0,0 +1,238 @@
|
||||
import { describe, expect, it } from "bun:test";
|
||||
import type { ReceiveOperation } from "@cashu/coco-core";
|
||||
import { Database } from "bun:sqlite";
|
||||
import {
|
||||
clearInterruptedReceiveReservations,
|
||||
deleteReceiveTokenReservation,
|
||||
groupExecutingReceives,
|
||||
hashReceiveToken,
|
||||
initReceiveDedupSchema,
|
||||
listProcessingReceiveTokens,
|
||||
reconcileExecutingReceives,
|
||||
reserveReceiveToken,
|
||||
updateReceiveToken,
|
||||
type ReceiveReconcileSource,
|
||||
} from "./receive-dedup";
|
||||
|
||||
function operation(
|
||||
id: string,
|
||||
secret: string,
|
||||
outputB: string,
|
||||
createdAt = 1,
|
||||
): Extract<ReceiveOperation, { state: "executing" }> {
|
||||
return {
|
||||
id,
|
||||
mintUrl: "https://mint.example",
|
||||
unit: "sat",
|
||||
amount: 1,
|
||||
state: "executing",
|
||||
createdAt,
|
||||
updatedAt: createdAt,
|
||||
fee: 0,
|
||||
inputProofs: [{ amount: 1, id: "keyset", secret, C: "02" }],
|
||||
outputData: {
|
||||
keep: [
|
||||
{
|
||||
blindedMessage: { amount: 1, id: "keyset", B_: outputB },
|
||||
blindingFactor: "1",
|
||||
secret: "00",
|
||||
},
|
||||
],
|
||||
send: [],
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
function source(options: {
|
||||
operations: ReceiveOperation[];
|
||||
states?: string[];
|
||||
restored?: string[];
|
||||
saved?: string[];
|
||||
}) {
|
||||
const rolledBack: Array<{ id: string; reason: string }> = [];
|
||||
let stateChecks = 0;
|
||||
let restoreChecks = 0;
|
||||
const reconcileSource: ReceiveReconcileSource = {
|
||||
listExecuting: async () => options.operations,
|
||||
checkProofStates: async (op) => {
|
||||
stateChecks++;
|
||||
return (options.states ?? op.inputProofs.map(() => "SPENT")).map((state) => ({ state }));
|
||||
},
|
||||
restoreOutputs: async (_mintUrl, outputs) => {
|
||||
restoreChecks++;
|
||||
const restored = new Set(options.restored ?? []);
|
||||
return outputs.filter((output) => restored.has(output.B_));
|
||||
},
|
||||
hasSavedOutputs: async (op) => (options.saved ?? []).includes(op.id),
|
||||
rollBack: async (op, reason) => {
|
||||
rolledBack.push({ id: op.id, reason });
|
||||
},
|
||||
};
|
||||
return {
|
||||
source: reconcileSource,
|
||||
rolledBack,
|
||||
stateChecks: () => stateChecks,
|
||||
restoreChecks: () => restoreChecks,
|
||||
};
|
||||
}
|
||||
|
||||
describe("receive token hashing and reservation", () => {
|
||||
it("hashes the exact encoded token deterministically", () => {
|
||||
expect(hashReceiveToken("cashuA")).toBe(hashReceiveToken("cashuA"));
|
||||
expect(hashReceiveToken("cashuA")).not.toBe(hashReceiveToken("cashuB"));
|
||||
});
|
||||
|
||||
it("atomically reserves an exact token once", () => {
|
||||
const database = new Database(":memory:");
|
||||
initReceiveDedupSchema(database);
|
||||
const first = reserveReceiveToken(database, "cashuA");
|
||||
const second = reserveReceiveToken(database, "cashuA");
|
||||
expect(first.acquired).toBe(true);
|
||||
expect(second).toMatchObject({
|
||||
acquired: false,
|
||||
existing: { state: "processing", operationId: null },
|
||||
});
|
||||
database.close();
|
||||
});
|
||||
|
||||
it("persists successful status and releases retryable reservations", () => {
|
||||
const database = new Database(":memory:");
|
||||
initReceiveDedupSchema(database);
|
||||
const reservation = reserveReceiveToken(database, "cashuA");
|
||||
updateReceiveToken(database, reservation.tokenHash, {
|
||||
state: "succeeded",
|
||||
operationId: "receive-1",
|
||||
});
|
||||
expect(reserveReceiveToken(database, "cashuA").existing).toMatchObject({
|
||||
state: "succeeded",
|
||||
operationId: "receive-1",
|
||||
});
|
||||
deleteReceiveTokenReservation(database, reservation.tokenHash);
|
||||
expect(reserveReceiveToken(database, "cashuA").acquired).toBe(true);
|
||||
database.close();
|
||||
});
|
||||
|
||||
it("clears only interrupted reservations without Coco operation ids", () => {
|
||||
const database = new Database(":memory:");
|
||||
initReceiveDedupSchema(database);
|
||||
const orphan = reserveReceiveToken(database, "orphan");
|
||||
const linked = reserveReceiveToken(database, "linked");
|
||||
updateReceiveToken(database, linked.tokenHash, {
|
||||
state: "processing",
|
||||
operationId: "receive-1",
|
||||
});
|
||||
expect(clearInterruptedReceiveReservations(database)).toBe(1);
|
||||
expect(listProcessingReceiveTokens(database)).toEqual([
|
||||
expect.objectContaining({ tokenHash: linked.tokenHash, operationId: "receive-1" }),
|
||||
]);
|
||||
expect(reserveReceiveToken(database, "orphan").tokenHash).toBe(orphan.tokenHash);
|
||||
database.close();
|
||||
});
|
||||
});
|
||||
|
||||
describe("legacy receive grouping", () => {
|
||||
it("groups identical input proofs independent of operation id", () => {
|
||||
const groups = groupExecutingReceives([
|
||||
operation("one", "same", "B1"),
|
||||
operation("two", "same", "B2"),
|
||||
operation("three", "different", "B3"),
|
||||
]);
|
||||
expect(groups.map((group) => group.map((op) => op.id))).toEqual([
|
||||
["one", "two"],
|
||||
["three"],
|
||||
]);
|
||||
});
|
||||
|
||||
it("checks an unspent duplicate group once and retains one operation", async () => {
|
||||
const mock = source({
|
||||
operations: [operation("one", "same", "B1", 1), operation("two", "same", "B2", 2)],
|
||||
states: ["UNSPENT"],
|
||||
});
|
||||
const result = await reconcileExecutingReceives(mock.source);
|
||||
expect(mock.stateChecks()).toBe(1);
|
||||
expect(mock.restoreChecks()).toBe(0);
|
||||
expect(mock.rolledBack.map((entry) => entry.id)).toEqual(["two"]);
|
||||
expect(result).toMatchObject({ rolledBack: 1, retained: 1, unresolved: 0 });
|
||||
});
|
||||
|
||||
it("rolls back all spent duplicates when restore finds no outputs", async () => {
|
||||
const mock = source({
|
||||
operations: [operation("one", "same", "B1"), operation("two", "same", "B2")],
|
||||
states: ["SPENT"],
|
||||
restored: [],
|
||||
});
|
||||
const result = await reconcileExecutingReceives(mock.source);
|
||||
expect(mock.stateChecks()).toBe(1);
|
||||
expect(mock.restoreChecks()).toBe(1);
|
||||
expect(mock.rolledBack.map((entry) => entry.id)).toEqual(["one", "two"]);
|
||||
expect(result.rolledBack).toBe(2);
|
||||
});
|
||||
|
||||
it("retains the operation whose output the mint can restore", async () => {
|
||||
const mock = source({
|
||||
operations: [operation("one", "same", "B1"), operation("two", "same", "B2")],
|
||||
states: ["SPENT"],
|
||||
restored: ["B2"],
|
||||
});
|
||||
const result = await reconcileExecutingReceives(mock.source);
|
||||
expect(mock.rolledBack.map((entry) => entry.id)).toEqual(["one"]);
|
||||
expect(result).toMatchObject({ rolledBack: 1, retained: 1 });
|
||||
});
|
||||
|
||||
it("leaves a group untouched when the mint returns incomplete states", async () => {
|
||||
const first = operation("one", "same", "B1");
|
||||
first.inputProofs.push({ amount: 1, id: "keyset", secret: "other", C: "03" });
|
||||
const second = operation("two", "same", "B2");
|
||||
second.inputProofs.push({ amount: 1, id: "keyset", secret: "other", C: "03" });
|
||||
const mock = source({
|
||||
operations: [first, second],
|
||||
states: ["SPENT"],
|
||||
});
|
||||
const result = await reconcileExecutingReceives(mock.source);
|
||||
expect(mock.rolledBack).toEqual([]);
|
||||
expect(mock.restoreChecks()).toBe(0);
|
||||
expect(result.unresolved).toBe(2);
|
||||
});
|
||||
|
||||
it("leaves a group untouched when proof states are mixed", async () => {
|
||||
const first = operation("one", "same", "B1");
|
||||
first.inputProofs.push({ amount: 1, id: "keyset", secret: "other", C: "03" });
|
||||
const second = operation("two", "same", "B2");
|
||||
second.inputProofs.push({ amount: 1, id: "keyset", secret: "other", C: "03" });
|
||||
const mock = source({
|
||||
operations: [first, second],
|
||||
states: ["SPENT", "UNSPENT"],
|
||||
});
|
||||
const result = await reconcileExecutingReceives(mock.source);
|
||||
expect(mock.rolledBack).toEqual([]);
|
||||
expect(result.unresolved).toBe(2);
|
||||
});
|
||||
|
||||
it("scales mint checks with unique input groups rather than row count", async () => {
|
||||
const operations: ReceiveOperation[] = [];
|
||||
for (let group = 0; group < 47; group++) {
|
||||
const repeats = group < 44 ? 41 : 40; // 1,924 rows total
|
||||
for (let index = 0; index < repeats; index++) {
|
||||
operations.push(
|
||||
operation(
|
||||
`operation-${group}-${index}`,
|
||||
`secret-${group}`,
|
||||
`B-${group}-${index}`,
|
||||
index,
|
||||
),
|
||||
);
|
||||
}
|
||||
}
|
||||
expect(operations).toHaveLength(1924);
|
||||
const mock = source({ operations, states: ["SPENT"], restored: [] });
|
||||
const result = await reconcileExecutingReceives(mock.source);
|
||||
expect(mock.stateChecks()).toBe(47);
|
||||
expect(mock.restoreChecks()).toBe(47);
|
||||
expect(result).toMatchObject({
|
||||
executing: 1924,
|
||||
uniqueGroups: 47,
|
||||
rolledBack: 1924,
|
||||
unresolved: 0,
|
||||
});
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,384 @@
|
||||
import { createHash } from "crypto";
|
||||
import type { ReceiveOperation } from "@cashu/coco-core";
|
||||
import type { Database } from "bun:sqlite";
|
||||
|
||||
const TOKEN_TABLE = "routstrd_receive_tokens";
|
||||
const MAINTENANCE_TABLE = "routstrd_wallet_maintenance";
|
||||
const RECONCILE_BACKUP_KEY = "receive_reconcile_v1_backup";
|
||||
|
||||
export type StoredReceiveTokenState = "processing" | "succeeded" | "failed";
|
||||
|
||||
export interface StoredReceiveTokenRow {
|
||||
tokenHash: string;
|
||||
state: StoredReceiveTokenState;
|
||||
operationId: string | null;
|
||||
error: string | null;
|
||||
}
|
||||
|
||||
export interface ReceiveTokenReservation {
|
||||
acquired: boolean;
|
||||
tokenHash: string;
|
||||
existing?: StoredReceiveTokenRow;
|
||||
}
|
||||
|
||||
export interface ReceiveReconcileResult {
|
||||
executing: number;
|
||||
uniqueGroups: number;
|
||||
rolledBack: number;
|
||||
retained: number;
|
||||
unresolved: number;
|
||||
recoveryOperationIds: string[];
|
||||
backupPath?: string;
|
||||
}
|
||||
|
||||
interface SerializedBlindedMessageLike {
|
||||
amount: number;
|
||||
id: string;
|
||||
B_: string;
|
||||
}
|
||||
|
||||
interface SerializedOutput {
|
||||
blindedMessage: SerializedBlindedMessageLike;
|
||||
secret: string;
|
||||
}
|
||||
|
||||
interface SerializedOutputDataLike {
|
||||
keep: SerializedOutput[];
|
||||
send: SerializedOutput[];
|
||||
}
|
||||
|
||||
type ExecutingReceive = Extract<ReceiveOperation, { state: "executing" }>;
|
||||
|
||||
export interface ReceiveReconcileSource {
|
||||
listExecuting(): Promise<ReceiveOperation[]>;
|
||||
checkProofStates(operation: ExecutingReceive): Promise<Array<{ state: string }>>;
|
||||
restoreOutputs(
|
||||
mintUrl: string,
|
||||
outputs: SerializedBlindedMessageLike[],
|
||||
): Promise<Array<{ B_: string }>>;
|
||||
hasSavedOutputs(operation: ExecutingReceive): Promise<boolean>;
|
||||
rollBack(operation: ExecutingReceive, reason: string): Promise<void>;
|
||||
}
|
||||
|
||||
export function hashReceiveToken(token: string): string {
|
||||
return createHash("sha256").update(token).digest("hex");
|
||||
}
|
||||
|
||||
export function initReceiveDedupSchema(database: Database): void {
|
||||
database.exec(`
|
||||
CREATE TABLE IF NOT EXISTS ${TOKEN_TABLE} (
|
||||
tokenHash TEXT PRIMARY KEY,
|
||||
state TEXT NOT NULL CHECK (state IN ('processing', 'succeeded', 'failed')),
|
||||
operationId TEXT,
|
||||
error TEXT,
|
||||
createdAt INTEGER NOT NULL,
|
||||
updatedAt INTEGER NOT NULL
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS ${MAINTENANCE_TABLE} (
|
||||
key TEXT PRIMARY KEY,
|
||||
value TEXT NOT NULL,
|
||||
updatedAt INTEGER NOT NULL
|
||||
);
|
||||
`);
|
||||
}
|
||||
|
||||
export function listProcessingReceiveTokens(database: Database): StoredReceiveTokenRow[] {
|
||||
return database
|
||||
.query(
|
||||
`SELECT tokenHash, state, operationId, error
|
||||
FROM ${TOKEN_TABLE} WHERE state = 'processing'`,
|
||||
)
|
||||
.all() as StoredReceiveTokenRow[];
|
||||
}
|
||||
|
||||
export function deleteReceiveTokenReservation(database: Database, tokenHash: string): void {
|
||||
database.query(`DELETE FROM ${TOKEN_TABLE} WHERE tokenHash = ?`).run(tokenHash);
|
||||
}
|
||||
|
||||
export function clearInterruptedReceiveReservations(database: Database): number {
|
||||
const result = database
|
||||
.query(
|
||||
`DELETE FROM ${TOKEN_TABLE}
|
||||
WHERE state = 'processing' AND operationId IS NULL`,
|
||||
)
|
||||
.run();
|
||||
return result.changes;
|
||||
}
|
||||
|
||||
export function reserveReceiveToken(database: Database, token: string): ReceiveTokenReservation {
|
||||
const tokenHash = hashReceiveToken(token);
|
||||
const now = Math.floor(Date.now() / 1000);
|
||||
const result = database
|
||||
.query(
|
||||
`INSERT OR IGNORE INTO ${TOKEN_TABLE}
|
||||
(tokenHash, state, operationId, error, createdAt, updatedAt)
|
||||
VALUES (?, 'processing', NULL, NULL, ?, ?)`,
|
||||
)
|
||||
.run(tokenHash, now, now);
|
||||
|
||||
if (result.changes > 0) return { acquired: true, tokenHash };
|
||||
|
||||
const existing = database
|
||||
.query(
|
||||
`SELECT tokenHash, state, operationId, error
|
||||
FROM ${TOKEN_TABLE} WHERE tokenHash = ?`,
|
||||
)
|
||||
.get(tokenHash) as StoredReceiveTokenRow | null;
|
||||
return { acquired: false, tokenHash, ...(existing ? { existing } : {}) };
|
||||
}
|
||||
|
||||
export function releaseReceiveToken(database: Database, tokenHash: string): void {
|
||||
database
|
||||
.query(`DELETE FROM ${TOKEN_TABLE} WHERE tokenHash = ? AND operationId IS NULL`)
|
||||
.run(tokenHash);
|
||||
}
|
||||
|
||||
export function updateReceiveToken(
|
||||
database: Database,
|
||||
tokenHash: string,
|
||||
update: { state: StoredReceiveTokenState; operationId?: string; error?: string },
|
||||
): void {
|
||||
database
|
||||
.query(
|
||||
`UPDATE ${TOKEN_TABLE}
|
||||
SET state = ?, operationId = COALESCE(?, operationId), error = ?, updatedAt = ?
|
||||
WHERE tokenHash = ?`,
|
||||
)
|
||||
.run(
|
||||
update.state,
|
||||
update.operationId ?? null,
|
||||
update.error ?? null,
|
||||
Math.floor(Date.now() / 1000),
|
||||
tokenHash,
|
||||
);
|
||||
}
|
||||
|
||||
export function receiveInputFingerprint(operation: ReceiveOperation): string {
|
||||
// Legacy cleanup is intentionally stricter than new-request token hashing:
|
||||
// only operations with the exact same proof material are grouped. A shared
|
||||
// secret alone is not enough evidence to discard an operation.
|
||||
const proofs = operation.inputProofs
|
||||
.map((proof) => ({
|
||||
amount: proof.amount,
|
||||
id: proof.id,
|
||||
secret: proof.secret,
|
||||
C: proof.C,
|
||||
witness: proof.witness ?? null,
|
||||
}))
|
||||
.sort((left, right) =>
|
||||
JSON.stringify(left).localeCompare(JSON.stringify(right)),
|
||||
);
|
||||
return createHash("sha256")
|
||||
.update(JSON.stringify([operation.mintUrl, operation.unit, proofs]))
|
||||
.digest("hex");
|
||||
}
|
||||
|
||||
export function groupExecutingReceives(
|
||||
operations: ReceiveOperation[],
|
||||
): ExecutingReceive[][] {
|
||||
const groups = new Map<string, ExecutingReceive[]>();
|
||||
for (const operation of operations) {
|
||||
if (operation.state !== "executing") continue;
|
||||
const key = receiveInputFingerprint(operation);
|
||||
const group = groups.get(key) ?? [];
|
||||
group.push(operation);
|
||||
groups.set(key, group);
|
||||
}
|
||||
return [...groups.values()].map((group) =>
|
||||
group.sort((left, right) => left.createdAt - right.createdAt || left.id.localeCompare(right.id)),
|
||||
);
|
||||
}
|
||||
|
||||
function outputData(operation: ExecutingReceive): SerializedOutputDataLike | null {
|
||||
const value = operation.outputData as unknown;
|
||||
if (!value || typeof value !== "object") return null;
|
||||
const candidate = value as Partial<SerializedOutputDataLike>;
|
||||
if (!Array.isArray(candidate.keep) || !Array.isArray(candidate.send)) return null;
|
||||
return { keep: candidate.keep, send: candidate.send };
|
||||
}
|
||||
|
||||
function outputBlindedMessages(operation: ExecutingReceive): SerializedBlindedMessageLike[] {
|
||||
const data = outputData(operation);
|
||||
if (!data) return [];
|
||||
return [...data.keep, ...data.send]
|
||||
.map((entry) => entry?.blindedMessage)
|
||||
.filter((message): message is SerializedBlindedMessageLike =>
|
||||
Boolean(
|
||||
message &&
|
||||
typeof message.B_ === "string" &&
|
||||
message.B_ &&
|
||||
typeof message.id === "string" &&
|
||||
Number.isFinite(message.amount),
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
async function rollBackOthers(
|
||||
source: ReceiveReconcileSource,
|
||||
group: ExecutingReceive[],
|
||||
retainedIds: Set<string>,
|
||||
reason: string,
|
||||
): Promise<number> {
|
||||
let count = 0;
|
||||
for (const operation of group) {
|
||||
if (retainedIds.has(operation.id)) continue;
|
||||
await source.rollBack(operation, reason);
|
||||
count++;
|
||||
}
|
||||
return count;
|
||||
}
|
||||
|
||||
/**
|
||||
* Reduce legacy executing receives before coco-core's per-operation recovery.
|
||||
*
|
||||
* No state is changed unless the evidence is conclusive:
|
||||
* - saved local outputs identify the successful operation;
|
||||
* - all inputs are unspent, so one deterministic retry is sufficient; or
|
||||
* - all inputs are spent and restore identifies an output owner (or confirms
|
||||
* that none of the stored deterministic outputs exist).
|
||||
*
|
||||
* Unreachable mints, malformed output data, partial restore ownership, and
|
||||
* mixed proof states are left untouched for normal recovery/manual review.
|
||||
*/
|
||||
export async function reconcileExecutingReceives(
|
||||
source: ReceiveReconcileSource,
|
||||
): Promise<ReceiveReconcileResult> {
|
||||
const executing = (await source.listExecuting()).filter(
|
||||
(operation): operation is ExecutingReceive => operation.state === "executing",
|
||||
);
|
||||
const groups = groupExecutingReceives(executing);
|
||||
const result: ReceiveReconcileResult = {
|
||||
executing: executing.length,
|
||||
uniqueGroups: groups.length,
|
||||
rolledBack: 0,
|
||||
retained: 0,
|
||||
unresolved: 0,
|
||||
recoveryOperationIds: [],
|
||||
};
|
||||
|
||||
for (const group of groups) {
|
||||
try {
|
||||
const savedOwners = new Set<string>();
|
||||
for (const operation of group) {
|
||||
if (await source.hasSavedOutputs(operation)) savedOwners.add(operation.id);
|
||||
}
|
||||
if (savedOwners.size > 0) {
|
||||
result.rolledBack += await rollBackOthers(
|
||||
source,
|
||||
group,
|
||||
savedOwners,
|
||||
"Duplicate receive: outputs already persisted by another operation",
|
||||
);
|
||||
result.retained += savedOwners.size;
|
||||
result.recoveryOperationIds.push(...savedOwners);
|
||||
continue;
|
||||
}
|
||||
|
||||
const states = await source.checkProofStates(group[0]!);
|
||||
const completeStates = states.length === group[0]!.inputProofs.length;
|
||||
const allUnspent = completeStates && states.every((state) => state.state === "UNSPENT");
|
||||
const allSpent = completeStates && states.every((state) => state.state === "SPENT");
|
||||
|
||||
if (allUnspent) {
|
||||
const canonical = group.find((operation) => outputBlindedMessages(operation).length > 0);
|
||||
if (!canonical) {
|
||||
result.unresolved += group.length;
|
||||
continue;
|
||||
}
|
||||
result.rolledBack += await rollBackOthers(
|
||||
source,
|
||||
group,
|
||||
new Set([canonical.id]),
|
||||
"Duplicate receive: identical inputs retained by canonical operation",
|
||||
);
|
||||
result.retained++;
|
||||
result.recoveryOperationIds.push(canonical.id);
|
||||
continue;
|
||||
}
|
||||
|
||||
if (!allSpent) {
|
||||
result.unresolved += group.length;
|
||||
continue;
|
||||
}
|
||||
|
||||
const outputOwners = new Map<
|
||||
string,
|
||||
{ message: SerializedBlindedMessageLike; owners: Set<string> }
|
||||
>();
|
||||
let malformed = false;
|
||||
for (const operation of group) {
|
||||
const outputs = outputBlindedMessages(operation);
|
||||
if (outputs.length === 0) {
|
||||
malformed = true;
|
||||
continue;
|
||||
}
|
||||
for (const output of outputs) {
|
||||
const entry = outputOwners.get(output.B_) ?? {
|
||||
message: output,
|
||||
owners: new Set<string>(),
|
||||
};
|
||||
entry.owners.add(operation.id);
|
||||
outputOwners.set(output.B_, entry);
|
||||
}
|
||||
}
|
||||
if (malformed || outputOwners.size === 0) {
|
||||
result.unresolved += group.length;
|
||||
continue;
|
||||
}
|
||||
|
||||
const restored = await source.restoreOutputs(
|
||||
group[0]!.mintUrl,
|
||||
[...outputOwners.values()].map((entry) => entry.message),
|
||||
);
|
||||
const restoredOwners = new Set<string>();
|
||||
for (const output of restored) {
|
||||
for (const owner of outputOwners.get(output.B_)?.owners ?? []) {
|
||||
restoredOwners.add(owner);
|
||||
}
|
||||
}
|
||||
|
||||
if (restoredOwners.size === 0) {
|
||||
for (const operation of group) {
|
||||
await source.rollBack(
|
||||
operation,
|
||||
"Recovered: input proofs spent without recoverable outputs",
|
||||
);
|
||||
result.rolledBack++;
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
||||
result.rolledBack += await rollBackOthers(
|
||||
source,
|
||||
group,
|
||||
restoredOwners,
|
||||
"Duplicate receive: recoverable outputs belong to another operation",
|
||||
);
|
||||
result.retained += restoredOwners.size;
|
||||
result.recoveryOperationIds.push(...restoredOwners);
|
||||
} catch {
|
||||
// A network, parsing, or persistence failure is not evidence that an
|
||||
// operation is safe to abandon. Leave the complete group untouched.
|
||||
result.unresolved += group.length;
|
||||
}
|
||||
}
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
export function getReceiveReconcileBackup(database: Database): string | null {
|
||||
const row = database
|
||||
.query(`SELECT value FROM ${MAINTENANCE_TABLE} WHERE key = ?`)
|
||||
.get(RECONCILE_BACKUP_KEY) as { value: string } | null;
|
||||
return row?.value ?? null;
|
||||
}
|
||||
|
||||
export function setReceiveReconcileBackup(database: Database, path: string): void {
|
||||
database
|
||||
.query(
|
||||
`INSERT INTO ${MAINTENANCE_TABLE} (key, value, updatedAt)
|
||||
VALUES (?, ?, ?)
|
||||
ON CONFLICT(key) DO UPDATE SET value = excluded.value, updatedAt = excluded.updatedAt`,
|
||||
)
|
||||
.run(RECONCILE_BACKUP_KEY, path, Math.floor(Date.now() / 1000));
|
||||
}
|
||||
Reference in New Issue
Block a user