Rename messages

This commit is contained in:
Balint Gabor 2018-01-12 14:01:21 +01:00
parent 852a6d669b
commit 23861e0fa8

View file

@ -82,9 +82,9 @@ type SubscribeMsg struct {
Priority uint8 // delivered on priority channel Priority uint8 // delivered on priority channel
} }
// UnsyncedKeysMsg is the protocol msg for offering to hand over a // OfferedHashesMsg is the protocol msg for offering to hand over a
// stream section // stream section
type UnsyncedKeysMsg struct { type OfferedHashesMsg struct {
Stream string // name of Stream Stream string // name of Stream
Key []byte // subtype or key Key []byte // subtype or key
From, To uint64 // peer and db-specific entry count From, To uint64 // peer and db-specific entry count
@ -104,22 +104,22 @@ type ChunkDeliveryMsg struct {
from Peer // [not serialised] protocol registers the requester from Peer // [not serialised] protocol registers the requester
} }
// String pretty prints UnsyncedKeysMsg // String pretty prints OfferedHashesMsg
func (self UnsyncedKeysMsg) String() string { func (self OfferedHashesMsg) String() string {
return fmt.Sprintf("Stream '%v' [%v-%v] (%v)", self.Stream, self.From, self.To, len(self.Hashes)/HashSize) return fmt.Sprintf("Stream '%v' [%v-%v] (%v)", self.Stream, self.From, self.To, len(self.Hashes)/HashSize)
} }
// WantedKeysMsg is the protocol msg data for signaling which hashes // WantedHashesMsg is the protocol msg data for signaling which hashes
// offered in UnsyncedKeysMsg downstream peer actually wants sent over // offered in OfferedHashesMsg downstream peer actually wants sent over
type WantedKeysMsg struct { type WantedHashesMsg struct {
Stream string // name of stream Stream string // name of stream
Key []byte // subtype or key Key []byte // subtype or key
Want []byte // bitvector indicating which keys of the batch needed Want []byte // bitvector indicating which keys of the batch needed
From, To uint64 // next interval offset - empty if not to be continued From, To uint64 // next interval offset - empty if not to be continued
} }
// String pretty prints WantedKeysMsg // String pretty prints WantedHashesMsg
func (self WantedKeysMsg) String() string { func (self WantedHashesMsg) String() string {
return fmt.Sprintf("Stream '%v', Want: %x, Next: [%v-%v]", self.Stream, self.Want, self.From, self.To) return fmt.Sprintf("Stream '%v', Want: %x, Next: [%v-%v]", self.Stream, self.Want, self.From, self.To)
} }
@ -410,7 +410,7 @@ func (self *StreamerPeer) setIncomingStreamer(s string, i IncomingStreamer, prio
priority: priority, priority: priority,
next: next, next: next,
} }
next <- struct{}{} // this is to allow wantedKeysMsg before first batch arrives next <- struct{}{} // this is to allow wantedHashesMsg before first batch arrives
return nil return nil
} }
@ -489,13 +489,13 @@ func (self *StreamerPeer) handleSubscribeMsg(req *SubscribeMsg) error {
if err != nil { if err != nil {
return nil return nil
} }
go self.SendUnsyncedKeys(os, req.From, req.To) go self.SendOfferedHashes(os, req.From, req.To)
return nil return nil
} }
// handleUnsyncedKeysMsg protocol msg handler calls the incoming streamer interface // handleOfferedHashesMsg protocol msg handler calls the incoming streamer interface
// Filter method // Filter method
func (self *StreamerPeer) handleUnsyncedKeysMsg(req *UnsyncedKeysMsg) error { func (self *StreamerPeer) handleOfferedHashesMsg(req *OfferedHashesMsg) error {
s, err := self.getIncomingStreamer(req.Stream) s, err := self.getIncomingStreamer(req.Stream)
if err != nil { if err != nil {
return err return err
@ -529,7 +529,7 @@ func (self *StreamerPeer) handleUnsyncedKeysMsg(req *UnsyncedKeysMsg) error {
} }
s.next <- struct{}{} s.next <- struct{}{}
}() }()
// only send wantedKeysMsg if all missing chunks of the previous batch arrived // only send wantedHashesMsg if all missing chunks of the previous batch arrived
// except // except
if s.live { if s.live {
s.sessionAt = req.From s.sessionAt = req.From
@ -538,7 +538,7 @@ func (self *StreamerPeer) handleUnsyncedKeysMsg(req *UnsyncedKeysMsg) error {
if from == to { if from == to {
return nil return nil
} }
msg := &WantedKeysMsg{ msg := &WantedHashesMsg{
Stream: req.Stream, Stream: req.Stream,
Want: want.Bytes(), Want: want.Bytes(),
From: from, From: from,
@ -555,17 +555,17 @@ func (self *StreamerPeer) handleUnsyncedKeysMsg(req *UnsyncedKeysMsg) error {
return nil return nil
} }
// handleWantedKeysMsg protocol msg handler // handleWantedHashesMsg protocol msg handler
// * sends the next batch of unsynced keys // * sends the next batch of unsynced keys
// * sends the actual data chunks as per WantedKeysMsg // * sends the actual data chunks as per WantedHashesMsg
func (self *StreamerPeer) handleWantedKeysMsg(req *WantedKeysMsg) error { func (self *StreamerPeer) handleWantedHashesMsg(req *WantedHashesMsg) error {
s, err := self.getOutgoingStreamer(req.Stream) s, err := self.getOutgoingStreamer(req.Stream)
if err != nil { if err != nil {
return err return err
} }
hashes := s.currentBatch hashes := s.currentBatch
// launch in go routine since GetBatch blocks until new hashes arrive // launch in go routine since GetBatch blocks until new hashes arrive
go self.SendUnsyncedKeys(s, req.From, req.To) go self.SendOfferedHashes(s, req.From, req.To)
l := len(hashes) / HashSize l := len(hashes) / HashSize
want, err := bv.NewFromBytes(req.Want, l) want, err := bv.NewFromBytes(req.Want, l)
if err != nil { if err != nil {
@ -611,14 +611,14 @@ func (self *StreamerPeer) SendPriority(msg interface{}, priority uint8) error {
return self.pq.Push(nil, msg, int(priority)) return self.pq.Push(nil, msg, int(priority))
} }
// UnsyncedKeys sends UnsyncedKeysMsg protocol msg // OfferedHashes sends OfferedHashesMsg protocol msg
func (self *StreamerPeer) SendUnsyncedKeys(s *outgoingStreamer, f, t uint64) error { func (self *StreamerPeer) SendOfferedHashes(s *outgoingStreamer, f, t uint64) error {
hashes, from, to, proof, err := s.SetNextBatch(f, t) hashes, from, to, proof, err := s.SetNextBatch(f, t)
if err != nil { if err != nil {
return err return err
} }
s.currentBatch = hashes s.currentBatch = hashes
msg := &UnsyncedKeysMsg{ msg := &OfferedHashesMsg{
HandoverProof: proof, HandoverProof: proof,
Hashes: hashes, Hashes: hashes,
From: from, From: from,
@ -634,8 +634,8 @@ var StreamerSpec = &protocols.Spec{
MaxMsgSize: 10 * 1024 * 1024, MaxMsgSize: 10 * 1024 * 1024,
Messages: []interface{}{ Messages: []interface{}{
HandshakeMsg{}, HandshakeMsg{},
UnsyncedKeysMsg{}, OfferedHashesMsg{},
WantedKeysMsg{}, WantedHashesMsg{},
TakeoverProofMsg{}, TakeoverProofMsg{},
SubscribeMsg{}, SubscribeMsg{},
}, },
@ -668,14 +668,14 @@ func (self *StreamerPeer) HandleMsg(msg interface{}) error {
case *SubscribeMsg: case *SubscribeMsg:
return self.handleSubscribeMsg(msg) return self.handleSubscribeMsg(msg)
case *UnsyncedKeysMsg: case *OfferedHashesMsg:
return self.handleUnsyncedKeysMsg(msg) return self.handleOfferedHashesMsg(msg)
case *TakeoverProofMsg: case *TakeoverProofMsg:
return self.handleTakeoverProofMsg(msg) return self.handleTakeoverProofMsg(msg)
case *WantedKeysMsg: case *WantedHashesMsg:
return self.handleWantedKeysMsg(msg) return self.handleWantedHashesMsg(msg)
case *ChunkDeliveryMsg: case *ChunkDeliveryMsg:
return self.handleChunkDeliveryMsg(msg) return self.handleChunkDeliveryMsg(msg)