Skip to content

Commit 0dcfe44

Browse files
committed
feat(fast-inbox): archiver syncs Inbox buckets (A-1379)
1 parent 7f41bee commit 0dcfe44

17 files changed

Lines changed: 712 additions & 37 deletions

yarn-project/archiver/src/archiver-sync.test.ts

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -259,6 +259,36 @@ describe('Archiver Sync', () => {
259259
expect(
260260
(await archiver.getCheckpoints({ from: CheckpointNumber(1), limit: 100 })).map(b => b.checkpoint.number),
261261
).toEqual([1, 2, 3]);
262+
263+
// Inbox buckets: each of the three L1 message blocks opened its own bucket, in insertion order.
264+
const t1 = fake.getTimestampAtL1Block(98n);
265+
const t2 = fake.getTimestampAtL1Block(2504n);
266+
const t3 = fake.getTimestampAtL1Block(2511n);
267+
268+
expect(await archiver.getInboxBucket(1n)).toMatchObject({
269+
seq: 1n,
270+
msgCount: 3,
271+
totalMsgCount: 3n,
272+
timestamp: t1,
273+
});
274+
expect(await archiver.getInboxBucket(3n)).toMatchObject({
275+
seq: 3n,
276+
msgCount: 3,
277+
totalMsgCount: 9n,
278+
timestamp: t3,
279+
isOpen: true,
280+
});
281+
expect((await archiver.getInboxBucket(1n))!.isOpen).toBe(false);
282+
283+
// At-or-before lookups resolve the latest bucket not opened after the given timestamp.
284+
expect((await archiver.getLatestInboxBucketAtOrBefore(t3))!.seq).toEqual(3n);
285+
expect((await archiver.getLatestInboxBucketAtOrBefore(t2))!.seq).toEqual(2n);
286+
expect(await archiver.getLatestInboxBucketAtOrBefore(t1 - 1n)).toBeUndefined();
287+
288+
// Messages between buckets, in insertion order.
289+
expect(await archiver.getL1ToL2MessagesBetweenBuckets(0n, 3n)).toEqual([...msgs1, ...msgs2, ...msgs3]);
290+
expect(await archiver.getL1ToL2MessagesBetweenBuckets(1n, 2n)).toEqual(msgs2);
291+
expect(await archiver.getL1ToL2MessagesBetweenBuckets(2n, 3n)).toEqual(msgs3);
262292
}, 30_000);
263293

264294
it('ignores checkpoint 3 because it has been pruned', async () => {

yarn-project/archiver/src/l1/data_retrieval.ts

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -387,6 +387,9 @@ function mapLogInboxMessage(log: MessageSentLog): InboxMessage {
387387
l1BlockHash: log.l1BlockHash,
388388
checkpointNumber: log.args.checkpointNumber,
389389
rollingHash: log.args.rollingHash,
390+
inboxRollingHash: log.args.inboxRollingHash,
391+
bucketSeq: log.args.bucketSeq,
392+
bucketTimestamp: log.l1BlockTimestamp,
390393
};
391394
}
392395

yarn-project/archiver/src/modules/data_source_base.ts

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -40,7 +40,7 @@ import {
4040
} from '@aztec/stdlib/epoch-helpers';
4141
import type { L2LogsSource } from '@aztec/stdlib/interfaces/server';
4242
import type { LogResult, PrivateLogsQuery, PublicLogsQuery } from '@aztec/stdlib/logs';
43-
import type { L1ToL2MessageSource, L2ToL1MembershipWitness } from '@aztec/stdlib/messaging';
43+
import type { InboxBucket, L1ToL2MessageSource, L2ToL1MembershipWitness } from '@aztec/stdlib/messaging';
4444
import { AppendOnlyTreeSnapshot } from '@aztec/stdlib/trees';
4545
import type { BlockHeader, IndexedTxEffect, TxHash } from '@aztec/stdlib/tx';
4646
import type { UInt64 } from '@aztec/stdlib/types';
@@ -324,6 +324,18 @@ export abstract class ArchiverDataSourceBase
324324
return this.stores.messages.getL1ToL2MessageIndex(l1ToL2Message);
325325
}
326326

327+
public getLatestInboxBucketAtOrBefore(timestamp: bigint): Promise<InboxBucket | undefined> {
328+
return this.stores.messages.getLatestInboxBucketAtOrBefore(timestamp);
329+
}
330+
331+
public getInboxBucket(seq: bigint): Promise<InboxBucket | undefined> {
332+
return this.stores.messages.getInboxBucket(seq);
333+
}
334+
335+
public getL1ToL2MessagesBetweenBuckets(fromExclusive: bigint, toInclusive: bigint): Promise<Fr[]> {
336+
return this.stores.messages.getL1ToL2MessagesBetweenBuckets(fromExclusive, toInclusive);
337+
}
338+
327339
private async getPublishedCheckpointFromCheckpointData(checkpoint: CheckpointData): Promise<PublishedCheckpoint> {
328340
const blocksForCheckpoint = await this.stores.blocks.getBlocksForCheckpoint(checkpoint.checkpointNumber);
329341
if (!blocksForCheckpoint) {

yarn-project/archiver/src/store/data_stores.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,7 @@ import { FunctionNamesCache } from './function_names_cache.js';
1313
import { LogStore } from './log_store.js';
1414
import { MessageStore } from './message_store.js';
1515

16-
export const ARCHIVER_DB_VERSION = 7;
16+
export const ARCHIVER_DB_VERSION = 8;
1717

1818
/**
1919
* Represents the latest L1 block processed by the archiver for various objects in L2.

yarn-project/archiver/src/store/message_store.test.ts

Lines changed: 162 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import { NUMBER_OF_L1_L2_MESSAGES_PER_ROLLUP } from '@aztec/constants';
22
import { CheckpointNumber } from '@aztec/foundation/branded-types';
33
import { Buffer16, Buffer32 } from '@aztec/foundation/buffer';
4+
import { Fr } from '@aztec/foundation/curves/bn254';
45
import { toArray } from '@aztec/foundation/iterable';
56
import { openTmpStore } from '@aztec/kv-store/lmdb-v2';
67
import { Checkpoint, type PublishedCheckpoint } from '@aztec/stdlib/checkpoint';
@@ -136,6 +137,7 @@ describe('MessageStore', () => {
136137
const msgs2 = makeInboxMessages(3, {
137138
initialCheckpointNumber: CheckpointNumber(20),
138139
initialHash: msgs1.at(-1)!.rollingHash,
140+
initialInboxHash: msgs1.at(-1)!.inboxRollingHash,
139141
});
140142

141143
await messageStore.addL1ToL2Messages(msgs1);
@@ -327,4 +329,164 @@ describe('MessageStore', () => {
327329
});
328330
});
329331
});
332+
333+
describe('Inbox buckets', () => {
334+
// Builds `count` consecutive valid messages in a single checkpoint, then reassigns their bucket sequence and
335+
// timestamp per the given per-message spec so we can exercise multi-message and rollover buckets.
336+
const makeBucketedMessages = (spec: { seq: bigint; timestamp: bigint }[]): InboxMessage[] => {
337+
const msgs = makeInboxMessages(spec.length, {
338+
initialCheckpointNumber: CheckpointNumber(1),
339+
messagesPerCheckpoint: spec.length,
340+
});
341+
msgs.forEach((msg, i) => {
342+
msg.bucketSeq = spec[i].seq;
343+
msg.bucketTimestamp = spec[i].timestamp;
344+
});
345+
return msgs;
346+
};
347+
348+
// Three buckets over six messages: bucket 1 = [0,1,2], bucket 2 = [3,4], bucket 3 = [5].
349+
const threeBucketSpec = [
350+
{ seq: 1n, timestamp: 100n },
351+
{ seq: 1n, timestamp: 100n },
352+
{ seq: 1n, timestamp: 100n },
353+
{ seq: 2n, timestamp: 200n },
354+
{ seq: 2n, timestamp: 200n },
355+
{ seq: 3n, timestamp: 300n },
356+
];
357+
358+
it('snapshots buckets as messages are inserted', async () => {
359+
const msgs = makeBucketedMessages(threeBucketSpec);
360+
await messageStore.addL1ToL2Messages(msgs);
361+
362+
expect(await messageStore.getInboxBucket(1n)).toEqual({
363+
seq: 1n,
364+
inboxRollingHash: msgs[2].inboxRollingHash,
365+
totalMsgCount: 3n,
366+
timestamp: 100n,
367+
msgCount: 3,
368+
lastMessageIndex: msgs[2].index,
369+
isOpen: false,
370+
});
371+
expect(await messageStore.getInboxBucket(2n)).toEqual({
372+
seq: 2n,
373+
inboxRollingHash: msgs[4].inboxRollingHash,
374+
totalMsgCount: 5n,
375+
timestamp: 200n,
376+
msgCount: 2,
377+
lastMessageIndex: msgs[4].index,
378+
isOpen: false,
379+
});
380+
expect(await messageStore.getInboxBucket(3n)).toEqual({
381+
seq: 3n,
382+
inboxRollingHash: msgs[5].inboxRollingHash,
383+
totalMsgCount: 6n,
384+
timestamp: 300n,
385+
msgCount: 1,
386+
lastMessageIndex: msgs[5].index,
387+
isOpen: true,
388+
});
389+
expect(await messageStore.getInboxBucket(4n)).toBeUndefined();
390+
});
391+
392+
it('continues a bucket that spans two insertion batches', async () => {
393+
const msgs = makeBucketedMessages(threeBucketSpec);
394+
await messageStore.addL1ToL2Messages(msgs.slice(0, 2));
395+
await messageStore.addL1ToL2Messages(msgs.slice(2));
396+
397+
// Bucket 1 keeps accumulating across the batch boundary rather than restarting its message count.
398+
expect(await messageStore.getInboxBucket(1n)).toMatchObject({ msgCount: 3, totalMsgCount: 3n });
399+
expect(await messageStore.getInboxBucket(3n)).toMatchObject({ msgCount: 1, totalMsgCount: 6n });
400+
});
401+
402+
it('throws if the consensus rolling hash is not correct', async () => {
403+
const msgs = makeInboxMessages(5);
404+
msgs[1].inboxRollingHash = Fr.random();
405+
await expect(messageStore.addL1ToL2Messages(msgs)).rejects.toThrow(MessageStoreError);
406+
});
407+
408+
it('resolves the latest bucket at or before a timestamp', async () => {
409+
await messageStore.addL1ToL2Messages(makeBucketedMessages(threeBucketSpec));
410+
411+
expect((await messageStore.getLatestInboxBucketAtOrBefore(100n))!.seq).toEqual(1n);
412+
expect((await messageStore.getLatestInboxBucketAtOrBefore(150n))!.seq).toEqual(1n);
413+
expect((await messageStore.getLatestInboxBucketAtOrBefore(300n))!.seq).toEqual(3n);
414+
expect((await messageStore.getLatestInboxBucketAtOrBefore(10_000n))!.seq).toEqual(3n);
415+
expect(await messageStore.getLatestInboxBucketAtOrBefore(99n)).toBeUndefined();
416+
});
417+
418+
it('resolves rollover buckets that share a timestamp to the highest sequence', async () => {
419+
// Buckets 2 and 3 share timestamp 200 (a full bucket rolling over within the same L1 block).
420+
const msgs = makeBucketedMessages([
421+
{ seq: 1n, timestamp: 100n },
422+
{ seq: 2n, timestamp: 200n },
423+
{ seq: 3n, timestamp: 200n },
424+
]);
425+
await messageStore.addL1ToL2Messages(msgs);
426+
427+
expect((await messageStore.getLatestInboxBucketAtOrBefore(200n))!.seq).toEqual(3n);
428+
});
429+
430+
it('returns messages between buckets in insertion order', async () => {
431+
const msgs = makeBucketedMessages(threeBucketSpec);
432+
await messageStore.addL1ToL2Messages(msgs);
433+
const leaves = msgs.map(m => m.leaf);
434+
435+
expect(await messageStore.getL1ToL2MessagesBetweenBuckets(0n, 3n)).toEqual(leaves);
436+
expect(await messageStore.getL1ToL2MessagesBetweenBuckets(1n, 2n)).toEqual(leaves.slice(3, 5));
437+
expect(await messageStore.getL1ToL2MessagesBetweenBuckets(2n, 3n)).toEqual(leaves.slice(5));
438+
// An empty (fromExclusive, toInclusive] range yields no messages.
439+
expect(await messageStore.getL1ToL2MessagesBetweenBuckets(3n, 3n)).toEqual([]);
440+
// Unknown upper bucket yields no messages.
441+
expect(await messageStore.getL1ToL2MessagesBetweenBuckets(0n, 9n)).toEqual([]);
442+
});
443+
444+
it('rewinds buckets when messages are removed', async () => {
445+
const msgs = makeBucketedMessages(threeBucketSpec);
446+
await messageStore.addL1ToL2Messages(msgs);
447+
448+
// Remove the last two messages (msgs[4] in bucket 2, msgs[5] in bucket 3), splitting bucket 2.
449+
await messageStore.removeL1ToL2Messages(msgs[4].index);
450+
451+
expect(await messageStore.getInboxBucket(3n)).toBeUndefined();
452+
expect(await messageStore.getInboxBucket(2n)).toEqual({
453+
seq: 2n,
454+
inboxRollingHash: msgs[3].inboxRollingHash,
455+
totalMsgCount: 4n,
456+
timestamp: 200n,
457+
msgCount: 1,
458+
lastMessageIndex: msgs[3].index,
459+
isOpen: true,
460+
});
461+
expect(await messageStore.getInboxBucket(1n)).toMatchObject({ msgCount: 3, totalMsgCount: 3n, isOpen: false });
462+
463+
// Bucket 3's timestamp index entry is gone, so an at-or-before lookup falls back to bucket 2.
464+
expect((await messageStore.getLatestInboxBucketAtOrBefore(300n))!.seq).toEqual(2n);
465+
});
466+
467+
it('rewinds a rollover bucket sharing a timestamp with the surviving boundary', async () => {
468+
const msgs = makeBucketedMessages([
469+
{ seq: 1n, timestamp: 100n },
470+
{ seq: 2n, timestamp: 200n },
471+
{ seq: 3n, timestamp: 200n },
472+
]);
473+
await messageStore.addL1ToL2Messages(msgs);
474+
475+
// Removing the last message deletes bucket 3, whose timestamp (200) is shared with the surviving bucket 2.
476+
await messageStore.removeL1ToL2Messages(msgs[2].index);
477+
478+
expect(await messageStore.getInboxBucket(3n)).toBeUndefined();
479+
expect((await messageStore.getLatestInboxBucketAtOrBefore(200n))!.seq).toEqual(2n);
480+
});
481+
482+
it('clears all buckets when every message is removed', async () => {
483+
const msgs = makeBucketedMessages(threeBucketSpec);
484+
await messageStore.addL1ToL2Messages(msgs);
485+
486+
await messageStore.removeL1ToL2Messages(msgs[0].index);
487+
488+
expect(await messageStore.getInboxBucket(1n)).toBeUndefined();
489+
expect(await messageStore.getLatestInboxBucketAtOrBefore(300n)).toBeUndefined();
490+
});
491+
});
330492
});

0 commit comments

Comments
 (0)