From c781045f7b8b01b8755bf287f42a4eea5336cbbf Mon Sep 17 00:00:00 2001 From: Rinse Date: Thu, 3 Sep 2026 03:05:45 +0000 Subject: [PATCH 1/4] fix(client): chase on admission, never delete a snapshot over an RPC blip MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two silent state-loss bugs found on the production 5chan seeder, which served a divergent tally on two topics for 13 days. `makeRootChaser`'s skip predicate was wired to blockstore membership, but admission lives in the CRDT. The two disagree exactly on the node that most needs the chase: one whose persistent blockstore outlived the state keyed on it. Such a node decoded every peer's checkpoint successfully (all blocks resolve locally!), skipped every bundle in it, admitted nothing, and could never converge again — restarting did not help, since the cold-start pull feeds the same chaser. The dep is now `isAdmitted`, wired to the engine's `#checks` map (the same admission map `#restoreSnapshot` consults). Closes #44. `#restoreSnapshot` wrapped the decode AND the per-bundle admission in one try/catch whose catch assumed a corrupt blob and deleted it — but the admission path reads the gating chain, so one rate-limit window during a seeder's boot burst (64 topics restoring at once) permanently discarded a topic's persisted votes. The decode is now the whole of the corruption test; admission runs in `#admitRestored`, where each bundle is independent and a transient failure (a throwing head read, a head that has not reached the bundle's sample bucket) backlogs it for a retry instead of dropping it. A non-empty backlog suppresses the snapshot write, so a partial view can never overwrite the good blob. And none of it is silent any more: a discarded blob, an incomplete restore and a write that keeps failing all surface as a `SnapshotError` on the contest's `error` event — which meant registering the view's engine listeners BEFORE the join, since the restore runs inside it. Closes #45. Regression tests pin both at the voter level: a node holding every block of a served checkpoint (bundle block included) still admits it, and a restore interrupted by a 429 keeps its blob, reports itself incomplete, refuses to overwrite the blob with the partial view, and admits the votes on retry. --- src/client/voter.test.ts | 191 +++++++++++++++++++++- src/client/voter.ts | 230 ++++++++++++++++++++++----- src/errors.ts | 40 +++++ src/test-fixtures.ts | 11 +- src/transport/chase.test.ts | 28 +++- src/transport/chase.ts | 18 ++- src/transport/integration/harness.ts | 7 +- 7 files changed, 467 insertions(+), 58 deletions(-) diff --git a/src/client/voter.test.ts b/src/client/voter.test.ts index 119b3e2..1ae8c9b 100644 --- a/src/client/voter.test.ts +++ b/src/client/voter.test.ts @@ -31,12 +31,14 @@ import { MissingChainClientError, MissingFetchError, MissingPubsubError, + SnapshotError, UnknownRuleError, VoteEvictedError, VoterDestroyedError } from "../errors.js"; import { topicFor } from "../topic.js"; import { blockForBytes, encodeCheckpoint } from "../checkpoint/codec.js"; +import { bundleCidForBytes, encodeBundle } from "../crdt/codec.js"; import { bizCriteria, bizGateRef, @@ -49,7 +51,8 @@ import { fakeChains, stubChains, fakeSigner, - realSigner + realSigner, + SECOND_TEST_PRIVATE_KEY } from "../test-fixtures.js"; /** A valid base58btc IPNS community key (VotesBundleSchema rejects non-keys, so a real one is needed). */ @@ -58,6 +61,8 @@ const VALID_KEY = "12D3KooWEyoppNCUx8Yx66oV9fVnrJmG92pTuY6zbLDaz8T5XCiL"; const OTHER_KEY = "12D3KooWQYV9dGMFoRzNStwpXztXaBUjtPqi6aU76ZgUriHhKust"; /** A one-community upvote used across the publish tests. */ const VOTE = [{ community: { publicKey: VALID_KEY }, vote: 1 }]; +/** A second community's upvote, for tests that need two distinguishable ballots. */ +const OTHER_VOTE = [{ community: { publicKey: OTHER_KEY }, vote: 1 }]; /** The same upvote carrying a resolvable community name (exercises the name pipeline). */ const NAMED_VOTE = [{ community: { name: "memes.bso", publicKey: VALID_KEY }, vote: 1 }]; @@ -865,6 +870,118 @@ describe("checkpoint snapshot persistence (dataPath)", () => { await voter.destroy(); }); + /** + * A chain client whose head read fails on demand — the production shape of issue #45's + * incident: a seeder boots, restores ~64 topics at once, and one rate-limit window makes the + * head read behind the restore's admissibility check throw. `readContract` (the gate read) + * keeps working, so only the restore path is disturbed. + */ + const flakyHeadChains = (): { chains: ChainClientFactory; fail: (on: boolean) => void } => { + let failing = true; + const client = { + getBlockNumber: async () => { + if (failing) throw new Error("429 Too Many Requests"); + return 43200n; + }, + getBlock: async () => ({ hash: `0x${"11".repeat(32)}` }), + readContract: async () => 1n + }; + return { chains: () => client as unknown as ChainClient, fail: (on) => (failing = on) }; + }; + + /** Session 1 of a restart test: publish one verified vote and persist it on `leave()`. */ + const persistOneVote = async (dataPath: string, signer = realSigner()): Promise => { + const voter = new PubsubVoter({ dataPath, helia: fakeHelia(), chains: countingChains().chains }); + await (await voter.createContestVote({ criteria: bizCriteria(), votes: VOTE, signer })).publish(); + const contest = await voter.createContest({ criteria: bizCriteria() }); + await vi.waitFor(async () => expect((await contest.getTally()).ranking[0]?.chainVerified).toBe(true)); + await voter.destroy(); + }; + + it("keeps the snapshot when a transient RPC failure interrupts the restore (issue #45)", async () => { + const dataPath = await tempDataPath(); + const topic = await topicFor(bizCriteria()); + await persistOneVote(dataPath); + + // Session 2 restores into an RPC outage: the head read the admissibility check makes + // throws for every bundle. The blob is NOT corrupt — nothing about it failed — so + // deleting it (as one try/catch around the whole restore used to) permanently discards + // votes the node verified itself and no peer may be online to re-serve. + const flaky = flakyHeadChains(); + const voterB = new PubsubVoter({ dataPath, helia: fakeHelia(), chains: flaky.chains }); + const contestB = await voterB.createContest({ criteria: bizCriteria() }); + const errors: unknown[] = []; + contestB.on("error", (e) => errors.push(e)); + await contestB.update(); + expect(contestB.tally?.ranking).toEqual([]); // nothing admitted this session... + // ...and the incompleteness is reported, not silent (the whole loss was invisible). + expect(errors.filter((e) => e instanceof SnapshotError && e.failure === "restore-incomplete")).toHaveLength(1); + await voterB.destroy(); + + // The blob survived — session 3, with a healthy RPC, restores the vote it holds. + const reader = (await import("../storage/node.js")).makeStorage({ dataPath }); + expect(await reader.openSnapshots().get(topic)).toBeDefined(); + await reader.destroy(); + + const voterC = new PubsubVoter({ dataPath, helia: fakeHelia(), chains: countingChains().chains }); + const contestC = await voterC.createContest({ criteria: bizCriteria() }); + await contestC.update(); + expect(contestC.tally?.ranking[0]?.community.publicKey).toBe(VALID_KEY); + await voterC.destroy(); + }); + + it("suppresses the snapshot write while a restore is incomplete (a partial view must not overwrite the blob)", async () => { + const dataPath = await tempDataPath(); + const topic = await topicFor(bizCriteria()); + await persistOneVote(dataPath); // wallet #1's vote, on the blob + + // Session 2's restore is interrupted, then the RPC recovers and a SECOND wallet votes + // here. That publish changes the winner set and arms the snapshot write — which must + // skip, because our view is knowably missing the votes the backlog still holds. Writing + // it would replace a blob holding wallet #1's vote with one that never mentions it. + const flaky = flakyHeadChains(); + const voterB = new PubsubVoter({ dataPath, helia: fakeHelia(), chains: flaky.chains }); + const contestB = await voterB.createContest({ criteria: bizCriteria() }); + await contestB.update(); + flaky.fail(false); + const second = realSigner(SECOND_TEST_PRIVATE_KEY); + await (await voterB.createContestVote({ criteria: bizCriteria(), votes: OTHER_VOTE, signer: second })).publish(); + await vi.waitFor(async () => expect((await contestB.getTally()).ranking[0]?.chainVerified).toBe(true)); + await voterB.destroy(); // the leave() flush must skip too + + // Session 3 still finds wallet #1's vote: the blob was never overwritten. + const voterC = new PubsubVoter({ dataPath, helia: fakeHelia(), chains: countingChains().chains }); + const contestC = await voterC.createContest({ criteria: bizCriteria() }); + await contestC.update(); + expect(contestC.tally?.ranking.map((row) => row.community.publicKey)).toEqual([VALID_KEY]); + await voterC.destroy(); + expect(topic).toBe(await topicFor(bizCriteria())); + }); + + it("retries the interrupted restore and admits its bundles once the RPC recovers", async () => { + const dataPath = await tempDataPath(); + await persistOneVote(dataPath); + + const flaky = flakyHeadChains(); + vi.useFakeTimers(); + try { + const voterB = new PubsubVoter({ dataPath, helia: fakeHelia(), chains: flaky.chains }); + const contestB = await voterB.createContest({ criteria: bizCriteria() }); + await contestB.update(); + expect(contestB.tally?.ranking).toEqual([]); + + // The RPC comes back. Nothing else happens on this topic — no peer, no publish — so + // only the restore's own retry can recover the vote (RESTORE_RETRY_MS, the first + // backoff step). Before it existed, the session simply ran on without those ballots. + flaky.fail(false); + await vi.advanceTimersByTimeAsync(30_000); + expect((await contestB.getTally()).ranking[0]?.community.publicKey).toBe(VALID_KEY); + await voterB.destroy(); + } finally { + vi.useRealTimers(); + } + }); + it("runs the flush safely on the in-memory backend (`dataPath: false` — nothing to persist, nothing thrown)", async () => { const voter = new PubsubVoter({ dataPath: false, helia: fakeHelia(), chains: stubChains() }); await (await voter.createContestVote({ criteria: bizCriteria(), votes: VOTE, signer: fakeSigner() })).publish(); @@ -873,6 +990,78 @@ describe("checkpoint snapshot persistence (dataPath)", () => { }); }); +/** + * The chase's skip predicate is ADMISSION, not block presence (issue #44). The two disagree + * exactly on the node that needs the chase most: one whose persistent blockstore outlived the + * state keyed on it. Keyed on the blockstore, such a node decoded every peer's checkpoint + * successfully (all blocks resolve locally!), skipped every bundle in it, admitted nothing, and + * stayed silently wedged — for 13 days on the production seeder, while four peers served the + * ballots it was missing. + */ +describe("chase convergence with a warm blockstore (issue #44)", () => { + it("admits a served checkpoint's bundles whose blocks we already hold but never admitted", async () => { + // One real, signature-verifiable bundle, encoded into a real checkpoint. + const publisher = new PubsubVoter({ dataPath: false, helia: fakeHelia(), chains: stubChains() }); + const ballot = await publisher.createContestVote({ criteria: bizCriteria(), votes: VOTE, signer: realSigner() }); + const { bundle } = await ballot.publish(); + await publisher.destroy(); + const { root, chunks, blocks } = await encodeCheckpoint([bundle]); + + // The wedged node: every block of that checkpoint is in its blockstore — the chunk blocks + // AND the bundle's own block, which is what the production seeder really had on disk (it + // received the vote live, verified it and stored it) — while nothing is admitted, because + // the state keyed on those blocks was lost with its snapshot on an earlier restart. + const held = new Map(blocks.map((block) => [block.cid.toString(), block.bytes])); + const bundleBytes = encodeBundle(bundle); + held.set((await bundleCidForBytes(bundleBytes)).toString(), bundleBytes); + const blockstore = { + get: async (cid: CID) => { + const bytes = held.get(cid.toString()); + if (bytes === undefined) throw new Error("no block"); + return bytes; + }, + put: async (cid: CID, bytes: Uint8Array) => { + held.set(cid.toString(), bytes); + return cid; + }, + has: async (cid: CID) => held.has(cid.toString()) + }; + + // One peer serving exactly the root the wedged node's blocks belong to. + let topic = ""; + const record: FetchRootRecord = { + version: ROOT_RECORD_VERSION, + root, + chunks, + count: 1, + sizeBytes: blocks.reduce((total, block) => total + block.bytes.length, 0) + }; + const subscriber = { toString: () => "peerA" }; + const pubsub: PubsubService = { + publish: async () => ({}) as never, + subscribe: () => {}, + unsubscribe: () => {}, + getSubscribers: () => [subscriber] as unknown as ReturnType, + addEventListener: () => {}, + removeEventListener: () => {}, + topicValidators: new Map() + }; + const helia = { + libp2p: { services: { pubsub, fetch: rootFetchService({ records: () => ({ [topic]: record }) }) } }, + blockstore + } as unknown as HeliaInstance; + + const voter = new PubsubVoter({ dataPath: false, helia, chains: stubChains() }); + const contest = await voter.createContest({ criteria: bizCriteria() }); + topic = contest.topic; + await contest.update(); + // The cold-start pull hears the divergent root and chases it; the bundle must be verified + // and ADMITTED, not skipped because its block happens to be on disk. + await vi.waitFor(async () => expect((await contest.getTally()).ranking[0]?.community.publicKey).toBe(VALID_KEY)); + await voter.destroy(); + }); +}); + describe("republish cadence helper (client-owned republishing)", () => { it("recommends half the expiry window, rounded up", () => { expect(republishIntervalBuckets(bizCriteria())).toBe(15); // voteExpiryBuckets: 30 diff --git a/src/client/voter.ts b/src/client/voter.ts index 6920bf9..82dc6ec 100644 --- a/src/client/voter.ts +++ b/src/client/voter.ts @@ -68,6 +68,7 @@ import { criteriaCid, TOPIC_PREFIX } from "../topic.js"; import { InvalidCommunityNameError, MissingChainClientError, + SnapshotError, UnknownRuleError, VoteEvictedError, VoterDestroyedError @@ -671,6 +672,22 @@ const COLD_START_REPULL_WINDOW_MS = HEARTBEAT_INTERVAL_MS; * `leave()` flushes whatever is still pending so a clean shutdown never loses the tail. */ const SNAPSHOT_DEBOUNCE_MS = 10_000; +/** + * How long to wait before re-attempting the bundles a snapshot restore could not admit for a + * TRANSIENT reason — an RPC outage during a seeder's boot burst, or a gating-chain head that has + * not yet reached a bundle's sample bucket (see {@link ContestEngine.#restoreSnapshot}). Backs off + * exponentially to {@link RESTORE_RETRY_CAP_MS} and keeps retrying while joined: giving up would + * mean either abandoning those votes or suppressing this topic's snapshot writes forever, and a + * retry costs one head read against a memoized bucket, so the patient option is also the cheap one. + */ +const RESTORE_RETRY_MS = 30_000; +const RESTORE_RETRY_CAP_MS = 600_000; +/** + * Snapshot write attempts (one per debounce window) before leaving it to the next winner-set + * change. A store that fails every write must not spin, but a topic quiet enough to never change + * again must not keep a stale snapshot forever because ONE write failed (issue #45). + */ +const SNAPSHOT_WRITE_ATTEMPTS = 3; /** Per-root chase deadline (ms): a multi-block directed-bitswap pull, coarser than one message. */ const CHASE_TIMEOUT_MS = 30_000; /** Concurrent root chases; a spray of divergent roots queues, never floods. */ @@ -1135,6 +1152,19 @@ class ContestEngine { #heartbeatTimer: ReturnType | undefined; /** The armed (debounced) snapshot-write timer; flushed by `leave()`. See {@link #writeSnapshot}. */ #snapshotTimer: ReturnType | undefined; + /** Consecutive failed snapshot writes in the current run (see {@link SNAPSHOT_WRITE_ATTEMPTS}). */ + #snapshotWriteFailures = 0; + /** + * Bundles a restore decoded but could not admit for a TRANSIENT reason, awaiting + * {@link #restoreTimer}. Non-empty means the restore is INCOMPLETE, which suppresses the + * snapshot write: our view is missing votes the blob holds, and persisting it would overwrite + * the good blob with the partial one — the loss the retry exists to prevent (issue #45). + */ + #restoreBacklog: readonly VotesBundle[] = []; + /** The armed restore-retry timer; cleared by `leave()`. */ + #restoreTimer: ReturnType | undefined; + /** Consecutive restore retries, driving the backoff (reset by a complete restore). */ + #restoreAttempts = 0; /** * Tears down the per-join `subscription-change` re-pull (listener + window timer); set by * {@link #armSubscriptionRepull}, cleared by `leave()` or by the window expiring. @@ -1205,9 +1235,8 @@ class ContestEngine { nameResolutionCache: deps.nameResolutionCache, readHead: ({ chain }) => this.#readHead({ chain }) }); - // The gate/transport are (re)built on join(); the store, crdt, caches, verifier, and - // background verifier are stable per contest, so they survive re-joins of the topic. - this.#store = store; + // The gate/transport are (re)built on join(); the crdt, caches, verifier, and background + // verifier are stable per contest, so they survive re-joins of the topic. this.#cache = makeVerdictCache(); this.#acceptedDedup = makeAcceptedDedup(this.#bucketMath); this.#verifier = verifier; @@ -1243,7 +1272,6 @@ class ContestEngine { }); } - readonly #store: ReturnType; /** Per gate leaf: hash of its canonical rule ref + chainId — its keyspace in the shared gate store. */ readonly #ruleIds: string[]; /** The gate's leaf refs in document order, aligned with {@link #ruleIds}. */ @@ -1699,7 +1727,12 @@ class ContestEngine { verifyOffline: (bundle) => this.#verifier.verifyOffline(bundle), cache: this.#cache, isEvaluableNow: (bundle) => this.#isEvaluableNow(bundle), - hasBundle: (cid) => this.#store.has(cid), + // ADMISSION, not block presence (issue #44): our blockstore may hold a bundle's block + // while the winner-set does not know it — after a restart that lost the snapshot, an + // eviction, or a prune. Keyed on the blockstore, such a node skipped that bundle on + // every chase and never converged; `#checks` is the same admission map + // `#restoreSnapshot` consults, and it is dropped in lockstep with the CRDT's membership. + isAdmitted: async (cid) => this.#checks.has(cid.toString()), // `verified: false` is a provisional admit (offline checks only) whose deferred gate // read + name resolution ride `deferVerify`; `true` means a cached terminal verdict // already covers the full pipeline. @@ -1927,12 +1960,21 @@ class ContestEngine { if (this.#heartbeatTimer !== undefined) clearTimeout(this.#heartbeatTimer); this.#heartbeatTimer = undefined; // Flush the debounced snapshot write so a clean shutdown persists the latest state - // (still skipped if checks are pending — the stale-but-good snapshot stays put). + // (still skipped if checks are pending, or a restore backlog is outstanding — the + // stale-but-good snapshot stays put). `retry: false`: a failed flush must not arm a + // timer that outlives the join. if (this.#snapshotTimer !== undefined) { clearTimeout(this.#snapshotTimer); this.#snapshotTimer = undefined; - await this.#writeSnapshot(); + await this.#writeSnapshot({ retry: false }); } + // Drop the restore backlog with its timer: a re-join re-reads the blob (still on disk, + // since the backlog suppressed every write that could have thinned it) and starts over. + if (this.#restoreTimer !== undefined) clearTimeout(this.#restoreTimer); + this.#restoreTimer = undefined; + this.#restoreBacklog = []; + this.#restoreAttempts = 0; + this.#snapshotWriteFailures = 0; // Pause the background verifier's retry timer; pending state survives for a re-join. this.#background.stop(); this.#disarmSubscriptionRepull(); @@ -2062,17 +2104,25 @@ class ContestEngine { * wire bytes plus the blocks it references (just re-put by the encode, read back from the * blockstore so no second copy is held in memory). Best-effort like every persistent-cache * write — a broken store degrades to the pre-persistence behavior (the cold-start pull), - * never an error. + * never an error — but not silent: a write that keeps failing surfaces as a `SnapshotError` + * on the contest's `error` event, because the state it fails to keep is invisible otherwise. + * + * Skipped while our view is knowably INCOMPLETE, in either of the two ways it can be, because + * a write in that window overwrites a good blob with a lossy one: * - * Skipped while ANY admitted bundle has a deferred check pending: the encoder serves only - * fully verified bundles, so writing mid-settlement would persist a snapshot that OMITS the - * pending ones — right after a restore (where every reloaded bundle is provisional) that - * would clobber a good snapshot with a near-empty one, and a crash in that window would lose - * the very votes persistence exists to keep. Every settlement and eviction re-marks the - * state changed, so the write that was skipped here is re-armed by the last one to land. - */ - async #writeSnapshot(): Promise { - if (this.#hasUnsettledChecks()) return; + * - ANY admitted bundle has a deferred check pending — the encoder serves only fully + * verified bundles, so writing mid-settlement persists a snapshot that OMITS the pending + * ones (right after a restore, where every reloaded bundle is provisional, that is a + * near-empty snapshot). Every settlement and eviction re-marks the state changed, so the + * skipped write is re-armed by the last one to land. + * - a restore is still carrying a backlog ({@link #restoreBacklog}) — the blob holds votes + * a transient failure kept us from admitting, and the retry owns them until it lands. + * + * `retry` is false on the `leave()` flush: a failed write there must not arm a timer that + * outlives the join. + */ + async #writeSnapshot({ retry = true }: { retry?: boolean } = {}): Promise { + if (this.#hasUnsettledChecks() || this.#restoreBacklog.length > 0) return; try { await this.rootRecord(); // (re-)encode so the cache reflects the current winner-set // Read record + blocks from the cache as one unit: a state change racing the encode @@ -2083,9 +2133,18 @@ class ContestEngine { this.topic, encodeSnapshot({ record: encodeRootRecord(cached.record), blocks: cached.blocks.map((block) => block.bytes) }) ); - } catch { - // Best-effort: a failed write leaves the previous snapshot in place; the next - // state change re-arms the debounce and retries. + this.#snapshotWriteFailures = 0; + } catch (error) { + // The previous snapshot stays in place (a restart restores older state, never none). + // Retry on the debounce a couple of times rather than waiting for a winner-set change + // that a quiet topic may never see — then report and stop, so a dead store cannot spin. + this.#snapshotWriteFailures += 1; + if (retry && this.#snapshotWriteFailures < SNAPSHOT_WRITE_ATTEMPTS) { + this.#scheduleSnapshotWrite(); + return; + } + this.#snapshotWriteFailures = 0; + this.#emitError(new SnapshotError("write-failed", this.topic, error)); } } @@ -2097,9 +2156,18 @@ class ContestEngine { * bundle re-passes the offline signature/constraint checks, and the deferred gate read + * name resolution ride the background verifier (mostly persisted-gate-cache hits on a * restart) — so the trust model is unchanged: this is the node's own previously-validated - * state, re-validated on load. A corrupt or version-mismatched blob is removed and the join - * proceeds empty, exactly as before persistence; the cold-start pull then still runs, so a - * stale snapshot self-heals by union with the live topic. + * state, re-validated on load. The cold-start pull still runs afterwards, so a stale snapshot + * self-heals by union with the live topic. + * + * **Only a provably bad blob is discarded** (issue #45). The decode is the whole of the + * corruption test, so it is the whole of what this catch covers: a truncated, mangled or + * version-mismatched blob is removed and the join proceeds empty, exactly as before + * persistence. Everything after it — the per-bundle admission, which reads the gating chain — + * runs under {@link #admitRestored}, where a transient failure yields a RETRY and never a + * delete. The distinction is the incident: the two catches used to be one, so a momentary RPC + * failure during a seeder's boot burst (64 topics restoring at once, one rate-limit window) + * permanently discarded a topic's persisted votes and the node came up empty on a topic every + * other peer still served. */ async #restoreSnapshot(): Promise { let blob: Uint8Array | undefined; @@ -2109,39 +2177,101 @@ class ContestEngine { return; // unreadable store — degrade to a plain cold join } if (blob === undefined) return; + let winners: VotesBundle[]; try { const snapshot = decodeSnapshot(blob); const record = decodeRootRecord(snapshot.record); const byCid = new Map(); for (const bytes of snapshot.blocks) byCid.set((await blockForBytes(bytes)).cid.toString(), bytes); - const winners = await decodeCheckpoint(record.root, async (cid) => byCid.get(cid.toString()), record.chunks); - - const pending: PendingBundle[] = []; - for (const bundle of winners) { - // Same two-stage admit as the chase (see transport/chase.ts): offline checks - // synchronously before admit, chain/name checks deferred and batched. Admission - // goes through `crdt.add` (which stores the block AND registers the bundle) — a - // merge-by-CID would re-read through the blockstore, which the restore must not - // depend on: it may be fresh (the incident's in-memory case) or hold the block - // already (a persistent one), neither of which says the CRDT knows the bundle. + winners = await decodeCheckpoint(record.root, async (cid) => byCid.get(cid.toString()), record.chunks); + } catch (error) { + // Corrupt, truncated, or version-mismatched blob: discard it and join empty — the + // pre-persistence behavior. Never let a bad snapshot block the join. + void this.#deps.snapshots.remove(this.topic).catch(() => {}); + this.#emitError(new SnapshotError("discarded", this.topic, error)); + return; + } + this.#restoreAttempts = 0; + await this.#admitRestored(winners); + } + + /** + * Admit a decoded snapshot's bundles, one independently of the next. Three outcomes per + * bundle, and the whole point of the method is that they stay distinct: + * + * - **admitted** — the same two-stage admit as the chase (transport/chase.ts): offline + * checks synchronously before admit, chain/name checks deferred and batched. Admission + * goes through `crdt.add` (which stores the block AND registers the bundle) — a + * merge-by-CID would re-read through the blockstore, which the restore must not depend + * on: it may be fresh (the incident's in-memory case) or hold the block already (a + * persistent one), neither of which says the CRDT knows the bundle. + * - **dropped** — the offline checks REFUSED it. That is a verdict on the bundle, not on + * the environment, and no retry changes it. + * - **backlogged** — a transient failure: the read that decides evaluability threw (an RPC + * outage), or our gating-chain head has not yet reached the bundle's sample bucket. Those + * go to {@link #restoreBacklog}, which suppresses the snapshot write and arms a retry; + * dropping them silently is how a restart's ballots disappeared with the blob intact. + */ + async #admitRestored(winners: readonly VotesBundle[]): Promise { + const pending: PendingBundle[] = []; + const backlog: VotesBundle[] = []; + for (const bundle of winners) { + try { const cid = await bundleCidForBytes(encodeBundle(bundle)); if (this.#checks.has(cid.toString())) continue; // already admitted (a re-join) - if (!(await this.#isEvaluableNow(bundle))) continue; + if (!(await this.#isEvaluableNow(bundle))) { + backlog.push(bundle); // our head is behind the ballot: retry, do not forget it + continue; + } const offline = await this.#verifier.verifyOffline(bundle); - if (!offline.valid) continue; + if (!offline.valid) continue; // the bundle's own fault — a sibling's admit stands await this.#crdt.add(bundle); this.#recordChecks(cid, bundle, false); pending.push({ cid, bundle }); + } catch { + // Infrastructure, not the bundle (a chain read, a blockstore write): keep it. + backlog.push(bundle); } - if (pending.length > 0) { - this.#background.enqueue(pending); - this.#onStateChanged(); - } - } catch { - // Corrupt, truncated, or version-mismatched blob: discard it and join empty — the - // pre-persistence behavior. Never let a bad snapshot block the join. - void this.#deps.snapshots.remove(this.topic).catch(() => {}); } + const wasIncomplete = this.#restoreBacklog.length > 0; + this.#restoreBacklog = backlog; + if (pending.length > 0) { + this.#background.enqueue(pending); + this.#onStateChanged(); + } + if (backlog.length === 0) { + // The backlog just drained: re-arm the write it was suppressing, so the now-complete + // view reaches disk even if the retry admitted nothing new (a re-join that raced us). + if (wasIncomplete) this.#scheduleSnapshotWrite(); + return; + } + // Report the incompleteness once per restore (not once per retry): an operator needs to + // know this topic's tally is missing persisted votes, not a line every backoff window. + if (!wasIncomplete) this.#emitError(new SnapshotError("restore-incomplete", this.topic)); + this.#armRestoreRetry(); + } + + /** + * Arm the next backlog retry (exponential backoff, capped). Kept armed for as long as the + * backlog is non-empty and we stay joined: the alternative — giving up — means either + * discarding the votes or suppressing this topic's snapshot writes forever. + */ + #armRestoreRetry(): void { + if (this.#restoreTimer !== undefined || !this.#joined) return; + const delay = Math.min(RESTORE_RETRY_MS * 2 ** this.#restoreAttempts, RESTORE_RETRY_CAP_MS); + this.#restoreAttempts += 1; + const timer = setTimeout(() => { + this.#restoreTimer = undefined; + const backlog = this.#restoreBacklog; + if (backlog.length === 0 || !this.#joined) return; + void this.#admitRestored(backlog).catch(() => { + // #admitRestored never throws; re-arm defensively so a backlog cannot strand. + this.#armRestoreRetry(); + }); + }, delay); + // Don't hold a Node process open; no-op in the browser. + (timer as { unref?: () => void }).unref?.(); + this.#restoreTimer = timer; } /** @@ -2464,10 +2594,22 @@ class ContestView implements Contest { async update(): Promise { if (this.#subscribed) return; - await this.#engine.join(); + // Registered BEFORE the join, not after: the join restores this contest's persisted + // snapshot, and every way that can go wrong (a blob discarded as corrupt, a restore an + // RPC outage left incomplete — `SnapshotError`) is reported on the `error` event from + // inside `join()`. Attaching afterwards made exactly those reports unobservable. this.#engine.addUpdateListener(this.#onEngineUpdate); this.#engine.addErrorListener(this.#onEngineError); this.#subscribed = true; + try { + await this.#engine.join(); + } catch (error) { + // A failed join leaves no subscription behind (a retry must be able to re-arm). + this.#engine.removeUpdateListener(this.#onEngineUpdate); + this.#engine.removeErrorListener(this.#onEngineError); + this.#subscribed = false; + throw error; + } // Populate `tally` and fire an initial `update` for the current state. await this.#engine.refreshTallyNow(); } diff --git a/src/errors.ts b/src/errors.ts index d1c2768..e1d551f 100644 --- a/src/errors.ts +++ b/src/errors.ts @@ -235,3 +235,43 @@ export class VoteEvictedError extends Error { this.name = "VoteEvictedError"; } } + +/** + * Emitted (never thrown) on a contest's `error` event when this node's OWN persisted checkpoint + * snapshot could not be kept ({@link SnapshotFailure} says which way): a blob discarded as + * unreadable, a restore left incomplete by a transient failure, or a write that kept failing. + * + * Persistence is best-effort by design — every one of these degrades to the pre-persistence + * behaviour (the cold-start pull), never to a broken join — but "best-effort" used to mean + * "silent", and silent state loss is exactly how a seeder served a divergent tally for 13 days + * without anyone noticing (issue #45). Operators get a fact and a topic instead of an unexplained + * gap; there is nothing for a client to do about it beyond logging or alerting. + */ +export class SnapshotError extends Error { + constructor( + /** Which snapshot operation failed. */ + readonly failure: SnapshotFailure, + /** The topic whose snapshot is affected. */ + readonly topic: string, + /** The underlying throw, when there was one (a `restore-incomplete` may have none). */ + readonly cause?: unknown + ) { + super(`Checkpoint snapshot ${failure} for topic ${topic}${cause === undefined ? "" : `: ${String(cause)}`}`); + this.name = "SnapshotError"; + } +} + +/** + * What went wrong with a persisted checkpoint snapshot: + * + * - `discarded` — the blob did not decode (corrupt, truncated, version-mismatched) and was + * removed. The votes it held are gone from disk; the topic joins empty and re-converges from + * peers. + * - `restore-incomplete` — the blob decoded, but a transient failure (an RPC outage during the + * boot burst, a head that has not reached a bundle's sample bucket) left some of its bundles + * un-admitted. The blob is KEPT and the restore retries; snapshot writes are suppressed for + * this topic meanwhile, so a partial view can never overwrite the good blob. + * - `write-failed` — the store rejected the snapshot write repeatedly. The previous snapshot + * stays in place, so a restart restores an older state rather than none. + */ +export type SnapshotFailure = "discarded" | "restore-incomplete" | "write-failed"; diff --git a/src/test-fixtures.ts b/src/test-fixtures.ts index 8d2e91a..803d24b 100644 --- a/src/test-fixtures.ts +++ b/src/test-fixtures.ts @@ -201,14 +201,21 @@ export function fakeSigner(address = "0x0000000000000000000000000000000000000001 }; } +/** Anvil/hardhat test account #2 — a SECOND wallet for {@link realSigner}. */ +export const SECOND_TEST_PRIVATE_KEY = "0x5de4111afa1a4b94908f83103eb1f1706367c2e68ca870fc3fb9a804cdab365a" as const; + /** * A REAL signer over the anvil/hardhat test account #1 (as in verify/bundle.test.ts and the * two-node integration test), for tests whose bundles must survive the verifier's signature * recovery — e.g. the checkpoint-snapshot restore, which re-runs `verifyOffline` on reload. * `fakeSigner`'s placeholder fails recovery by design. + * + * `privateKey` selects a different wallet ({@link SECOND_TEST_PRIVATE_KEY} is anvil account #2) — + * the LWW winner-set is keyed by wallet, so a test that needs two coexisting bundles needs two of + * these, not one signing twice. */ -export function realSigner(): VoteSigner { - const account = privateKeyToAccount("0x59c6995e998f97a5a0044966f0945389dc9e86dae88c7a8412f4603b6b78690d"); +export function realSigner(privateKey: `0x${string}` = "0x59c6995e998f97a5a0044966f0945389dc9e86dae88c7a8412f4603b6b78690d"): VoteSigner { + const account = privateKeyToAccount(privateKey); return { address: () => account.address, signBallot: async (typedData) => ({ signature: await account.signTypedData(typedData), type: EIP712_SIGNATURE_TYPE }) diff --git a/src/transport/chase.test.ts b/src/transport/chase.test.ts index baed3fb..c191f60 100644 --- a/src/transport/chase.test.ts +++ b/src/transport/chase.test.ts @@ -72,7 +72,7 @@ function harness(over: Partial & Pick ({ valid: true }), cache: makeVerdictCache(), - hasBundle: async () => false, + isAdmitted: async () => false, admit: async ({ cid, verified }) => { admitted.push({ cid, verified }); }, @@ -117,13 +117,13 @@ describe("makeRootChaser", () => { expect(h.chaser.inFlight()).toBe(0); }); - it("reports every CID a checkpoint contained — including bundles already held, which never reach admit", async () => { + it("reports every CID a checkpoint contained — including bundles already admitted, which never reach admit", async () => { const winners = [bundle("0x1"), bundle("0x2")]; const { root, getBlock } = await checkpointOf(winners); - // Everything in this checkpoint is already in our store: the publisher's own bundle is + // Everything in this checkpoint is already in our winner-set: the publisher's own bundle is // always in exactly this position, so if `onCheckpointContents` fired from the admit path // it would never mention the one CID a publisher actually wants to hear about. - const h = harness({ getBlock, hasBundle: async () => true }); + const h = harness({ getBlock, isAdmitted: async () => true }); h.chaser.chase(root); await h.settle(); expect(h.admitted).toHaveLength(0); @@ -168,13 +168,13 @@ describe("makeRootChaser", () => { expect(fetched).toContain(root.toString()); }); - it("skips bundles we already hold without re-verifying", async () => { + it("skips bundles we already admitted without re-verifying", async () => { const winners = [bundle("0x1")]; const { root, getBlock } = await checkpointOf(winners); let verifies = 0; const h = harness({ getBlock, - hasBundle: async () => true, + isAdmitted: async () => true, verifyOffline: async () => { verifies++; return { valid: true }; @@ -188,6 +188,22 @@ describe("makeRootChaser", () => { expect(h.merged()).toBe(0); }); + it("admits a bundle whose blocks are already local but which was never admitted (issue #44's wedge)", async () => { + // The skip predicate is ADMISSION, not block presence. A node holding every block of a + // checkpoint locally — a persistent blockstore that outlived the state keyed on it — used + // to skip every bundle of every chase and could never converge again, silently, while + // peers served the votes. Here `getBlock` answers entirely from the local block map (no + // want leaves the node) and nothing is admitted yet: the chase must verify and admit. + const winners = [bundle("0x1"), bundle("0x2")]; + const { root, map } = await checkpointOf(winners); + const localOnly: RootChaserDeps["getBlock"] = async (cid) => map.get(cid.toString()); + const h = harness({ getBlock: localOnly, isAdmitted: async () => false }); + h.chaser.chase(root); + await h.settle(); + expect(h.admitted).toHaveLength(2); + expect(h.merged()).toBe(1); + }); + it("drops a liar's offline-invalid bundle before admit but keeps the honest one (per-bundle trust)", async () => { const good = bundle("0x1"); const forged = bundle("0xbad"); diff --git a/src/transport/chase.ts b/src/transport/chase.ts index 0f92c36..245d641 100644 --- a/src/transport/chase.ts +++ b/src/transport/chase.ts @@ -110,8 +110,18 @@ export interface RootChaserDeps { cache: VerdictCache; /** The gate's freshness guard (see gossip-validator.ts); omitted ⇒ no check. */ isEvaluableNow?: (bundle: VotesBundle) => Promise; - /** Skip bundles we already hold (their CID is in the store) without re-verifying. */ - hasBundle: (cid: CID) => Promise; + /** + * Have we already ADMITTED this bundle (is it in the CRDT / winner-set)? Admitted bundles are + * skipped without re-verifying. + * + * Admission, never block presence: this was wired to the blockstore once, and the two disagree + * exactly when it matters (issue #44). A node whose persistent blockstore holds a bundle's + * block but whose admission state was lost — a restart that dropped the snapshot, an eviction, + * a prune — skipped that bundle on EVERY chase and could never converge again, silently, while + * every peer served it. Re-verifying a locally-held block costs no network (the bytes are the + * ones we just decoded), so the safe side of the disagreement is cheap. + */ + isAdmitted: (cid: CID) => Promise; /** * Store an offline-valid bundle's block bytes and admit its CID into the CRDT (idempotent). * `verified: true` means a cached terminal verdict already covers the FULL pipeline (the @@ -178,7 +188,7 @@ export function makeRootChaser(deps: RootChaserDeps): RootChaser { verifyOffline, cache, isEvaluableNow, - hasBundle, + isAdmitted, admit, deferVerify, onMerged, @@ -257,7 +267,7 @@ export function makeRootChaser(deps: RootChaserDeps): RootChaser { const bytes = encodeBundle(bundle); const cid = await bundleCidForBytes(bytes); contained.push(cid); // recorded BEFORE the skip below — see onCheckpointContents - if (await hasBundle(cid)) continue; // already held — nothing to verify + if (await isAdmitted(cid)) continue; // already in the winner-set — nothing to verify const cached = cache.get(cid); if (cached) { if (!cached.valid) continue; // known bad — skip diff --git a/src/transport/integration/harness.ts b/src/transport/integration/harness.ts index af51fe5..3c37caa 100644 --- a/src/transport/integration/harness.ts +++ b/src/transport/integration/harness.ts @@ -200,9 +200,13 @@ export async function makeVoteNode(topic: string, options: VoteNodeOptions = {}) checkGates: async () => ({ kind: "leaf", leaf: 0, satisfied: true, score: 1n, penalize: false }) }; + // Mirrors the voter's admission map (`#checks`): what the CRDT knows, NOT what the blockstore + // holds — the chase's skip predicate must be admission (voter.ts `isAdmitted`, issue #44). + const admitted = new Set(); const admit = async ({ cid, bytes }: { cid: CID; bytes: Uint8Array; bundle: VotesBundle }): Promise => { await blockstore.put(cid, bytes); await crdt.merge([cid]); + admitted.add(cid.toString()); }; const acceptedBundles: Uint8Array[] = []; @@ -255,7 +259,7 @@ export async function makeVoteNode(topic: string, options: VoteNodeOptions = {}) }, verifyOffline: (bundle) => verifier.verifyOffline(bundle), cache, - hasBundle: (cid) => store.has(cid), + isAdmitted: async (cid) => admitted.has(cid.toString()), admit, // The harness runs the whole swapped pipeline in `verifyOffline` above, so nothing is // left deferred; the background verifier has its own unit tests. @@ -346,6 +350,7 @@ export async function makeVoteNode(topic: string, options: VoteNodeOptions = {}) const cid = await bundleCidForBytes(bytes); await blockstore.put(cid, bytes); await crdt.merge([cid]); + admitted.add(cid.toString()); }, stop: async () => { await transport.stop(); From 8559c13b970f30a8875fed7d593670765f4da22f Mon Sep 17 00:00:00 2001 From: Rinse Date: Thu, 3 Sep 2026 03:06:58 +0000 Subject: [PATCH 2/4] docs: record the two silent state-loss bugs (#44, #45) and their invariants --- AGENTS.md | 2 +- DESIGN.md | 6 +++++- README.md | 2 +- 3 files changed, 7 insertions(+), 3 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index c59f1f4..07b4785 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -1,6 +1,6 @@ # Agent Instructions for @bitsocial/pubsub-voting -Trustless pubsub voting library that runs on a host's shared libp2p/Helia node. The **engine and reactive facade are implemented** (schemas, canonical encoding, topic derivation, the CRDT, transport, verify, tally, chain reads, and the `PubsubVoter` / `Contest` (`createContest`) / `ContestVote` (`createContestVote`) facade; a contest is addressed by its full criteria document and **identity is per-ballot** — the `VoteSigner` is required by `createContestVote({ criteria, votes, signer })` and exposed as `ContestVote.signer`, never injected into the voter, so one client on the host's shared node publishes for as many wallets as the host holds keys for; there is no `readOnly` flag and no `ReadOnlyError` (a ballot that cannot sign has nothing to render — reading is `createContest`'s job). There is no manifest option, no `start()`, and the checkpoint fetch responder registers lazily on the first topic join, unregistering on the last leave). **Keeping a live vote from decaying is the consuming client's job** — the library publishes each vote once and exports `republishIntervalBuckets(criteria)` for the client to schedule its own refreshes; there is no republish scheduler and no library-side vote persistence (see DESIGN.md "Republishing is the client's job"). The **live-delta transport is implemented too**: the pubsub payload is a two-kind union — one inline bundle per message, or a tiny **root record** — with root-record checkpoint sync (on-demand encode, suppressed topic heartbeat, libp2p-fetch pull with the `MissingFetchError` construction guard, divergent roots chased through advertiser-seeded bitswap sessions — targeted wants at the peers whose last advertised root matches, one router query per root via the `maxProviders` headroom, broadcast `get` fallback, feature-detected `BlockstoreLike.createSession` — see DESIGN.md "Block pull") and a binary bundle block encoding. **Chain verification is two-stage** (see DESIGN.md "Background chain verification"): the live gossip gate still verifies fully before forwarding, but cold-join checkpoint bundles (and own publishes) admit on the synchronous offline checks and settle the deferred gate read + name resolution in the background verifier (`src/verify/background.ts`, batched per bucket via the rules' optional `evaluateMany`/multicall3 — 200 reads per `aggregate3`, and every chain client is wrapped by the voter in the cross-contest read coalescer `src/chain/coalescer.ts`, which merges parallel pinned-block `readContract`s into shared multicalls under one per-client in-flight budget of 2; free public RPCs throttle bigger bursts, see DESIGN.md) — the tally carries per-row `chainVerified`/`nameResolved` flags, checkpoints serve only fully verified bundles, the local verdict doubles as the publisher's rejection feedback (gossipsub has none: `publish()` preflights community names — `src/verify/name-preflight.ts`, `InvalidCommunityNameError`, resolver outages fail open — and a background eviction of an OWN bundle surfaces as `VoteEvictedError` on the vote's and contest's `error` events, flipping `publishingState` to `failed` post hoc), and the benchmark's cold joiner runs against a mock ETH gateway (`benchmark/rpc-gateway.ts`) so those reads are measured, not free — or, with `BENCH_RPC_URL=https://mainnet.base.org`, against the REAL Base mainnet at the real bucket sample block (the probe rule in `benchmark/signing.ts` keeps the reads real while admitting the bench's empty wallets). The **provider-record announcer is implemented** (`src/transport/announce/`, Node-only — the browser build swaps in an inert stub via the same package.json `browser`-field remap as `src/storage/`): a seeder sets `PubsubVoterOptions.httpRouterUrls` and the voter PUTs one **IPIP-0526 signed** batched record per router (the routers are assumed to run [pkcprotocol/pkc-http-router](https://github.com/pkcprotocol/pkc-http-router), which verifies signatures and 403s an unsigned or unstamped record — issue #38; the `Signature` is made by the node's own libp2p key over the payload bytes **as serialized into the body**, with a fresh `Payload.Timestamp` per tick, and a node exposing no key fails construction with `MissingPrivateKeyError`) (every joined contest's criteria CID + checkpoint root + chunk CIDs; addrs filtered client-side to public IP/DNS **plus the exactly-unspecified `0.0.0.0`/`::` sentinels the production router rewrites to the PUT's source IP** — and when the filter comes up empty the announcer synthesizes those sentinels from the node's listen ports, because libp2p never reports a wildcard and withholds unconfirmed public interface addrs pending AutoNAT; only a loopback-only node announces nothing) hourly, debounced 10 s on checkpoint changes / topic joins, and on `self:peer:update` — absent/empty means never announce (plain clients are not dialable), querying still rides the injected node's `libp2p.contentRouting` with no URLs in this library, and the cold-join bench exercises announce→router→discover→dial end-to-end (run.mjs hosts the mock router and reverse-tunnels it to the seeder; no hardcoded provider record). **Three pieces of state persist across restarts** under `PubsubVoterOptions.dataPath` (default `{cwd}/.bitsocial-pubsub-voting` on Node, IndexedDB in the browser, `false` → in-memory; the pkc-js storage stack — better-sqlite3 / localforage — lives in `src/storage/`, swapped by a package.json `browser`-field remap of one module): gate results keyed `(ruleHash, wallet, sampleBlock)` with a deterministic expiry purge, name resolutions under pkc-js's exact `NameResolutionCache` rule (LRU 5000, per-call max-age, 3600s at both verify call sites, failures never cached), and **each joined contest's checkpoint snapshot** (`src/checkpoint/snapshot.ts`, non-LRU, one atomic blob per topic in `{dataPath}/checkpoints.db`) — the node's own last fully-verified winner-set, written debounced 10 s on winner-set changes but SKIPPED while any deferred check is pending (else a restore-window write would clobber a good snapshot with a near-empty one), flushed on `leave()`, and reloaded at `join()` through the chase's decode+offline-verify+background-settle path, so a seeder restart with no other peer online no longer empties the tally (issue #14) — see DESIGN.md "Persistent caches". Tests and benchmarks pass `dataPath: false`; a stray `.bitsocial-pubsub-voting/` in the repo root means one forgot to. The **two-node and three-node-relay gossipsub integration tests are implemented** (`src/transport/integration/`, real `@libp2p/gossipsub` `17.0.1`, run via `npm run test:integration`, excluded from the unit `npm test`), and the **pkc-js host contract is pinned at both levels**: an offline unit test (`src/transport/pkc-js-host.test.ts`, in `npm test`) builds a stock `PKC({ libp2pJsClientsOptions })` instance and asserts its shared Helia node passes the voter's construction guards and the `adaptBlockstore` round-trip, and an e2e test (`integration/pkc-js-host.integration.test.ts`) runs three pkc-js-hosted voters through live publish → forward-gate verify → cold-join checkpoint pull on the stock host config. The host side has caught up: pkc-js registers gossipsub + `@libp2p/fetch` on the shared node as of `0.0.63` (pkc-js#183 closed 2026-07) and exposes the node through the public, semver-covered `Libp2pJsClient.heliaNode` accessor as of `0.0.72` (pkc-js#221 / PR #223) — the pinned devDependency and both host tests use it; the remaining pkc-js work is gossipsub score tuning. Read [DESIGN.md](./DESIGN.md) before changing anything. +Trustless pubsub voting library that runs on a host's shared libp2p/Helia node. The **engine and reactive facade are implemented** (schemas, canonical encoding, topic derivation, the CRDT, transport, verify, tally, chain reads, and the `PubsubVoter` / `Contest` (`createContest`) / `ContestVote` (`createContestVote`) facade; a contest is addressed by its full criteria document and **identity is per-ballot** — the `VoteSigner` is required by `createContestVote({ criteria, votes, signer })` and exposed as `ContestVote.signer`, never injected into the voter, so one client on the host's shared node publishes for as many wallets as the host holds keys for; there is no `readOnly` flag and no `ReadOnlyError` (a ballot that cannot sign has nothing to render — reading is `createContest`'s job). There is no manifest option, no `start()`, and the checkpoint fetch responder registers lazily on the first topic join, unregistering on the last leave). **Keeping a live vote from decaying is the consuming client's job** — the library publishes each vote once and exports `republishIntervalBuckets(criteria)` for the client to schedule its own refreshes; there is no republish scheduler and no library-side vote persistence (see DESIGN.md "Republishing is the client's job"). The **live-delta transport is implemented too**: the pubsub payload is a two-kind union — one inline bundle per message, or a tiny **root record** — with root-record checkpoint sync (on-demand encode, suppressed topic heartbeat, libp2p-fetch pull with the `MissingFetchError` construction guard, divergent roots chased through advertiser-seeded bitswap sessions — targeted wants at the peers whose last advertised root matches, one router query per root via the `maxProviders` headroom, broadcast `get` fallback, feature-detected `BlockstoreLike.createSession` — see DESIGN.md "Block pull") and a binary bundle block encoding. **Chain verification is two-stage** (see DESIGN.md "Background chain verification"): the live gossip gate still verifies fully before forwarding, but cold-join checkpoint bundles (and own publishes) admit on the synchronous offline checks and settle the deferred gate read + name resolution in the background verifier (`src/verify/background.ts`, batched per bucket via the rules' optional `evaluateMany`/multicall3 — 200 reads per `aggregate3`, and every chain client is wrapped by the voter in the cross-contest read coalescer `src/chain/coalescer.ts`, which merges parallel pinned-block `readContract`s into shared multicalls under one per-client in-flight budget of 2; free public RPCs throttle bigger bursts, see DESIGN.md) — the tally carries per-row `chainVerified`/`nameResolved` flags, checkpoints serve only fully verified bundles, the local verdict doubles as the publisher's rejection feedback (gossipsub has none: `publish()` preflights community names — `src/verify/name-preflight.ts`, `InvalidCommunityNameError`, resolver outages fail open — and a background eviction of an OWN bundle surfaces as `VoteEvictedError` on the vote's and contest's `error` events, flipping `publishingState` to `failed` post hoc), and the benchmark's cold joiner runs against a mock ETH gateway (`benchmark/rpc-gateway.ts`) so those reads are measured, not free — or, with `BENCH_RPC_URL=https://mainnet.base.org`, against the REAL Base mainnet at the real bucket sample block (the probe rule in `benchmark/signing.ts` keeps the reads real while admitting the bench's empty wallets). The **provider-record announcer is implemented** (`src/transport/announce/`, Node-only — the browser build swaps in an inert stub via the same package.json `browser`-field remap as `src/storage/`): a seeder sets `PubsubVoterOptions.httpRouterUrls` and the voter PUTs one **IPIP-0526 signed** batched record per router (the routers are assumed to run [pkcprotocol/pkc-http-router](https://github.com/pkcprotocol/pkc-http-router), which verifies signatures and 403s an unsigned or unstamped record — issue #38; the `Signature` is made by the node's own libp2p key over the payload bytes **as serialized into the body**, with a fresh `Payload.Timestamp` per tick, and a node exposing no key fails construction with `MissingPrivateKeyError`) (every joined contest's criteria CID + checkpoint root + chunk CIDs; addrs filtered client-side to public IP/DNS **plus the exactly-unspecified `0.0.0.0`/`::` sentinels the production router rewrites to the PUT's source IP** — and when the filter comes up empty the announcer synthesizes those sentinels from the node's listen ports, because libp2p never reports a wildcard and withholds unconfirmed public interface addrs pending AutoNAT; only a loopback-only node announces nothing) hourly, debounced 10 s on checkpoint changes / topic joins, and on `self:peer:update` — absent/empty means never announce (plain clients are not dialable), querying still rides the injected node's `libp2p.contentRouting` with no URLs in this library, and the cold-join bench exercises announce→router→discover→dial end-to-end (run.mjs hosts the mock router and reverse-tunnels it to the seeder; no hardcoded provider record). **Three pieces of state persist across restarts** under `PubsubVoterOptions.dataPath` (default `{cwd}/.bitsocial-pubsub-voting` on Node, IndexedDB in the browser, `false` → in-memory; the pkc-js storage stack — better-sqlite3 / localforage — lives in `src/storage/`, swapped by a package.json `browser`-field remap of one module): gate results keyed `(ruleHash, wallet, sampleBlock)` with a deterministic expiry purge, name resolutions under pkc-js's exact `NameResolutionCache` rule (LRU 5000, per-call max-age, 3600s at both verify call sites, failures never cached), and **each joined contest's checkpoint snapshot** (`src/checkpoint/snapshot.ts`, non-LRU, one atomic blob per topic in `{dataPath}/checkpoints.db`) — the node's own last fully-verified winner-set, written debounced 10 s on winner-set changes but SKIPPED while our view is knowably incomplete — any deferred check pending, or a restore backlog outstanding (either way a write would clobber a good snapshot with a near-empty one) — flushed on `leave()`, and reloaded at `join()` through the chase's decode+offline-verify+background-settle path, so a seeder restart with no other peer online no longer empties the tally (issue #14). **Only a provably bad blob is discarded**: the decode is the whole corruption test, while the per-bundle admission (which reads the chain) backlogs a transient failure for an exponential-backoff retry instead of deleting the blob — one try/catch over both meant an RPC blip during a seeder's 64-topic boot burst permanently discarded a topic's votes (issue #45) — and a discarded blob / incomplete restore / repeatedly-failing write each surface as a `SnapshotError` on the contest's `error` event (hence the view registers its engine listeners BEFORE `join()`). The chase's skip predicate is **admission** (`isAdmitted` → the engine's `#checks` map), never blockstore membership: a node whose persistent blockstore outlived its admission state skipped every bundle of every chase and could never converge again (issue #44). See DESIGN.md "Persistent caches" and "Block pull". Tests and benchmarks pass `dataPath: false`; a stray `.bitsocial-pubsub-voting/` in the repo root means one forgot to. The **two-node and three-node-relay gossipsub integration tests are implemented** (`src/transport/integration/`, real `@libp2p/gossipsub` `17.0.1`, run via `npm run test:integration`, excluded from the unit `npm test`), and the **pkc-js host contract is pinned at both levels**: an offline unit test (`src/transport/pkc-js-host.test.ts`, in `npm test`) builds a stock `PKC({ libp2pJsClientsOptions })` instance and asserts its shared Helia node passes the voter's construction guards and the `adaptBlockstore` round-trip, and an e2e test (`integration/pkc-js-host.integration.test.ts`) runs three pkc-js-hosted voters through live publish → forward-gate verify → cold-join checkpoint pull on the stock host config. The host side has caught up: pkc-js registers gossipsub + `@libp2p/fetch` on the shared node as of `0.0.63` (pkc-js#183 closed 2026-07) and exposes the node through the public, semver-covered `Libp2pJsClient.heliaNode` accessor as of `0.0.72` (pkc-js#221 / PR #223) — the pinned devDependency and both host tests use it; the remaining pkc-js work is gossipsub score tuning. Read [DESIGN.md](./DESIGN.md) before changing anything. **Who may vote is `criteria.gate`: one rule, or a boolean tree of them** (`{ rule }` / `{ all: [...] }` / `{ any: [...] }`, schema in `src/schema/criteria.ts`, fold in `src/rules/gate.ts`). A rule answers one question about one wallet and never sees the gate it sits in — composition is the document's, so "the Pass, or a moderator, and not banned" is bytes rather than a bespoke rule. The shape rules are canonicity, not style — a redundant spelling is a topic fork, so the schema rejects them: a branch needs **≥ 2 children**, may not **repeat a child** (canonical-bytes comparison, SIBLINGS ONLY — a rule may repeat across branches, because `{any:[{all:[A,B]},{all:[A,C]},{all:[B,C]}]}` ("any 2 of 3") has no repetition-free spelling), and may not nest a branch of its **own kind** (`{all:[{all:[A,B]},C]}` ≡ `{all:[A,B,C]}` — min and `some` are associative). This is not a normal form and does not claim to be: child order stays significant, absorption survives (`{all:[{any:[A,B]},A]}` ≡ `A`) and distribution is untouched, so it refuses the accidental fork rather than guaranteeing one spelling per meaning (DESIGN.md "What canonicity here does NOT claim"). Because a rule may repeat, a leaf's `ruleId` is NOT unique within a gate — `EligibilityCheck.leaf` (its position) is the render key, `ruleId` is which question it asks and which memo it shares. The depth/leaf caps are enforced on the RAW value before the recursive schema descends, or a pathological document stack-overflows out of `safeParse` before any cap fires. Normalizing instead of rejecting was refused: validation is a check, never a transform, or the shipped document stops being the document whose bytes are the topic. Child ORDER is significant by design (it is the lazy gate's evaluation order). The leaf is **wrapped** (`RuleRefSchema` is loose, so a rule option named `all`/`any` would make a bare leaf ambiguous), depth ≤ 4 and leaves ≤ 8 (attacker-supplied input every peer parses), and no rule names a chain at all: the contest names ONE (`criteria.bucketChainId`, a numeric chain id — there is no `requires.chains` map and no ticker) and every rule reads it, so "this leaf answered about another chain's history" is inexpressible rather than validated. The host's `ChainClientFactory` is keyed `{ chainId }`. Multi-chain gating is future work and is blocked on the semantics, not the field — see DESIGN.md "Open questions". The fold reads structure and discriminants, never rule identity: score = min across `all` / max across satisfied `any`; the **blame set** is the failures that EXPLAIN a refusal (a leaf failing inside a satisfied `any` blames nothing — never `checks.filter(failed)`); and **`penalize` is not a simple OR** — an `all` is attributable if ANY failing child is, an `any` only if EVERY child is. The inline gate evaluates lazily (it may under-report attributability, the fail-safe direction); the background verifier and `checkEligibility` score every leaf, the first because its batching axis is the rule (one `evaluateMany` per leaf per round), the second because naming each failure is the feature. A leaf whose chain read **throws** aborts a `verify` (a verdict reached on a read that never happened would evict a vote over an outage) but is folded as **unknown** by `checkEligibility`, which still answers if another branch decided the gate and re-throws only when it cannot — hence the three-valued fold. `Contest.checkEligibility` returns `{ eligible, score | error, checks, failures, gate }` — key rows by `leaf` (its position), never by `type` or `ruleId`, both of which may repeat within one gate. Leaves asking the SAME question (identical canonical ref) are evaluated ONCE, not once per position (`dedupeLeaves`): one batched call per distinct question in the background verifier, one shared promise per wallet inline. The verdict carries no gate score at all — nothing consumed it, and a min-across-`all` fold over unrelated rules means nothing. diff --git a/DESIGN.md b/DESIGN.md index e2f72a8..a92ae96 100644 --- a/DESIGN.md +++ b/DESIGN.md @@ -303,6 +303,8 @@ The record is **deliberately unauthenticated — it carries zero authority**: th **Block pull: directed bitswap at the peers that advertised the root.** The chunk blocks behind a root are fetched by CID via bitswap **from the connected peers whose record named that root** — they provably hold what they just advertised, so this is a directed transfer over existing connections, never provider-discovery toward an unreachable publisher (the failure mode that removed bitswap from the live path). When the fetch-protocol response supplied the **chunk index** (see "The root record"), the joiner has the chunk CIDs up front — verified against `root` locally — so it pulls all chunks **in parallel** and never fetches the root-manifest block at all; only a heartbeat-only divergence (no index) or an index that fails the local root check falls back to fetching the manifest first, then its chunks. Matching roots across peers make the pull cheap: download the DAG once, spreading chunk wants across the agreeing advertisers; a divergent root is pulled separately and unioned. Bulk serving rides bitswap's engine — per-peer ledgers, want-queue caps, requester-paced flow control — rather than a response burst a custom protocol would have to defend by itself. The chase side is bounded like everything else at the gate: a concurrency cap and a per-root deadline, so a spray of bogus roots queues bounded work; a bogus advertised root costs one directed attempt against its advertiser, and a peer whose roots repeatedly fail to resolve or verify is locally deprioritized. +**What "we already have this bundle" means: admission, never block presence.** A chased checkpoint's bundles are skipped when we already hold them — and "hold" must mean *admitted into the winner-set*, not *present in the blockstore*. The two disagree on exactly the node that most needs the chase: one whose **persistent** blockstore outlived the state keyed on it (a restart that lost the snapshot, an eviction, a prune). Keyed on the blockstore, such a node decoded every peer's checkpoint perfectly — every block resolved *locally* — skipped every bundle in it, admitted nothing, and was **permanently wedged** on that topic: restarting did not help, since the cold-start pull feeds the same chaser, and nothing on the path logs, so the node looked healthy while silently not counting votes every other peer served. That is issue #44, found after a production seeder served a divergent tally on two topics for 13 days. The predicate is now the engine's admission map (`isAdmitted`, the same map the snapshot restore consults). Erring the other way is cheap by construction: the bytes were just decoded, so a re-verify of a locally-held block costs no network — and it restores the re-verification the chase otherwise guarantees. + **Directed sessions, not broadcast wants (shipped 2026-07).** "Directed at the advertisers" is enforced mechanically, not just by connection topology: each chased root's blocks are pulled through `helia.blockstore.createSession(root, { providers, maxProviders: providers.length + 1 })`, seeded with that root's advertisers (`src/transport/chase.ts` `openSession` seam, wired in `src/client/voter.ts`). The plain per-block `blockstore.get` it replaced had two pathologies in the pinned helia (`@helia/bitswap`'s `want()`): an **unconditional routing query per block** — a `findProviders` fan-out to every HTTP router on the host node, pure waste while nothing announces chunk CIDs (a 63-contest directory cold load fired 100+) — and a **WANT-BLOCK broadcast to every connected peer**, chatter proportional to blocks × connections (the pathology pkc-js fixed in PR #191; the per-`get` `ProviderOptions.providers` shortcut is a dead end — `Bitswap.want()` never consults it in the pinned version). Provider selection is ranked by how provable "holds these bytes" is: (1) **the hint's sender** — cold start's `pull(peer)` has the `PeerId` in hand, and the gossip gate's `onRootRecord(record, from)` sender is no longer dropped at the voter's wiring; (2) **every still-connected peer whose last advertised root matches** the chased root, from a small bounded per-contest peer→root map fed by heartbeats, divergence responses, and cold-start pulls — the topic subscribers who provably converged on the exact state being pulled, so same-state subscribers seed the session organically over the heartbeat cycle; a late advertiser of an in-flight root joins the running session via `addPeer` instead of being dropped by the in-flight dedup; (3) **router-announced providers, in parallel**: the one-slot `maxProviders` headroom keeps the session's background `findNewProviders` running, so the routers are queried **once per root** (not once per block) — empty until seeders announce root/chunk provider records (see "Provider-record announces"), after which it automatically covers what subscribers cannot: a seeder pinning checkpoint blocks without joining the topic, and a joiner surviving its advertiser dropping mid-chase. The raw subscriber set is deliberately **not** seeded — seeding everyone reapproaches the broadcast being eliminated, and a peer on a *different* root only probabilistically holds the chunks — but it stays reachable as the safety net: with no advertiser in hand the session is never opened (an unseeded session fail-fasts with `InsufficientProvidersError` where a broadcast want succeeds via any connected topic peer), and on the first failed session want (a converged advertiser holds all blocks or none, so one miss marks the session dead) the chase falls back to the plain broadcast `get` for the remaining blocks. The seam is feature-detected end to end: `BlockstoreLike.createSession` is optional, `adaptBlockstore` exposes it only when the injected node can make sessions, and a plain blockstore (the unit tests' mocks) simply takes the broadcast path. No wire change. **Cold-join pull.** A joining peer subscribes and asks up to `k` peers for their current root record over the fetch protocol (heartbeats also arrive passively, but the active ask gives an immediate answer and `k` independent ones). Peers come from **two discovery sources raced concurrently, neither blocking the other**: (1) the peers gossipsub already knows subscribe to this topic (`getSubscribers`), and (2) the **providers of the criteria CID from the host's HTTP content router(s)** — `libp2p.contentRouting.findProviders`, the pkc-js delegated-Routing-V1-over-HTTP pattern (no DHT), where the criteria CID doubles as the routing key so a provider record for it means "I run this contest." Source (2) is what makes a *fresh* join fast: on a cold start `getSubscribers` is typically empty until subscription gossip propagates (several seconds over a WAN link), whereas the router names a provider directly — so the peer dials it and fetches the root record **immediately**, without waiting to *learn* who is subscribed. Source (1) additionally stays **armed on gossipsub's `subscription-change` for one heartbeat interval after the join**: the initial `getSubscribers` call is an instantaneous snapshot, and a joiner that dials a seeder and joins immediately (the normal browser boot order) races subscription gossip — at the instant of `join()` it sees zero subscribers, pulls nothing, and without this re-trigger would idle until the topic heartbeat (issue #15; measured 90+ s vs ~4 s). Each peer whose subscription becomes visible inside the window is pulled exactly once (the pull's `seen` set dedups, the per-peer fetch budget bounds concurrency), which closes the gap for router-less clients too; past the window the passive heartbeat covers divergence detection, so a long-lived node does not pay one fetch per churning subscriber forever. It then pulls the blocks for each distinct root via directed bitswap from its advertisers; recovers and checks every EIP-712 signature **offline**; **unions + LWW-merges** every offline-valid bundle into one winner-set (admitted *provisionally* — the deferred gate reads and name resolutions run batched in the background and confirm or evict after the fact, see "Background chain verification"); and keeps merging live gossip on top. **Restart is a snapshot reload plus a cold join**: the winner-set itself is in-memory, but each joined contest's last fully-verified checkpoint is persisted under `dataPath` and reloaded at `join()` through this same decode-and-verify path (see "Persistent caches", checkpoint snapshots) — so a seeder that restarts while no other peer is online serves its previous state immediately instead of losing it (the incident behind issue #14) — and this pull then still runs on top, converging the reloaded snapshot with the live topic. Root agreement across the `k` peers is a **download-once optimization, never an acceptance quorum**: when the roots match, pull the blocks once and know every advertiser held identical bytes; when they differ, pull *every* distinct root and union — a root served by only one peer is still pulled and merged, because trust is per-bundle (each signature self-verifies), not per-checkpoint (majority). A quorum-acceptance rule would hand a colluding majority of seeders exactly the vote-hiding power union exists to deny. Divergent checkpoints are never "chosen between" — they are unioned, so a seeder that omits a vote cannot hide it as long as one honest peer includes it, and one that injects a forgery is caught by the local signature check. **No checkpoint proves completeness** — union across independent peers is the guarantee, not any single root CID. @@ -397,7 +399,9 @@ Three pieces of engine state persist across restarts, under the voter's `dataPat - **Rule caches** (`(ruleNamespace, key, epoch) → value`, `src/rules/cache.ts`). The keys and the epochs are the rule's (see "What a rule owns, and what the pipeline owns"); what persists is whatever it memoized. A read pinned to a historical block is free of staleness by construction, so the v1 gate keys its fallback leg by that block and it stays valid forever, while its head leg is keyed by a coarse head window — which is what bounds how long a *negative* lingers. In the pinned case: a restart re-serves every settled gate read from the store instead of the RPC, which is what moves `START→ALL-VERIFIED` on a warm re-join. The namespace (sha256 of the canonical dag-cbor rule reference + chainId) is exactly the sharing boundary — two contests over one gate (a 5chan-style directory of boards on the same Pass) share each other's reads; different gates cannot collide. Caching sits at this *semantic* layer rather than raw RPC because the background verifier batches through multicall3: batched calldata varies with batch composition, so RPC-level keys would rarely repeat. Eviction is deterministic, not heuristic: a score at bucket B is only ever consulted while bundles from B are admissible, so an entry whose epoch has rolled past is unreachable, and only the rule can say when that is, so the rule purges (`RuleCache.purgeBelow`): the v1 gate drops its head-keyed entries as the window advances and leaves its pinned ones alone. The store's LRU bound is the backstop for anything a rule never purges. An in-memory FIFO front keeps the verify hot path off the store; a broken store read or write degrades to a live chain read, never an error. - **Name resolutions** (`src/verify/name-resolution-cache.ts`) — a port of pkc-js's `NameResolutionCache`, same rule so a host reasons about one policy across both libraries: LRU-bounded (5000, pkc-js's cap), **no stored TTL**, freshness is a per-call max-age (`Cache-Control: max-age` semantics), only successful resolutions are stored (a failure is retried live, never negatively cached), keyed `{name}::{resolverKey}::{sha256(provider)}`. Both verify call sites (inline forward-gate and background verifier) read with **max-age 3600s** — pkc-js's own background-resolution value — so the RPC cost of popular names is bounded while a re-pointed name is honored within at most an hour, strictly tighter than the republish-cycle window the re-point analysis in "Tally" already accepts. This cache exists here because pkc-js caches at *its* call sites, not inside the injected resolver — without it every verify pays a live registry read pkc-js would have served from cache. v1 resolves at head; once pinned-block resolution lands (see "Open questions"), a resolution becomes deterministic like a gate read and the max-age can go. -- **Checkpoint snapshots** (`src/checkpoint/snapshot.ts` + `SnapshotStorage` in `src/storage/`, one blob per topic): the node's **own last fully-verified checkpoint** — the root record's wire bytes plus the blocks it references — written **debounced** (10 s, announcer-style: the first change in a window arms the timer, the rest coalesce) on the same winner-set transition that dirties the on-demand encode and re-announces provider records, **flushed on `leave()`**, and **reloaded at `join()`** before any network activity. This is what makes a seeder restart survivable with no other peer online (issue #14): previously the state lived only in process memory and the host blockstore (typically in-memory), so a restart emptied the tally unless some online peer happened to re-advertise its checkpoint within the heartbeat window. The trust model is unchanged — the blob is the node's own previously-validated state, re-validated on load by the SAME pipeline as a chased remote checkpoint: block CIDs are re-derived from the bytes (a flipped bit self-invalidates), each bundle re-passes the offline signature/constraint checks, and the deferred gate reads + name resolutions ride the background verifier — which on a restart mostly hits the persisted gate-result store, so the re-verification reads disk, not RPC. Admission goes through `crdt.add` (never a merge-by-CID, which would re-read through the blockstore — fresh or already-holding, neither implies the CRDT knows the bundle). Two guards: the write is **skipped while any admitted bundle's deferred checks are pending** — the encoder serves only fully verified bundles, so writing mid-settlement (right after a restore, when every reloaded bundle is provisional) would clobber a good snapshot with a near-empty one, and a crash in that window would lose the very votes persistence exists to keep; the last settlement (or eviction) to land re-marks the state and re-arms the write. And a corrupt or version-mismatched blob is removed and the join proceeds empty — never blocked. Deliberately NOT an LRU store: eviction would silently corrupt a snapshot; each topic's blob is replaced whole (one atomic sqlite row / IndexedDB value), so a torn snapshot cannot exist. The cold-start pull still runs after the reload, so a stale snapshot self-heals by union with the live topic. +- **Checkpoint snapshots** (`src/checkpoint/snapshot.ts` + `SnapshotStorage` in `src/storage/`, one blob per topic): the node's **own last fully-verified checkpoint** — the root record's wire bytes plus the blocks it references — written **debounced** (10 s, announcer-style: the first change in a window arms the timer, the rest coalesce) on the same winner-set transition that dirties the on-demand encode and re-announces provider records, **flushed on `leave()`**, and **reloaded at `join()`** before any network activity. This is what makes a seeder restart survivable with no other peer online (issue #14): previously the state lived only in process memory and the host blockstore (typically in-memory), so a restart emptied the tally unless some online peer happened to re-advertise its checkpoint within the heartbeat window. The trust model is unchanged — the blob is the node's own previously-validated state, re-validated on load by the SAME pipeline as a chased remote checkpoint: block CIDs are re-derived from the bytes (a flipped bit self-invalidates), each bundle re-passes the offline signature/constraint checks, and the deferred gate reads + name resolutions ride the background verifier — which on a restart mostly hits the persisted gate-result store, so the re-verification reads disk, not RPC. Admission goes through `crdt.add` (never a merge-by-CID, which would re-read through the blockstore — fresh or already-holding, neither implies the CRDT knows the bundle). The write is **skipped while our view is knowably incomplete**, in either of the two ways it can be, because a write in that window overwrites a good blob with a lossy one: (a) while any admitted bundle's deferred checks are pending — the encoder serves only fully verified bundles, so writing mid-settlement (right after a restore, when every reloaded bundle is provisional) would clobber a good snapshot with a near-empty one, and a crash in that window would lose the very votes persistence exists to keep; the last settlement (or eviction) to land re-marks the state and re-arms the write — and (b) while a restore still carries a backlog (below). + + **Only a provably bad blob is discarded.** The decode is the whole of the corruption test, so it is the whole of what the discard covers: a truncated, mangled or version-mismatched blob is removed and the join proceeds empty — never blocked. Everything after it — the per-bundle admission, which reads the gating chain to check the bundle's bucket is reachable — is a different question, and a failure there says nothing about the blob. Both used to sit in one try/catch whose catch deleted, so a momentary RPC failure during a seeder's boot burst (64 topics restoring at once, one rate-limit window) **permanently discarded a topic's persisted votes** and the node came up empty on a topic every other peer still served (issue #45; the same incident as #44, and the selection effect is worth naming: topics whose voters republish on schedule self-heal within days, so the visible casualties are precisely the ballots of voters who went offline — the ones persistence exists to protect). Now the two are separate: admission runs per bundle, and a bundle that fails for a **transient** reason (a throwing head read, a head that has not yet reached the bundle's sample bucket) goes to a **backlog** that suppresses the snapshot write and is retried on an exponential backoff (30 s → 10 min) for as long as the contest stays joined — giving up would mean either abandoning those votes or suppressing this topic's writes forever, and a retry costs one head read against a memoized bucket. A bundle the **offline checks** refuse is dropped and its siblings still admit: that is a verdict on the bundle, and no retry changes it. None of this is silent any more, either — a discarded blob, an incomplete restore and a write that keeps failing each surface as a `SnapshotError` on the contest's `error` event (which is why the reactive view registers its engine listeners *before* the join: the restore runs inside it). Persistence stays best-effort — every one of these degrades to the pre-persistence behaviour, never to a broken join — but "best-effort" no longer means "invisible". Deliberately NOT an LRU store: eviction would silently corrupt a snapshot; each topic's blob is replaced whole (one atomic sqlite row / IndexedDB value), so a torn snapshot cannot exist. The cold-start pull still runs after the reload, so a stale snapshot self-heals by union with the live topic. ### Tally diff --git a/README.md b/README.md index e23afc1..c612bb4 100644 --- a/README.md +++ b/README.md @@ -94,7 +94,7 @@ Construction throws `MissingPubsubError`, `MissingBlockstoreError`, or `MissingF ```ts const contest = await voter.createContest({ criteria }); // criteria: the contest's full document (strictly validated here) contest.on("update", () => render(contest.tally)); // tally rides the object; recomputed before each emit -contest.on("error", (err) => showConnectivityWarning(err)); // tally chain read failed, the background verifier's RPC/resolver is down (retrying), or a deferred check evicted THIS wallet's own vote (VoteEvictedError) +contest.on("error", (err) => showConnectivityWarning(err)); // tally chain read failed, the background verifier's RPC/resolver is down (retrying), a deferred check evicted THIS wallet's own vote (VoteEvictedError), or this node's persisted checkpoint could not be kept (SnapshotError) await contest.update(); // join the topic, cold-start, begin emitting // const fresh = await contest.getTally(); // or force a fresh read, bypassing the cache // await contest.stop(); // leave the topic From 9c53631f1d6499a9c99dd0a4204f74c30da6c1ea Mon Sep 17 00:00:00 2001 From: Rinse Date: Thu, 3 Sep 2026 03:21:33 +0000 Subject: [PATCH 3/4] docs(bench): record the 2026-09-03 cold-join re-measure (parity with a master control) --- benchmark/RESULTS.md | 12 ++++++++++++ 1 file changed, 12 insertions(+) diff --git a/benchmark/RESULTS.md b/benchmark/RESULTS.md index 8141753..f8e1c97 100644 --- a/benchmark/RESULTS.md +++ b/benchmark/RESULTS.md @@ -141,6 +141,18 @@ within jitter: `START→TALLY` 2.46s / 2.48s / 2.48s / 2.66s / 4.80s and `START 3.09s / 3.27s / 6.08s for N=1…1000, with `router` 1.00s, `connect` 1.58–1.59s, `fetch` 0.51–1.09s and `gate-RPC` 3 / 3 / 3 / 3 / 9 unchanged. The table above is left at the 2026-08-09 numbers.* +*Re-measured 2026-09-03 (median of 3) with the **admission-keyed chase skip and the split snapshot +restore** (the chaser's skip predicate moved from blockstore membership to the engine's admission map; +the snapshot restore's decode and per-bundle admission got separate failure handling, with a retry +backlog). Neither touches the measured path here (the joiner starts empty, so nothing is skipped and +there is no snapshot to restore): `START→TALLY` 3.18s / 3.12s / 3.13s / 3.28s / 5.82s and +`START→VERIFIED` 3.79s / 3.73s / 3.74s / 3.89s / 7.09s for N=1…1000, against a back-to-back **master +control** on the same link window of 3.14s / 3.12s / 3.13s / 3.37s / 5.61s and 3.74s / 3.73s / 3.74s / +3.98s / 6.88s — every column within jitter (`connect` 1.91–1.97s vs 1.91–1.93s, `fetch` 0.84–1.79s vs +0.84–1.44s, `verify+merge` 0.31–1.98s vs 0.31–2.05s, `gate-RPC` 3 / 9 and inner `reads` N+1 in both). +Both runs again sit ~0.35s above the table's `connect` 1.56s, as the 2026-08-31 control did — the link +window, not the change. The table above is left at the 2026-08-09 numbers.* + *†The `N=10000` row (median of 3, one rep timed out on WAN jitter) is a separate single-contest run from the **previous baseline** (2026-07-08: instant fake chain, inline verification — before the mock ETH gateway and background chain verification existed, hence no gate-RPC/VERIFIED values) — a realistic From ab52bde4659a1efa8f2562390945fcc3ea46b207 Mon Sep 17 00:00:00 2001 From: Rinse Date: Thu, 3 Sep 2026 03:39:33 +0000 Subject: [PATCH 4/4] fix(client): drop per-bundle check state at the join-time prune too MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `computeTally`'s prune deleted each removed CID's `#checks` entry and own-bundle tracking; `join()`'s prune discarded the returned CIDs. The asymmetry was harmless while `#checks` was only bookkeeping, but it now decides two things: `#hasUnsettledChecks` (an orphaned pending entry suppresses every snapshot write for the contest) and the chase's `isAdmitted` (an orphan makes it skip a bundle we no longer hold). Both prunes now run one method. Reachable in one restart: `verifyOffline` is deliberately expiry-blind, so a snapshot older than the expiry window restores its bundles and the join-time prune drops them again moments later. The tally refresh's own prune usually cleans up first — but only when an update listener is registered, which a publish-driven join has none of, and the regression test joins that way. Reported by CodeRabbit on #49. --- src/client/voter.test.ts | 34 ++++++++++++++++++++++++++++++++-- src/client/voter.ts | 26 ++++++++++++++++++++------ 2 files changed, 52 insertions(+), 8 deletions(-) diff --git a/src/client/voter.test.ts b/src/client/voter.test.ts index 1ae8c9b..90bf10c 100644 --- a/src/client/voter.test.ts +++ b/src/client/voter.test.ts @@ -84,11 +84,11 @@ function advancingChains(currentBlock: () => bigint): ChainClientFactory { * test can observe the provisional (chainVerified: false) window deterministically before the * background verifier settles it. */ -function gatedChains(): { chains: ChainClientFactory; release: () => void } { +function gatedChains(head = 43200n): { chains: ChainClientFactory; release: () => void } { let release!: () => void; const gate = new Promise((resolve) => (release = resolve)); const client = { - getBlockNumber: async () => 43200n, + getBlockNumber: async () => head, getBlock: async () => ({ hash: `0x${"11".repeat(32)}` }), readContract: async () => { await gate; @@ -982,6 +982,36 @@ describe("checkpoint snapshot persistence (dataPath)", () => { } }); + it("drops the check state of bundles the join-time prune removes (no orphan left to block writes)", async () => { + // A snapshot older than the expiry window: `verifyOffline` is deliberately expiry-blind + // (verify/bundle.ts), so the restore admits its bundles and the join-time prune drops them + // again moments later. Their per-bundle check state must go with them — an orphan reads as + // "an admitted bundle whose deferred checks are pending", which suppresses every snapshot + // write for this contest, and as "already admitted" to the chase, for a bundle that is no + // longer in the winner-set. This join runs off `publish()`, with no update listener + // registered, so nothing kicks the tally refresh whose own prune would otherwise clean up + // first: the join-time prune is the only one that runs. + const dataPath = await tempDataPath(); + const signer = realSigner(); + const voterA = new PubsubVoter({ dataPath, helia: fakeHelia(), chains: countingChains().chains }); + const ballotA = await voterA.createContestVote({ criteria: bizCriteria(), votes: VOTE, signer }); + const { cid: decayed } = await ballotA.publish(); // sampled at bucket 1 (head 43200) + const contestA = await voterA.createContest({ criteria: bizCriteria() }); + await vi.waitFor(async () => expect((await contestA.getTally()).ranking[0]?.chainVerified).toBe(true)); + await voterA.destroy(); + + // Session 2 comes up 40 buckets later — past `voteExpiryBuckets: 30`. The gate reads never + // settle here, so nothing can quietly clean the orphan up after the fact. + const { chains, release } = gatedChains(43200n * 40n); + const voterB = new PubsubVoter({ dataPath, helia: fakeHelia(), chains }); + const contestB = await voterB.createContest({ criteria: bizCriteria() }); + const second = realSigner(SECOND_TEST_PRIVATE_KEY); + await (await voterB.createContestVote({ criteria: bizCriteria(), votes: OTHER_VOTE, signer: second })).publish(); + expect(contestB.checksFor(decayed)).toBeUndefined(); + await voterB.destroy(); + release(); + }); + it("runs the flush safely on the in-memory backend (`dataPath: false` — nothing to persist, nothing thrown)", async () => { const voter = new PubsubVoter({ dataPath: false, helia: fakeHelia(), chains: stubChains() }); await (await voter.createContestVote({ criteria: bizCriteria(), votes: VOTE, signer: fakeSigner() })).publish(); diff --git a/src/client/voter.ts b/src/client/voter.ts index 82dc6ec..50ff42f 100644 --- a/src/client/voter.ts +++ b/src/client/voter.ts @@ -1592,6 +1592,24 @@ class ContestEngine { if (i >= 0) this.#errorListeners.splice(i, 1); } + /** + * Drop decayed and superseded bundles from the CRDT **and** the per-bundle state keyed by + * their CIDs — one method because the two must not drift. `#checks` is not bookkeeping: it is + * what {@link #hasUnsettledChecks} reads (a leaked pending entry suppresses this contest's + * snapshot writes for a bundle that is no longer in the winner-set) and what the chase's + * `isAdmitted` reads (a leaked entry makes it skip a bundle we no longer hold). The join-time + * prune used to discard the removed CIDs — reachable in one restart: `verifyOffline` is + * deliberately expiry-blind (verify/bundle.ts), so a snapshot older than the expiry window + * restores its bundles and this prune removes them again moments later. + */ + async #pruneDecayed(): Promise { + for (const removed of await this.#crdt.prune(this.#currentBucketCache)) { + const key = removed.toString(); + this.#checks.delete(key); + this.#forgetOwnBundle(key); // expiry is decay, not an eviction — no error + } + } + /** Compute the current ranking fresh (refreshing the bucket + pruning when state is present). */ async computeTally(): Promise { // With state present, refresh the bucket so the tally's `current()` filters expiry against @@ -1600,11 +1618,7 @@ class ContestEngine { // chain reads" property). if (this.#crdt.nodeCount() > 0) { await this.#refreshBucket(); - for (const removed of await this.#crdt.prune(this.#currentBucketCache)) { - const key = removed.toString(); - this.#checks.delete(key); - this.#forgetOwnBundle(key); // expiry is decay, not an eviction — no error - } + await this.#pruneDecayed(); } return this.#tally.compute(); } @@ -1781,7 +1795,7 @@ class ContestEngine { // constant-weight tally" property. if (this.#crdt.nodeCount() > 0) { await this.#refreshBucket(); - await this.#crdt.prune(this.#currentBucketCache); + await this.#pruneDecayed(); } }