diff --git a/yarn-project/pxe/src/events/event_service.test.ts b/yarn-project/pxe/src/events/event_service.test.ts index 0b605e35879..bbd8f857422 100644 --- a/yarn-project/pxe/src/events/event_service.test.ts +++ b/yarn-project/pxe/src/events/event_service.test.ts @@ -36,6 +36,8 @@ describe('validateAndStoreEvents', () => { beforeEach(async () => { const store = await openTmpStore('test'); privateEventStore = new PrivateEventStore(store); + // Leave a change set open for the tests to operate under: every store operation requires one. + privateEventStore.beginChangeSet('test'); contractAddress = await AztecAddress.random(); recipient = await AztecAddress.random(); @@ -91,6 +93,7 @@ describe('validateAndStoreEvents', () => { await eventService.validateAndStoreEvents([request], recipient, map); await privateEventStore.commitChangeSet('test'); + privateEventStore.beginChangeSet('test'); } it('should throw when tx does not exist or has no effects', async () => { @@ -113,12 +116,7 @@ describe('validateAndStoreEvents', () => { await runStoreEvent({ eventContent: otherContent, eventCommitment: otherCommitment }); - const result = await privateEventStore.getPrivateEvents(eventSelector, { - contractAddress, - fromBlock: blockNumber, - toBlock: blockNumber + 1, - scopes: [recipient], - }); + const result = await readEvents(); expect(result.length).toEqual(0); expect(logger.warn).toHaveBeenCalledWith(expect.stringMatching(/commitment is not present in its tx/)); @@ -128,12 +126,7 @@ describe('validateAndStoreEvents', () => { // Commitment is legitimately present in the tx, but the provided content does not hash to it. await runStoreEvent({ eventContent: [Fr.random(), Fr.random()] }); - const result = await privateEventStore.getPrivateEvents(eventSelector, { - contractAddress, - fromBlock: blockNumber, - toBlock: blockNumber + 1, - scopes: [recipient], - }); + const result = await readEvents(); expect(result.length).toEqual(0); expect(logger.warn).toHaveBeenCalledWith(expect.stringMatching(/content does not hash to the provided commitment/)); @@ -143,12 +136,7 @@ describe('validateAndStoreEvents', () => { await runStoreEvent(); // I should be able to retrieve the private event I just saved using getPrivateEvents - const result = await privateEventStore.getPrivateEvents(eventSelector, { - contractAddress, - fromBlock: blockNumber, - toBlock: blockNumber + 1, - scopes: [recipient], - }); + const result = await readEvents(); expect(result.length).toEqual(1); expect(result[0].packedEvent).toEqual(eventContent); @@ -157,4 +145,13 @@ describe('validateAndStoreEvents', () => { function defaultValidationTxDataMap() { return new Map([[txEffect.txHash.toString(), validationTxData]]); } + + /** Reads the fixture's events through the change set the tests operate under. */ + function readEvents() { + return privateEventStore.getPrivateEvents( + eventSelector, + { contractAddress, fromBlock: blockNumber, toBlock: blockNumber + 1, scopes: [recipient] }, + 'test', + ); + } }); diff --git a/yarn-project/pxe/src/operation_lifecycle.test.ts b/yarn-project/pxe/src/operation_lifecycle.test.ts index 4a6532241b6..baab21264ca 100644 --- a/yarn-project/pxe/src/operation_lifecycle.test.ts +++ b/yarn-project/pxe/src/operation_lifecycle.test.ts @@ -23,6 +23,7 @@ describe('runOperation', () => { discarded = []; const recordingStore: StagedStore = { storeName: 'recording_store', + beginChangeSet: () => {}, commitChangeSet: id => { committed.push(id); return Promise.resolve(); @@ -141,6 +142,7 @@ describe('runOperation', () => { stagedStores: [ { storeName: 'undiscardable_store', + beginChangeSet: () => {}, commitChangeSet: () => Promise.resolve(), discardChangeSet: () => { throw new Error('cannot discard'); diff --git a/yarn-project/pxe/src/operation_queue.test.ts b/yarn-project/pxe/src/operation_queue.test.ts index 486f7649fda..274c934bd1e 100644 --- a/yarn-project/pxe/src/operation_queue.test.ts +++ b/yarn-project/pxe/src/operation_queue.test.ts @@ -29,6 +29,7 @@ describe('OperationQueue', () => { discarded = []; const recordingStore: StagedStore = { storeName: 'recording_store', + beginChangeSet: () => {}, commitChangeSet: id => { committed.push(id); return Promise.resolve(); diff --git a/yarn-project/pxe/src/pxe.test.ts b/yarn-project/pxe/src/pxe.test.ts index df3f705dd6c..a18430788f5 100644 --- a/yarn-project/pxe/src/pxe.test.ts +++ b/yarn-project/pxe/src/pxe.test.ts @@ -425,6 +425,8 @@ describe('PXE', () => { scope = await AztecAddress.random(); privateEventStore = new PrivateEventStore(kvStore); + // Leave a change set open for the tests to operate under: every store operation requires one. + privateEventStore.beginChangeSet('test'); }); let eventCounter = 0; diff --git a/yarn-project/pxe/src/pxe.ts b/yarn-project/pxe/src/pxe.ts index 41ce45b1e77..afc3e6ef09b 100644 --- a/yarn-project/pxe/src/pxe.ts +++ b/yarn-project/pxe/src/pxe.ts @@ -1366,15 +1366,8 @@ export class PXE { * Defaults to the latest known block to PXE + 1. * @returns - The packed events with block and tx metadata. */ - public async getPrivateEvents( - eventSelector: EventSelector, - filter: PrivateEventFilter, - ): Promise { - let anchorBlockNumber: BlockNumber; - - await this.operationQueue.runSynced(async ({ changeSetId, anchorBlockHeader }) => { - anchorBlockNumber = anchorBlockHeader.getBlockNumber(); - + public getPrivateEvents(eventSelector: EventSelector, filter: PrivateEventFilter): Promise { + return this.operationQueue.runSynced(async ({ changeSetId, anchorBlockHeader }) => { const contractFunctionSimulator = this.#getSimulatorForTx(); await this.contractSyncService.ensureContractSynced({ @@ -1394,16 +1387,15 @@ export class PXE { scopes: filter.scopes, triggeredBy: undefined, }); - }); - // anchorBlockNumber is set during the operation and fixed to whatever it is after a block sync - const sanitizedFilter = new PrivateEventFilterValidator(anchorBlockNumber!).validate(filter); + const sanitizedFilter = new PrivateEventFilterValidator(anchorBlockHeader.getBlockNumber()).validate(filter); - this.log.debug( - `Getting private events for ${sanitizedFilter.contractAddress.toString()} from ${sanitizedFilter.fromBlock} to ${sanitizedFilter.toBlock}`, - ); + this.log.debug( + `Getting private events for ${sanitizedFilter.contractAddress.toString()} from ${sanitizedFilter.fromBlock} to ${sanitizedFilter.toBlock}`, + ); - return this.privateEventStore.getPrivateEvents(eventSelector, sanitizedFilter); + return this.privateEventStore.getPrivateEvents(eventSelector, sanitizedFilter, changeSetId); + }); } /** diff --git a/yarn-project/pxe/src/storage/backwards_compatibility_tests/schema_tests.ts b/yarn-project/pxe/src/storage/backwards_compatibility_tests/schema_tests.ts index 6b33ce90703..a5147654d3e 100644 --- a/yarn-project/pxe/src/storage/backwards_compatibility_tests/schema_tests.ts +++ b/yarn-project/pxe/src/storage/backwards_compatibility_tests/schema_tests.ts @@ -436,6 +436,7 @@ export const SCHEMA_TESTS: readonly SchemaTest[] = [ const privateEventStore = new PrivateEventStore(kvStore); const changeSetId = 'fixture-change-set'; + privateEventStore.beginChangeSet(changeSetId); // Two (contract, selector) pairs and two block numbers so each multimap exhibits both a multi-value row // (contractA/selectorA → {e1, e2} and blockN1 → {e1, e2}) and a contrasting single-value row. diff --git a/yarn-project/pxe/src/storage/private_event_store/private_event_store.test.ts b/yarn-project/pxe/src/storage/private_event_store/private_event_store.test.ts index 28831992d6d..21059dcbee9 100644 --- a/yarn-project/pxe/src/storage/private_event_store/private_event_store.test.ts +++ b/yarn-project/pxe/src/storage/private_event_store/private_event_store.test.ts @@ -10,7 +10,7 @@ import { TxHash } from '@aztec-labs/stdlib/tx'; import type { PackedPrivateEvent } from '../../pxe.js'; import type { ChangeSetId } from '../staged_write_coordinator.js'; -import { PrivateEventStore } from './private_event_store.js'; +import { PrivateEventStore, type PrivateEventStoreFilter } from './private_event_store.js'; const getRandomMsgContent = () => { return [Fr.random(), Fr.random(), Fr.random()]; @@ -33,6 +33,8 @@ describe('PrivateEventStore', () => { beforeEach(async () => { kvStore = await openTmpStore('private_event_store_test'); privateEventStore = new PrivateEventStore(kvStore); + // Leave a change set open for the tests to operate under: every store operation requires one. + privateEventStore.beginChangeSet('test'); contractAddress = await AztecAddress.random(); scope = await AztecAddress.random(); msgContent = getRandomMsgContent(); @@ -74,10 +76,10 @@ describe('PrivateEventStore', () => { }, 'test', ); - await privateEventStore.commitChangeSet('test'); + await cycleChangeSet(); } - const events = await privateEventStore.getPrivateEvents(eventSelector, { + const events = await readEvents({ contractAddress, fromBlock: l2BlockNumber, toBlock: l2BlockNumber + 1, @@ -115,9 +117,9 @@ describe('PrivateEventStore', () => { metadata, 'test', ); - await privateEventStore.commitChangeSet('test'); + await cycleChangeSet(); - const events = await privateEventStore.getPrivateEvents(eventSelector, { + const events = await readEvents({ contractAddress, fromBlock: l2BlockNumber, toBlock: l2BlockNumber + 1, @@ -163,10 +165,10 @@ describe('PrivateEventStore', () => { }, 'test', ); - await privateEventStore.commitChangeSet('test'); + await cycleChangeSet(); } - const events = await privateEventStore.getPrivateEvents(eventSelector, { + const events = await readEvents({ contractAddress, fromBlock: l2BlockNumber, toBlock: l2BlockNumber + 1, @@ -232,10 +234,10 @@ describe('PrivateEventStore', () => { }, 'test', ); - await privateEventStore.commitChangeSet('test'); + await cycleChangeSet(); } - const events = await privateEventStore.getPrivateEvents(eventSelector, { + const events = await readEvents({ contractAddress, fromBlock: 150, toBlock: 150 + 100, @@ -281,10 +283,10 @@ describe('PrivateEventStore', () => { }, 'test', ); - await privateEventStore.commitChangeSet('test'); + await cycleChangeSet(); } - const events = await privateEventStore.getPrivateEvents(eventSelector, { + const events = await readEvents({ contractAddress, fromBlock: l2BlockNumber, toBlock: l2BlockNumber + 1, @@ -319,20 +321,28 @@ describe('PrivateEventStore', () => { 'test', ); - await privateEventStore.commitChangeSet('test'); + await cycleChangeSet(); const filter = { contractAddress, fromBlock: l2BlockNumber, toBlock: l2BlockNumber + 1 }; - const eventsScope1 = await privateEventStore.getPrivateEvents(eventSelector, { ...filter, scopes: [scope1] }); + const eventsScope1 = await privateEventStore.getPrivateEvents( + eventSelector, + { ...filter, scopes: [scope1] }, + 'test', + ); expect(eventsScope1).toHaveLength(1); expect(eventsScope1[0].packedEvent).toEqual(msgContent); - const eventsScope2 = await privateEventStore.getPrivateEvents(eventSelector, { ...filter, scopes: [scope2] }); + const eventsScope2 = await privateEventStore.getPrivateEvents( + eventSelector, + { ...filter, scopes: [scope2] }, + 'test', + ); expect(eventsScope2).toHaveLength(1); expect(eventsScope2[0].packedEvent).toEqual(msgContent); // Querying with both scopes returns the event once - const eventsBoth = await privateEventStore.getPrivateEvents(eventSelector, { + const eventsBoth = await readEvents({ ...filter, scopes: [scope1, scope2], }); @@ -340,7 +350,7 @@ describe('PrivateEventStore', () => { }); it('returns empty array when no events match criteria', async () => { - const events = await privateEventStore.getPrivateEvents(eventSelector, { + const events = await readEvents({ contractAddress, fromBlock: l2BlockNumber, toBlock: l2BlockNumber + 1, @@ -413,10 +423,10 @@ describe('PrivateEventStore', () => { }, 'test', ); - await privateEventStore.commitChangeSet('test'); + await cycleChangeSet(); } - const events = await privateEventStore.getPrivateEvents(eventSelector, { + const events = await readEvents({ contractAddress, fromBlock: 0, toBlock: 0 + 1000, @@ -481,10 +491,10 @@ describe('PrivateEventStore', () => { }, 'test', ); - await privateEventStore.commitChangeSet('test'); + await cycleChangeSet(); } - const events = await privateEventStore.getPrivateEvents(eventSelector, { + const events = await readEvents({ contractAddress, fromBlock: 0, toBlock: 1000, @@ -550,10 +560,10 @@ describe('PrivateEventStore', () => { }, 'test', ); - await privateEventStore.commitChangeSet('test'); + await cycleChangeSet(); } - const events = await privateEventStore.getPrivateEvents(eventSelector, { + const events = await readEvents({ contractAddress, fromBlock: 0, toBlock: 1000, @@ -598,16 +608,7 @@ describe('PrivateEventStore', () => { await kvStore.transactionAsync(() => privateEventStore.rollbackToBlock(9)); // Block 9 event survives; block 10 event is gone. - expect(await privateEventStore.eventIdsAtBlock(9)).toEqual([eventAt9.toString()]); - expect(await privateEventStore.eventIdsAtBlock(10)).toHaveLength(0); - - // getPrivateEvents confirms only the block-9 event is retrievable. - const events = await privateEventStore.getPrivateEvents(eventSelector, { - contractAddress, - fromBlock: 9, - toBlock: 11, - scopes: [scope], - }); + const events = await readEventsAfterRollback(9, 11); expect(events).toHaveLength(1); expect(events[0].l2BlockHash.equals(BLOCK_HASH_9)).toBe(true); }); @@ -625,9 +626,9 @@ describe('PrivateEventStore', () => { await kvStore.transactionAsync(() => privateEventStore.rollbackToBlock(9)); // Block 9 survives; both 10 and the non-contiguous 12 are swept. - expect(await privateEventStore.eventIdsAtBlock(9)).toEqual([eventAt9.toString()]); - expect(await privateEventStore.eventIdsAtBlock(10)).toHaveLength(0); - expect(await privateEventStore.eventIdsAtBlock(12)).toHaveLength(0); + const events = await readEventsAfterRollback(9, 13); + expect(events).toHaveLength(1); + expect(events[0].l2BlockHash.equals(BLOCK_HASH_9)).toBe(true); }); it('is idempotent — re-running an already-applied rollback is a no-op', async () => { @@ -642,8 +643,9 @@ describe('PrivateEventStore', () => { // Re-running over the already-truncated tail must not throw and must not change anything. await kvStore.transactionAsync(() => privateEventStore.rollbackToBlock(9)); - expect(await privateEventStore.eventIdsAtBlock(9)).toEqual([eventAt9.toString()]); - expect(await privateEventStore.eventIdsAtBlock(10)).toHaveLength(0); + const events = await readEventsAfterRollback(9, 11); + expect(events).toHaveLength(1); + expect(events[0].l2BlockHash.equals(BLOCK_HASH_9)).toBe(true); }); it('allows re-adding an event after rollback', async () => { @@ -651,24 +653,18 @@ describe('PrivateEventStore', () => { // siloed event commitment (this store's key). Rollback must leave no residue behind, so re-storing that // commitment succeeds with no key collision and the event becomes retrievable again. const commitment = Fr.random(); - const readBack = () => - privateEventStore.getPrivateEvents(eventSelector, { - contractAddress, - fromBlock: 9, - toBlock: 11, - scopes: [scope], - }); await storeEventAt(commitment, 10, BLOCK_HASH_10); await privateEventStore.commitChangeSet('test'); await kvStore.transactionAsync(() => privateEventStore.rollbackToBlock(9)); - expect(await readBack()).toHaveLength(0); + expect(await readEventsAfterRollback(9, 11)).toHaveLength(0); // Re-add the same commitment, as happens when the tx is re-included after the reorg. + privateEventStore.beginChangeSet('test'); await storeEventAt(commitment, 10, BLOCK_HASH_10); await privateEventStore.commitChangeSet('test'); - expect(await readBack()).toHaveLength(1); + expect(await readEventsAfterRollback(9, 11)).toHaveLength(1); }); it('handles rollback with no events to remove', async () => { @@ -679,133 +675,30 @@ describe('PrivateEventStore', () => { // Rolling back to a block above every stored event removes nothing. await kvStore.transactionAsync(() => privateEventStore.rollbackToBlock(20)); - expect(await privateEventStore.eventIdsAtBlock(10)).toEqual([eventAt10.toString()]); + expect(await readEventsAfterRollback(10, 11)).toHaveLength(1); }); - it('throws when rollback is called while staged writes are pending', async () => { - // Stage an event under a change set but never commit it, so the store still holds in-flight staged data. - await privateEventStore.storePrivateEventLog( + /** Reads in a change set opened and closed around the read: the store rejects a rollback while one is open. */ + async function readEventsAfterRollback(fromBlock: number, toBlock: number) { + privateEventStore.beginChangeSet('read-change-set'); + const events = await privateEventStore.getPrivateEvents( eventSelector, - randomness, - msgContent, - Fr.random(), - { - contractAddress, - scope, - txHash: TxHash.random(), - l2BlockNumber: BlockNumber(10), - l2BlockHash, - txIndexInBlock: 0, - eventIndexInTx: 0, - }, - 'uncommitted-change-set', + { contractAddress, fromBlock, toBlock, scopes: [scope] }, + 'read-change-set', ); - - await expect(kvStore.transactionAsync(() => privateEventStore.rollbackToBlock(0))).rejects.toThrow( - 'PXE private event store rollback is not allowed while staged writes are pending', - ); - - privateEventStore.discardChangeSet('uncommitted-change-set'); - - await expect(kvStore.transactionAsync(() => privateEventStore.rollbackToBlock(0))).resolves.not.toThrow(); - }); - }); - - describe('eventIdsAtBlock', () => { - it('returns the event id of an event stored at a given block', async () => { - await privateEventStore.storePrivateEventLog( - eventSelector, - randomness, - msgContent, - siloedEventCommitment, - { - contractAddress, - scope, - txHash, - l2BlockNumber, - l2BlockHash, - txIndexInBlock: 0, - eventIndexInTx: 0, - }, - 'test', - ); - await privateEventStore.commitChangeSet('test'); - - const ids = await privateEventStore.eventIdsAtBlock(l2BlockNumber); - expect(ids).toContain(siloedEventCommitment.toString()); - }); - - it('returns all event ids when multiple events are stored at the same block', async () => { - const siloedEventCommitment2 = Fr.random(); - - await privateEventStore.storePrivateEventLog( - eventSelector, - randomness, - msgContent, - siloedEventCommitment, - { - contractAddress, - scope, - txHash, - l2BlockNumber, - l2BlockHash, - txIndexInBlock: 0, - eventIndexInTx: 0, - }, - 'test', - ); - await privateEventStore.storePrivateEventLog( - eventSelector, - randomness, - getRandomMsgContent(), - siloedEventCommitment2, - { - contractAddress, - scope, - txHash: TxHash.random(), - l2BlockNumber, - l2BlockHash, - txIndexInBlock: 0, - eventIndexInTx: 1, - }, - 'test', - ); - await privateEventStore.commitChangeSet('test'); - - const ids = await privateEventStore.eventIdsAtBlock(l2BlockNumber); - expect(new Set(ids)).toEqual(new Set([siloedEventCommitment.toString(), siloedEventCommitment2.toString()])); - }); + privateEventStore.discardChangeSet('read-change-set'); + return events; + } }); describe('change-set', () => { - it('stages events without affecting committed storage', async () => { - const commitChangeSetId: ChangeSetId = 'commit-change-set'; + it('sees its own staged events before they are committed', async () => { const stagedChangeSetId: ChangeSetId = 'staged'; - - const committedEventRandomness = Fr.random(); const stagedEventRandomness = Fr.random(); - - // Store committed event - await privateEventStore.storePrivateEventLog( - eventSelector, - committedEventRandomness, - msgContent, - Fr.random(), - { - contractAddress, - scope, - txHash, - l2BlockNumber, - l2BlockHash, - txIndexInBlock: randomInt(100), - eventIndexInTx: randomInt(100), - }, - commitChangeSetId, - ); - await privateEventStore.commitChangeSet(commitChangeSetId); - - // Store staged event (not committed) const stagedMsgContent = getRandomMsgContent(); + + privateEventStore.discardChangeSet('test'); + privateEventStore.beginChangeSet(stagedChangeSetId); await privateEventStore.storePrivateEventLog( eventSelector, stagedEventRandomness, @@ -814,7 +707,7 @@ describe('PrivateEventStore', () => { { contractAddress, scope, - txHash: TxHash.random(), + txHash, l2BlockNumber, l2BlockHash, txIndexInBlock: randomInt(100), @@ -823,108 +716,89 @@ describe('PrivateEventStore', () => { stagedChangeSetId, ); - // With a fresh changeSetId, should only see committed event - const events = await privateEventStore.getPrivateEvents(eventSelector, { + const eventFilter = { contractAddress, fromBlock: l2BlockNumber, toBlock: l2BlockNumber + 1, scopes: [scope], - }); + }; + + const events = await privateEventStore.getPrivateEvents(eventSelector, eventFilter, stagedChangeSetId); expect(events).toHaveLength(1); - expect(events[0].packedEvent).toEqual(msgContent); + expect(events[0].packedEvent).toEqual(stagedMsgContent); }); - it('commit promotes staged events to main storage', async () => { - const stagedChangeSetId: ChangeSetId = 'staged'; - const stagedEventRandomness = Fr.random(); - const stagedMsgContent = getRandomMsgContent(); - + it('returns a committed event once when the open change set adds a scope to it', async () => { + const otherScope = await AztecAddress.random(); + const commitment = Fr.random(); + const metadata = { + contractAddress, + scope, + txHash, + l2BlockNumber, + l2BlockHash, + txIndexInBlock: 0, + eventIndexInTx: 0, + }; + + await privateEventStore.storePrivateEventLog(eventSelector, randomness, msgContent, commitment, metadata, 'test'); + await cycleChangeSet(); + + // Re-storing the committed event under a second scope stages it while it is also in the committed index. await privateEventStore.storePrivateEventLog( eventSelector, - stagedEventRandomness, - stagedMsgContent, - Fr.random(), - { - contractAddress, - scope, - txHash, - l2BlockNumber, - l2BlockHash, - txIndexInBlock: randomInt(100), - eventIndexInTx: randomInt(100), - }, - stagedChangeSetId, + randomness, + msgContent, + commitment, + { ...metadata, scope: otherScope }, + 'test', ); - await privateEventStore.commitChangeSet(stagedChangeSetId); - - // Now should see the event with a fresh changeSetId - const events = await privateEventStore.getPrivateEvents(eventSelector, { - contractAddress, - fromBlock: l2BlockNumber, - toBlock: l2BlockNumber + 1, - scopes: [scope], - }); - expect(events).toHaveLength(1); - expect(events[0].packedEvent).toEqual(stagedMsgContent); + const events = await privateEventStore.getPrivateEvents( + eventSelector, + { contractAddress, fromBlock: l2BlockNumber, toBlock: l2BlockNumber + 1, scopes: [otherScope] }, + 'test', + ); + expect(events).toEqual([expectedEvent]); }); - it('discardChangeSet removes staged events without affecting main', async () => { - const commitChangeSetId: ChangeSetId = 'commit-change-set'; - const stagedChangeSetId: ChangeSetId = 'staged'; - const committedEventRandomness = Fr.random(); - const stagedEventRandomness = Fr.random(); + it('excludes staged events belonging to another contract', async () => { + const otherContractAddress = await AztecAddress.random(); - // Store committed event await privateEventStore.storePrivateEventLog( eventSelector, - committedEventRandomness, + randomness, msgContent, Fr.random(), { - contractAddress, + contractAddress: otherContractAddress, scope, txHash, l2BlockNumber, l2BlockHash, - txIndexInBlock: randomInt(100), - eventIndexInTx: randomInt(100), + txIndexInBlock: 0, + eventIndexInTx: 0, }, - commitChangeSetId, + 'test', ); - await privateEventStore.commitChangeSet(commitChangeSetId); - // Store staged event (not committed) - const stagedMsgContent = getRandomMsgContent(); - await privateEventStore.storePrivateEventLog( + const events = await privateEventStore.getPrivateEvents( eventSelector, - stagedEventRandomness, - stagedMsgContent, - Fr.random(), - { - contractAddress, - scope, - txHash: TxHash.random(), - l2BlockNumber, - l2BlockHash, - txIndexInBlock: randomInt(100), - eventIndexInTx: randomInt(100), - }, - stagedChangeSetId, + { contractAddress, fromBlock: l2BlockNumber, toBlock: l2BlockNumber + 1, scopes: [scope] }, + 'test', ); - - // Discard change set - privateEventStore.discardChangeSet(stagedChangeSetId); - - // Should only see committed event - const events = await privateEventStore.getPrivateEvents(eventSelector, { - contractAddress, - fromBlock: l2BlockNumber, - toBlock: l2BlockNumber + 1, - scopes: [scope], - }); - expect(events).toHaveLength(1); - expect(events[0].packedEvent).toEqual(msgContent); + expect(events).toEqual([]); }); }); + + /** Reads through the change set the tests operate under. */ + function readEvents(filter: PrivateEventStoreFilter) { + return privateEventStore.getPrivateEvents(eventSelector, filter, 'test'); + } + + /** Commits the open change set and opens a fresh one, so everything staged so far lands in the committed index. */ + async function cycleChangeSet() { + await privateEventStore.commitChangeSet('test'); + privateEventStore.beginChangeSet('test'); + } }); diff --git a/yarn-project/pxe/src/storage/private_event_store/private_event_store.ts b/yarn-project/pxe/src/storage/private_event_store/private_event_store.ts index e39827937a6..5dc8257df25 100644 --- a/yarn-project/pxe/src/storage/private_event_store/private_event_store.ts +++ b/yarn-project/pxe/src/storage/private_event_store/private_event_store.ts @@ -2,15 +2,14 @@ import { BlockNumber } from '@aztec-labs/foundation/branded-types'; import { Fr } from '@aztec-labs/foundation/curves/bn254'; import { createLogger } from '@aztec-labs/foundation/log'; import { allToCompletion } from '@aztec-labs/foundation/promise'; -import { Semaphore } from '@aztec-labs/foundation/queue'; import type { AztecAsyncKVStore, AztecAsyncMap, AztecAsyncMultiMap } from '@aztec-labs/kv-store'; import type { EventSelector } from '@aztec-labs/stdlib/abi'; import type { AztecAddress } from '@aztec-labs/stdlib/aztec-address'; import type { InTx, TxHash } from '@aztec-labs/stdlib/tx'; import type { PackedPrivateEvent } from '../../pxe.js'; -import type { Rollbackable } from '../rollbackable.js'; -import type { ChangeSetId, StagedStore } from '../staged_write_coordinator.js'; +import { BaseStagingStore, type ReadonlyDb } from '../base_staging_store.js'; +import type { ChangeSetId } from '../staged_write_coordinator.js'; import { StoredPrivateEvent } from './stored_private_event.js'; export type PrivateEventStoreFilter = { @@ -42,33 +41,20 @@ type StoredEventBuffer = Buffer; * Append-only: events are never deleted during normal operation. Reorgs are handled by delete-on-prune, which removes * every event originating on a reorg'd block. */ -export class PrivateEventStore implements StagedStore, Rollbackable { - readonly storeName: string = 'private_event'; - - #store: AztecAsyncKVStore; - /** Actual private event log entries, keyed by siloedEventCommitment */ - #events: AztecAsyncMap; - /** Multi-map from contractAddress_eventSelector to siloedEventCommitment for efficient lookup */ - #eventsByContractAndEventSelector: AztecAsyncMultiMap; - /** Multi-map from block number to siloedEventCommitment, for delete-on-prune. */ - #eventsByBlockNumber: AztecAsyncMultiMap; - - /** changeSetId => eventId (event siloed nullifier) => StoredPrivateEvent */ - #eventsForChangeSet: Map>; - - /** Per-change-set locks to prevent concurrent writes from affecting each other. */ - #changeSetLocks: Map; - +export class PrivateEventStore extends BaseStagingStore { logger = createLogger('private_event_store'); constructor(store: AztecAsyncKVStore) { - this.#store = store; - this.#events = this.#store.openMap('private_event_logs'); - this.#eventsByContractAndEventSelector = this.#store.openMultiMap('events_by_contract_selector'); - this.#eventsByBlockNumber = this.#store.openMultiMap('events_by_block_number'); - - this.#eventsForChangeSet = new Map(); - this.#changeSetLocks = new Map(); + super({ + storeName: 'private_event', + store, + buildChangeSet: () => new Map(), + buildDb: db => ({ + events: db.openMap('private_event_logs'), + eventsByContractAndEventSelector: db.openMultiMap('events_by_contract_selector'), + eventsByBlockNumber: db.openMultiMap('events_by_block_number'), + }), + }); } /** @@ -91,45 +77,42 @@ export class PrivateEventStore implements StagedStore, Rollbackable { metadata: PrivateEventMetadata, changeSetId: ChangeSetId, ) { - return this.#withChangeSetLock(changeSetId, () => - this.#store.transactionAsync(async () => { - const { contractAddress, scope, txHash, l2BlockNumber, l2BlockHash, txIndexInBlock, eventIndexInTx } = metadata; - const eventId = siloedEventCommitment.toString(); + return this.withChangeSetAndDb(changeSetId, async (changeSet, db) => { + const { contractAddress, scope, txHash, l2BlockNumber, l2BlockHash, txIndexInBlock, eventIndexInTx } = metadata; + const eventId = siloedEventCommitment.toString(); - this.logger.verbose('storing private event log (staged)', { - eventId, - contractAddress, - scope, - msgContent, - l2BlockNumber, - }); + this.logger.verbose('storing private event log (staged)', { + eventId, + contractAddress, + scope, + msgContent, + l2BlockNumber, + }); - const existing = await this.#readEvent(eventId, changeSetId); - - if (existing) { - // If we already stored this event, we still want to make sure to track it for the given scope - existing.addScope(scope.toString()); - this.#writeEvent(eventId, existing, changeSetId); - } else { - this.#writeEvent( - eventId, - new StoredPrivateEvent( - randomness, - msgContent, - l2BlockNumber, - l2BlockHash, - txHash, - txIndexInBlock, - eventIndexInTx, - contractAddress, - eventSelector, - new Set([scope.toString()]), - ), - changeSetId, - ); - } - }), - ); + const existing = await this.#readEvent(changeSet, db, eventId); + + if (existing) { + // If we already stored this event, we still want to make sure to track it for the given scope + existing.addScope(scope.toString()); + changeSet.set(eventId, existing); + } else { + changeSet.set( + eventId, + new StoredPrivateEvent( + randomness, + msgContent, + l2BlockNumber, + l2BlockHash, + txHash, + txIndexInBlock, + eventIndexInTx, + contractAddress, + eventSelector, + new Set([scope.toString()]), + ), + ); + } + }); } /** @@ -140,27 +123,40 @@ export class PrivateEventStore implements StagedStore, Rollbackable { * fromBlock: The block number to search from (inclusive). * toBlock: The block number to search upto (exclusive). * scope: - The addresses that decrypted the logs. + * @param changeSetId - the change set to read staged data from. * @returns - The event log contents, augmented with metadata about the transaction and block in which the event was * included. */ public getPrivateEvents( eventSelector: EventSelector, filter: PrivateEventStoreFilter, + changeSetId: ChangeSetId, ): Promise { - return this.#store.transactionAsync(async () => { + return this.withChangeSetAndDb(changeSetId, async (changeSet, db) => { const key = this.#keyFor(filter.contractAddress, eventSelector); const targetScopes = new Set(filter.scopes.map(s => s.toString())); - // Map from eventId to the promise that reads the event buffer. + // Map from eventId to the promise that reads the event. // We start reads during iteration to keep DB requests pending and avoid IndexedDB auto-commit. - const eventReadPromises: Map> = new Map(); + const eventReadPromises: Map> = new Map(); - for await (const eventId of this.#eventsByContractAndEventSelector.getValuesAsync(key)) { - eventReadPromises.set(eventId, this.#events.getAsync(eventId)); + // Committed events indexed by contract address and event selector + for await (const eventId of db.eventsByContractAndEventSelector.getValuesAsync(key)) { + eventReadPromises.set(eventId, this.#readEvent(changeSet, db, eventId)); } + // Staged events have no row in the committed index yet, so add them here. Skip any the loop above already + // picked up. + [...changeSet.entries()] + .filter( + ([eventId, stagedEvent]) => + !eventReadPromises.has(eventId) && + this.#keyFor(stagedEvent.contractAddress, stagedEvent.eventSelector) === key, + ) + .forEach(([eventId, stagedEvent]) => eventReadPromises.set(eventId, Promise.resolve(stagedEvent))); + const eventIds = [...eventReadPromises.keys()]; - const eventBuffers = await allToCompletion([...eventReadPromises.values()]); + const storedEvents = await allToCompletion([...eventReadPromises.values()]); const events: Array<{ l2BlockNumber: number; @@ -171,18 +167,16 @@ export class PrivateEventStore implements StagedStore, Rollbackable { for (let i = 0; i < eventIds.length; i++) { const eventId = eventIds[i]; - const eventBuffer = eventBuffers[i]; + const storedPrivateEvent = storedEvents[i]; - // Defensive, if it happens, there's a problem with how we're handling #eventsByContractAndEventSelector - if (!eventBuffer) { + // Defensive, if it happens, there's a problem with how we're handling db.eventsByContractAndEventSelector + if (!storedPrivateEvent) { this.logger.verbose( `EventId ${eventId} does not exist in main index but it is referenced from contract event selector index`, ); continue; } - const storedPrivateEvent = StoredPrivateEvent.fromBuffer(eventBuffer); - // Filter by block range if (storedPrivateEvent.l2BlockNumber < filter.fromBlock || storedPrivateEvent.l2BlockNumber >= filter.toBlock) { continue; @@ -227,85 +221,46 @@ export class PrivateEventStore implements StagedStore, Rollbackable { }); } - /** Returns the ids (siloed event commitments) of all events emitted at the given block number. Used by delete-on-prune. */ - public async eventIdsAtBlock(blockNumber: number): Promise { - const eventIds: string[] = []; - for await (const eventId of this.#eventsByBlockNumber.getValuesAsync(blockNumber)) { - eventIds.push(eventId); - } - return eventIds; - } - /** - * Rolls the store back to `toBlock`: deletes every event anchored to a block strictly above it, as if nothing past - * that block height ever happened. Used by the reorg (`chain-pruned`) path to truncate the orphaned tail. Scanning - * from `toBlock + 1` upward covers everything above the rollback target without needing to know the chain tip. - * - * Must be called inside a transaction owned by the caller (it issues no `transactionAsync` of its own, the reorg path - * wraps it together with the anchor update, and IndexedDB has no nested transactions). Throws if any change set has - * uncommitted staged writes, since rolling back mid-change-set could later re-introduce events anchored to deleted - * blocks. + * Deletes every event originating on a block strictly above `toBlock`, as if nothing past that block height ever + * happened, truncating the orphaned tail on a reorg. Scanning from `toBlock + 1` upward covers everything above the + * rollback target without needing to know the chain tip. */ - public async rollbackToBlock(toBlock: number): Promise { - if (this.#eventsForChangeSet.size > 0) { - throw new Error('PXE private event store rollback is not allowed while staged writes are pending'); - } + protected async applyRollback(toBlock: number, db: PrivateEventStoreDb): Promise { // Snapshot before mutating so we never delete from the multimap we are iterating. - const orphaned: { block: number; eventId: string }[] = []; - for await (const [block, eventId] of this.#eventsByBlockNumber.entriesAsync({ start: toBlock + 1 })) { + const orphaned: { block: BlockNum; eventId: EventId }[] = []; + for await (const [block, eventId] of db.eventsByBlockNumber.entriesAsync({ start: toBlock + 1 })) { orphaned.push({ block, eventId }); } let removedCount = 0; for (const { block, eventId } of orphaned) { - const buf = await this.#events.getAsync(eventId); + const buf = await db.events.getAsync(eventId); if (!buf) { throw new Error(`Event not found for eventId ${eventId}`); } const stored = StoredPrivateEvent.fromBuffer(buf); - await this.#events.delete(eventId); - await this.#eventsByContractAndEventSelector.deleteValue( + await db.events.delete(eventId); + await db.eventsByContractAndEventSelector.deleteValue( this.#keyFor(stored.contractAddress, stored.eventSelector), eventId, ); - await this.#eventsByBlockNumber.deleteValue(block, eventId); + await db.eventsByBlockNumber.deleteValue(block, eventId); removedCount++; } this.logger.verbose('rolled back private events', { removedCount, toBlock }); } - /** - * Commits in-memory staged data to persistent storage. - * - * Called by StagedWriteCoordinator when an operation completes successfully. - * - * Note: StagedWriteCoordinator wraps all commits in a single transaction, so we don't need our own transactionAsync - * here (and using one would throw on IndexedDB as it does not support nested txs). - * - * @param changeSetId - The changeSetId identifying which staged data to commit - */ - async commitChangeSet(changeSetId: ChangeSetId): Promise { - // Note: Don't use #withChangeSetLock here - commit runs within StagedWriteCoordinator's transactionAsync, - // and awaiting the lock would create a microtask boundary with no pending DB request, - // causing IndexedDB to auto-commit the transaction. - for (const [eventId, entry] of this.#getEventsForChangeSet(changeSetId).entries()) { + protected async flushChangeSet(changeSet: PrivateEventStoreChangeSet, db: PrivateEventStoreDb): Promise { + for (const [eventId, entry] of changeSet.entries()) { const lookupKey = this.#keyFor(entry.contractAddress, entry.eventSelector); this.logger.verbose('storing private event log', { eventId, lookupKey }); await allToCompletion([ - this.#events.set(eventId, entry.toBuffer()), - this.#eventsByContractAndEventSelector.set(lookupKey, eventId), - this.#eventsByBlockNumber.set(entry.l2BlockNumber, eventId), + db.events.set(eventId, entry.toBuffer()), + db.eventsByContractAndEventSelector.set(lookupKey, eventId), + db.eventsByBlockNumber.set(entry.l2BlockNumber, eventId), ]); } - - this.#clearChangeSetData(changeSetId); - } - - /** - * Discards in-memory staged data without persisting it. - */ - discardChangeSet(changeSetId: ChangeSetId): void { - this.#clearChangeSetData(changeSetId); } /** @@ -313,71 +268,38 @@ export class PrivateEventStore implements StagedStore, Rollbackable { * * Returns undefined if the event does not exist in the store overall. */ - async #readEvent(eventId: string, changeSetId: ChangeSetId): Promise { + async #readEvent( + changeSet: PrivateEventStoreChangeSet, + db: ReadonlyDb, + eventId: EventId, + ): Promise { // Always issue DB read to keep IndexedDB transaction alive (they auto-commit when a new micro-task starts and there // are no pending read requests). The staged value still takes precedence if it exists. - const buffer = await this.#events.getAsync(eventId); - const eventForChangeSet = this.#getEventsForChangeSet(changeSetId).get(eventId); + const buffer = await db.events.getAsync(eventId); + const eventForChangeSet = changeSet.get(eventId); return eventForChangeSet ?? (buffer ? StoredPrivateEvent.fromBuffer(buffer) : undefined); } /** - * Writes an event to in-memory staged data. + * Returns a string key based on @param contractAddress and @param eventSelector. * - * Writes are only allowed in a change set context. Events modified while staged will only be persisted when `commit` - * is called. + * The returned key is meant to be used when interacting with the db.eventsByContractAndEventSelector index. */ - #writeEvent(eventId: string, entry: StoredPrivateEvent, changeSetId: ChangeSetId) { - this.#getEventsForChangeSet(changeSetId).set(eventId, entry); + #keyFor(contractAddress: AztecAddress, eventSelector: EventSelector): ContractAndSelectorKey { + return `${contractAddress.toString()}_${eventSelector.toString()}`; } +} - /** - * Get in-memory data only visible to @param changeSetId - */ - #getEventsForChangeSet(changeSetId: ChangeSetId): Map { - let eventsForChangeSet = this.#eventsForChangeSet.get(changeSetId); - if (eventsForChangeSet === undefined) { - eventsForChangeSet = new Map(); - this.#eventsForChangeSet.set(changeSetId, eventsForChangeSet); - } - return eventsForChangeSet; - } +/** A change set's staged data: the events it has stored, keyed by their siloed event commitment. */ +type PrivateEventStoreChangeSet = Map; - /** - * Clear data structures supporting a specific change set. - */ - #clearChangeSetData(changeSetId: ChangeSetId) { - this.#eventsForChangeSet.delete(changeSetId); - this.#changeSetLocks.delete(changeSetId); - } +type PrivateEventStoreDb = { + /** Actual private event log entries, keyed by siloedEventCommitment */ + events: AztecAsyncMap; - /** - * Ensures a function can only run once it acquires a unique per-change-set lock, and handles proper lock release - * after it runs. - * - * This primitive allows concurrent writes on this store without risking data corruption due to unsound write - * interleaving. - */ - async #withChangeSetLock(changeSetId: ChangeSetId, fn: () => Promise): Promise { - let lock = this.#changeSetLocks.get(changeSetId); - if (!lock) { - lock = new Semaphore(1); - this.#changeSetLocks.set(changeSetId, lock); - } - await lock.acquire(); - try { - return await fn(); - } finally { - lock.release(); - } - } + /** Multi-map from contractAddress_eventSelector to siloedEventCommitment for efficient lookup */ + eventsByContractAndEventSelector: AztecAsyncMultiMap; - /** - * Returns a string key based on @param contractAddress and @param eventSelector. - * - * The returned key is meant to be used when interacting with index #eventsByContractAndEventSelector. - */ - #keyFor(contractAddress: AztecAddress, eventSelector: EventSelector): string { - return `${contractAddress.toString()}_${eventSelector.toString()}`; - } -} + /** Multi-map from block number to siloedEventCommitment, for delete-on-prune. */ + eventsByBlockNumber: AztecAsyncMultiMap; +}; diff --git a/yarn-project/pxe/src/storage/staged_write_coordinator.test.ts b/yarn-project/pxe/src/storage/staged_write_coordinator.test.ts index 7c3b0d6fe5b..d4f375783c9 100644 --- a/yarn-project/pxe/src/storage/staged_write_coordinator.test.ts +++ b/yarn-project/pxe/src/storage/staged_write_coordinator.test.ts @@ -109,6 +109,7 @@ describe('StagedWriteCoordinator', () => { const committed: { changeSetId: ChangeSetId; inTransaction: boolean }[] = []; const mockStore = makeStagedStore({ storeName: 'mock_store', + beginChangeSet: () => {}, commitChangeSet: changeSetId => { committed.push({ changeSetId, inTransaction }); return Promise.resolve(); @@ -155,6 +156,7 @@ describe('StagedWriteCoordinator', () => { const discardChangeSetMock = jest.fn<() => void>(); const mockStore = makeStagedStore({ storeName: 'mock_store', + beginChangeSet: () => {}, commitChangeSet: commitMock, discardChangeSet: discardChangeSetMock, }); @@ -235,6 +237,7 @@ describe('StagedWriteCoordinator', () => { const discardChangeSetMock = jest.fn<() => void>(); const mockStore = makeStagedStore({ storeName: 'mock_store', + beginChangeSet: () => {}, commitChangeSet: commitMock, discardChangeSet: discardChangeSetMock, }); @@ -247,6 +250,11 @@ describe('StagedWriteCoordinator', () => { /** A staged store that does nothing on commit and discard, so a test only spells out the part it exercises. */ function makeStagedStore(overrides: Partial & Pick): StagedStore { - return { commitChangeSet: () => Promise.resolve(), discardChangeSet: () => {}, ...overrides }; + return { + beginChangeSet: () => {}, + commitChangeSet: () => Promise.resolve(), + discardChangeSet: () => {}, + ...overrides, + }; } }); diff --git a/yarn-project/pxe/src/storage/staged_write_coordinator.ts b/yarn-project/pxe/src/storage/staged_write_coordinator.ts index 82fb90d8820..9ecabb40a11 100644 --- a/yarn-project/pxe/src/storage/staged_write_coordinator.ts +++ b/yarn-project/pxe/src/storage/staged_write_coordinator.ts @@ -24,12 +24,9 @@ export interface StagedStore { /** * Notifies the store that a change set has been opened. * - * TODO: make it required once every staged store extends `BaseStagingStore`. It is optional only while - * they migrate to per-change-set staging. - * * @param changeSetId - The change set identifier */ - beginChangeSet?(changeSetId: ChangeSetId): void; + beginChangeSet(changeSetId: ChangeSetId): void; /** * Commits staged data to persistent storage. Will be called within a db transaction for atomicity, alongside the @@ -168,7 +165,7 @@ export class StagedWriteCoordinator { const begunStores: StagedStore[] = []; try { for (const store of this.#stagedStores.values()) { - store.beginChangeSet?.(changeSetId); + store.beginChangeSet(changeSetId); begunStores.push(store); } } catch (err) { diff --git a/yarn-project/txe/src/oracle/interfaces.ts b/yarn-project/txe/src/oracle/interfaces.ts index f4fefd29fb8..773a87ee6a3 100644 --- a/yarn-project/txe/src/oracle/interfaces.ts +++ b/yarn-project/txe/src/oracle/interfaces.ts @@ -87,7 +87,12 @@ export interface ITxeExecutionOracle { setAuthorizeAllUtilityCallTargets(authorizeAll: boolean): void; getLastBlockTimestamp(): Promise; getLastTxEffects(): Promise; - getPrivateEvents(selector: EventSelector, contractAddress: AztecAddress, scope: AztecAddress): Promise; + getPrivateEvents( + selector: EventSelector, + contractAddress: AztecAddress, + scope: AztecAddress, + changeSetId: ChangeSetId, + ): Promise; privateCallNewFlow( from: AztecAddress | undefined, targetContractAddress: AztecAddress, diff --git a/yarn-project/txe/src/oracle/txe_oracle_top_level_context.ts b/yarn-project/txe/src/oracle/txe_oracle_top_level_context.ts index fcd2430cde8..d9ae2c5cfc1 100644 --- a/yarn-project/txe/src/oracle/txe_oracle_top_level_context.ts +++ b/yarn-project/txe/src/oracle/txe_oracle_top_level_context.ts @@ -252,14 +252,23 @@ export class TXEOracleTopLevelContext implements IMiscOracle, ITxeExecutionOracl }); } - async getPrivateEvents(selector: EventSelector, contractAddress: AztecAddress, scope: AztecAddress) { + async getPrivateEvents( + selector: EventSelector, + contractAddress: AztecAddress, + scope: AztecAddress, + changeSetId: ChangeSetId, + ) { const events = ( - await this.privateEventStore.getPrivateEvents(selector, { - contractAddress, - scopes: [scope], - fromBlock: 0, - toBlock: (await this.getLastBlockNumber()) + 1, - }) + await this.privateEventStore.getPrivateEvents( + selector, + { + contractAddress, + scopes: [scope], + fromBlock: 0, + toBlock: (await this.getLastBlockNumber()) + 1, + }, + changeSetId, + ) ).map(e => e.packedEvent); if (events.length > MAX_PRIVATE_EVENTS_PER_TXE_QUERY) { diff --git a/yarn-project/txe/src/txe_session.ts b/yarn-project/txe/src/txe_session.ts index 6cd4564f1b0..469100b61d7 100644 --- a/yarn-project/txe/src/txe_session.ts +++ b/yarn-project/txe/src/txe_session.ts @@ -660,7 +660,7 @@ export class TXESession implements TXESessionStateHandler { await handler.syncContractNonOracleMethod(contractAddress, scope, this.currentChangeSetId); // Cycle the change set to commit the stores after the contract sync. await this.cycleOperation(); - return handler.getPrivateEvents(selector, contractAddress, scope); + return handler.getPrivateEvents(selector, contractAddress, scope, this.currentChangeSetId); } private handlerAsTxe(): ITxeExecutionOracle {