Skip to content

Commit b573584

Browse files
committed
fix(sync): sync distant P2P data heads by range
1 parent c27a9fd commit b573584

3 files changed

Lines changed: 114 additions & 1 deletion

File tree

‎pkg/sync/sync_service_test.go‎

Lines changed: 86 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,9 +5,14 @@ import (
55
cryptoRand "crypto/rand"
66
"math/rand"
77
"path/filepath"
8+
"sync/atomic"
89
"testing"
910
"time"
1011

12+
goheader "github.com/celestiaorg/go-header"
13+
"github.com/celestiaorg/go-header/headertest"
14+
goheaderlocal "github.com/celestiaorg/go-header/local"
15+
goheadersync "github.com/celestiaorg/go-header/sync"
1116
"github.com/ipfs/go-datastore"
1217
"github.com/ipfs/go-datastore/sync"
1318
"github.com/libp2p/go-libp2p/core/crypto"
@@ -26,6 +31,87 @@ import (
2631
"github.com/evstack/ev-node/types"
2732
)
2833

34+
type countingP2PDataGetter struct {
35+
goheader.Getter[*types.P2PData]
36+
getByHeightCalls atomic.Uint64
37+
rangeCalls atomic.Uint64
38+
}
39+
40+
func (g *countingP2PDataGetter) GetByHeight(ctx context.Context, height uint64) (*types.P2PData, error) {
41+
g.getByHeightCalls.Add(1)
42+
return g.Getter.GetByHeight(ctx, height)
43+
}
44+
45+
func (g *countingP2PDataGetter) GetRangeByHeight(
46+
ctx context.Context,
47+
from *types.P2PData,
48+
to uint64,
49+
) ([]*types.P2PData, error) {
50+
g.rangeCalls.Add(1)
51+
return g.Getter.GetRangeByHeight(ctx, from, to)
52+
}
53+
54+
type verifierCapturingP2PDataSubscriber struct {
55+
*headertest.Subscriber[*types.P2PData]
56+
verifier func(context.Context, *types.P2PData) error
57+
}
58+
59+
func (s *verifierCapturingP2PDataSubscriber) SetVerifier(
60+
verifier func(context.Context, *types.P2PData) error,
61+
) error {
62+
s.verifier = verifier
63+
return nil
64+
}
65+
66+
func TestDataSyncerDistantHeadUsesRangeSync(t *testing.T) {
67+
const targetHeight uint64 = 64
68+
69+
ctx, cancel := context.WithTimeout(t.Context(), 5*time.Second)
70+
defer cancel()
71+
72+
chain := make([]*types.P2PData, 0, targetHeight)
73+
blockTime := time.Now().Add(-time.Duration(targetHeight) * time.Millisecond)
74+
var previousHash types.Hash
75+
for height := uint64(1); height <= targetHeight; height++ {
76+
_, data := types.GetRandomBlock(height, 1, "data-sync-catchup")
77+
data.Metadata.Time = uint64(blockTime.Add(time.Duration(height) * time.Millisecond).UnixNano())
78+
data.LastDataHash = previousHash
79+
previousHash = data.Hash()
80+
chain = append(chain, &types.P2PData{Data: data})
81+
}
82+
83+
remoteStore := &headertest.Store[*types.P2PData]{Headers: make(map[uint64]*types.P2PData)}
84+
require.NoError(t, remoteStore.Append(ctx, chain...))
85+
localStore := &headertest.Store[*types.P2PData]{Headers: make(map[uint64]*types.P2PData)}
86+
require.NoError(t, localStore.Append(ctx, chain[0]))
87+
88+
getter := &countingP2PDataGetter{Getter: goheaderlocal.NewExchange(remoteStore)}
89+
subscriber := &verifierCapturingP2PDataSubscriber{
90+
Subscriber: &headertest.Subscriber[*types.P2PData]{},
91+
}
92+
syncer, err := goheadersync.NewSyncer(
93+
getter,
94+
localStore,
95+
subscriber,
96+
goheadersync.WithBlockTime(time.Second),
97+
)
98+
require.NoError(t, err)
99+
require.NoError(t, syncer.Start(ctx))
100+
t.Cleanup(func() { _ = syncer.Stop(context.Background()) })
101+
102+
getter.getByHeightCalls.Store(0)
103+
getter.rangeCalls.Store(0)
104+
require.NoError(t, subscriber.verifier(ctx, chain[targetHeight-1]))
105+
require.LessOrEqual(t, getter.getByHeightCalls.Load(), uint64(1),
106+
"validating a distant head must not fetch intermediate blocks one by one")
107+
108+
require.Eventually(t, func() bool {
109+
return syncer.State().ToHeight == targetHeight && localStore.Height() == targetHeight
110+
}, time.Second, time.Millisecond)
111+
require.Equal(t, targetHeight, localStore.Height())
112+
require.Positive(t, getter.rangeCalls.Load())
113+
}
114+
29115
func TestHeaderSyncServiceStartForPublishingWithPeers(t *testing.T) {
30116
mainKV := sync.MutexWrap(datastore.NewMapDatastore())
31117
pk, _, err := crypto.GenerateEd25519Key(cryptoRand.Reader)

‎types/p2p_envelope.go‎

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -109,8 +109,12 @@ func (p *P2PData) DAHint() uint64 {
109109
return p.DAHeightHint
110110
}
111111

112-
// Verify verifies against untrusted data.
112+
// Verify verifies the data hash linkage for adjacent data.
113113
func (p *P2PData) Verify(untrusted *P2PData) error {
114+
if p.Height()+1 != untrusted.Height() {
115+
return nil
116+
}
117+
114118
return p.Data.Verify(untrusted.Data)
115119
}
116120

‎types/p2p_envelope_test.go‎

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ import (
55
"testing"
66
"time"
77

8+
goheader "github.com/celestiaorg/go-header"
89
"github.com/libp2p/go-libp2p/core/crypto"
910
"github.com/stretchr/testify/assert"
1011
"github.com/stretchr/testify/require"
@@ -42,6 +43,28 @@ func TestP2PEnvelope_MarshalUnmarshal(t *testing.T) {
4243
assert.Equal(t, envelope.Txs, newEnvelope.Txs)
4344
}
4445

46+
func TestP2PDataVerifyAdjacentHeads(t *testing.T) {
47+
now := time.Now()
48+
_, trustedData := GetRandomBlock(10, 1, "test-chain")
49+
trustedData.Metadata.Time = uint64(now.UnixNano())
50+
51+
_, validData := GetRandomBlock(11, 1, "test-chain")
52+
validData.Metadata.Time = uint64(now.Add(time.Second).UnixNano())
53+
validData.LastDataHash = trustedData.Hash()
54+
require.NoError(t, goheader.Verify(
55+
&P2PData{Data: trustedData},
56+
&P2PData{Data: validData},
57+
))
58+
59+
_, invalidData := GetRandomBlock(11, 1, "test-chain")
60+
invalidData.Metadata.Time = uint64(now.Add(time.Second).UnixNano())
61+
invalidData.LastDataHash = bytes.Repeat([]byte{0x1}, 32)
62+
require.Error(t, goheader.Verify(
63+
&P2PData{Data: trustedData},
64+
&P2PData{Data: invalidData},
65+
))
66+
}
67+
4568
func TestP2PSignedHeader_MarshalUnmarshal(t *testing.T) {
4669
_, pubKey, err := crypto.GenerateEd25519Key(nil)
4770
require.NoError(t, err)

0 commit comments

Comments
 (0)