swarm/network/stream: remove HandoverProof

This commit is contained in:
Anton Evangelatov 2019-05-31 12:39:34 +02:00
parent 654a51a88a
commit e9f4d1aa60
7 changed files with 21 additions and 68 deletions

View file

@ -72,7 +72,6 @@ func TestStreamerUpstreamRetrieveRequestMsgExchangeWithoutStore(t *testing.T) {
{ //to which the peer responds with offered hashes { //to which the peer responds with offered hashes
Code: 1, Code: 1,
Msg: &OfferedHashesMsg{ Msg: &OfferedHashesMsg{
HandoverProof: nil,
Hashes: nil, Hashes: nil,
From: 0, From: 0,
To: 0, To: 0,

View file

@ -339,7 +339,7 @@ func (s *testExternalServer) SessionIndex() (uint64, error) {
return s.sessionAt, nil return s.sessionAt, nil
} }
func (s *testExternalServer) SetNextBatch(from uint64, to uint64) ([]byte, uint64, uint64, *HandoverProof, error) { func (s *testExternalServer) SetNextBatch(from uint64, to uint64) ([]byte, uint64, uint64, error) {
if to > s.maxKeys { if to > s.maxKeys {
to = s.maxKeys to = s.maxKeys
} }
@ -347,7 +347,7 @@ func (s *testExternalServer) SetNextBatch(from uint64, to uint64) ([]byte, uint6
for i := from; i <= to; i++ { for i := from; i <= to; i++ {
s.keyFunc(b[(i-from)*HashSize:(i-from+1)*HashSize], i) s.keyFunc(b[(i-from)*HashSize:(i-from+1)*HashSize], i)
} }
return b, from, to, nil, nil return b, from, to, nil
} }
func (s *testExternalServer) GetData(context.Context, []byte) ([]byte, error) { func (s *testExternalServer) GetData(context.Context, []byte) ([]byte, error) {

View file

@ -186,7 +186,6 @@ type OfferedHashesMsg struct {
Stream Stream // name of Stream Stream Stream // name of Stream
From, To uint64 // peer and db-specific entry count From, To uint64 // peer and db-specific entry count
Hashes []byte // stream of hashes (128) Hashes []byte // stream of hashes (128)
*HandoverProof // HandoverProof
} }
// String pretty prints OfferedHashesMsg // String pretty prints OfferedHashesMsg
@ -385,12 +384,6 @@ type Handover struct {
Root []byte // Root hash for indexed segment inclusion proofs Root []byte // Root hash for indexed segment inclusion proofs
} }
// HandoverProof represents a signed statement that the upstream peer handed over the stream section
type HandoverProof struct {
Sig []byte // Sign(Hash(Serialisation(Handover)))
*Handover
}
// Takeover represents a statement that downstream peer took over (stored all data) // Takeover represents a statement that downstream peer took over (stored all data)
// handed over // handed over
type Takeover Handover type Takeover Handover

View file

@ -183,7 +183,7 @@ func (p *Peer) SendOfferedHashes(s *server, f, t uint64) error {
defer metrics.GetOrRegisterResettingTimer("send.offered.hashes", nil).UpdateSince(time.Now()) defer metrics.GetOrRegisterResettingTimer("send.offered.hashes", nil).UpdateSince(time.Now())
hashes, from, to, proof, err := s.setNextBatch(f, t) hashes, from, to, err := s.setNextBatch(f, t)
if err != nil { if err != nil {
return err return err
} }
@ -191,14 +191,8 @@ func (p *Peer) SendOfferedHashes(s *server, f, t uint64) error {
if len(hashes) == 0 { if len(hashes) == 0 {
return nil return nil
} }
if proof == nil {
proof = &HandoverProof{
Handover: &Handover{},
}
}
s.currentBatch = hashes s.currentBatch = hashes
msg := &OfferedHashesMsg{ msg := &OfferedHashesMsg{
HandoverProof: proof,
Hashes: hashes, Hashes: hashes,
From: from, From: from,
To: to, To: to,

View file

@ -474,7 +474,7 @@ type server struct {
// setNextBatch adjusts passed interval based on session index and whether // setNextBatch adjusts passed interval based on session index and whether
// stream is live or history. It calls Server SetNextBatch with adjusted // stream is live or history. It calls Server SetNextBatch with adjusted
// interval and returns batch hashes and their interval. // interval and returns batch hashes and their interval.
func (s *server) setNextBatch(from, to uint64) ([]byte, uint64, uint64, *HandoverProof, error) { func (s *server) setNextBatch(from, to uint64) ([]byte, uint64, uint64, error) {
if s.stream.Live { if s.stream.Live {
if from == 0 { if from == 0 {
from = s.sessionIndex from = s.sessionIndex
@ -484,7 +484,7 @@ func (s *server) setNextBatch(from, to uint64) ([]byte, uint64, uint64, *Handove
} }
} else { } else {
if (to < from && to != 0) || from > s.sessionIndex { if (to < from && to != 0) || from > s.sessionIndex {
return nil, 0, 0, nil, nil return nil, 0, 0, nil
} }
if to == 0 || to > s.sessionIndex { if to == 0 || to > s.sessionIndex {
to = s.sessionIndex to = s.sessionIndex
@ -500,7 +500,7 @@ type Server interface {
// Based on this index, live and history stream intervals // Based on this index, live and history stream intervals
// will be adjusted before calling SetNextBatch. // will be adjusted before calling SetNextBatch.
SessionIndex() (uint64, error) SessionIndex() (uint64, error)
SetNextBatch(uint64, uint64) (hashes []byte, from uint64, to uint64, proof *HandoverProof, err error) SetNextBatch(uint64, uint64) (hashes []byte, from uint64, to uint64, err error)
GetData(context.Context, []byte) ([]byte, error) GetData(context.Context, []byte) ([]byte, error)
Close() Close()
} }

View file

@ -127,8 +127,8 @@ func (s *testServer) SessionIndex() (uint64, error) {
return s.sessionIndex, nil return s.sessionIndex, nil
} }
func (self *testServer) SetNextBatch(from uint64, to uint64) ([]byte, uint64, uint64, *HandoverProof, error) { func (self *testServer) SetNextBatch(from uint64, to uint64) ([]byte, uint64, uint64, error) {
return make([]byte, HashSize), from + 1, to + 1, nil, nil return make([]byte, HashSize), from + 1, to + 1, nil
} }
func (self *testServer) GetData(context.Context, []byte) ([]byte, error) { func (self *testServer) GetData(context.Context, []byte) ([]byte, error) {
@ -179,9 +179,6 @@ func TestStreamerDownstreamSubscribeUnsubscribeMsgExchange(t *testing.T) {
{ {
Code: 1, Code: 1,
Msg: &OfferedHashesMsg{ Msg: &OfferedHashesMsg{
HandoverProof: &HandoverProof{
Handover: &Handover{},
},
Hashes: hashes, Hashes: hashes,
From: 5, From: 5,
To: 8, To: 8,
@ -264,9 +261,6 @@ func TestStreamerUpstreamSubscribeUnsubscribeMsgExchange(t *testing.T) {
Code: 1, Code: 1,
Msg: &OfferedHashesMsg{ Msg: &OfferedHashesMsg{
Stream: stream, Stream: stream,
HandoverProof: &HandoverProof{
Handover: &Handover{},
},
Hashes: make([]byte, HashSize), Hashes: make([]byte, HashSize),
From: 6, From: 6,
To: 9, To: 9,
@ -330,9 +324,6 @@ func TestStreamerUpstreamSubscribeUnsubscribeMsgExchangeLive(t *testing.T) {
Code: 1, Code: 1,
Msg: &OfferedHashesMsg{ Msg: &OfferedHashesMsg{
Stream: stream, Stream: stream,
HandoverProof: &HandoverProof{
Handover: &Handover{},
},
Hashes: make([]byte, HashSize), Hashes: make([]byte, HashSize),
From: 1, From: 1,
To: 0, To: 0,
@ -441,9 +432,6 @@ func TestStreamerUpstreamSubscribeLiveAndHistory(t *testing.T) {
Code: 1, Code: 1,
Msg: &OfferedHashesMsg{ Msg: &OfferedHashesMsg{
Stream: NewStream("foo", "", false), Stream: NewStream("foo", "", false),
HandoverProof: &HandoverProof{
Handover: &Handover{},
},
Hashes: make([]byte, HashSize), Hashes: make([]byte, HashSize),
From: 6, From: 6,
To: 9, To: 9,
@ -454,9 +442,6 @@ func TestStreamerUpstreamSubscribeLiveAndHistory(t *testing.T) {
Code: 1, Code: 1,
Msg: &OfferedHashesMsg{ Msg: &OfferedHashesMsg{
Stream: stream, Stream: stream,
HandoverProof: &HandoverProof{
Handover: &Handover{},
},
From: 11, From: 11,
To: 0, To: 0,
Hashes: make([]byte, HashSize), Hashes: make([]byte, HashSize),
@ -514,9 +499,6 @@ func TestStreamerDownstreamCorruptHashesMsgExchange(t *testing.T) {
{ {
Code: 1, Code: 1,
Msg: &OfferedHashesMsg{ Msg: &OfferedHashesMsg{
HandoverProof: &HandoverProof{
Handover: &Handover{},
},
Hashes: corruptHashes, Hashes: corruptHashes,
From: 5, From: 5,
To: 8, To: 8,
@ -579,9 +561,6 @@ func TestStreamerDownstreamOfferedHashesMsgExchange(t *testing.T) {
{ {
Code: 1, Code: 1,
Msg: &OfferedHashesMsg{ Msg: &OfferedHashesMsg{
HandoverProof: &HandoverProof{
Handover: &Handover{},
},
Hashes: hashes, Hashes: hashes,
From: 5, From: 5,
To: 8, To: 8,
@ -670,9 +649,6 @@ func TestStreamerRequestSubscriptionQuitMsgExchange(t *testing.T) {
Code: 1, Code: 1,
Msg: &OfferedHashesMsg{ Msg: &OfferedHashesMsg{
Stream: NewStream("foo", "", false), Stream: NewStream("foo", "", false),
HandoverProof: &HandoverProof{
Handover: &Handover{},
},
Hashes: make([]byte, HashSize), Hashes: make([]byte, HashSize),
From: 6, From: 6,
To: 9, To: 9,
@ -683,9 +659,6 @@ func TestStreamerRequestSubscriptionQuitMsgExchange(t *testing.T) {
Code: 1, Code: 1,
Msg: &OfferedHashesMsg{ Msg: &OfferedHashesMsg{
Stream: stream, Stream: stream,
HandoverProof: &HandoverProof{
Handover: &Handover{},
},
From: 11, From: 11,
To: 0, To: 0,
Hashes: make([]byte, HashSize), Hashes: make([]byte, HashSize),
@ -787,9 +760,6 @@ func TestMaxPeerServersWithUnsubscribe(t *testing.T) {
Code: 1, Code: 1,
Msg: &OfferedHashesMsg{ Msg: &OfferedHashesMsg{
Stream: stream, Stream: stream,
HandoverProof: &HandoverProof{
Handover: &Handover{},
},
Hashes: make([]byte, HashSize), Hashes: make([]byte, HashSize),
From: 1, From: 1,
To: 0, To: 0,
@ -891,9 +861,6 @@ func TestMaxPeerServersWithoutUnsubscribe(t *testing.T) {
Code: 1, Code: 1,
Msg: &OfferedHashesMsg{ Msg: &OfferedHashesMsg{
Stream: stream, Stream: stream,
HandoverProof: &HandoverProof{
Handover: &Handover{},
},
Hashes: make([]byte, HashSize), Hashes: make([]byte, HashSize),
From: 1, From: 1,
To: 0, To: 0,

View file

@ -92,7 +92,7 @@ func (s *SwarmSyncerServer) SessionIndex() (uint64, error) {
// chunk addresses. If at least one chunk is added to the batch and no new chunks // chunk addresses. If at least one chunk is added to the batch and no new chunks
// are added in batchTimeout period, the batch will be returned. This function // are added in batchTimeout period, the batch will be returned. This function
// will block until new chunks are received from localstore pull subscription. // will block until new chunks are received from localstore pull subscription.
func (s *SwarmSyncerServer) SetNextBatch(from, to uint64) ([]byte, uint64, uint64, *HandoverProof, error) { func (s *SwarmSyncerServer) SetNextBatch(from, to uint64) ([]byte, uint64, uint64, error) {
batchStart := time.Now() batchStart := time.Now()
descriptors, stop := s.netStore.SubscribePull(context.Background(), s.po, from, to) descriptors, stop := s.netStore.SubscribePull(context.Background(), s.po, from, to)
defer stop() defer stop()
@ -131,7 +131,7 @@ func (s *SwarmSyncerServer) SetNextBatch(from, to uint64) ([]byte, uint64, uint6
if err != nil { if err != nil {
metrics.GetOrRegisterCounter("syncer.set-next-batch.set-sync-err", nil).Inc(1) metrics.GetOrRegisterCounter("syncer.set-next-batch.set-sync-err", nil).Inc(1)
log.Debug("syncer pull subscription - err setting chunk as synced", "correlateId", s.correlateId, "err", err) log.Debug("syncer pull subscription - err setting chunk as synced", "correlateId", s.correlateId, "err", err)
return nil, 0, 0, nil, err return nil, 0, 0, err
} }
batchSize++ batchSize++
if batchStartID == nil { if batchStartID == nil {
@ -171,7 +171,7 @@ func (s *SwarmSyncerServer) SetNextBatch(from, to uint64) ([]byte, uint64, uint6
// if batch start id is not set, return 0 // if batch start id is not set, return 0
batchStartID = new(uint64) batchStartID = new(uint64)
} }
return batch, *batchStartID, batchEndID, nil, nil return batch, *batchStartID, batchEndID, nil
} }
// SwarmSyncerClient // SwarmSyncerClient