diff --git a/.env.example b/.env.example index 91e22de7..d56f9b64 100644 --- a/.env.example +++ b/.env.example @@ -49,6 +49,8 @@ X402_APPRAISAL_URL= # ROUND_CONTRACT_ID=C… # KEEPER_DRY_RUN=true # ROUND_ID=1 +# KEEPER_CHECKPOINT_PATH=.keeper-checkpoint.json # durable watch cursor +# KEEPER_STORE_PATH=.keeper-store.json # watched round queue # WATCH_POLL_MS=15000 # WATCH_ROUND_IDS=1,2,5 # WATCH_FROM=1 diff --git a/.gitignore b/.gitignore index f1a41d3e..d4efb572 100644 --- a/.gitignore +++ b/.gitignore @@ -35,6 +35,8 @@ deployments/*.local.json # Keeper .keeper-store.json .keeper-store.json.corrupted.* +.keeper-checkpoint.json +.keeper-checkpoint.json.corrupted.* # Coverage coverage/ diff --git a/docs/THREAT_MODEL.md b/docs/THREAT_MODEL.md index ff7312d5..18f8a368 100644 --- a/docs/THREAT_MODEL.md +++ b/docs/THREAT_MODEL.md @@ -46,7 +46,7 @@ | --- | --- | --- | | Winner doesn't pay | Escrow locked at commit; settle pulls from escrow | Requires valid reveal | | Drand never delivers R | `void` after grace refunds all escrow | Grace window must be configured | -| Double settle | Idempotent settle skips settled bids; keeper watchers take an exclusive per-round lease so only one process reveals or settles | Proven in e2e | +| Double settle | Idempotent settle skips settled bids; keeper watch checkpoint records completed steps (hash-verified on startup) so a restart cannot rebroadcast | Proven in e2e | ### Identity privacy diff --git a/services/keeper/README.md b/services/keeper/README.md index 5d0fc2e2..ddff1e13 100644 --- a/services/keeper/README.md +++ b/services/keeper/README.md @@ -116,6 +116,53 @@ Response shape (typed in `@sub-rosa/sdk` as `KeeperStatusResponse`): 6. **Failure states are visible, not hidden.** A round whose on-chain lookup fails is surfaced with `status: "Unknown"` or `status: "NotFound"` and the `lastError` field populated. The process does not crash on upstream errors. 7. **Secrets in responses: none.** The status API never emits secret keys, signed transactions, or bidder private data. Bidder *addresses* (public on-chain identifiers) are included so dashboards can show bidder counts. +## Watch Checkpoint (restart safety) + +The in-memory settlement guard is lost on restart. The **watch checkpoint** is the durable version: a small local JSON file (default `.keeper-checkpoint.json`, override with `KEEPER_CHECKPOINT_PATH`) that records how far each round got, so a crash after a confirmed reveal, clear, or settle does not broadcast that step again. + +### Checkpoint format + +```json +{ + "version": 1, + "network": "Test SDF Network ; September 2015", + "contractId": "C...", + "rounds": { + "1": { + "roundId": "1", + "completedSteps": ["open-reveal", "reveal", "clear", "settle"], + "lastCompletedStep": "settle", + "lastTransactionHash": "0x…", + "stepHashes": { "settle": "0x…" }, + "updatedAt": "2026-09-30T00:00:00.000Z" + } + } +} +``` + +| Field | Meaning | +|-------|---------| +| `network` / `contractId` | The deployment the cursor was recorded for. | +| `completedSteps` | Steps observed complete, in completion order. | +| `lastCompletedStep` | The cursor position — the last step this process finished. | +| `lastTransactionHash` | Transaction hash of the last step, when the SDK exposes one. | +| `stepHashes` | Per-step hashes, re-verified on every startup. | +| `updatedAt` | ISO-8601 timestamp of the last write. | + +Steps tracked: `open-reveal`, `reveal`, `clear`, `settle`, `void`. + +### Startup validation + +1. **Binding.** If `network` or `contractId` on disk does not match the process configuration, the keeper refuses to start (`KeeperCheckpointMismatchError`) instead of replaying a cursor from another deployment. Point `KEEPER_CHECKPOINT_PATH` at a per-deployment file, or delete the file when you switch networks. +2. **Hash verification.** Every recorded `stepHashes` entry is re-checked (when a verifier is wired in). A hash that comes back `failed` or `missing` is rolled back so the step is retried; a `confirmed` hash is trusted even if the RPC replica still reports the pre-step status. +3. **Chain reconciliation.** Cursor entries with no hash to verify are checked against the on-chain status. If the chain cannot prove the step happened, the entry is dropped rather than stranding the round. + +`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. + +### 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`). + ## Persisted Queue / Store ## Persisted Queue / Store diff --git a/services/keeper/package.json b/services/keeper/package.json index 71205d79..8e4d9669 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/watch-lease.test.ts src/settlement-guard.test.ts src/status.test.ts src/status-server.test.ts src/queue-replay.test.ts", + "test": "node --import tsx --test src/checkpoint.test.ts src/checkpoint-restart.test.ts 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", "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 new file mode 100644 index 00000000..297c752d --- /dev/null +++ b/services/keeper/src/checkpoint-restart.test.ts @@ -0,0 +1,548 @@ +// Copyright (c) 2026 Sub Rosa contributors +// checkpoint-restart.test.ts +// +// Restart-resilience coverage for the keeper watch cursor. +// +// Two crash points are exercised for every irreversible step: +// +// crash BEFORE the checkpoint write — the step is retried on restart, the +// contract's idempotent rejection downgrades it to a skip, and the keeper +// records it so the *next* restart is quiet. +// +// crash AFTER the checkpoint write — the step is not re-broadcast, even when +// the RPC replica the keeper reads is still behind (a "lagging replica" +// Cleared status must not re-settle a settled round). +// +// Every scenario runs on a fake clock and an in-memory chain: no RPC, no Drand, +// no signing material. + +import assert from "node:assert/strict"; +import * as fs from "node:fs"; +import * as os from "node:os"; +import * as path from "node:path"; +import { afterEach, beforeEach, describe, test } from "node:test"; + +import type { Round, SubRosaClient } from "@sub-rosa/sdk"; +import { createFakeTime, type FakeClock, type FakeScheduler } from "@sub-rosa/time"; + +import { + CHECKPOINT_VERSION, + KeeperCheckpointStore, + type KeeperCheckpointFile, + type KeeperStep, + type TransactionHashStatus, +} from "./checkpoint.js"; +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"; + +const NETWORK = "Test SDF Network ; September 2015"; +const CONTRACT = "CTESTCONTRACT"; +/** 2023-11-14T22:13:20Z — comfortably past the fixture reveal deadline. */ +const NOW_MS = 1_700_000_000_000; +const REVEAL_DEADLINE = 1_699_900_000n; +const ROUND_ID = 1n; + +const silentLogger = { + debug: () => {}, + info: () => {}, + warn: () => {}, + error: () => {}, +}; + +let dir: string; +let file: string; +let clock: FakeClock; +let scheduler: FakeScheduler; + +beforeEach(() => { + dir = fs.mkdtempSync(path.join(os.tmpdir(), "keeper-restart-")); + file = path.join(dir, "checkpoint.json"); + ({ clock, scheduler } = createFakeTime(NOW_MS)); +}); + +afterEach(() => { + fs.rmSync(dir, { recursive: true, force: true }); +}); + +// ── In-memory chain ──────────────────────────────────────────────────────── + +interface FakeChainOptions { + status: string; + bidders?: string[]; + /** Keep the reported status fixed even after mutations (lagging replica). */ + stickyStatus?: boolean; +} + +/** + * Minimal stand-in for `SubRosaClient`. Mutation methods record every + * broadcast so a test can assert that a step was *not* re-submitted, and + * `settleOnce`-style flags model the contract's idempotent rejections. + */ +class FakeChain { + calls: string[] = []; + status: string; + bidderList: string[]; + private readonly stickyStatus: boolean; + + constructor(options: FakeChainOptions) { + this.status = options.status; + this.bidderList = options.bidders ?? []; + this.stickyStatus = options.stickyStatus ?? false; + } + + private advance(status: string): void { + if (!this.stickyStatus) this.status = status; + } + + async getRound(roundId: bigint | number = ROUND_ID): Promise { + this.calls.push("getRound"); + if (BigInt(roundId) !== ROUND_ID) throw new Error("HostError: RoundNotFound(1)"); + return { + auditor_pubkey: Buffer.alloc(32), + bidders: this.bidderList, + clearing_rule: { tag: "HighestBid", values: undefined }, + commit_deadline: REVEAL_DEADLINE - 3_600n, + item_ref: Buffer.alloc(32), + operator: "GOPERATOR", + reveal_deadline: REVEAL_DEADLINE, + reveal_round: 42n, + status: { tag: this.status, values: undefined }, + winner: undefined, + winning_bid: 0n, + } as unknown as Round; + } + + async *bidders(): AsyncGenerator { + this.calls.push("bidders"); + for (const bidder of this.bidderList) yield bidder; + } + + async getBidState(): Promise { + this.calls.push("getBidState"); + throw new Error("getBidState must not run when the cursor records reveals"); + } + + async getSeal(): Promise { + this.calls.push("getSeal"); + throw new Error("getSeal must not run when the cursor records reveals"); + } + + async openReveal(): Promise { + this.calls.push("openReveal"); + this.advance("Revealing"); + } + + async reveal(): Promise { + this.calls.push("reveal"); + } + + async clear(): Promise { + this.calls.push("clear"); + this.advance("Cleared"); + return "GWINNER"; + } + + async settle(): Promise { + this.calls.push("settle"); + this.advance("Settled"); + } + + async void(): Promise { + this.calls.push("void"); + this.advance("Voided"); + } + + count(call: string): number { + return this.calls.filter((c) => c === call).length; + } +} + +function asSdk(chain: FakeChain): SubRosaClient { + return chain as unknown as SubRosaClient; +} + +function newStore(options: { dryRun?: boolean } = {}): KeeperCheckpointStore { + return new KeeperCheckpointStore({ + path: file, + network: NETWORK, + contractId: CONTRACT, + clock, + logger: silentLogger, + ...options, + }); +} + +/** Seed a cursor as a previous process would have left it on disk. */ +function seedCheckpoint( + steps: KeeperStep[], + hashes: Partial> = {}, +): void { + const rounds: KeeperCheckpointFile["rounds"] = {}; + for (const [index, step] of steps.entries()) { + rounds[ROUND_ID.toString()] = { + roundId: ROUND_ID.toString(), + completedSteps: steps.slice(0, index + 1), + lastCompletedStep: step, + lastTransactionHash: hashes[steps[steps.length - 1]] ?? null, + stepHashes: Object.fromEntries( + steps.slice(0, index + 1).flatMap((s) => (hashes[s] ? [[s, hashes[s]]] : [])), + ), + updatedAt: clock.toISOString(), + }; + } + fs.writeFileSync( + file, + JSON.stringify( + { version: CHECKPOINT_VERSION, network: NETWORK, contractId: CONTRACT, rounds }, + null, + 2, + ), + "utf-8", + ); +} + +function readCheckpoint(): KeeperCheckpointFile { + return JSON.parse(fs.readFileSync(file, "utf-8")) as KeeperCheckpointFile; +} + +const confirmAll = async (): Promise => "confirmed"; + +// ── Crash AFTER the checkpoint write ─────────────────────────────────────── + +describe("restart after the checkpoint was written", () => { + 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. + const chain = new FakeChain({ status: "Cleared" }); + seedCheckpoint(["open-reveal", "reveal", "clear", "settle"], { + settle: "0xsethash", + }); + const checkpoint = newStore(); + + const log: string[] = []; + await resumeCheckpoint({ + checkpoint, + sdk: asSdk(chain), + log: (m) => log.push(m), + verifyTransaction: confirmAll, + }); + + const tick = await watchRound( + { + sdk: asSdk(chain), + drand: {} as never, + log: () => {}, + time: { clock, scheduler }, + checkpoint, + }, + ROUND_ID, + ); + + assert.equal(chain.count("settle"), 0, "settle must not be re-broadcast"); + assert.ok(log.some((line) => line.includes("confirmed settle for round 1"))); + assert.ok( + tick.close?.skipped.includes("settle already complete (checkpoint)"), + `expected a checkpoint skip, got ${JSON.stringify(tick.close?.skipped)}`, + ); + assert.equal(tick.close?.settled, false); + assert.equal(tick.finalStatus, "Cleared"); + }); + + test("a confirmed clear is not submitted again", async () => { + const chain = new FakeChain({ status: "Revealing" }); + seedCheckpoint(["open-reveal", "reveal", "clear"], { clear: "0xclearhash" }); + const checkpoint = newStore(); + await checkpoint.verifyHashes(confirmAll); + + const result = await closeRound( + { + sdk: asSdk(chain), + drand: {} as never, + log: () => {}, + time: { clock, scheduler }, + checkpoint, + }, + ROUND_ID, + ); + + assert.equal(chain.count("clear"), 0, "clear must not be re-broadcast"); + assert.ok(result.skipped.includes("clear already complete (checkpoint)")); + assert.equal(result.cleared, false); + }); + + test("a confirmed void is not submitted again", async () => { + const chain = new FakeChain({ status: "Open" }); + seedCheckpoint(["void"], { void: "0xvoidhash" }); + const checkpoint = newStore(); + await checkpoint.verifyHashes(confirmAll); + + const result = await voidIfStale( + { + sdk: asSdk(chain), + drand: {} as never, + log: () => {}, + time: { clock, scheduler }, + checkpoint, + }, + ROUND_ID, + ); + + assert.equal(chain.count("void"), 0, "void must not be re-broadcast"); + assert.ok(result.skipped.includes("void already complete (checkpoint)")); + assert.equal(result.voided, false); + }); + + test("recorded reveals skip the per-bidder pass entirely", async () => { + const chain = new FakeChain({ status: "Revealing", bidders: ["G1", "G2"] }); + seedCheckpoint(["open-reveal", "reveal"], { reveal: "0xrevealhash" }); + const checkpoint = newStore(); + await checkpoint.verifyHashes(confirmAll); + + const result = await keepRound( + { + sdk: asSdk(chain), + drand: {} as never, + log: () => {}, + time: { clock, scheduler }, + checkpoint, + }, + ROUND_ID, + ); + + assert.equal(chain.count("bidders"), 0, "no seal reads on a resumed cursor"); + assert.equal(chain.count("getBidState"), 0); + assert.equal(chain.count("reveal"), 0); + assert.deepEqual(result.revealed, []); + assert.ok( + result.skipped.some((s) => s.reason === "reveals complete (checkpoint)"), + ); + }); +}); + +// ── Crash BEFORE the checkpoint write ────────────────────────────────────── + +describe("restart before the checkpoint was written", () => { + test("the step is retried, downgraded to a skip, and then recorded", async () => { + const chain = new FakeChain({ status: "Cleared", stickyStatus: true }); + // The contract is idempotent: a second settle is rejected, not executed. + const sdk = { + ...chain, + getRound: () => chain.getRound(), + settle: async () => { + chain.calls.push("settle"); + if (chain.count("settle") > 1) throw new Error("HostError: AlreadySettled(1)"); + }, + } as unknown as SubRosaClient; + + // Pass 1: settle lands, then the process dies before the checkpoint write. + const crashing = newStore({ dryRun: true }); + const first = await closeRound( + { sdk, drand: {} as never, log: () => {}, time: { clock, scheduler }, checkpoint: crashing }, + ROUND_ID, + ); + assert.equal(first.settled, true); + assert.equal(chain.count("settle"), 1); + assert.equal(fs.existsSync(file), false, "crash-before leaves no checkpoint"); + + // Pass 2 (restart): the cursor knows nothing, so settle is broadcast again + // and the contract's idempotent rejection is recorded as completion. + const restarted = newStore(); + const second = await closeRound( + { sdk, drand: {} as never, log: () => {}, time: { clock, scheduler }, checkpoint: restarted }, + ROUND_ID, + ); + assert.equal(second.settled, false); + assert.equal(chain.count("settle"), 2); + assert.ok( + second.skipped.some((line) => line.includes("AlreadySettled")), + `expected an idempotent skip, got ${JSON.stringify(second.skipped)}`, + ); + assert.equal(restarted.isComplete(ROUND_ID, "settle"), true); + assert.deepEqual(readCheckpoint().rounds["1"].completedSteps, ["settle"]); + + // Pass 3 (restart again): the cursor now stops the broadcast for good. + const third = newStore(); + await third.verifyHashes(confirmAll); + await closeRound( + { sdk, drand: {} as never, log: () => {}, time: { clock, scheduler }, checkpoint: third }, + ROUND_ID, + ); + assert.equal(chain.count("settle"), 2, "settle is broadcast at most twice"); + }); + + test("a lost cursor is rebuilt from the chain, not trusted blindly", async () => { + // A cursor written before the crash, with no transaction hash to confirm, on + // a replica that says the round never settled: retry instead of stranding. + const chain = new FakeChain({ status: "Cleared" }); + seedCheckpoint(["settle"]); + const checkpoint = newStore(); + const log: string[] = []; + + await resumeCheckpoint({ + checkpoint, + sdk: asSdk(chain), + log: (m) => log.push(m), + verifyTransaction: confirmAll, + }); + + assert.equal(checkpoint.isComplete(ROUND_ID, "settle"), false); + assert.ok(log.some((line) => line.includes("retrying settle"))); + + await closeRound( + { sdk: asSdk(chain), drand: {} as never, log: () => {}, time: { clock, scheduler }, checkpoint }, + ROUND_ID, + ); + assert.equal(chain.count("settle"), 1, "the unprovable step is retried"); + }); + + test("a hash that failed on the network is retried on the next pass", async () => { + const chain = new FakeChain({ status: "Cleared" }); + seedCheckpoint(["clear", "settle"], { clear: "0xclearhash", settle: "0xsethash" }); + const checkpoint = newStore(); + + const verifications = await checkpoint.verifyHashes(async (hash) => + hash === "0xsethash" ? "failed" : "confirmed", + ); + + assert.deepEqual( + verifications.map((v) => [v.step, v.retained]), + [ + ["clear", true], + ["settle", false], + ], + ); + assert.equal(checkpoint.isComplete(ROUND_ID, "settle"), false); + assert.equal(checkpoint.isComplete(ROUND_ID, "clear"), true); + + await closeRound( + { sdk: asSdk(chain), drand: {} as never, log: () => {}, time: { clock, scheduler }, checkpoint }, + ROUND_ID, + ); + assert.equal(chain.count("settle"), 1, "the failed step is retried"); + }); + + test("an unreadable round during reconciliation keeps the cursor", async () => { + seedCheckpoint(["settle"], { settle: "0xsethash" }); + const checkpoint = newStore(); + const log: string[] = []; + + await resumeCheckpoint({ + checkpoint, + sdk: { + getRound: async () => { + throw new Error("HostError: ConnectionError"); + }, + } as unknown as SubRosaClient, + log: (m) => log.push(m), + verifyTransaction: confirmAll, + }); + + assert.equal(checkpoint.isComplete(ROUND_ID, "settle"), true); + assert.ok(log.some((line) => line.includes("could not read round 1"))); + }); +}); + +// ── Recording ────────────────────────────────────────────────────────────── + +describe("cursor recording", () => { + test("a cleared round records reveal, clear, and settle with their hashes", async () => { + const chain = new FakeChain({ status: "Revealing" }); + const checkpoint = newStore(); + + const result = await watchRound( + { + sdk: asSdk(chain), + drand: {} as never, + log: () => {}, + time: { clock, scheduler }, + checkpoint, + }, + ROUND_ID, + ); + + assert.equal(result.close?.cleared, true); + assert.equal(result.close?.settled, true); + const written = readCheckpoint().rounds["1"]; + assert.deepEqual(written.completedSteps, ["reveal", "clear", "settle"]); + assert.equal(written.lastCompletedStep, "settle"); + assert.equal(written.lastTransactionHash, null, "the SDK returns no tx hash"); + assert.equal(written.updatedAt, clock.toISOString()); + }); + + test("a round with no bids still records the reveal step", async () => { + const chain = new FakeChain({ status: "Revealing", bidders: [] }); + const checkpoint = newStore(); + + await keepRound( + { + sdk: asSdk(chain), + drand: {} as never, + log: () => {}, + time: { clock, scheduler }, + checkpoint, + }, + ROUND_ID, + ); + + assert.equal(checkpoint.isComplete(ROUND_ID, "reveal"), true); + assert.equal(chain.count("reveal"), 0); + }); + + test("a stale round records the void step", async () => { + const chain = new FakeChain({ status: "Open" }); + const checkpoint = newStore(); + + const result = await voidIfStale( + { + sdk: asSdk(chain), + drand: {} as never, + log: () => {}, + time: { clock, scheduler }, + checkpoint, + }, + ROUND_ID, + ); + + assert.equal(result.voided, true); + assert.equal(checkpoint.isComplete(ROUND_ID, "void"), true); + }); +}); + +// ── Watch loop integration ───────────────────────────────────────────────── + +describe("watch loop restart", () => { + test("a resumed watch loop does not settle a round the cursor already settled", async () => { + const chain = new FakeChain({ status: "Cleared" }); + seedCheckpoint(["open-reveal", "reveal", "clear", "settle"], { + settle: "0xsethash", + }); + const store = new KeeperStore(path.join(dir, "queue.json"), silentLogger); + store.addRound(ROUND_ID, { contractId: CONTRACT, network: NETWORK, lastStatus: "Cleared" }); + + let stopChecks = 0; + const logs: string[] = []; + + await runWatchLoop({ + sdk: asSdk(chain), + drand: {} as never, + log: (m) => logs.push(m), + pollMs: 0, + contractId: CONTRACT, + network: NETWORK, + store, + settlementGuard: createSettlementGuard(clock), + checkpoint: newStore(), + verifyTransaction: confirmAll, + isStopping: () => stopChecks++ >= 3, + time: { clock, scheduler }, + }); + + assert.equal(chain.count("settle"), 0, "the watch loop must not resettle"); + assert.ok( + logs.some((line) => line.includes("confirmed settle for round 1")), + `expected a confirmed-settle log line, got ${JSON.stringify(logs)}`, + ); + }); +}); diff --git a/services/keeper/src/checkpoint.test.ts b/services/keeper/src/checkpoint.test.ts new file mode 100644 index 00000000..e3d43a7c --- /dev/null +++ b/services/keeper/src/checkpoint.test.ts @@ -0,0 +1,424 @@ +// Copyright (c) 2026 Sub Rosa contributors +// checkpoint.test.ts +// +// Unit coverage for the durable watch cursor: on-disk schema, startup binding +// validation, transaction-hash verification, chain reconciliation, and the +// no-write dry-run mode. + +import assert from "node:assert/strict"; +import * as fs from "node:fs"; +import * as os from "node:os"; +import * as path from "node:path"; +import { afterEach, beforeEach, describe, test } from "node:test"; + +import { createFakeTime } from "@sub-rosa/time"; + +import { + CHECKPOINT_SKIP_STEPS, + CHECKPOINT_VERSION, + KeeperCheckpointError, + KeeperCheckpointMismatchError, + KeeperCheckpointStore, + checkpointBindingMismatch, + checkpointHasStep, + planCheckpointRollback, + planCheckpointStep, + readCheckpointFile, + type KeeperCheckpointFile, + type KeeperStep, +} from "./checkpoint.js"; + +const NETWORK = "Test SDF Network ; September 2015"; +const OTHER_NETWORK = "Public Global Stellar Network ; September 2015"; +const CONTRACT = "CTESTCONTRACT"; +const OTHER_CONTRACT = "COTHERCONTRACT"; +const START_MS = Date.parse("2026-09-30T00:00:00.000Z"); + +const silentLogger = { + debug: () => {}, + info: () => {}, + warn: () => {}, + error: () => {}, +}; + +let dir: string; +let file: string; + +function makeStore( + overrides: Partial[0]> = {}, +): KeeperCheckpointStore { + return new KeeperCheckpointStore({ + path: file, + network: NETWORK, + contractId: CONTRACT, + clock: createFakeTime(START_MS).clock, + logger: silentLogger, + ...overrides, + }); +} + +function readFile(): KeeperCheckpointFile { + return JSON.parse(fs.readFileSync(file, "utf-8")) as KeeperCheckpointFile; +} + +function writeFile(contents: unknown): void { + fs.writeFileSync(file, JSON.stringify(contents, null, 2), "utf-8"); +} + +beforeEach(() => { + dir = fs.mkdtempSync(path.join(os.tmpdir(), "keeper-checkpoint-")); + file = path.join(dir, "checkpoint.json"); +}); + +afterEach(() => { + fs.rmSync(dir, { recursive: true, force: true }); +}); + +describe("KeeperCheckpointStore persistence", () => { + test("persists round id, last completed step, transaction hash, and network", () => { + const store = makeStore(); + store.markComplete(7n, "open-reveal"); + store.markComplete(7n, "reveal"); + store.markComplete(7n, "clear"); + store.markComplete(7n, "settle", "0xsethash"); + + const written = readFile(); + assert.equal(written.version, CHECKPOINT_VERSION); + assert.equal(written.network, NETWORK); + assert.equal(written.contractId, CONTRACT); + assert.deepEqual(written.rounds["7"], { + roundId: "7", + completedSteps: ["open-reveal", "reveal", "clear", "settle"], + lastCompletedStep: "settle", + lastTransactionHash: "0xsethash", + stepHashes: { settle: "0xsethash" }, + updatedAt: "2026-09-30T00:00:00.000Z", + }); + }); + + test("resumes the cursor across a restart", () => { + makeStore().markComplete(3n, "clear"); + + const resumed = makeStore(); + assert.equal(resumed.isComplete(3n, "clear"), true); + assert.equal(resumed.isComplete(3n, "settle"), false); + assert.equal(resumed.isComplete(4n, "clear"), false); + assert.deepEqual(resumed.listRoundIds(), [3n]); + assert.equal(resumed.read(3n)?.lastCompletedStep, "clear"); + }); + + test("records a step once even when it is replayed", () => { + const store = makeStore(); + store.markComplete(1n, "settle", "0xabc"); + store.markComplete(1n, "settle", "0xabc"); + assert.deepEqual(store.read(1n)?.completedSteps, ["settle"]); + assert.deepEqual(store.read(1n)?.stepHashes, { settle: "0xabc" }); + }); + + test("orders multiple round cursors numerically", () => { + const store = makeStore(); + for (const id of [10n, 2n, 7n]) store.markComplete(id, "void"); + assert.deepEqual(store.listRoundIds(), [2n, 7n, 10n]); + }); + + test("starts empty when no file exists yet", () => { + const store = makeStore(); + assert.deepEqual(store.listRoundIds(), []); + assert.equal(fs.existsSync(file), false); + }); + + test("backs up a corrupted file instead of replaying garbage", () => { + fs.writeFileSync(file, "{ not json", "utf-8"); + const store = makeStore(); + assert.deepEqual(store.listRoundIds(), []); + store.markComplete(1n, "settle"); + const backups = fs.readdirSync(dir).filter((f) => f.includes(".corrupted.")); + assert.equal(backups.length, 1); + }); + + test("drops malformed round entries but keeps valid ones", () => { + writeFile({ + version: CHECKPOINT_VERSION, + network: NETWORK, + contractId: CONTRACT, + rounds: { + "1": { + roundId: "1", + completedSteps: ["settle", "not-a-step"], + lastCompletedStep: "settle", + stepHashes: { settle: 42, void: "0xvoid" }, + }, + "2": "nonsense", + "3": null, + }, + }); + + const store = makeStore(); + assert.deepEqual(store.listRoundIds(), [1n]); + assert.deepEqual(store.read(1n)?.completedSteps, ["settle"]); + assert.deepEqual(store.read(1n)?.stepHashes, { void: "0xvoid" }); + }); +}); + +describe("KeeperCheckpointStore startup validation", () => { + test("refuses to start when the checkpoint contract id does not match", () => { + writeFile({ + version: CHECKPOINT_VERSION, + network: NETWORK, + contractId: OTHER_CONTRACT, + rounds: { + "1": { + roundId: "1", + completedSteps: ["settle"], + lastCompletedStep: "settle", + lastTransactionHash: null, + stepHashes: {}, + updatedAt: "2026-09-30T00:00:00.000Z", + }, + }, + }); + + assert.throws( + () => makeStore(), + (error: unknown) => { + assert.ok(error instanceof KeeperCheckpointMismatchError); + assert.equal(error.field, "contractId"); + assert.equal(error.expected, CONTRACT); + assert.equal(error.actual, OTHER_CONTRACT); + assert.match(error.message, /refusing to start/); + return true; + }, + ); + }); + + test("refuses to start when the checkpoint network does not match", () => { + writeFile({ + version: CHECKPOINT_VERSION, + network: OTHER_NETWORK, + contractId: CONTRACT, + rounds: {}, + }); + + assert.throws( + () => makeStore(), + (error: unknown) => { + assert.ok(error instanceof KeeperCheckpointMismatchError); + assert.equal(error.field, "network"); + return true; + }, + ); + }); + + test("refuses to start on an unsupported schema version", () => { + writeFile({ version: 99, network: NETWORK, contractId: CONTRACT, rounds: {} }); + assert.throws(() => makeStore(), KeeperCheckpointError); + }); + + test("resumes when the binding matches", () => { + writeFile({ + version: CHECKPOINT_VERSION, + network: NETWORK, + contractId: CONTRACT, + rounds: { + "5": { + roundId: "5", + completedSteps: ["settle"], + lastCompletedStep: "settle", + lastTransactionHash: "0xdead", + stepHashes: { settle: "0xdead" }, + updatedAt: "2026-09-30T00:00:00.000Z", + }, + }, + }); + assert.equal(makeStore().isComplete(5n, "settle"), true); + }); +}); + +describe("KeeperCheckpointStore hash verification", () => { + test("keeps steps whose transaction hash is confirmed", async () => { + const store = makeStore(); + store.markComplete(1n, "settle", "0xgood"); + + const results = await store.verifyHashes(async () => "confirmed"); + + assert.deepEqual(results, [ + { roundId: "1", step: "settle", transactionHash: "0xgood", status: "confirmed", retained: true }, + ]); + assert.equal(store.isComplete(1n, "settle"), true); + }); + + test("rolls back a step whose transaction failed", async () => { + const store = makeStore(); + store.markComplete(1n, "clear", "0xbad"); + store.markComplete(1n, "settle", "0xgood"); + + const results = await store.verifyHashes(async (hash) => + hash === "0xbad" ? "failed" : "confirmed", + ); + + assert.equal(results.find((r) => r.step === "clear")?.retained, false); + assert.equal(store.isComplete(1n, "clear"), false); + assert.equal(store.isComplete(1n, "settle"), true); + assert.equal(store.read(1n)?.lastCompletedStep, "settle"); + assert.deepEqual(readFile().rounds["1"].completedSteps, ["settle"]); + }); + + test("rolls back a step whose transaction the network never saw", async () => { + const store = makeStore(); + store.markComplete(2n, "settle", "0xghost"); + + const results = await store.verifyHashes(async () => "missing"); + + assert.equal(results[0].status, "missing"); + assert.equal(results[0].retained, false); + assert.equal(store.isComplete(2n, "settle"), false); + assert.deepEqual(store.read(2n)?.completedSteps, []); + assert.equal(store.read(2n)?.lastCompletedStep, null); + assert.equal(store.read(2n)?.lastTransactionHash, null); + }); + + test("is a no-op when no verifier is available", async () => { + const store = makeStore(); + store.markComplete(1n, "settle", "0xgood"); + assert.deepEqual(await store.verifyHashes(), []); + assert.equal(store.isComplete(1n, "settle"), true); + }); + + test("keeps the cursor when the hash lookup itself fails", async () => { + const store = makeStore(); + store.markComplete(1n, "settle", "0xgood"); + + const results = await store.verifyHashes(async () => { + throw new Error("rpc unreachable"); + }); + + assert.equal(results[0].retained, true); + assert.equal(store.isComplete(1n, "settle"), true); + }); +}); + +describe("KeeperCheckpointStore chain reconciliation", () => { + test("drops hashless steps the on-chain status cannot prove", () => { + const store = makeStore(); + store.markComplete(1n, "settle"); + + const dropped = store.reconcile(1n, "Cleared"); + + assert.deepEqual(dropped, ["settle"]); + assert.equal(store.isComplete(1n, "settle"), false); + }); + + test("keeps hashless steps the on-chain status proves", () => { + const store = makeStore(); + store.markComplete(1n, "clear"); + + assert.deepEqual(store.reconcile(1n, "Settled"), []); + assert.equal(store.isComplete(1n, "clear"), true); + }); + + test("keeps hash-confirmed steps even when the replica lags behind", () => { + const store = makeStore(); + store.markComplete(1n, "settle", "0xgood"); + + assert.deepEqual(store.reconcile(1n, "Cleared"), []); + assert.equal(store.isComplete(1n, "settle"), true); + }); + + test("ignores rounds without cursor state", () => { + const store = makeStore(); + assert.deepEqual(store.reconcile(9n, "Settled"), []); + }); +}); + +describe("KeeperCheckpointStore dry-run mode", () => { + test("tracks progress in memory and never writes to disk", () => { + const store = makeStore({ dryRun: true }); + store.markComplete(1n, "settle", "0xgood"); + + assert.equal(store.isComplete(1n, "settle"), true); + assert.equal(fs.existsSync(file), false, "dry-run must not create the file"); + + store.markComplete(1n, "clear"); + assert.deepEqual(store.read(1n)?.completedSteps, ["settle", "clear"]); + assert.equal(fs.existsSync(file), false); + }); +}); + +describe("checkpoint helpers", () => { + test("checkpointHasStep reads the cursor", () => { + const entry = planCheckpointStep( + { version: 1, network: NETWORK, contractId: CONTRACT, rounds: {} }, + 1n, + "void", + ).rounds["1"]; + assert.equal(checkpointHasStep(entry, "void"), true); + assert.equal(checkpointHasStep(entry, "settle"), false); + assert.equal(checkpointHasStep(undefined, "void"), false); + }); + + test("checkpointBindingMismatch reports the conflicting field", () => { + const file2 = { network: NETWORK, contractId: CONTRACT }; + assert.equal( + checkpointBindingMismatch(file2, { network: NETWORK, contractId: CONTRACT }), + null, + ); + assert.equal( + checkpointBindingMismatch(file2, { network: OTHER_NETWORK, contractId: CONTRACT }), + "network", + ); + assert.equal( + checkpointBindingMismatch(file2, { network: NETWORK, contractId: OTHER_CONTRACT }), + "contractId", + ); + }); + + test("planCheckpointStep is pure", () => { + const base: KeeperCheckpointFile = { + version: 1, + network: NETWORK, + contractId: CONTRACT, + rounds: {}, + }; + const next = planCheckpointStep(base, 1n, "clear", { + transactionHash: "0xabc", + at: "2026-09-30T00:00:00.000Z", + }); + assert.deepEqual(base.rounds, {}, "input file must not be mutated"); + assert.equal(next.rounds["1"].lastTransactionHash, "0xabc"); + }); + + test("planCheckpointRollback keeps the remaining cursor intact", () => { + let file2: KeeperCheckpointFile = { + version: 1, + network: NETWORK, + contractId: CONTRACT, + rounds: {}, + }; + file2 = planCheckpointStep(file2, 1n, "clear", { transactionHash: "0xclear" }); + file2 = planCheckpointStep(file2, 1n, "settle", { transactionHash: "0xsettle" }); + const rolled = planCheckpointRollback(file2, 1n, "settle"); + + assert.deepEqual(rolled.rounds["1"].completedSteps, ["clear"]); + assert.equal(rolled.rounds["1"].lastCompletedStep, "clear"); + assert.equal(rolled.rounds["1"].lastTransactionHash, "0xclear"); + assert.deepEqual(rolled.rounds["1"].stepHashes, { clear: "0xclear" }); + // The original object is untouched. + assert.deepEqual(file2.rounds["1"].completedSteps, ["clear", "settle"]); + }); + + test("readCheckpointFile never throws on malformed input", () => { + assert.equal(readCheckpointFile(path.join(dir, "absent.json")), undefined); + fs.writeFileSync(file, "{ nope", "utf-8"); + assert.equal(readCheckpointFile(file, silentLogger), undefined); + writeFile({ version: 1, network: NETWORK, contractId: CONTRACT }); + assert.equal(readCheckpointFile(file, silentLogger)?.rounds, undefined); + }); + + test("open-reveal is not skip-eligible, the irreversible steps are", () => { + const steps: KeeperStep[] = ["open-reveal", "reveal", "clear", "settle", "void"]; + assert.ok(!CHECKPOINT_SKIP_STEPS.includes("open-reveal")); + for (const step of steps.filter((s) => s !== "open-reveal")) { + assert.ok(CHECKPOINT_SKIP_STEPS.includes(step), `${step} should be skip-eligible`); + } + }); +}); diff --git a/services/keeper/src/checkpoint.ts b/services/keeper/src/checkpoint.ts new file mode 100644 index 00000000..57086ced --- /dev/null +++ b/services/keeper/src/checkpoint.ts @@ -0,0 +1,588 @@ +import { normalizeError } from "@sub-rosa/logging/errors"; +// Copyright (c) 2026 Sub Rosa contributors +// checkpoint.ts +// +// Durable watch cursor for the keeper. +// +// The in-memory settlement guard is not enough to survive a restart: if the +// process dies after broadcasting a settle (or clear) transaction but before it +// re-reads the chain, the next process has no idea the step is already in +// flight and happily broadcasts it again. This module persists that cursor to a +// small local JSON file so a restart cannot resettle. +// +// File format (one file per keeper process, bound to one contract + network): +// +// { +// "version": 1, +// "network": "Test SDF Network ; September 2015", +// "contractId": "C...", +// "rounds": { +// "1": { +// "roundId": "1", +// "completedSteps": ["open-reveal", "reveal", "clear", "settle"], +// "lastCompletedStep": "settle", +// "lastTransactionHash": "0x…", +// "stepHashes": { "settle": "0x…" }, +// "updatedAt": "2026-09-30T00:00:00.000Z" +// } +// } +// } +// +// Safety rules, in order of precedence: +// +// 1. Binding. The file records the network and contract id it was written +// for. A mismatch with the process config refuses to start — replaying a +// cursor recorded against a different contract or network is worse than +// not having a cursor at all. +// 2. Hash verification. On startup every recorded transaction hash is +// re-checked. A step whose hash is `failed` or `missing` is rolled back so +// the step is retried; a `confirmed` hash is trusted even when the RPC +// replica the keeper reads is still behind. +// 3. Chain reconciliation. Steps with no transaction hash to verify fall back +// to the on-chain status: if the round never reached the state that proves +// the step happened, the cursor entry is dropped and the step is retried. +// +// Every mutation is a pure function (`planCheckpointStep`) so the dry-run +// planner can print exactly what a live run would write without touching disk. + +import { createLogger, type Logger } from "@sub-rosa/logging"; +import * as fs from "fs"; +import * as path from "path"; +import { systemClock, type Clock } from "@sub-rosa/time"; + +import { normalizeRoundId, type RoundIdInput } from "./store.js"; + +const diagnostics = createLogger("services.keeper.src.checkpoint"); + +/** Schema version of the on-disk checkpoint file. */ +export const CHECKPOINT_VERSION = 1; + +/** Default checkpoint path, overridable with `KEEPER_CHECKPOINT_PATH`. */ +export const DEFAULT_CHECKPOINT_PATH = ".keeper-checkpoint.json"; + +/** + * Keeper lifecycle steps tracked by the cursor. + * + * `open-reveal` is recorded for auditability but is never used to skip work: + * whether the reveal window is open is already authoritative on-chain (the round + * status is `Open` or `Revealing`), so trusting the cursor there could strand a + * round whose opening transaction never landed. + */ +export type KeeperStep = "open-reveal" | "reveal" | "clear" | "settle" | "void"; + +/** Steps whose completion is durable in the checkpoint and must not be + * re-broadcast after a restart. `open-reveal` is deliberately excluded. */ +export const CHECKPOINT_SKIP_STEPS: readonly KeeperStep[] = [ + "reveal", + "clear", + "settle", + "void", +]; + +/** + * On-chain statuses that prove each step actually took effect. Used to roll + * back cursor entries that carry no verifiable transaction hash. + */ +export const STEP_SATISFIED_BY_STATUS: Record = { + "open-reveal": ["Revealing", "Cleared", "Settled", "Voided"], + reveal: ["Revealing", "Cleared", "Settled", "Voided"], + clear: ["Cleared", "Settled", "Voided"], + settle: ["Settled", "Voided"], + void: ["Voided"], +}; + +export interface KeeperCheckpoint { + roundId: string; + /** Steps observed complete, in the order the keeper completed them. */ + completedSteps: KeeperStep[]; + /** Most recently completed step — the "last completed step" of the cursor. */ + lastCompletedStep: KeeperStep | null; + /** Transaction hash of the last completed step, when one was available. */ + lastTransactionHash: string | null; + /** Per-step transaction hashes, used to re-verify the cursor on startup. */ + stepHashes: Partial>; + /** ISO-8601 timestamp of the last checkpoint write for this round. */ + updatedAt: string; +} + +export interface KeeperCheckpointFile { + version: number; + network: string; + contractId: string; + rounds: Record; +} + +/** Result of re-checking a recorded transaction hash. */ +export type TransactionHashStatus = "confirmed" | "failed" | "missing"; + +/** Injectable hash lookup (an RPC `getTransaction` wrapper in production). */ +export type TransactionHashVerifier = ( + transactionHash: string, +) => Promise; + +export interface CheckpointVerification { + roundId: string; + step: KeeperStep; + transactionHash: string; + status: TransactionHashStatus; + /** False when the step was rolled back because the hash did not confirm. */ + retained: boolean; +} + +/** Cursor surface consumed by the keeper phases. */ +export interface WatchCheckpoint { + isComplete(roundId: bigint | number, step: KeeperStep): boolean; + markComplete( + roundId: bigint | number, + step: KeeperStep, + transactionHash?: string | null, + ): void; + read(roundId: bigint | number): KeeperCheckpoint | undefined; +} + +/** Cursor plus the startup resume surface used by the watch loop. */ +export interface ResumableCheckpoint extends WatchCheckpoint { + listRoundIds(): bigint[]; + verifyHashes(verify?: TransactionHashVerifier): Promise; + reconcile(roundId: bigint | number, observedStatus: string): KeeperStep[]; +} + +/** Base class for checkpoint failures that must stop the process. */ +export class KeeperCheckpointError extends Error { + constructor(message: string) { + super(message); + this.name = "KeeperCheckpointError"; + } +} + +/** Raised when the checkpoint file was written for another contract/network. */ +export class KeeperCheckpointMismatchError extends KeeperCheckpointError { + constructor( + readonly field: "network" | "contractId", + readonly expected: string, + readonly actual: string, + ) { + super( + `refusing to start: checkpoint ${field} ${JSON.stringify(actual)} does not match process config ${JSON.stringify(expected)}`, + ); + this.name = "KeeperCheckpointMismatchError"; + } +} + +/** + * Compare a checkpoint file's binding against the process configuration. + * Returns the mismatching field, or `null` when the cursor may be resumed. + */ +export function checkpointBindingMismatch( + file: Pick, + config: { network: string; contractId: string }, +): "network" | "contractId" | null { + if (file.network !== config.network) return "network"; + if (file.contractId !== config.contractId) return "contractId"; + return null; +} + +export interface KeeperCheckpointStoreOptions { + /** File path. Defaults to `KEEPER_CHECKPOINT_PATH`, then `.keeper-checkpoint.json`. */ + path?: string; + network: string; + contractId: string; + clock?: Clock; + logger?: Logger; + /** Track progress in memory only — never write to disk (dry-run mode). */ + dryRun?: boolean; +} + +function emptyFile( + network: string, + contractId: string, +): KeeperCheckpointFile { + return { version: CHECKPOINT_VERSION, network, contractId, rounds: {} }; +} + +/** Placeholder for entries written before a timestamp was recorded. */ +const UNKNOWN_TIMESTAMP = "1970-01-01T00:00:00.000Z"; + +function isKeeperStep(value: unknown): value is KeeperStep { + return ( + typeof value === "string" && + Object.prototype.hasOwnProperty.call(STEP_SATISFIED_BY_STATUS, value) + ); +} + +/** Parse one round entry, dropping malformed cursor data. */ +function parseCheckpointEntry( + roundId: string, + value: unknown, + logger: Logger, +): KeeperCheckpoint | undefined { + if (!value || typeof value !== "object" || Array.isArray(value)) { + logger.warn( + "checkpoint-dropping-malformed-entry", + `[Checkpoint] Dropping malformed entry for round ${roundId}: expected an object`, + ); + return undefined; + } + const stored = value as Partial; + const steps = Array.isArray(stored.completedSteps) + ? stored.completedSteps.filter(isKeeperStep) + : []; + const stepHashes: Partial> = {}; + if (stored.stepHashes && typeof stored.stepHashes === "object") { + for (const [step, hash] of Object.entries(stored.stepHashes)) { + if (isKeeperStep(step) && typeof hash === "string" && hash) { + stepHashes[step] = hash; + } + } + } + const lastCompletedStep = + steps.length > 0 ? steps[steps.length - 1] : null; + return { + roundId, + completedSteps: steps, + lastCompletedStep, + lastTransactionHash: + typeof stored.lastTransactionHash === "string" + ? stored.lastTransactionHash + : null, + stepHashes, + updatedAt: + typeof stored.updatedAt === "string" + ? stored.updatedAt + : UNKNOWN_TIMESTAMP, + }; +} + +/** + * Pure cursor update. Returns a new file object; never mutates the input and + * never touches the filesystem, so the dry-run planner can reuse it verbatim. + */ +export function planCheckpointStep( + file: KeeperCheckpointFile, + roundId: RoundIdInput, + step: KeeperStep, + options: { transactionHash?: string | null; at?: string } = {}, +): KeeperCheckpointFile { + const id = normalizeRoundId(roundId); + const transactionHash = options.transactionHash ?? null; + const previous = file.rounds[id]; + const completedSteps = previous + ? [...previous.completedSteps] + : []; + if (!completedSteps.includes(step)) completedSteps.push(step); + + const stepHashes: Partial> = { + ...(previous?.stepHashes ?? {}), + }; + if (transactionHash) stepHashes[step] = transactionHash; + else delete stepHashes[step]; + + const updatedAt = options.at ?? previous?.updatedAt ?? UNKNOWN_TIMESTAMP; + const entry: KeeperCheckpoint = { + roundId: id, + completedSteps, + lastCompletedStep: step, + lastTransactionHash: transactionHash, + stepHashes, + updatedAt, + }; + + return { + version: file.version || CHECKPOINT_VERSION, + network: file.network, + contractId: file.contractId, + rounds: { ...file.rounds, [id]: entry }, + }; +} + +/** Drop a step from the cursor (used when its transaction hash fails). */ +export function planCheckpointRollback( + file: KeeperCheckpointFile, + roundId: RoundIdInput, + step: KeeperStep, +): KeeperCheckpointFile { + const id = normalizeRoundId(roundId); + const previous = file.rounds[id]; + if (!previous || !previous.completedSteps.includes(step)) return file; + const completedSteps = previous.completedSteps.filter((s) => s !== step); + const stepHashes = { ...previous.stepHashes }; + delete stepHashes[step]; + const entry: KeeperCheckpoint = { + ...previous, + completedSteps, + stepHashes, + lastCompletedStep: + completedSteps.length > 0 ? completedSteps[completedSteps.length - 1] : null, + lastTransactionHash: + completedSteps.length > 0 + ? (stepHashes[completedSteps[completedSteps.length - 1]] ?? null) + : null, + }; + return { + ...file, + rounds: { ...file.rounds, [id]: entry }, + }; +} + +/** True when the cursor records `step` as complete for the round. */ +export function checkpointHasStep( + checkpoint: KeeperCheckpoint | undefined, + step: KeeperStep, +): boolean { + return checkpoint?.completedSteps.includes(step) ?? false; +} + +/** + * Read a checkpoint file without validating it against any process config. + * Returns `undefined` when the file is missing, empty, or unparseable — used by + * the dry-run planner, which must never mutate state. + */ +export function readCheckpointFile( + filePath: string, + logger: Logger = diagnostics, +): KeeperCheckpointFile | undefined { + if (!fs.existsSync(filePath)) return undefined; + try { + const content = fs.readFileSync(filePath, "utf-8"); + if (!content.trim()) return undefined; + const parsed = JSON.parse(content) as Partial; + if (!parsed.rounds || typeof parsed.rounds !== "object") return undefined; + const rounds: Record = {}; + for (const [key, value] of Object.entries(parsed.rounds)) { + const entry = parseCheckpointEntry(key, value, logger); + if (entry) rounds[key] = entry; + } + return { + version: typeof parsed.version === "number" ? parsed.version : CHECKPOINT_VERSION, + network: typeof parsed.network === "string" ? parsed.network : "", + contractId: typeof parsed.contractId === "string" ? parsed.contractId : "", + rounds, + }; + } catch (e) { + logger.warn( + "checkpoint-failed-to-parse", + `[Checkpoint] Failed to parse ${filePath}: ${normalizeError(e).message}`, + ); + return undefined; + } +} + +/** + * Durable, per-round watch cursor. + * + * The constructor refuses to load a file whose network or contract id does not + * match the process configuration — a cursor recorded against a different + * deployment is meaningless and unsafe to replay. + */ +export class KeeperCheckpointStore implements ResumableCheckpoint { + readonly filePath: string; + readonly dryRun: boolean; + private readonly clock: Clock; + private readonly logger: Logger; + private readonly network: string; + private readonly contractId: string; + private data: KeeperCheckpointFile; + + constructor(options: KeeperCheckpointStoreOptions) { + this.filePath = + options.path || + process.env.KEEPER_CHECKPOINT_PATH || + DEFAULT_CHECKPOINT_PATH; + this.network = options.network; + this.contractId = options.contractId; + this.clock = options.clock ?? systemClock; + this.logger = options.logger ?? diagnostics; + this.dryRun = options.dryRun ?? false; + this.data = this.load(); + } + + /** Bind check against the process config, then normalize the round entries. */ + private load(): KeeperCheckpointFile { + if (!fs.existsSync(this.filePath)) { + return emptyFile(this.network, this.contractId); + } + + let parsed: KeeperCheckpointFile; + try { + const content = fs.readFileSync(this.filePath, "utf-8"); + parsed = content.trim() ? (JSON.parse(content) as KeeperCheckpointFile) : emptyFile(this.network, this.contractId); + } catch (e) { + this.logger.warn( + "checkpoint-failed-to-parse", + `[Checkpoint] Failed to parse ${this.filePath}. Backing up the corrupted file and starting fresh.`, + ); + try { + fs.renameSync( + this.filePath, + `${this.filePath}.corrupted.${this.clock.nowMs()}`, + ); + } catch (backupErr) { + this.logger.error( + "checkpoint-could-not-backup-corrupted-file", + "[Checkpoint] Could not back up the corrupted checkpoint file:", + { "backupErr_0": normalizeError(backupErr) }, + ); + } + return emptyFile(this.network, this.contractId); + } + + if (parsed.network !== this.network) { + throw new KeeperCheckpointMismatchError( + "network", + this.network, + parsed.network, + ); + } + if (parsed.contractId !== this.contractId) { + throw new KeeperCheckpointMismatchError( + "contractId", + this.contractId, + parsed.contractId, + ); + } + if (parsed.version !== CHECKPOINT_VERSION) { + throw new KeeperCheckpointError( + `unsupported checkpoint version ${JSON.stringify(parsed.version)} in ${this.filePath} (expected ${CHECKPOINT_VERSION})`, + ); + } + this.logger.info( + "checkpoint-resumed", + `[Checkpoint] resuming ${this.filePath}: ${Object.keys(parsed.rounds ?? {}).length} round cursor(s)`, + ); + + const rounds: Record = {}; + for (const [key, value] of Object.entries(parsed.rounds ?? {})) { + const entry = parseCheckpointEntry(key, value, this.logger); + if (entry) rounds[key] = entry; + } + return { version: CHECKPOINT_VERSION, network: this.network, contractId: this.contractId, rounds }; + } + + private save(): void { + if (this.dryRun) { + this.logger.info( + "checkpoint-dry-run-no-write", + `[Checkpoint] dry-run: not writing ${this.filePath}`, + ); + return; + } + try { + const dir = path.dirname(this.filePath); + if (dir !== ".") fs.mkdirSync(dir, { recursive: true }); + fs.writeFileSync(this.filePath, JSON.stringify(this.data, null, 2), "utf-8"); + } catch (e) { + this.logger.error( + "checkpoint-failed-to-save", + `[Checkpoint] Failed to save checkpoint to ${this.filePath}:`, + { "e_0": normalizeError(e) }, + ); + } + } + + isComplete(roundId: bigint | number, step: KeeperStep): boolean { + return checkpointHasStep(this.data.rounds[normalizeRoundId(roundId)], step); + } + + markComplete( + roundId: bigint | number, + step: KeeperStep, + transactionHash?: string | null, + ): void { + this.data = planCheckpointStep(this.data, roundId, step, { + transactionHash, + at: this.clock.toISOString(), + }); + this.save(); + } + + read(roundId: bigint | number): KeeperCheckpoint | undefined { + return this.data.rounds[normalizeRoundId(roundId)]; + } + + /** Round ids that carry cursor state, ascending. */ + listRoundIds(): bigint[] { + return Object.keys(this.data.rounds) + .sort((a, b) => (BigInt(a) < BigInt(b) ? -1 : BigInt(a) > BigInt(b) ? 1 : 0)) + .map((id) => BigInt(id)); + } + + /** Snapshot of the file exactly as a live run would write it. */ + snapshot(): KeeperCheckpointFile { + return JSON.parse(JSON.stringify(this.data)) as KeeperCheckpointFile; + } + + /** + * Re-verify every recorded transaction hash. Steps whose hash is not + * `confirmed` are rolled back so the next tick retries them. With no + * verifier available the recorded hashes are trusted as-is and nothing is + * rolled back — hash confirmation is a gate, never a requirement to run. + */ + async verifyHashes( + verify?: TransactionHashVerifier, + ): Promise { + const results: CheckpointVerification[] = []; + if (!verify) return results; + + let changed = false; + for (const [id, entry] of Object.entries(this.data.rounds)) { + for (const [step, hash] of Object.entries(entry.stepHashes)) { + if (!isKeeperStep(step)) continue; + let status: TransactionHashStatus; + try { + status = await verify(hash); + } catch (e) { + // An unreachable RPC must not discard durable progress: treat it as + // "unverifiable" and keep the cursor. + this.logger.warn( + "checkpoint-hash-lookup-failed", + `[Checkpoint] hash lookup failed for round ${id} step ${step}: ${normalizeError(e).message}`, + ); + results.push({ roundId: id, step, transactionHash: hash, status: "missing", retained: true }); + continue; + } + const retained = status === "confirmed"; + results.push({ roundId: id, step, transactionHash: hash, status, retained }); + if (!retained) { + this.data = planCheckpointRollback(this.data, id, step); + changed = true; + this.logger.warn( + "checkpoint-rolled-back-step", + `[Checkpoint] rolling back round ${id} step ${step}: transaction ${hash} is ${status}`, + ); + } + } + } + if (changed) this.save(); + return results; + } + + /** + * Chain reconciliation for cursor entries with no transaction hash to verify. + * If the observed on-chain status cannot prove a recorded step happened, the + * step is dropped so the keeper retries it instead of stranding the round. + */ + reconcile( + roundId: bigint | number, + observedStatus: string, + ): KeeperStep[] { + const id = normalizeRoundId(roundId); + const entry = this.data.rounds[id]; + if (!entry) return []; + const dropped: KeeperStep[] = []; + for (const step of entry.completedSteps) { + if (entry.stepHashes[step]) continue; // hash-confirmed: trust the cursor + if (STEP_SATISFIED_BY_STATUS[step].includes(observedStatus)) continue; + dropped.push(step); + } + if (dropped.length === 0) return dropped; + for (const step of dropped) { + this.data = planCheckpointRollback(this.data, roundId, step); + } + this.save(); + this.logger.warn( + "checkpoint-reconciled-with-chain", + `[Checkpoint] round ${roundId} is ${observedStatus}; dropped unverified cursor steps: ${dropped.join(", ")}`, + ); + return dropped; + } +} diff --git a/services/keeper/src/dry-run.test.ts b/services/keeper/src/dry-run.test.ts index 2119c2cc..4cd0d7ce 100644 --- a/services/keeper/src/dry-run.test.ts +++ b/services/keeper/src/dry-run.test.ts @@ -1,13 +1,20 @@ // Copyright (c) 2026 Sub Rosa contributors import assert from "node:assert/strict"; +import * as fs from "node:fs"; +import * as os from "node:os"; +import * as path from "node:path"; import { describe, test } from "node:test"; import type { BidState, Round } from "@sub-rosa/sdk"; +import { createFakeTime } from "@sub-rosa/time"; +import { KeeperCheckpointStore } from "./checkpoint.js"; import { + DRY_RUN_PHASE_STEP, buildKeeperDryRunSummary, decideKeeperDryRunAction, parseKeeperRunConfig, + planDryRunCheckpoint, type KeeperDryRunReader, } from "./dry-run.js"; @@ -204,6 +211,12 @@ describe("buildKeeperDryRunSummary", () => { sdk as KeeperDryRunReader, 7n, 500, + { + checkpointPath: "/tmp/dry-run-checkpoint.json", + network: "Test SDF Network ; September 2015", + contractId: CONTRACT_ID, + nowIso: "2026-09-30T00:00:00.000Z", + }, ); assert.deepEqual(summary, { @@ -216,6 +229,29 @@ describe("buildKeeperDryRunSummary", () => { currentPhase: "revealing", nextAction: "reveal 1 pending bidder", transactionsSubmitted: 0, + checkpoint: { + path: "/tmp/dry-run-checkpoint.json", + network: "Test SDF Network ; September 2015", + contractId: CONTRACT_ID, + proposedStep: "reveal", + mismatch: null, + filesWritten: 0, + proposedFile: { + version: 1, + network: "Test SDF Network ; September 2015", + contractId: CONTRACT_ID, + rounds: { + "7": { + roundId: "7", + completedSteps: ["reveal"], + lastCompletedStep: "reveal", + lastTransactionHash: null, + stepHashes: {}, + updatedAt: "2026-09-30T00:00:00.000Z", + }, + }, + }, + }, }); assert.deepEqual(reads, ["round:7", "bid:7:G1", "bid:7:G2"]); assert.equal(mutations, 0); @@ -245,3 +281,116 @@ describe("buildKeeperDryRunSummary", () => { assert.equal(summary.transactionsSubmitted, 0); }); }); + +describe("dry-run checkpoint preview", () => { + test("writes nothing to the checkpoint path and submits nothing", async () => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), "keeper-dry-run-")); + const checkpointPath = path.join(dir, "checkpoint.json"); + const mutations: string[] = []; + const reader = { + async getRound() { + return { ...baseRound, status: { tag: "Cleared", values: undefined } }; + }, + async getBidState() { + return bidState(true); + }, + async settle() { + mutations.push("settle"); + }, + async clear() { + mutations.push("clear"); + }, + } as unknown as KeeperDryRunReader; + + try { + const summary = await buildKeeperDryRunSummary(reader, 3n, 5_000, { + checkpointPath, + network: "Test SDF Network ; September 2015", + contractId: CONTRACT_ID, + nowIso: "2026-09-30T00:00:00.000Z", + }); + + assert.equal(summary.checkpoint.proposedStep, "settle"); + assert.equal(summary.checkpoint.filesWritten, 0); + assert.equal(summary.transactionsSubmitted, 0); + assert.equal(fs.existsSync(checkpointPath), false, "dry-run must not create the file"); + assert.deepEqual(mutations, []); + assert.deepEqual(fs.readdirSync(dir), []); + } finally { + fs.rmSync(dir, { recursive: true, force: true }); + } + }); + + test("proposes the same file content the live keeper store would write", () => { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), "keeper-dry-run-")); + const checkpointPath = path.join(dir, "checkpoint.json"); + const options = { + checkpointPath, + network: "Test SDF Network ; September 2015", + contractId: CONTRACT_ID, + nowIso: "2026-09-30T00:00:00.000Z", + }; + + try { + const preview = planDryRunCheckpoint(5n, "clear", options); + + const store = new KeeperCheckpointStore({ + path: checkpointPath, + network: options.network, + contractId: options.contractId, + clock: createFakeTime(Date.parse(options.nowIso)).clock, + }); + store.markComplete(5n, "clear"); + + assert.deepEqual(preview.proposedFile, store.snapshot()); + } finally { + fs.rmSync(dir, { recursive: true, force: true }); + } + }); + + test("reports a binding conflict that would stop a live keeper", () => { + const current = { + version: 1 as const, + network: "Public Global Stellar Network ; September 2015", + contractId: CONTRACT_ID, + rounds: {}, + }; + + assert.equal( + planDryRunCheckpoint(1n, "settle", { + network: "Test SDF Network ; September 2015", + contractId: CONTRACT_ID, + currentCheckpoint: current, + }).mismatch, + "network", + ); + assert.equal( + planDryRunCheckpoint(1n, "settle", { + network: current.network, + contractId: "CDIFFERENT", + currentCheckpoint: current, + }).mismatch, + "contractId", + ); + assert.equal( + planDryRunCheckpoint(1n, "settle", { + network: current.network, + contractId: current.contractId, + currentCheckpoint: current, + }).mismatch, + null, + ); + }); + + test("proposes no checkpoint change for terminal phases", () => { + for (const phase of ["awaiting-clear", "complete"] as const) { + assert.equal(DRY_RUN_PHASE_STEP[phase], null); + const preview = planDryRunCheckpoint(1n, DRY_RUN_PHASE_STEP[phase], { + network: "Test SDF Network ; September 2015", + contractId: CONTRACT_ID, + }); + assert.equal(preview.proposedStep, null); + assert.deepEqual(preview.proposedFile.rounds, {}); + } + }); +}); diff --git a/services/keeper/src/dry-run.ts b/services/keeper/src/dry-run.ts index 4097b26d..6a2ae050 100644 --- a/services/keeper/src/dry-run.ts +++ b/services/keeper/src/dry-run.ts @@ -3,6 +3,14 @@ import type { BidState, Round, SubRosaClient } from "@sub-rosa/sdk"; import { systemClock } from "@sub-rosa/time"; import { VOID_GRACE_SECONDS } from "./keeper.js"; +import { + CHECKPOINT_VERSION, + DEFAULT_CHECKPOINT_PATH, + checkpointBindingMismatch, + planCheckpointStep, + type KeeperCheckpointFile, + type KeeperStep, +} from "./checkpoint.js"; const DEFAULT_RPC_URL = "https://soroban-testnet.stellar.org"; const DEFAULT_NETWORK_PASSPHRASE = "Test SDF Network ; September 2015"; @@ -26,6 +34,17 @@ export type KeeperDryRunPhase = | "ready-to-settle" | "complete"; +/** Cursor step each dry-run phase would record on the next live pass. */ +export const DRY_RUN_PHASE_STEP: Record = { + "awaiting-drand": "open-reveal", + "stale-open": "void", + revealing: "reveal", + "awaiting-clear": null, + "ready-to-clear": "clear", + "ready-to-settle": "settle", + complete: null, +}; + export interface KeeperDryRunDecision { currentPhase: KeeperDryRunPhase; nextAction: string; @@ -39,6 +58,8 @@ export interface KeeperDryRunSummary extends KeeperDryRunDecision { bidderCount: number; revealedCount: number | null; transactionsSubmitted: 0; + /** Checkpoint preview — never written to disk by a dry run. */ + checkpoint: KeeperDryRunCheckpoint; } export type KeeperDryRunReader = Pick< @@ -46,6 +67,35 @@ export type KeeperDryRunReader = Pick< "getRound" | "getBidState" >; +/** + * The checkpoint a live run *would* write for the next step, plus the binding + * conflict that would stop it. Nothing here is persisted: dry-run neither + * submits a transaction nor touches the checkpoint file. + */ +export interface KeeperDryRunCheckpoint { + path: string; + network: string; + contractId: string; + /** Step the next live pass would record, or null when nothing is pending. */ + proposedStep: KeeperStep | null; + /** Field whose stored value conflicts with the process config, if any. */ + mismatch: "network" | "contractId" | null; + /** Exact file content a live run would write. */ + proposedFile: KeeperCheckpointFile; + /** Always 0 — dry-run writes nothing. */ + filesWritten: 0; +} + +export interface KeeperDryRunOptions { + checkpointPath?: string; + network?: string; + contractId?: string; + /** Current on-disk checkpoint, read-only, so the preview includes history. */ + currentCheckpoint?: KeeperCheckpointFile; + /** ISO timestamp recorded in the preview. Defaults to the system clock. */ + nowIso?: string; +} + function requiredEnv( env: Record, name: string, @@ -173,10 +223,56 @@ async function countRevealedBids( } } +/** + * Build the checkpoint a live run would write, without writing it. + * + * Pure: the returned file object is produced by the same planner the keeper + * store uses, so a dry run shows byte-for-byte what a live pass would persist. + */ +export function planDryRunCheckpoint( + roundId: bigint | number, + step: KeeperStep | null, + options: KeeperDryRunOptions = {}, +): KeeperDryRunCheckpoint { + const path = options.checkpointPath ?? DEFAULT_CHECKPOINT_PATH; + const network = options.network ?? ""; + const contractId = options.contractId ?? ""; + const current = options.currentCheckpoint; + const bound = network !== "" && contractId !== ""; + const mismatch = + current && bound + ? checkpointBindingMismatch(current, { network, contractId }) + : null; + + const base: KeeperCheckpointFile = current + ? { + version: current.version || CHECKPOINT_VERSION, + network: current.network, + contractId: current.contractId, + rounds: { ...current.rounds }, + } + : { version: CHECKPOINT_VERSION, network, contractId, rounds: {} }; + + const proposedFile = step + ? planCheckpointStep(base, roundId, step, { at: options.nowIso }) + : base; + + return { + path, + network, + contractId, + proposedStep: step, + mismatch, + proposedFile, + filesWritten: 0, + }; +} + export async function buildKeeperDryRunSummary( reader: KeeperDryRunReader, roundId: bigint | number, nowSeconds = systemClock.nowSeconds(), + options: KeeperDryRunOptions = {}, ): Promise { const rid = BigInt(roundId); const round = await reader.getRound(rid); @@ -202,5 +298,9 @@ export async function buildKeeperDryRunSummary( revealedCount, ...decision, transactionsSubmitted: 0, + checkpoint: planDryRunCheckpoint(rid, DRY_RUN_PHASE_STEP[decision.currentPhase], { + ...options, + nowIso: options.nowIso ?? systemClock.toISOString(), + }), }; } diff --git a/services/keeper/src/index.ts b/services/keeper/src/index.ts index aaf00fc2..606b8b02 100644 --- a/services/keeper/src/index.ts +++ b/services/keeper/src/index.ts @@ -29,12 +29,39 @@ export { buildKeeperDryRunSummary, decideKeeperDryRunAction, parseKeeperRunConfig, + planDryRunCheckpoint, + DRY_RUN_PHASE_STEP, + type KeeperDryRunCheckpoint, type KeeperDryRunDecision, + type KeeperDryRunOptions, type KeeperDryRunPhase, type KeeperDryRunReader, type KeeperDryRunSummary, type KeeperRunConfig, } from "./dry-run.js"; +export { + checkpointBindingMismatch, + checkpointHasStep, + planCheckpointRollback, + planCheckpointStep, + readCheckpointFile, + CHECKPOINT_SKIP_STEPS, + CHECKPOINT_VERSION, + DEFAULT_CHECKPOINT_PATH, + KeeperCheckpointError, + KeeperCheckpointMismatchError, + KeeperCheckpointStore, + STEP_SATISFIED_BY_STATUS, + type CheckpointVerification, + type KeeperCheckpoint, + type KeeperCheckpointFile, + type KeeperCheckpointStoreOptions, + type KeeperStep, + type ResumableCheckpoint, + type TransactionHashStatus, + type TransactionHashVerifier, + type WatchCheckpoint, +} from "./checkpoint.js"; export { createStatusServer, createStatusHandler, @@ -56,6 +83,9 @@ export { type KeeperServiceHealth, type KeeperStatusResponse, } from "./status.js"; -export { runWatchLoop, type RunWatchLoopParams } from "./watch-loop.js"; -export { KeeperQueue, type KeeperQueueOptions } from "./queue.js"; -export { KeeperStore, type WatchedRound, type RoundIdInput } from "./store.js"; +export { + resumeCheckpoint, + runWatchLoop, + type ResumeCheckpointParams, + type RunWatchLoopParams, +} from "./watch-loop.js"; diff --git a/services/keeper/src/keeper.ts b/services/keeper/src/keeper.ts index 5d3b6aed..ec72db37 100644 --- a/services/keeper/src/keeper.ts +++ b/services/keeper/src/keeper.ts @@ -19,6 +19,11 @@ import { normalizeError } from "@sub-rosa/logging/errors"; import type { SubRosaClient } from "@sub-rosa/sdk"; import { openBid, fetchRoundSignature, type DrandClient } from "@sub-rosa/tlock"; import { compareRoundIds } from "./store.js"; +import { + CHECKPOINT_SKIP_STEPS, + type KeeperStep, + type WatchCheckpoint, +} from "./checkpoint.js"; import { resolveTimeContext, systemTime, @@ -41,7 +46,11 @@ export interface KeeperDeps { pollMs?: number; /** Injectable wall clock and scheduler. Default: systemTime. */ time?: PartialTimeContext; - settlementGuard?: SettlementGuard; + /** + * Durable watch cursor. When present, steps it records as complete are not + * re-broadcast — that is what stops a restart from resettling a round. + */ + checkpoint?: WatchCheckpoint; } export interface SkipRecord { @@ -62,6 +71,9 @@ export interface KeeperResult { const IDEMPOTENT_OPEN = ["RevealAlreadyOpen", "WrongStatus", "AlreadyCleared"]; const IDEMPOTENT_REVEAL = ["AlreadyRevealed"]; +/** Skip reasons that still mean "this bid is revealed on chain". */ +const REVEALED_REASONS = new Set(["already revealed", "already revealed (race)"]); + export function errorName(e: unknown): string { return normalizeError(e).message; } @@ -75,6 +87,30 @@ function keeperTime(deps: KeeperDeps): TimeContext { return resolveTimeContext(systemTime, deps.time); } +/** + * True when the durable checkpoint says this step already completed and must not + * be broadcast again. Only skip-eligible steps qualify: `open-reveal` is derived + * from the on-chain status, so trusting the cursor there could strand a round. + */ +function stepAlreadyDone( + deps: KeeperDeps, + roundId: bigint, + step: KeeperStep, +): boolean { + if (!deps.checkpoint || !CHECKPOINT_SKIP_STEPS.includes(step)) return false; + return deps.checkpoint.isComplete(roundId, step); +} + +/** Record a completed step in the durable cursor (no-op without a checkpoint). */ +function recordStep( + deps: KeeperDeps, + roundId: bigint, + step: KeeperStep, + transactionHash?: string | null, +): void { + deps.checkpoint?.markComplete(roundId, step, transactionHash); +} + /** Wait until Drand round R should be published. Returns false if R is still in * the future after `maxWaitSeconds`. */ export async function waitForRound( @@ -154,65 +190,83 @@ export async function keepRound( 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"); round = await sdk.getRound(rid); } // ── Phase B: decrypt every seal and reveal it ───────────────────────── if (round.status.tag === "Revealing") { - const bidders: string[] = []; - for await (const addr of sdk.bidders(rid)) bidders.push(addr); - log(`revealing ${bidders.length} bidder(s)`); + if (stepAlreadyDone(deps, rid, "reveal")) { + // A previous process already revealed every bid and recorded it. Re-decrypting + // and re-broadcasting would be pure waste, so trust the cursor. + log(`reveals skipped: checkpoint records every bid revealed`); + result.skipped.push({ bidder: "*", reason: "reveals complete (checkpoint)" }); + } else { + const bidders: string[] = []; + for await (const addr of sdk.bidders(rid)) bidders.push(addr); + log(`revealing ${bidders.length} bidder(s)`); + + for (const bidder of bidders) { + let state; + try { + state = await sdk.getBidState(rid, bidder); + } catch (e) { + result.skipped.push({ bidder, reason: `state read failed: ${errorName(e)}` }); + continue; + } + // Option None decodes as null/undefined; a revealed bid is a bigint. + if (state.revealed_value != null) { + result.skipped.push({ bidder, reason: "already revealed" }); + continue; + } - for (const bidder of bidders) { - let state; - try { - state = await sdk.getBidState(rid, bidder); - } catch (e) { - result.skipped.push({ bidder, reason: `state read failed: ${errorName(e)}` }); - continue; - } - // Option None decodes as null/undefined; a revealed bid is a bigint. - if (state.revealed_value != null) { - result.skipped.push({ bidder, reason: "already revealed" }); - continue; - } + const seal = await sdk.getSeal(rid, bidder); + if (!seal) { + result.skipped.push({ bidder, reason: "seal expired/absent" }); + continue; + } - const seal = await sdk.getSeal(rid, bidder); - if (!seal) { - result.skipped.push({ bidder, reason: "seal expired/absent" }); - continue; - } + let opened; + try { + opened = await openBid(new Uint8Array(seal.ciphertext), drand); + } catch (e) { + result.skipped.push({ bidder, reason: `decrypt failed: ${errorName(e)}` }); + continue; + } - let opened; - try { - opened = await openBid(new Uint8Array(seal.ciphertext), drand); - } catch (e) { - result.skipped.push({ bidder, reason: `decrypt failed: ${errorName(e)}` }); - continue; + try { + await sdk.reveal({ + roundId: rid, + bidder, + value: opened.value, + nonce: opened.nonce, + }); + result.revealed.push(bidder); + log(`revealed ${bidder} = ${opened.value}`); + } catch (e) { + if (errorMatches(e, IDEMPOTENT_REVEAL)) { + result.skipped.push({ bidder, reason: "already revealed (race)" }); + } else if (errorMatches(e, ["HashMismatch"])) { + // A reveal that does not hash to H is rejected by the contract; the + // canonical value is whatever we decrypted, so this only happens for a + // corrupt seal — record and move on. + result.skipped.push({ bidder, reason: "hash mismatch (corrupt seal)" }); + } else if (errorMatches(e, ["RevealWindowClosed"])) { + result.skipped.push({ bidder, reason: "reveal window closed" }); + } else { + throw e; + } + } } - try { - await sdk.reveal({ - roundId: rid, - bidder, - value: opened.value, - nonce: opened.nonce, - }); - result.revealed.push(bidder); - log(`revealed ${bidder} = ${opened.value}`); - } catch (e) { - if (errorMatches(e, IDEMPOTENT_REVEAL)) { - result.skipped.push({ bidder, reason: "already revealed (race)" }); - } else if (errorMatches(e, ["HashMismatch"])) { - // A reveal that does not hash to H is rejected by the contract; the - // canonical value is whatever we decrypted, so this only happens for a - // corrupt seal — record and move on. - result.skipped.push({ bidder, reason: "hash mismatch (corrupt seal)" }); - } else if (errorMatches(e, ["RevealWindowClosed"])) { - result.skipped.push({ bidder, reason: "reveal window closed" }); - } else { - throw e; - } + // Only advance the cursor when every bidder ended up revealed. A seal we + // could not decrypt, or a window that closed mid-pass, must stay retryable. + if (result.skipped.every((s) => REVEALED_REASONS.has(s.reason))) { + recordStep(deps, rid, "reveal"); + } else { + log(`reveals incomplete; cursor left before the reveal step`); } } round = await sdk.getRound(rid); @@ -259,27 +313,38 @@ export async function closeRound( // ── Phase C: clear once the reveal window has closed ────────────────── if (round.status.tag === "Revealing") { - const now = clock.nowSeconds(); - if (now <= Number(round.reveal_deadline)) { - result.skipped.push(`reveal window open until ${round.reveal_deadline}`); - result.finalStatus = round.status.tag; - return result; - } - 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}`); + if (stepAlreadyDone(deps, rid, "clear")) { + // The cursor says a previous process cleared this round. Do not re-broadcast + // clear: the winner is already fixed and the escrow is already committed. + log(`clear skipped: checkpoint records the round as cleared`); + result.skipped.push("clear already complete (checkpoint)"); + } else { + const now = clock.nowSeconds(); + if (now <= Number(round.reveal_deadline)) { + result.skipped.push(`reveal window open until ${round.reveal_deadline}`); + result.finalStatus = round.status.tag; + return result; } - } catch (e) { - if (errorMatches(e, ["AlreadyCleared", "RevealStillOpen", "WrongStatus", "RoundVoided"])) { - result.skipped.push(`clear skipped: ${errorName(e)}`); - } else { - throw e; + 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; + } } } round = await sdk.getRound(rid); @@ -287,28 +352,28 @@ export async function closeRound( // ── Phase D: settle a cleared round (real SAC transfers) ────────────── if (round.status.tag === "Cleared") { - if (deps.settlementGuard) { - const check = deps.settlementGuard.canSettle(rid); - if (!check.allowed) { - result.skipped.push(`settle skipped: duplicate (${check.event.skippedDuplicateReason})`); - round = await sdk.getRound(rid); - result.finalStatus = round.status.tag; - return result; - } - deps.settlementGuard.markSubmitted(rid); - } - try { - await sdk.settle(rid); - result.settled = true; - deps.settlementGuard?.markTerminal(rid, "settled on-chain"); - log(`settled round ${rid}`); - } catch (e) { - if (errorMatches(e, ["AlreadySettled", "NotCleared", "WrongStatus"])) { - deps.settlementGuard?.markTerminal(rid, `skipped: ${errorName(e)}`); - result.skipped.push(`settle skipped: ${errorName(e)}`); - } else { - deps.settlementGuard?.markRetryable(rid, errorName(e)); - throw e; + if (stepAlreadyDone(deps, rid, "settle")) { + // This is the whole point of the cursor: a restart after a confirmed + // settle must not pay the escrow out a second time. + log(`settle skipped: checkpoint records the round as settled`); + result.skipped.push("settle already complete (checkpoint)"); + } 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; + } } } round = await sdk.getRound(rid); @@ -357,6 +422,13 @@ export async function voidIfStale( return result; } + if (stepAlreadyDone(deps, rid, "void")) { + log(`void skipped: checkpoint records the round as voided`); + result.skipped.push("void already complete (checkpoint)"); + result.finalStatus = round.status.tag; + return result; + } + const now = clock.nowSeconds(); const voidAfter = Number(round.reveal_deadline) + VOID_GRACE_SECONDS; if (now <= voidAfter) { @@ -369,6 +441,7 @@ export async function voidIfStale( await sdk.void(rid); result.voided = true; log(`voided round ${rid} (Drand liveness / grace elapsed)`); + recordStep(deps, rid, "void"); } catch (e) { if (errorMatches(e, ["NotVoidable", "WrongStatus", "AlreadyCleared"])) { result.skipped.push(errorName(e)); diff --git a/services/keeper/src/run.ts b/services/keeper/src/run.ts index 34a084d6..f91a3c23 100644 --- a/services/keeper/src/run.ts +++ b/services/keeper/src/run.ts @@ -6,17 +6,23 @@ const diagnostics = createLogger("services.keeper.src.run"); // all) and prints the result. Re-running is safe: completed work is skipped. // // Env: -// ROUND_CONTRACT_ID deployed Round contract id (C…) -// ROUND_ID round to keep (default 1) -// KEEPER_DRY_RUN true prints a read-only preflight summary and exits -// KEEPER_SECRET funded signer secret (S…); not required for dry-run -// MAX_WAIT_SECONDS how long to wait for round R (default 0) -// RPC_URL default https://soroban-testnet.stellar.org -// NETWORK_PASSPHRASE default testnet +// ROUND_CONTRACT_ID deployed Round contract id (C…) +// ROUND_ID round to keep (default 1) +// KEEPER_DRY_RUN true prints a read-only preflight summary and exits +// KEEPER_SECRET funded signer secret (S…); not required for dry-run +// MAX_WAIT_SECONDS how long to wait for round R (default 0) +// RPC_URL default https://soroban-testnet.stellar.org +// NETWORK_PASSPHRASE default testnet +// KEEPER_CHECKPOINT_PATH default .keeper-checkpoint.json import { SubRosaClient } from "@sub-rosa/sdk"; import { quicknet } from "@sub-rosa/tlock"; +import { + DEFAULT_CHECKPOINT_PATH, + KeeperCheckpointStore, + readCheckpointFile, +} from "./checkpoint.js"; import { buildKeeperDryRunSummary, parseKeeperRunConfig, @@ -32,9 +38,29 @@ async function main() { networkPassphrase: config.networkPassphrase, contractId: config.contractId, }); - const summary = await buildKeeperDryRunSummary(reader, config.roundId); + // Dry run reads the checkpoint for context but never writes it and never + // builds a transaction. + const checkpointPath = + process.env.KEEPER_CHECKPOINT_PATH ?? DEFAULT_CHECKPOINT_PATH; + const summary = await buildKeeperDryRunSummary( + reader, + config.roundId, + undefined, + { + checkpointPath, + network: config.networkPassphrase, + contractId: config.contractId, + currentCheckpoint: readCheckpointFile(checkpointPath), + }, + ); diagnostics.info("keeper-dry-run-summary", "keeper dry-run summary:"); diagnostics.info("progress", JSON.stringify(summary, bigintReplacer, 2)); + if (summary.checkpoint.mismatch) { + diagnostics.warn( + "keeper-dry-run-checkpoint-mismatch", + `dry-run: a live keeper would refuse to start — checkpoint ${summary.checkpoint.mismatch} does not match this process config.`, + ); + } return; } @@ -45,12 +71,20 @@ async function main() { 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 result = await keepRound( { sdk, drand: quicknet(), log: (m) => diagnostics.info("progress-2", `· ${m}`), maxWaitSeconds: config.maxWaitSeconds, + checkpoint, }, config.roundId, ); diff --git a/services/keeper/src/serve.ts b/services/keeper/src/serve.ts index 4b893f1e..8d6c1217 100644 --- a/services/keeper/src/serve.ts +++ b/services/keeper/src/serve.ts @@ -29,6 +29,7 @@ import { Keypair } from "@stellar/stellar-sdk"; import { SubRosaClient } from "@sub-rosa/sdk"; import { quicknet } from "@sub-rosa/tlock"; +import { KeeperCheckpointStore } from "./checkpoint.js"; import { createSettlementGuard } from "./settlement-guard.js"; import { createStatusServer, withGracefulShutdown } from "./status-server.js"; import { KeeperStore } from "./store.js"; @@ -65,6 +66,9 @@ async function main() { const log = (m: string) => diagnostics.info("progress", `· ${m}`); 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 }); @@ -121,6 +125,7 @@ async function main() { store, queue, settlementGuard, + checkpoint, isStopping: () => stopping, owner: process.env.KEEPER_OWNER?.trim() || generateLeaseOwner(), leaseMs: parseLeaseMs(process.env.KEEPER_LEASE_MS), diff --git a/services/keeper/src/watch-loop.ts b/services/keeper/src/watch-loop.ts index 69fc8af6..11a39321 100644 --- a/services/keeper/src/watch-loop.ts +++ b/services/keeper/src/watch-loop.ts @@ -24,7 +24,10 @@ import { import type { SettlementGuard } from "./settlement-guard.js"; import type { KeeperLogger } from "./keeper.js"; import { KeeperStore } from "./store.js"; -import { KeeperQueue } from "./queue.js"; +import type { + ResumableCheckpoint, + TransactionHashVerifier, +} from "./checkpoint.js"; export interface RunWatchLoopParams { sdk: SubRosaClient; @@ -36,6 +39,10 @@ export interface RunWatchLoopParams { store: KeeperStore; settlementGuard: SettlementGuard; isStopping: () => boolean; + /** Durable watch cursor. When omitted the loop still runs, just without a cursor. */ + checkpoint?: ResumableCheckpoint; + /** Optional transaction-hash lookup used to confirm recorded cursor steps. */ + verifyTransaction?: TransactionHashVerifier; /** Injectable wall clock and scheduler. Default: systemTime. */ time?: PartialTimeContext; queue?: KeeperQueue; @@ -66,56 +73,55 @@ async function resolveRoundIds(reader: SubRosaClient): Promise { }); } +export interface ResumeCheckpointParams { + checkpoint?: ResumableCheckpoint; + sdk: Pick; + log: KeeperLogger; + verifyTransaction?: TransactionHashVerifier; +} + /** - * Bounds in-flight round execution during shutdown by the given scheduler and timeout. + * Startup resume gate for the watch cursor. + * + * 1. Re-check every recorded transaction hash. A hash that is not `confirmed` + * is rolled back so the step is retried. + * 2. Reconcile the remaining (hashless) cursor entries against the on-chain + * status. If the chain cannot prove a step happened, drop it rather than + * stranding the round. */ -async function waitForInFlight( - promise: Promise, - opts: { - isStopping: () => boolean; - scheduler: Scheduler; - timeoutMs: number; - roundId: bigint; - }, -): Promise { - const { isStopping, scheduler, timeoutMs, roundId } = opts; - let done = false; +export async function resumeCheckpoint( + params: ResumeCheckpointParams, +): Promise { + const { checkpoint, sdk, log, verifyTransaction } = params; + if (!checkpoint) return; - const timeoutPromise = new Promise((_, reject) => { - const triggerTimeout = () => { - scheduler.setTimeout(() => { - if (!done) { - reject(new Error(`Shutdown timeout (${timeoutMs}ms) waiting for round ${roundId}`)); - } - }, timeoutMs); - }; + for (const verification of await checkpoint.verifyHashes(verifyTransaction)) { + log( + `checkpoint ${verification.retained ? "confirmed" : "rolled back"} ` + + `${verification.step} for round ${verification.roundId} ` + + `(tx ${verification.transactionHash}: ${verification.status})`, + ); + } - if (isStopping()) { - triggerTimeout(); - return; + for (const roundId of checkpoint.listRoundIds()) { + let status: string; + try { + status = (await sdk.getRound(roundId)).status.tag; + } catch (e) { + // An unreadable round is not proof of anything — keep the cursor and retry + // reconciliation on the next tick. + log( + `checkpoint: could not read round ${roundId} to reconcile: ${normalizeError(e).message}`, + ); + continue; } - - const checkInterval = 25; - let pollHandle: ReturnType | undefined; - - const check = () => { - if (done) return; - if (isStopping()) { - triggerTimeout(); - return; - } - pollHandle = scheduler.setTimeout(check, checkInterval); - }; - - pollHandle = scheduler.setTimeout(check, checkInterval); - }); - - return Promise.race([ - promise.finally(() => { - done = true; - }), - timeoutPromise, - ]); + const dropped = checkpoint.reconcile(roundId, status); + if (dropped.length > 0) { + log( + `checkpoint: round ${roundId} is ${status}; retrying ${dropped.join(", ")}`, + ); + } + } } export async function runWatchLoop(params: RunWatchLoopParams): Promise { @@ -129,6 +135,8 @@ export async function runWatchLoop(params: RunWatchLoopParams): Promise { store, settlementGuard, isStopping, + checkpoint, + verifyTransaction, time, owner: explicitOwner, leaseMs, @@ -140,11 +148,9 @@ export async function runWatchLoop(params: RunWatchLoopParams): Promise { const resolvedTime = resolveTimeContext(systemTime, time); const { clock, scheduler } = resolvedTime; - const deps: KeeperDeps = { sdk, drand, log, time: resolvedTime, settlementGuard }; - const queue = params.queue ?? new KeeperQueue(store, { contractId, network }); - const shutdownTimeoutMs = params.shutdownTimeoutMs ?? 30000; + const deps: KeeperDeps = { sdk, drand, log, time: resolvedTime, checkpoint }; - const shouldStop = () => isStopping() || queue.isStopping(); + await resumeCheckpoint({ checkpoint, sdk, log, verifyTransaction }); while (!shouldStop()) { const started = clock.nowMs(); diff --git a/services/keeper/src/watch.ts b/services/keeper/src/watch.ts index 8405dab8..eda053ae 100644 --- a/services/keeper/src/watch.ts +++ b/services/keeper/src/watch.ts @@ -21,6 +21,7 @@ import { Keypair } from "@stellar/stellar-sdk"; import { SubRosaClient } from "@sub-rosa/sdk"; import { quicknet } from "@sub-rosa/tlock"; +import { KeeperCheckpointStore } from "./checkpoint.js"; import { createSettlementGuard } from "./settlement-guard.js"; import { KeeperStore } from "./store.js"; import { KeeperQueue } from "./queue.js"; @@ -59,6 +60,9 @@ 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 }); @@ -77,6 +81,7 @@ async function main() { store, queue, settlementGuard, + checkpoint, isStopping: () => stopping, owner: process.env.KEEPER_OWNER?.trim() || generateLeaseOwner(), leaseMs: parseLeaseMs(process.env.KEEPER_LEASE_MS),