Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
73 changes: 63 additions & 10 deletions sei-tendermint/internal/autobahn/avail/inner.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,55 @@ import (
"github.com/sei-protocol/sei-chain/sei-tendermint/libs/utils"
)

// blockQueue is a per-lane block queue.
type blockQueue struct {
queue[types.BlockNumber, *types.Signed[*types.LaneProposal]]
// last is None, or this node's last pushed proposal at height >= first-1.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[suggestion] last is only ever read for local block production (ProduceLocalBlock, state.go:664) — PushBlock deliberately does not consult it — yet pushBack sets it for every lane, so retentionFloor and unpersistedLast keep and may flush one extra signed proposal (full payload) per foreign lane, in memory and in the WAL, with no consumer. On a large committee that is one retained proposal per validator held past the anchor for the life of the process.

Related: the doc comment says "this node's last pushed proposal", but PushBlock sets it from a remote producer's block too. A reader could reasonably conclude foreign lanes don't retain anything here.

If gating is wanted, the choke point is addLane/newBlockQueue (which is where the local key could be compared against lane.Validator) rather than a condition at each pushBack caller. Otherwise, at minimum reword the comment to match what the code does.

last utils.Option[*types.Signed[*types.LaneProposal]]
}

func newBlockQueue() *blockQueue {
return &blockQueue{queue: *newQueue[types.BlockNumber, *types.Signed[*types.LaneProposal]]()}
}

func (q *blockQueue) pushBack(p *types.Signed[*types.LaneProposal]) {
q.queue.pushBack(p)
q.last = utils.Some(p)
}

// prune drops [first, newFirst). last is kept when newFirst <= next and
// cleared when newFirst > next.
func (q *blockQueue) prune(newFirst types.BlockNumber) {
if newFirst <= q.first {
return
}
if newFirst > q.next {
// TODO: seed last from a non-empty LaneRange LastHash at Next()-1.
// Empty ranges carry a zero LastHash, so they cannot replace a local last.
q.last = utils.None[*types.Signed[*types.LaneProposal]]()
}
q.queue.prune(newFirst)
}

// unpersistedLast returns the last block once it has left the active range and
// block persistence has not reached it.
func (q *blockQueue) unpersistedLast(nextToPersist types.BlockNumber) utils.Option[*types.Signed[*types.LaneProposal]] {
if p, ok := q.last.Get(); ok {
if n := p.Msg().Block().Header().BlockNumber(); n < q.first && nextToPersist <= n {
return utils.Some(p)
}
}
return utils.None[*types.Signed[*types.LaneProposal]]()
}

// retentionFloor returns the lowest block number the WAL must still hold.
func (q *blockQueue) retentionFloor() types.BlockNumber {
if p, ok := q.last.Get(); ok {
return min(q.first, p.Msg().Block().Header().BlockNumber())
}
return q.first
}

// inner holds roads and per-LaneID block/vote maps.
type inner struct {
persistedCommitQC utils.AtomicSend[utils.Option[*types.CommitQC]] // latest persisted CommitQC
Expand All @@ -25,7 +74,7 @@ type inner struct {
// When it lags applied, epochForVote falls back to this committee for
// departing-lane voters.
anchorEpoch utils.Option[*types.Epoch]
blocks map[types.LaneID]*queue[types.BlockNumber, *types.Signed[*types.LaneProposal]]
blocks map[types.LaneID]*blockQueue
votes map[types.LaneID]*queue[types.BlockNumber, *blockVotes]
// nextBlockToPersist tracks per-lane how far block persistence has progressed.
// RecvBatch only yields blocks below this cursor for voting.
Expand Down Expand Up @@ -61,7 +110,7 @@ func newInner(ep *types.Epoch, first types.RoadIndex) *inner {
persistedCommitQC: utils.NewAtomicSend(utils.None[*types.CommitQC]()),
consensusSpec: utils.NewAtomicSend(types.ConsensusSpec{CommitQC: utils.None[*types.CommitQC](), Epoch: ep}),
roads: roads,
blocks: map[types.LaneID]*queue[types.BlockNumber, *types.Signed[*types.LaneProposal]]{},
blocks: map[types.LaneID]*blockQueue{},
votes: map[types.LaneID]*queue[types.BlockNumber, *blockVotes]{},
nextBlockToPersist: map[types.LaneID]types.BlockNumber{},
}
Expand All @@ -85,16 +134,18 @@ func (i *inner) restoreBlocks(blocks map[types.LaneID][]persist.LoadedBlock) err
return fmt.Errorf("lane %s: loaded %d blocks exceeds capacity %d", lane, len(bs), BlocksPerLane)
}
if b.Number < q.next {
// Certified. Restore last from the proposal at first-1 when present.
if b.Number+1 == q.first {
q.last = utils.Some(b.Proposal)
}
continue
}
if b.Number != q.next {
return fmt.Errorf("lane %s: non-contiguous persisted blocks: expected %d, got %d", lane, q.next, b.Number)
}
// We check the parent hash only for the blocks above the anchor, because:
// * node can cast LaneVote for the block of the lane without checking the parent hash,
// in case the previous block was already (executed and) pruned from memory.
// * current WAL implementation is lazily pruning on disk, so old executed blocks might be loaded on startup.
if q.Len() > 0 {
// Parent is checked only inside [first, next). last restored from
// first-1 is for local production, not this check.
if q.first < q.next {
ph := b.Proposal.Msg().Block().Header().ParentHash()
if q.q[q.next-1].Msg().Block().Header().Hash() != ph {
return fmt.Errorf("lane %s: parent hash mismatch at block %d", lane, b.Number)
Expand Down Expand Up @@ -192,7 +243,7 @@ func (i *inner) addLane(lane types.LaneID) bool {
if _, ok := i.blocks[lane]; ok {
return false
}
i.blocks[lane] = newQueue[types.BlockNumber, *types.Signed[*types.LaneProposal]]()
i.blocks[lane] = newBlockQueue()
i.votes[lane] = newQueue[types.BlockNumber, *blockVotes]()
i.nextBlockToPersist[lane] = 0
return true
Expand Down Expand Up @@ -248,8 +299,10 @@ func (i *inner) prune(anchor data.Anchor) int {
bq := i.blocks[lane]
vq.prune(lr.Next())
bq.prune(lr.Next())
if i.nextBlockToPersist[lane] < lr.Next() {
i.nextBlockToPersist[lane] = lr.Next()
// A lagging cursor stops at retentionFloor so an unflushed last can still
// be written. The cursor is never rewound: already past last means it is on disk.
if floor := bq.retentionFloor(); i.nextBlockToPersist[lane] < floor {
i.nextBlockToPersist[lane] = floor
}
}
if i.roads.Len() == 0 {
Expand Down
202 changes: 200 additions & 2 deletions sei-tendermint/internal/autobahn/avail/inner_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,166 @@ func contiguousBlocks(key types.SecretKey, lane types.LaneID, n int, rng utils.R
return bs
}

func TestBlockQueueRetainsLast(t *testing.T) {
rng := utils.TestRng()
key := types.GenSecretKey(rng)
lane := types.LaneID{Validator: key.Public(), Joined: 0}
blocks := contiguousBlocks(key, lane, 3, rng)
q := newBlockQueue()
for _, b := range blocks {
q.pushBack(b.Proposal)
}

q.prune(2)
require.Equal(t, types.BlockNumber(2), q.first)
require.Equal(t, types.BlockNumber(3), q.next)
require.Equal(t, utils.Some(blocks[2].Proposal), q.last)
// Block 2 is still active and carries the chain, so nothing below first is needed.
require.Equal(t, types.BlockNumber(2), q.retentionFloor())

q.prune(3)
require.Equal(t, types.BlockNumber(3), q.first)
require.Equal(t, types.BlockNumber(3), q.next)
require.Equal(t, utils.Some(blocks[2].Proposal), q.last)
require.Equal(t, types.BlockNumber(2), q.retentionFloor())

q.prune(5)
require.Equal(t, types.BlockNumber(5), q.first)
require.Equal(t, types.BlockNumber(5), q.next)
require.Equal(t, utils.None[*types.Signed[*types.LaneProposal]](), q.last)
require.Equal(t, types.BlockNumber(5), q.retentionFloor())
}

func TestInnerPruneViaQCRetainsLast(t *testing.T) {
rng := utils.TestRng()
registry, keys := epoch.GenRegistry(rng, 3)
ep := registry.MustEpoch(0)
lane := ep.Committee().Lane(keys[0].Public()).OrPanic("lane")
blocks := contiguousBlocks(keys[0], lane, 3, rng)
i := newInner(ep, 0)
for _, b := range blocks {
i.blocks[lane].pushBack(b.Proposal)
}

last := blocks[len(blocks)-1].Proposal.Msg().Block().Header()
qc := types.BuildCommitQC(ep, keys, utils.None[*types.CommitQC](), map[types.LaneID]*types.LaneQC{
lane: types.NewLaneQC(makeLaneVotes(keys, last)),
})
lr := qc.LaneRange(lane)
require.Equal(t, types.BlockNumber(3), lr.Next())
require.Equal(t, last.Hash(), lr.LastHash())

i.prune(data.Anchor{
CommitQC: qc,
AppQC: data.TestAppQC(keys, types.NewAppProposal(qc.Proposal(), types.AppHash{})),
Epoch: ep,
})
q := i.blocks[lane]
require.Equal(t, types.BlockNumber(3), q.first)
require.Equal(t, types.BlockNumber(3), q.next)
require.Equal(t, utils.Some(blocks[2].Proposal), q.last)
require.Equal(t, types.BlockNumber(2), q.retentionFloor())
require.Equal(t, types.BlockNumber(2), i.nextBlockToPersist[lane])
p, ok := q.unpersistedLast(i.nextBlockToPersist[lane]).Get()
require.True(t, ok)
require.Equal(t, blocks[2].Proposal, p)
}

func TestInnerPruneKeepsLastAheadOfQC(t *testing.T) {
rng := utils.TestRng()
registry, keys := epoch.GenRegistry(rng, 3)
ep := registry.MustEpoch(0)
lane := ep.Committee().Lane(keys[0].Public()).OrPanic("lane")
blocks := contiguousBlocks(keys[0], lane, 3, rng)
i := newInner(ep, 0)
for _, b := range blocks {
i.blocks[lane].pushBack(b.Proposal)
}

qc := types.BuildCommitQC(ep, keys, utils.None[*types.CommitQC](), map[types.LaneID]*types.LaneQC{
lane: types.NewLaneQC(makeLaneVotes(keys, blocks[0].Proposal.Msg().Block().Header())),
})
require.Equal(t, types.BlockNumber(1), qc.LaneRange(lane).Next())

i.prune(data.Anchor{
CommitQC: qc,
AppQC: data.TestAppQC(keys, types.NewAppProposal(qc.Proposal(), types.AppHash{})),
Epoch: ep,
})
q := i.blocks[lane]
require.Equal(t, types.BlockNumber(1), q.first)
require.Equal(t, types.BlockNumber(3), q.next)
require.Equal(t, utils.Some(blocks[2].Proposal), q.last)
}

func TestBlockQueueLastSurvivesRestart(t *testing.T) {
rng := utils.TestRng()
registry, keys := epoch.GenRegistry(rng, 3)
key := keys[0]
lane := registry.MustEpoch(0).Committee().Lane(key.Public()).OrPanic("lane")
block := contiguousBlocks(key, lane, 1, rng)[0]
dir := t.TempDir()

persister, _, err := persist.NewBlockPersister(utils.Some(dir))
require.NoError(t, err)
require.NoError(t, persister.PruneAndPersist(
lane,
0,
[]*types.Signed[*types.LaneProposal]{block.Proposal},
))
q := newBlockQueue()
q.pushBack(block.Proposal)
q.prune(1)
require.NoError(t, persister.PruneAndPersist(lane, q.retentionFloor(), nil))
require.NoError(t, persister.Close())

persister, loaded, err := persist.NewBlockPersister(utils.Some(dir))
require.NoError(t, err)
require.NoError(t, persister.Close())
i := newInner(registry.MustEpoch(0), 0)
i.blocks[lane].prune(1)
require.NoError(t, i.restoreBlocks(loaded))
last, ok := i.blocks[lane].last.Get()
require.True(t, ok)
require.Equal(t, block.Proposal.Msg().Block().Header().Hash(), last.Msg().Block().Header().Hash())
}

func TestBlockQueueUnflushedLastSurvivesRestart(t *testing.T) {
rng := utils.TestRng()
registry, keys := epoch.GenRegistry(rng, 3)
key := keys[0]
lane := registry.MustEpoch(0).Committee().Lane(key.Public()).OrPanic("lane")
block := contiguousBlocks(key, lane, 1, rng)[0]
dir := t.TempDir()

persister, _, err := persist.NewBlockPersister(utils.Some(dir))
require.NoError(t, err)
q := newBlockQueue()
q.pushBack(block.Proposal)
q.prune(1)

p, ok := q.unpersistedLast(0).Get()
require.True(t, ok)
require.Equal(t, block.Proposal, p)
require.NoError(t, persister.PruneAndPersist(
lane,
q.retentionFloor(),
[]*types.Signed[*types.LaneProposal]{p},
))
require.False(t, q.unpersistedLast(1).IsPresent())
require.NoError(t, persister.Close())

persister, loaded, err := persist.NewBlockPersister(utils.Some(dir))
require.NoError(t, err)
require.NoError(t, persister.Close())
i := newInner(registry.MustEpoch(0), 0)
i.blocks[lane].prune(1)
require.NoError(t, i.restoreBlocks(loaded))
last, ok := i.blocks[lane].last.Get()
require.True(t, ok)
require.Equal(t, block.Proposal.Msg().Block().Header().Hash(), last.Msg().Block().Header().Hash())
}

func TestRestoreInner_Empty(t *testing.T) {
rng := utils.TestRng()
registry, _ := epoch.GenRegistry(rng, 4)
Expand Down Expand Up @@ -90,6 +250,44 @@ func TestRestoreInner_LoadedBlocks(t *testing.T) {
require.Equal(t, types.BlockNumber(0), q.next)
})

t.Run("last below anchor", func(t *testing.T) {
rng := utils.TestRng()
registry, keys := epoch.GenRegistry(rng, 4)
lane := registry.MustEpoch(0).Committee().Lane(keys[0].Public()).OrPanic("keys[0]")
blocks := contiguousBlocks(keys[0], lane, 2, rng)
i := newInner(registry.MustEpoch(0), 0)
i.blocks[lane].prune(2)

err := i.restoreBlocks(map[types.LaneID][]persist.LoadedBlock{lane: blocks})
require.NoError(t, err)
q := i.blocks[lane]
require.Equal(t, types.BlockNumber(2), q.first)
require.Equal(t, types.BlockNumber(2), q.next)
require.Equal(t, utils.Some(blocks[1].Proposal), q.last)
require.Equal(t, types.BlockNumber(1), q.retentionFloor())
require.Equal(t, types.BlockNumber(2), i.nextBlockToPersist[lane])
})

t.Run("leftover below first is not parent-checked", func(t *testing.T) {
rng := utils.TestRng()
registry, keys := epoch.GenRegistry(rng, 4)
lane := registry.MustEpoch(0).Committee().Lane(keys[0].Public()).OrPanic("keys[0]")
old := testSignedBlock(keys[0], lane, 0, types.BlockHeaderHash{}, rng)
live := testSignedBlock(keys[0], lane, 1, types.GenBlockHeaderHash(rng), rng)
i := newInner(registry.MustEpoch(0), 0)
i.blocks[lane].prune(1)

err := i.restoreBlocks(map[types.LaneID][]persist.LoadedBlock{lane: {
{Number: 0, Proposal: old},
{Number: 1, Proposal: live},
}})
require.NoError(t, err)
q := i.blocks[lane]
require.Equal(t, types.BlockNumber(1), q.first)
require.Equal(t, types.BlockNumber(2), q.next)
require.Equal(t, utils.Some(live), q.last)
})

t.Run("foreign loaded lane does not touch committee queues", func(t *testing.T) {
rng := utils.TestRng()
registry, _ := epoch.GenRegistry(rng, 4)
Expand Down Expand Up @@ -218,7 +416,7 @@ func TestAddLane_ReportsNewLaneForEachMembershipPeriod(t *testing.T) {
a := types.GenSecretKey(rng)

i := &inner{
blocks: map[types.LaneID]*queue[types.BlockNumber, *types.Signed[*types.LaneProposal]]{},
blocks: map[types.LaneID]*blockQueue{},
votes: map[types.LaneID]*queue[types.BlockNumber, *blockVotes]{},
nextBlockToPersist: map[types.LaneID]types.BlockNumber{},
}
Expand Down Expand Up @@ -249,7 +447,7 @@ func TestRefreshConsensusSpec_WithholdsTipUntilNextViewEpochApplied(t *testing.T
persistedCommitQC: utils.NewAtomicSend(utils.None[*types.CommitQC]()),
consensusSpec: utils.NewAtomicSend(types.ConsensusSpec{CommitQC: utils.None[*types.CommitQC](), Epoch: ep0}),
roads: newQueue[types.RoadIndex, *road](),
blocks: map[types.LaneID]*queue[types.BlockNumber, *types.Signed[*types.LaneProposal]]{},
blocks: map[types.LaneID]*blockQueue{},
votes: map[types.LaneID]*queue[types.BlockNumber, *blockVotes]{},
nextBlockToPersist: map[types.LaneID]types.BlockNumber{},
}
Expand Down
Loading
Loading