From 2d1e9c8a232fabfe6a3dca71666c9a412659aa1e Mon Sep 17 00:00:00 2001 From: bug_lord <287525836+OpensrcLord@users.noreply.github.com> Date: Thu, 1 Oct 2026 20:16:25 +0100 Subject: [PATCH] fix: reconcile in-flight keeper transactions --- .gitignore | 1 + packages/sdk/src/client.test.ts | 182 +++++++++++- packages/sdk/src/client.ts | 276 +++++++++++++++++- packages/sdk/src/index.ts | 7 + packages/sdk/src/public-api-snapshot.test.ts | 1 + packages/sdk/src/submission.ts | 39 +++ services/keeper/README.md | 4 + services/keeper/package.json | 2 +- .../keeper/src/checkpoint-restart.test.ts | 81 +++++ services/keeper/src/index.ts | 6 + services/keeper/src/keeper.ts | 148 +++++++--- services/keeper/src/run.ts | 19 +- services/keeper/src/serve.ts | 5 +- services/keeper/src/submission-journal.ts | 110 +++++++ 14 files changed, 815 insertions(+), 66 deletions(-) create mode 100644 packages/sdk/src/submission.ts create mode 100644 services/keeper/src/submission-journal.ts diff --git a/.gitignore b/.gitignore index d4efb572..abd0c13c 100644 --- a/.gitignore +++ b/.gitignore @@ -37,6 +37,7 @@ deployments/*.local.json .keeper-store.json.corrupted.* .keeper-checkpoint.json .keeper-checkpoint.json.corrupted.* +.keeper-submissions.json # Coverage coverage/ diff --git a/packages/sdk/src/client.test.ts b/packages/sdk/src/client.test.ts index cfe6b135..39cf0b02 100644 --- a/packages/sdk/src/client.test.ts +++ b/packages/sdk/src/client.test.ts @@ -1,7 +1,8 @@ // Copyright (c) 2026 Sub Rosa contributors import assert from "node:assert/strict"; import { describe, it } from "node:test"; -import { Keypair, rpc, StrKey } from "@stellar/stellar-sdk"; +import { Keypair, rpc, StrKey, xdr } from "@stellar/stellar-sdk"; +import { createFakeTime } from "@sub-rosa/time"; import { SubRosaClient } from "./client.js"; import { @@ -14,6 +15,7 @@ import type { SubmitSignedTransactionParams, TransactionSubmitter, } from "./submitter.js"; +import { submissionKey, type SubmissionIdentity, type SubmissionJournal, type SubmissionRecord } from "./submission.js"; import { sealFixture, fixtureBinding } from "./testing/seal-fixture.js"; const BASE_CONFIG = { @@ -312,6 +314,183 @@ describe("SubRosaClient external submitter failures", () => { }); }); +describe("durable in-flight submission recovery", () => { + const identity: SubmissionIdentity = { + operation: "settle", + roundId: "9", + network: BASE_CONFIG.networkPassphrase, + contractId: BASE_CONFIG.contractId, + }; + + class MemoryJournal implements SubmissionJournal { + readonly records = new Map(); + async get(key: SubmissionIdentity) { + const value = this.records.get(submissionKey(key)); + return value ? { ...value } : undefined; + } + async put(record: SubmissionRecord) { + this.records.set(submissionKey(record), { ...record }); + } + } + + function makeClient(options: { + journal: MemoryJournal; + nowMs: number; + send: () => Promise; + lookup: () => Promise; + }) { + let builds = 0; + const fakeTime = createFakeTime(options.nowMs); + const server = { + getNetwork: async () => ({ passphrase: BASE_CONFIG.networkPassphrase, protocolVersion: "23" }), + getLedgerEntries: async () => ({ entries: [{}], latestLedger: 123 }), + sendTransaction: options.send, + getTransaction: options.lookup, + } as unknown as rpc.Server; + const client = new SubRosaClient({ + ...BASE_CONFIG, + secretKey: Keypair.random().secret(), + _server: server, + submissionJournal: options.journal, + confirmTimeout: 1_000, + pollInterval: 100, + time: { clock: fakeTime.clock, scheduler: fakeTime.scheduler }, + }); + const signed = { + hash: () => Buffer.alloc(32, 7), + toXDR: () => "signed-xdr", + timeBounds: { minTime: "0", maxTime: String(Math.floor(options.nowMs / 1000) + 5) }, + }; + const assembled = { + signed, + async sign() {}, + options: { parseResultXdr: () => ({ unwrap: () => undefined }) }, + }; + Object.defineProperty(client.contract, "settle", { + configurable: true, + value: async () => { + builds += 1; + return assembled; + }, + }); + return { client, fakeTime, server, buildCount: () => builds }; + } + + it("reconciles a confirmed hash before asking the contract client to build a retry", async () => { + const nowMs = 4_000_000; + const journal = new MemoryJournal(); + await journal.put({ + ...identity, + hash: "0xconfirmed-before-build", + state: "pending", + submittedAtMs: nowMs - 1_000, + expiresAtMs: nowMs + 60_000, + }); + let sends = 0; + const { client, buildCount } = makeClient({ + journal, + nowMs, + send: async () => { + sends += 1; + return { status: "PENDING" }; + }, + lookup: async () => ({ + status: rpc.Api.GetTransactionStatus.SUCCESS, + returnValue: xdr.ScVal.scvVoid(), + }), + }); + + await client.settle(9); + assert.equal(buildCount(), 0); + assert.equal(sends, 0); + assert.equal((await journal.get(identity))?.state, "confirmed"); + }); + + it("journals before an ambiguous send, then reuses the same hash when it confirms", async () => { + const nowMs = 1_000_000; + const journal = new MemoryJournal(); + let sends = 0; + let found = false; + const { client } = makeClient({ + journal, + nowMs, + send: async () => { + sends += 1; + assert.equal((await journal.get(identity))?.state, "pending"); + throw new Error("response lost after RPC accepted transaction"); + }, + lookup: async () => ({ + status: found ? rpc.Api.GetTransactionStatus.SUCCESS : rpc.Api.GetTransactionStatus.NOT_FOUND, + ...(found ? { returnValue: xdr.ScVal.scvVoid() } : {}), + }), + }); + + await assert.rejects(client.settle(9), /did not finalize it in time/); + const pending = await journal.get(identity); + assert.equal(pending?.state, "pending"); + assert.equal(pending?.operation, "settle"); + assert.equal(pending?.roundId, "9"); + assert.equal(pending?.network, BASE_CONFIG.networkPassphrase); + assert.equal(sends, 1); + + found = true; + const recovered = await client.reconcileSubmission({ operation: "settle", roundId: "9" }); + assert.deepEqual(recovered, { hash: pending?.hash, state: "confirmed" }); + assert.equal(sends, 1, "reconciliation polls the old hash and never resubmits it"); + }); + + it("replaces an expired hash with one submission, not a retry burst", async () => { + const nowMs = 2_000_000; + const journal = new MemoryJournal(); + await journal.put({ + ...identity, + hash: "expired-hash", + state: "pending", + submittedAtMs: nowMs - 10_000, + expiresAtMs: nowMs - 1, + }); + let sends = 0; + let wasSubmitted = false; + const { client } = makeClient({ + journal, + nowMs, + send: async () => { + sends += 1; + wasSubmitted = true; + return { status: "PENDING" }; + }, + lookup: async () => wasSubmitted + ? { status: rpc.Api.GetTransactionStatus.SUCCESS, returnValue: xdr.ScVal.scvVoid() } + : { status: rpc.Api.GetTransactionStatus.NOT_FOUND }, + }); + + await client.settle(9); + assert.equal(sends, 1); + assert.notEqual((await journal.get(identity))?.hash, "expired-hash"); + assert.equal((await journal.get(identity))?.state, "confirmed"); + }); + + it("records a definitive network failure and does not treat it as an expiry", async () => { + const journal = new MemoryJournal(); + let sends = 0; + const { client } = makeClient({ + journal, + nowMs: 3_000_000, + send: async () => { + sends += 1; + return { status: "PENDING" }; + }, + lookup: async () => ({ status: rpc.Api.GetTransactionStatus.FAILED }), + }); + + await assert.rejects(client.settle(9), /ended with status FAILED/); + assert.equal((await journal.get(identity))?.state, "failed"); + const recovered = await client.reconcileSubmission({ operation: "settle", roundId: "9" }); + assert.equal(recovered?.state, "failed"); + assert.equal(sends, 1, "a definitive contract failure is not replaced as if it expired"); + }); +}); + describe("SubRosaClient passkey session binding", () => { const SWAPPED_CONTRACT_ID = StrKey.encodeContract(Buffer.alloc(32, 2)); const PUBLIC_PASSPHRASE = "Public Global Stellar Network ; September 2015"; @@ -430,4 +609,3 @@ describe("SubRosaClient passkey session binding", () => { } }); }); - diff --git a/packages/sdk/src/client.ts b/packages/sdk/src/client.ts index 8a74a865..df049bd6 100644 --- a/packages/sdk/src/client.ts +++ b/packages/sdk/src/client.ts @@ -6,7 +6,7 @@ // contract Spec embedded in the generated bindings, so the bytes on the wire are // exactly what the contract expects. -import { Keypair, rpc } from "@stellar/stellar-sdk"; +import { Keypair, rpc, scValToNative, xdr } from "@stellar/stellar-sdk"; import { normalizeError } from "@sub-rosa/logging/errors"; import type { AssembledTransaction, @@ -34,6 +34,11 @@ import { assertSealedBid } from "./encrypted-blob.js"; import type { SealedBidBinding } from "./encrypted-blob.js"; import { networkFingerprint } from "./receipt.js"; import type { TransactionSubmitter } from "./submitter.js"; +import { + type SubmissionIdentity, + type SubmissionJournal, + type SubmissionRecord, +} from "./submission.js"; import { evaluatePreflight, classifyPreflightBuildError, @@ -95,14 +100,16 @@ export interface SubRosaClientConfig { allowHttp?: boolean; /** Optional external submitter. Direct Soroban RPC remains the default. */ submitter?: TransactionSubmitter; + /** Durable transaction journal used to reconcile submissions across retries/restarts. */ + submissionJournal?: SubmissionJournal; /** - * How long (ms) to poll RPC for transaction finality when using an external - * submitter. Must be at least 1_000. Default: 60_000. + * How long (ms) to poll RPC for transaction finality when a durable journal + * is enabled or an external submitter is used. Must be at least 1_000. Default: 60_000. */ confirmTimeout?: number; /** - * How long (ms) to wait between polling RPC for transaction status when - * using an external submitter. Must be at least 100. Default: 1_500. + * How long (ms) to wait between polling RPC for transaction status. Must be + * at least 100. Default: 1_500. */ pollInterval?: number; /** Injectable wall clock and scheduler. Default: systemTime. */ @@ -201,6 +208,7 @@ export class SubRosaClient { readonly #rpcUrl: string; readonly #allowHttp: boolean; readonly #submitter?: TransactionSubmitter; + readonly #submissionJournal?: SubmissionJournal; readonly #confirmTimeout: number; readonly #pollInterval: number; readonly #assetConfig?: import("./asset-config.js").AssetConfig; @@ -247,6 +255,7 @@ export class SubRosaClient { this.#rpcUrl = config.rpcUrl; this.#allowHttp = allowHttp; this.#submitter = config.submitter; + this.#submissionJournal = config.submissionJournal; this.#confirmTimeout = confirmTimeout; this.#pollInterval = pollInterval; this.#assetConfig = config.assetConfig; @@ -365,7 +374,11 @@ export class SubRosaClient { } - async #sendUnwrap(tx: AssembledTransaction>): Promise { + async #sendUnwrap( + tx: AssembledTransaction>, + identity: SubmissionIdentity, + ): Promise { + if (this.#submissionJournal) return this.#sendTracked(tx, identity); if (!this.#submitter) { try { const sent = await tx.signAndSend(); @@ -425,8 +438,211 @@ export class SubRosaClient { }); } + /** Poll a previously journaled transaction before the keeper builds a retry. */ + async reconcileSubmission( + identity: Omit, + ): Promise | undefined> { + if (!this.#submissionJournal) return undefined; + const fullIdentity: SubmissionIdentity = { + ...identity, + network: this.networkPassphrase, + contractId: this.contractId, + }; + const record = await this.#submissionJournal.get(fullIdentity); + if (!record || record.state !== "pending") { + return record ? { hash: record.hash, state: record.state } : undefined; + } + const status = await this.#pollTracked(record); + return { hash: record.hash, state: status.state }; + } + + async #sendTracked( + tx: AssembledTransaction>, + identity: SubmissionIdentity, + ): Promise { + await tx.sign(); + if (!tx.signed) throw new SubRosaSubmitError("transaction was not signed"); + + const signed = tx.signed; + const existing = await this.#submissionJournal!.get(identity); + if (existing?.state === "pending") { + const recovered = await this.#pollTracked(existing, tx); + if (recovered.state === "confirmed") { + if (!recovered.hasResult) throw new SubRosaMissingReturnValueError(existing.hash); + return recovered.value as T; + } + if (recovered.state === "failed") { + throw new SubRosaTransactionError(existing.hash, "FAILED"); + } + if (recovered.state === "pending") { + throw new SubRosaTimeoutError({ + hash: existing.hash, + submitter: this.#submitter?.name ?? "Soroban RPC", + lastStatus: "NOT_FOUND", + timeoutMs: this.#confirmTimeout, + pollIntervalMs: this.#pollInterval, + }); + } + } else if (existing?.state === "confirmed" && existing.resultXdr) { + return tx.options.parseResultXdr(xdr.ScVal.fromXDR(existing.resultXdr, "base64")).unwrap(); + } else if (existing?.state === "confirmed") { + throw new SubRosaMissingReturnValueError(existing.hash); + } else if (existing?.state === "failed") { + throw new SubRosaTransactionError(existing.hash, "FAILED"); + } + + const now = this.#clock.nowMs(); + const hash = Buffer.from(signed.hash()).toString("hex"); + const timeBounds = signed.timeBounds; + const expiresAtMs = timeBounds?.maxTime + ? Number(timeBounds.maxTime) * 1000 + : now + Math.max(300_000, this.#confirmTimeout); + let record: SubmissionRecord = { + ...identity, + hash, + state: "pending", + submittedAtMs: now, + expiresAtMs, + }; + // Persist the deterministic hash before any network submission call. If the + // RPC accepts the transaction but drops its response, the next run can poll + // this same hash instead of constructing a second transaction. + await this.#submissionJournal!.put(record); + + if (this.#submitter) { + try { + const accepted = await this.#submitter.submitSignedTransaction({ + signedTransactionXdr: signed.toXDR(), + contractId: this.contractId, + networkPassphrase: this.networkPassphrase, + rpcUrl: this.#rpcUrl, + }); + if (accepted.hash && accepted.hash !== record.hash) { + record = { ...record, hash: accepted.hash }; + await this.#submissionJournal!.put(record); + } + if (accepted.relayerTransactionId) { + record = { ...record, relayerTransactionId: accepted.relayerTransactionId }; + await this.#submissionJournal!.put(record); + } + } catch { + // A submitter timeout is ambiguous: still reconcile the signed tx hash. + } + } else { + try { + const sent = await this.#server.sendTransaction(signed); + if (sent.status === "ERROR") { + const failed = { ...record, state: "failed" as const, failure: "RPC rejected submission" }; + await this.#submissionJournal!.put(failed); + throw new SubRosaTransactionError(record.hash, "FAILED"); + } + } catch (error) { + if (error instanceof SubRosaTransactionError) throw error; + // The send response may have been lost after acceptance; poll the hash. + } + } + + const terminal = await this.#pollTracked(record, tx); + if (terminal.state === "confirmed") { + if (!terminal.hasResult) throw new SubRosaMissingReturnValueError(record.hash); + return terminal.value as T; + } + if (terminal.state === "failed") throw new SubRosaTransactionError(record.hash, "FAILED"); + throw new SubRosaTimeoutError({ + hash: record.hash, + submitter: this.#submitter?.name ?? "Soroban RPC", + lastStatus: terminal.state === "expired" ? "EXPIRED" : "NOT_FOUND", + timeoutMs: this.#confirmTimeout, + pollIntervalMs: this.#pollInterval, + }); + } + + async #pollTracked( + record: SubmissionRecord, + tx?: AssembledTransaction>, + ): Promise<{ state: SubmissionRecord["state"]; value?: T; hasResult?: boolean }> { + const deadline = this.#clock.nowMs() + this.#confirmTimeout; + let lastStatus = "NOT_FOUND"; + while (true) { + try { + const response = await this.#server.getTransaction(record.hash); + lastStatus = response.status; + if (response.status === rpc.Api.GetTransactionStatus.SUCCESS) { + const resultXdr = response.returnValue?.toXDR("base64"); + const confirmed: SubmissionRecord = { + ...record, + state: "confirmed", + ...(resultXdr ? { resultXdr } : {}), + }; + await this.#submissionJournal!.put(confirmed); + return { + state: "confirmed", + hasResult: Boolean(resultXdr), + ...(tx && resultXdr + ? { value: tx.options.parseResultXdr(xdr.ScVal.fromXDR(resultXdr, "base64")).unwrap() as T } + : {}), + }; + } + if (response.status === rpc.Api.GetTransactionStatus.FAILED) { + await this.#submissionJournal!.put({ + ...record, + state: "failed", + failure: "Soroban transaction failed", + }); + return { state: "failed" }; + } + } catch (error) { + if (error instanceof SubRosaTransactionError) throw error; + // RPC lookup failures do not clear the durable pending record. + } + if (lastStatus === rpc.Api.GetTransactionStatus.NOT_FOUND && this.#clock.nowMs() >= record.expiresAtMs) { + await this.#submissionJournal!.put({ ...record, state: "expired" }); + return { state: "expired" }; + } + if (this.#clock.nowMs() >= deadline) return { state: "pending" }; + await this.#sleep(this.#pollInterval); + } + } + #sleep: (ms: number) => Promise = (ms) => this.#scheduler.sleep(ms); + #submissionIdentity( + operation: string, + roundId: number | bigint | null, + discriminator?: string, + ): SubmissionIdentity { + return { + operation, + roundId: roundId === null ? null : normalizeRoundId(roundId), + ...(discriminator ? { discriminator } : {}), + network: this.networkPassphrase, + contractId: this.contractId, + }; + } + + async #recoverBeforeBuild(identity: SubmissionIdentity): Promise { + if (!this.#submissionJournal) return undefined; + const record = await this.#submissionJournal.get(identity); + if (!record) return undefined; + if (record.state === "confirmed") return record; + if (record.state === "failed") throw new SubRosaTransactionError(record.hash, "FAILED"); + if (record.state === "expired") return undefined; + + const status = await this.#pollTracked(record); + if (status.state === "confirmed") { + return (await this.#submissionJournal.get(identity)) ?? { ...record, state: "confirmed" }; + } + if (status.state === "failed") throw new SubRosaTransactionError(record.hash, "FAILED"); + if (status.state === "expired") return undefined; + throw new SubRosaTimeoutError({ + hash: record.hash, + submitter: this.#submitter?.name ?? "Soroban RPC", + lastStatus: "NOT_FOUND", + timeoutMs: this.#confirmTimeout, + pollIntervalMs: this.#pollInterval, + }); + } + // ── State-changing calls (sign + submit over RPC) ────────────────────── /** Build the on-chain asset_config argument from SDK params. */ #buildAssetConfig(params: CreateRoundParams): RoundAssetConfig { if (!params.assetConfig) { @@ -450,6 +666,25 @@ export class SubRosaClient { async createRound(params: CreateRoundParams): Promise { const operator = params.operator ?? this.#requireSource("operator"); + const submissionIdentity = this.#submissionIdentity( + "create_round", + null, + JSON.stringify({ + operator, + itemRef: toHex(params.itemRef), + revealRound: String(params.revealRound), + commitDeadline: String(params.commitDeadline), + revealDeadline: String(params.revealDeadline), + auditorPubkey: toHex(params.auditorPubkey), + clearingRule: params.clearingRule ?? "HighestBid", + assetConfig: params.assetConfig ?? null, + }), + ); + const recovered = await this.#recoverBeforeBuild(submissionIdentity); + if (recovered) { + if (!recovered.resultXdr) throw new SubRosaMissingReturnValueError(recovered.hash); + return BigInt(String(scValToNative(xdr.ScVal.fromXDR(recovered.resultXdr, "base64")))); + } const clearing_rule = { tag: params.clearingRule ?? "HighestBid", values: undefined, @@ -486,7 +721,7 @@ export class SubRosaClient { } as Parameters[0]), ); - return this.#sendUnwrap(tx); + return this.#sendUnwrap(tx, submissionIdentity); } async commit(params: CommitParams): Promise { @@ -513,6 +748,8 @@ export class SubRosaClient { ); } const seal_round = toBigInt(rawSealRound); + const submissionIdentity = this.#submissionIdentity("commit", params.roundId, bidder); + if (await this.#recoverBeforeBuild(submissionIdentity)) return; const tx = await this.#validatedContractCall(() => this.contract.commit({ round_id: normalizeRoundId(params.roundId), @@ -524,23 +761,27 @@ export class SubRosaClient { seal_round, }), ); - await this.#sendUnwrap(tx); + await this.#sendUnwrap(tx, submissionIdentity); } async openReveal( roundId: number | bigint, drandSignature: Uint8Array, ): Promise { + const submissionIdentity = this.#submissionIdentity("open_reveal", roundId); + if (await this.#recoverBeforeBuild(submissionIdentity)) return; const tx = await this.#validatedContractCall(() => this.contract.open_reveal({ round_id: normalizeRoundId(roundId), drand_signature: toBuffer(drandSignature), }), ); - await this.#sendUnwrap(tx); + await this.#sendUnwrap(tx, submissionIdentity); } async reveal(params: RevealParams): Promise { + const submissionIdentity = this.#submissionIdentity("reveal", params.roundId, params.bidder); + if (await this.#recoverBeforeBuild(submissionIdentity)) return; const tx = await this.#validatedContractCall(() => this.contract.reveal({ round_id: normalizeRoundId(params.roundId), @@ -549,31 +790,40 @@ export class SubRosaClient { nonce: toBuffer(params.nonce), }), ); - await this.#sendUnwrap(tx); + await this.#sendUnwrap(tx, submissionIdentity); } /** Clear a round. Returns the winning address, or undefined if the round was * voided for having no valid bids. */ async clear(roundId: number | bigint): Promise { + const submissionIdentity = this.#submissionIdentity("clear", roundId); + if (await this.#recoverBeforeBuild(submissionIdentity)) { + const round = await this.getRound(roundId); + return round.winner ?? undefined; + } const tx = await this.#validatedContractCall(() => this.contract.clear({ round_id: normalizeRoundId(roundId) }), ); - const winner = await this.#sendUnwrap(tx); + const winner = await this.#sendUnwrap(tx, submissionIdentity); return winner ?? undefined; } async settle(roundId: number | bigint): Promise { + const submissionIdentity = this.#submissionIdentity("settle", roundId); + if (await this.#recoverBeforeBuild(submissionIdentity)) return; const tx = await this.#validatedContractCall(() => this.contract.settle({ round_id: normalizeRoundId(roundId) }), ); - await this.#sendUnwrap(tx); + await this.#sendUnwrap(tx, submissionIdentity); } async void(roundId: number | bigint): Promise { + const submissionIdentity = this.#submissionIdentity("void", roundId); + if (await this.#recoverBeforeBuild(submissionIdentity)) return; const tx = await this.#validatedContractCall(() => this.contract.void({ round_id: normalizeRoundId(roundId) }), ); - await this.#sendUnwrap(tx); + await this.#sendUnwrap(tx, submissionIdentity); } // ── Preflight simulation (no signing/submission) ───────────────────── diff --git a/packages/sdk/src/index.ts b/packages/sdk/src/index.ts index 0d7d15cb..f97b84b7 100644 --- a/packages/sdk/src/index.ts +++ b/packages/sdk/src/index.ts @@ -8,6 +8,13 @@ export { type ClearingRuleTag, } from "./client.js"; export { normalizeRoundId, normalizeSorobanContractId } from "./ids.js"; +export { + submissionKey, + type SubmissionIdentity, + type SubmissionJournal, + type SubmissionRecord, + type SubmissionState, +} from "./submission.js"; export { type PreflightOperation, type PreflightResult, diff --git a/packages/sdk/src/public-api-snapshot.test.ts b/packages/sdk/src/public-api-snapshot.test.ts index 2c7d5583..10618e64 100644 --- a/packages/sdk/src/public-api-snapshot.test.ts +++ b/packages/sdk/src/public-api-snapshot.test.ts @@ -95,6 +95,7 @@ const EXPECTED_EXPORTS = [ "runMainnetReadiness", "summarizeDeploymentValue", "serializeReceipt", + "submissionKey", "tryDecodeBase64", "tryDecodeHex", "validateAssetConfig", diff --git a/packages/sdk/src/submission.ts b/packages/sdk/src/submission.ts new file mode 100644 index 00000000..f0ceb5a1 --- /dev/null +++ b/packages/sdk/src/submission.ts @@ -0,0 +1,39 @@ +// SPDX-License-Identifier: MIT + +/** A stable identity for one logical contract operation. */ +export interface SubmissionIdentity { + operation: string; + roundId: string | null; + network: string; + contractId: string; + /** Distinguishes repeated operations in the same round (for example bidder reveals). */ + discriminator?: string; +} + +export type SubmissionState = "pending" | "confirmed" | "expired" | "failed"; + +export interface SubmissionRecord extends SubmissionIdentity { + hash: string; + state: SubmissionState; + submittedAtMs: number; + expiresAtMs: number; + relayerTransactionId?: string | null; + resultXdr?: string; + failure?: string; +} + +/** Durable storage supplied by the host application; never store signed XDR or secrets. */ +export interface SubmissionJournal { + get(identity: SubmissionIdentity): Promise; + put(record: SubmissionRecord): Promise; +} + +export function submissionKey(identity: SubmissionIdentity): string { + return JSON.stringify([ + identity.network, + identity.contractId, + identity.operation, + identity.roundId, + identity.discriminator ?? "", + ]); +} diff --git a/services/keeper/README.md b/services/keeper/README.md index 74ea2026..0a6f02ea 100644 --- a/services/keeper/README.md +++ b/services/keeper/README.md @@ -163,6 +163,10 @@ Steps tracked: `open-reveal`, `reveal`, `clear`, `settle`, `void`. `open-reveal` is recorded but never used to skip work: whether the reveal window is open is already authoritative on-chain, and trusting the cursor there could strand an Open round. An unreadable or corrupted checkpoint file is backed up (`*.corrupted.`) and the keeper starts from a fresh cursor rather than guessing. +### In-flight transaction recovery + +The keeper also uses a separate durable submission journal (default `.keeper-submissions.json`, override with `KEEPER_SUBMISSION_PATH`). The SDK writes the hash, operation, round, network, submission time, and expiry before the RPC send call. On restart, the keeper polls a pending hash before constructing a replacement transaction. Confirmed hashes are recorded in the completion checkpoint; definitive failures stop the pass; expired hashes permit one replacement. The journal stores no signed XDR or secret key and refuses to open under a different network or contract binding. + ### Dry-run `KEEPER_DRY_RUN=true npm run start` prints the checkpoint it *would* write — path, binding, proposed step, and the exact file content — inside the dry-run summary. It submits no transactions (`transactionsSubmitted: 0`) and writes nothing (`checkpoint.filesWritten: 0`). If the existing checkpoint would block a live run, the summary reports it as `checkpoint.mismatch` (`network` or `contractId`). diff --git a/services/keeper/package.json b/services/keeper/package.json index 80676a89..fd750053 100644 --- a/services/keeper/package.json +++ b/services/keeper/package.json @@ -15,7 +15,7 @@ "watch": "node --import tsx src/watch.ts", "serve": "node --import tsx src/serve.ts", "queue": "node --import tsx src/queue.ts", - "test": "node --import tsx --test src/dry-run.test.ts src/keeper.test.ts src/watch.test.ts src/store.test.ts src/queue-cli.test.ts src/watch-queue.test.ts src/settlement-guard.test.ts src/status.test.ts src/status-server.test.ts src/queue-replay.test.ts scripts/mainnet-settle.test.ts", + "test": "node --import tsx --test src/dry-run.test.ts src/keeper.test.ts src/checkpoint-restart.test.ts src/watch.test.ts src/store.test.ts src/queue-cli.test.ts src/watch-queue.test.ts src/settlement-guard.test.ts src/status.test.ts src/status-server.test.ts src/queue-replay.test.ts scripts/mainnet-settle.test.ts", "typecheck": "tsc --noEmit -p tsconfig.json" }, "dependencies": { diff --git a/services/keeper/src/checkpoint-restart.test.ts b/services/keeper/src/checkpoint-restart.test.ts index 297c752d..c8323ff7 100644 --- a/services/keeper/src/checkpoint-restart.test.ts +++ b/services/keeper/src/checkpoint-restart.test.ts @@ -36,6 +36,7 @@ import { closeRound, keepRound, voidIfStale, watchRound } from "./keeper.js"; import { createSettlementGuard } from "./settlement-guard.js"; import { KeeperStore } from "./store.js"; import { resumeCheckpoint, runWatchLoop } from "./watch-loop.js"; +import { FileSubmissionJournal } from "./submission-journal.js"; const NETWORK = "Test SDF Network ; September 2015"; const CONTRACT = "CTESTCONTRACT"; @@ -66,6 +67,35 @@ afterEach(() => { fs.rmSync(dir, { recursive: true, force: true }); }); +describe("durable submission journal", () => { + test("persists operation metadata and pending hashes across store instances", async () => { + const journalPath = path.join(dir, "submissions.json"); + const journal = new FileSubmissionJournal({ network: NETWORK, contractId: CONTRACT, path: journalPath }); + const identity = { operation: "settle", roundId: "1", network: NETWORK, contractId: CONTRACT }; + const record = { + ...identity, + hash: "0xpendinghash", + state: "pending" as const, + submittedAtMs: NOW_MS, + expiresAtMs: NOW_MS + 60_000, + }; + await journal.put(record); + + const restarted = new FileSubmissionJournal({ network: NETWORK, contractId: CONTRACT, path: journalPath }); + assert.deepEqual(await restarted.get(identity), record); + assert.doesNotMatch(fs.readFileSync(journalPath, "utf8"), /signed-xdr|SECRET/); + }); + + test("refuses to resume a journal bound to another network", () => { + const journalPath = path.join(dir, "wrong-network.json"); + new FileSubmissionJournal({ network: NETWORK, contractId: CONTRACT, path: journalPath }); + assert.throws( + () => new FileSubmissionJournal({ network: "Public Network", contractId: CONTRACT, path: journalPath }), + /does not match configured/, + ); + }); +}); + // ── In-memory chain ──────────────────────────────────────────────────────── interface FakeChainOptions { @@ -212,6 +242,57 @@ const confirmAll = async (): Promise => "confirmed"; // ── Crash AFTER the checkpoint write ─────────────────────────────────────── describe("restart after the checkpoint was written", () => { + test("recovers a confirmed settle hash written before a crash, without rebroadcasting", async () => { + const chain = new FakeChain({ status: "Cleared" }); + const sdk = asSdk(chain); + Object.defineProperty(sdk, "reconcileSubmission", { + value: async () => ({ hash: "0xaccepted-settle", state: "confirmed" }), + }); + const checkpoint = newStore(); + + const result = await closeRound( + { sdk, drand: {} as never, log: () => {}, time: { clock, scheduler }, checkpoint }, + ROUND_ID, + ); + + assert.equal(result.settled, true); + assert.equal(chain.count("settle"), 0); + assert.equal(checkpoint.isComplete(ROUND_ID, "settle"), true); + assert.equal(checkpoint.read(ROUND_ID)?.stepHashes.settle, "0xaccepted-settle"); + }); + + test("keeps an unresolved settle pending instead of creating another transaction", async () => { + const chain = new FakeChain({ status: "Cleared" }); + const sdk = asSdk(chain); + Object.defineProperty(sdk, "reconcileSubmission", { + value: async () => ({ hash: "0xstill-pending", state: "pending" }), + }); + + await assert.rejects( + closeRound( + { sdk, drand: {} as never, log: () => {}, time: { clock, scheduler }, checkpoint: newStore() }, + ROUND_ID, + ), + /still pending/, + ); + assert.equal(chain.count("settle"), 0); + }); + + test("replaces an expired settle hash with exactly one new submission", async () => { + const chain = new FakeChain({ status: "Cleared" }); + const sdk = asSdk(chain); + Object.defineProperty(sdk, "reconcileSubmission", { + value: async () => ({ hash: "0xexpired-settle", state: "expired" }), + }); + + const result = await closeRound( + { sdk, drand: {} as never, log: () => {}, time: { clock, scheduler }, checkpoint: newStore() }, + ROUND_ID, + ); + assert.equal(result.settled, true); + assert.equal(chain.count("settle"), 1); + }); + test("a confirmed settle is not submitted again by a lagging replica", async () => { // The round reads "Cleared" because our RPC replica is behind — the settle // transaction is already on the network, and the cursor proves it. diff --git a/services/keeper/src/index.ts b/services/keeper/src/index.ts index 65711cbb..871a002f 100644 --- a/services/keeper/src/index.ts +++ b/services/keeper/src/index.ts @@ -108,3 +108,9 @@ export { type ResumeCheckpointParams, type RunWatchLoopParams, } from "./watch-loop.js"; +export { + DEFAULT_SUBMISSION_JOURNAL_PATH, + FileSubmissionJournal, + SubmissionJournalBindingError, + type FileSubmissionJournalOptions, +} from "./submission-journal.js"; diff --git a/services/keeper/src/keeper.ts b/services/keeper/src/keeper.ts index c193e8fc..baccdda0 100644 --- a/services/keeper/src/keeper.ts +++ b/services/keeper/src/keeper.ts @@ -36,8 +36,6 @@ import { type TimeContext, } from "@sub-rosa/time"; -import type { SettlementGuard } from "./settlement-guard.js"; - export type KeeperLogger = (msg: string) => void; export interface KeeperDeps { @@ -127,6 +125,39 @@ function recordStep( deps.checkpoint?.markComplete(roundId, step, transactionHash); } +/** Resolve a durable in-flight hash before the caller builds another tx. */ +async function recoverBeforeRetry( + deps: KeeperDeps, + roundId: bigint, + operation: string, + step: KeeperStep, + discriminator?: string, +): Promise { + if (typeof deps.sdk.reconcileSubmission !== "function") return false; + const recovered = await deps.sdk.reconcileSubmission({ + operation, + roundId: roundId.toString(), + ...(discriminator ? { discriminator } : {}), + }); + if (!recovered) return false; + if (recovered.state === "pending") { + throw new Error(`round ${roundId} ${operation} transaction ${recovered.hash} is still pending`); + } + if (recovered.state === "failed") { + throw new Error(`round ${roundId} ${operation} transaction ${recovered.hash} definitively failed`); + } + if (recovered.state === "confirmed") { + // Each reveal journal entry is per bidder; the aggregate cursor is only + // complete after every bidder has been checked in this pass. + if (step !== "reveal" || !discriminator) { + recordStep(deps, roundId, step, recovered.hash); + } + return true; + } + // An expired hash is terminal; exactly this pass may now build one replacement. + return false; +} + /** Wait until Drand round R should be published. Returns false if R is still in * the future after `maxWaitSeconds`. */ export async function waitForRound( @@ -195,20 +226,25 @@ export async function keepRound( return result; } - try { - await sdk.openReveal(rid, signature); - result.openedReveal = true; - log(`open_reveal OK (round ${rid} via Drand R=${R})`); - } catch (e) { - if (errorMatches(e, IDEMPOTENT_OPEN)) { - log(`open_reveal already done (${errorName(e)}); continuing`); - } else { - throw e; + const recoveredOpen = await recoverBeforeRetry(deps, rid, "open_reveal", "open-reveal"); + if (recoveredOpen) { + log(`open_reveal transaction recovered as confirmed for round ${rid}`); + } else { + try { + await sdk.openReveal(rid, signature); + result.openedReveal = true; + log(`open_reveal OK (round ${rid} via Drand R=${R})`); + } catch (e) { + if (errorMatches(e, IDEMPOTENT_OPEN)) { + log(`open_reveal already done (${errorName(e)}); continuing`); + } else { + throw e; + } } } // The reveal window is open from here on (we opened it, or the contract told // us it already was) — advance the cursor before revealing any bid. - recordStep(deps, rid, "open-reveal"); + if (!recoveredOpen) recordStep(deps, rid, "open-reveal"); round = await sdk.getRound(rid); } @@ -225,6 +261,12 @@ export async function keepRound( log(`revealing ${bidders.length} bidder(s)`); for (const bidder of bidders) { + const recoveredReveal = await recoverBeforeRetry(deps, rid, "reveal", "reveal", bidder); + if (recoveredReveal) { + result.revealed.push(bidder); + log(`reveal transaction recovered as confirmed for ${bidder}`); + continue; + } let state; try { state = await sdk.getBidState(rid, bidder); @@ -343,25 +385,31 @@ export async function closeRound( result.finalStatus = round.status.tag; return result; } - try { - const winner = await sdk.clear(rid); + const recoveredClear = await recoverBeforeRetry(deps, rid, "clear", "clear"); + if (recoveredClear) { result.cleared = true; - result.winner = winner; - if (winner === undefined) { - result.voided = true; - log(`cleared → no valid bids; round voided + refunded`); - } else { - log(`cleared → winner ${winner}`); - } - recordStep(deps, rid, "clear"); - } catch (e) { - if (errorMatches(e, ["AlreadyCleared", "RevealStillOpen", "WrongStatus", "RoundVoided"])) { - result.skipped.push(`clear skipped: ${errorName(e)}`); - if (errorMatches(e, ["AlreadyCleared"])) { - recordStep(deps, rid, "clear"); + log(`clear transaction recovered as confirmed for round ${rid}`); + } else { + try { + const winner = await sdk.clear(rid); + result.cleared = true; + result.winner = winner; + if (winner === undefined) { + result.voided = true; + log(`cleared → no valid bids; round voided + refunded`); + } else { + log(`cleared → winner ${winner}`); + } + recordStep(deps, rid, "clear"); + } catch (e) { + if (errorMatches(e, ["AlreadyCleared", "RevealStillOpen", "WrongStatus", "RoundVoided"])) { + result.skipped.push(`clear skipped: ${errorName(e)}`); + if (errorMatches(e, ["AlreadyCleared"])) { + recordStep(deps, rid, "clear"); + } + } else { + throw e; } - } else { - throw e; } } } @@ -376,21 +424,27 @@ export async function closeRound( log(`settle skipped: checkpoint records the round as settled`); result.skipped.push("settle already complete (checkpoint)"); } else { - try { - await sdk.settle(rid); + const recoveredSettle = await recoverBeforeRetry(deps, rid, "settle", "settle"); + if (recoveredSettle) { result.settled = true; - log(`settled round ${rid}`); - recordStep(deps, rid, "settle"); - } catch (e) { - if (errorMatches(e, ["AlreadySettled", "NotCleared", "WrongStatus"])) { - result.skipped.push(`settle skipped: ${errorName(e)}`); - // The contract telling us it is already settled is proof the step landed — - // record it so the next restart does not broadcast it again. - if (errorMatches(e, ["AlreadySettled"])) { - recordStep(deps, rid, "settle"); + log(`settle transaction recovered as confirmed for round ${rid}`); + } else { + try { + await sdk.settle(rid); + result.settled = true; + log(`settled round ${rid}`); + recordStep(deps, rid, "settle"); + } catch (e) { + if (errorMatches(e, ["AlreadySettled", "NotCleared", "WrongStatus"])) { + result.skipped.push(`settle skipped: ${errorName(e)}`); + // The contract telling us it is already settled is proof the step landed — + // record it so the next restart does not broadcast it again. + if (errorMatches(e, ["AlreadySettled"])) { + recordStep(deps, rid, "settle"); + } + } else { + throw e; } - } else { - throw e; } } } @@ -512,6 +566,16 @@ export async function voidIfStale( return result; } + const recoveredVoid = await recoverBeforeRetry(deps, rid, "void", "void"); + if (recoveredVoid) { + result.voided = true; + guard?.markTerminal(rid, "voided on-chain"); + log(`void transaction recovered as confirmed for round ${rid}`); + const after = await sdk.getRound(rid); + result.finalStatus = after.status.tag; + return result; + } + // ── A void is on the table: verify the refund set before dispatching ─── if (guard) { const view = await readSettlementView(sdk, rid, now); diff --git a/services/keeper/src/run.ts b/services/keeper/src/run.ts index f91a3c23..390accd3 100644 --- a/services/keeper/src/run.ts +++ b/services/keeper/src/run.ts @@ -28,6 +28,7 @@ import { parseKeeperRunConfig, } from "./dry-run.js"; import { keepRound } from "./keeper.js"; +import { FileSubmissionJournal } from "./submission-journal.js"; async function main() { const config = parseKeeperRunConfig(); @@ -64,19 +65,23 @@ async function main() { return; } - const sdk = new SubRosaClient({ - rpcUrl: config.rpcUrl, - networkPassphrase: config.networkPassphrase, - contractId: config.contractId, - secretKey: config.keeperSecret!, - }); - // Throws KeeperCheckpointMismatchError when the on-disk cursor was recorded // for another network or contract — better to stop than to replay it. const checkpoint = new KeeperCheckpointStore({ network: config.networkPassphrase, contractId: config.contractId, }); + const submissionJournal = new FileSubmissionJournal({ + network: config.networkPassphrase, + contractId: config.contractId, + }); + const sdk = new SubRosaClient({ + rpcUrl: config.rpcUrl, + networkPassphrase: config.networkPassphrase, + contractId: config.contractId, + secretKey: config.keeperSecret!, + submissionJournal, + }); const result = await keepRound( { diff --git a/services/keeper/src/serve.ts b/services/keeper/src/serve.ts index 6a9fa9a2..d91626eb 100644 --- a/services/keeper/src/serve.ts +++ b/services/keeper/src/serve.ts @@ -30,6 +30,7 @@ import { SubRosaClient } from "@sub-rosa/sdk"; import { quicknet } from "@sub-rosa/tlock"; import { KeeperCheckpointStore } from "./checkpoint.js"; +import { FileSubmissionJournal } from "./submission-journal.js"; import { createSettlementGuard } from "./settlement-guard.js"; import { createStatusServer, withGracefulShutdown } from "./status-server.js"; import { KeeperStore } from "./store.js"; @@ -50,11 +51,14 @@ async function main() { process.env.NETWORK_PASSPHRASE ?? "Test SDF Network ; September 2015"; const keeperSecret = reqEnv("KEEPER_SECRET"); + const checkpoint = new KeeperCheckpointStore({ network: networkPassphrase, contractId }); + const submissionJournal = new FileSubmissionJournal({ network: networkPassphrase, contractId }); const sdk = new SubRosaClient({ rpcUrl, networkPassphrase, contractId, secretKey: keeperSecret, + submissionJournal, }); const reader = new SubRosaClient({ rpcUrl, @@ -68,7 +72,6 @@ async function main() { const store = new KeeperStore(); // Durable watch cursor. Refuses to start when the file on disk was recorded // for another network or contract id. - const checkpoint = new KeeperCheckpointStore({ network: networkPassphrase, contractId }); const settlementGuard = createSettlementGuard(); const queue = new KeeperQueue(store, { contractId, network: networkPassphrase }); diff --git a/services/keeper/src/submission-journal.ts b/services/keeper/src/submission-journal.ts new file mode 100644 index 00000000..837607d9 --- /dev/null +++ b/services/keeper/src/submission-journal.ts @@ -0,0 +1,110 @@ +// Copyright (c) 2026 Sub Rosa contributors +// Durable recovery journal shared with the SDK submit path. +import * as fs from "node:fs"; +import * as path from "node:path"; +import type { + SubmissionIdentity, + SubmissionJournal, + SubmissionRecord, +} from "@sub-rosa/sdk"; +import { submissionKey } from "@sub-rosa/sdk"; + +export const DEFAULT_SUBMISSION_JOURNAL_PATH = ".keeper-submissions.json"; + +interface JournalFile { + version: 1; + network: string; + contractId: string; + records: Record; +} + +export class SubmissionJournalBindingError extends Error { + constructor(field: "network" | "contractId", expected: string, actual: string) { + super(`submission journal ${field} ${JSON.stringify(actual)} does not match configured ${JSON.stringify(expected)}`); + this.name = "SubmissionJournalBindingError"; + } +} + +export interface FileSubmissionJournalOptions { + network: string; + contractId: string; + path?: string; +} + +/** File-backed journal; records hashes and operation metadata only, never XDR or keys. */ +export class FileSubmissionJournal implements SubmissionJournal { + readonly filePath: string; + private readonly network: string; + private readonly contractId: string; + private data: JournalFile; + + constructor(options: FileSubmissionJournalOptions) { + this.filePath = options.path ?? process.env.KEEPER_SUBMISSION_PATH ?? DEFAULT_SUBMISSION_JOURNAL_PATH; + this.network = options.network; + this.contractId = options.contractId; + this.data = this.load(); + } + + private load(): JournalFile { + if (!fs.existsSync(this.filePath)) { + return { version: 1, network: this.network, contractId: this.contractId, records: {} }; + } + const parsed = JSON.parse(fs.readFileSync(this.filePath, "utf8")) as Partial; + if (parsed.network !== this.network) { + throw new SubmissionJournalBindingError("network", this.network, String(parsed.network ?? "")); + } + if (parsed.contractId !== this.contractId) { + throw new SubmissionJournalBindingError("contractId", this.contractId, String(parsed.contractId ?? "")); + } + if (parsed.version !== 1 || !parsed.records || typeof parsed.records !== "object" || Array.isArray(parsed.records)) { + throw new Error(`unsupported or malformed submission journal at ${this.filePath}`); + } + const records: Record = {}; + for (const [key, value] of Object.entries(parsed.records)) { + if (!value || typeof value !== "object" || Array.isArray(value)) { + throw new Error(`malformed submission record at ${this.filePath}`); + } + const record = value as SubmissionRecord; + if ( + typeof record.operation !== "string" || + !(typeof record.roundId === "string" || record.roundId === null) || + typeof record.hash !== "string" || !record.hash || + !["pending", "confirmed", "expired", "failed"].includes(record.state) || + !Number.isFinite(record.submittedAtMs) || + !Number.isFinite(record.expiresAtMs) || + record.network !== this.network || + record.contractId !== this.contractId || + key !== submissionKey(record) + ) { + throw new Error(`malformed or mismatched submission record at ${this.filePath}`); + } + records[key] = { ...record }; + } + return { version: 1, network: this.network, contractId: this.contractId, records }; + } + + async get(identity: SubmissionIdentity): Promise { + this.assertBinding(identity); + const record = this.data.records[submissionKey(identity)]; + return record ? { ...record } : undefined; + } + + async put(record: SubmissionRecord): Promise { + this.assertBinding(record); + this.data.records[submissionKey(record)] = { ...record }; + const directory = path.dirname(this.filePath); + if (directory !== ".") fs.mkdirSync(directory, { recursive: true }); + const temporaryPath = `${this.filePath}.${process.pid}.tmp`; + fs.writeFileSync(temporaryPath, JSON.stringify(this.data, null, 2), { encoding: "utf8", mode: 0o600 }); + fs.renameSync(temporaryPath, this.filePath); + } + + private assertBinding(identity: SubmissionIdentity): void { + if (identity.network !== this.network) { + throw new SubmissionJournalBindingError("network", this.network, identity.network); + } + if (identity.contractId !== this.contractId) { + throw new SubmissionJournalBindingError("contractId", this.contractId, identity.contractId); + } + } +}