swarm/network/stream: disambiguate chunk delivery messages (retrieval vs syncing)

This commit is contained in:
Fabio Barone 2018-10-15 19:48:56 -05:00
parent 6c313fff7b
commit 492095764a
4 changed files with 46 additions and 12 deletions

View file

@ -173,7 +173,8 @@ func (d *Delivery) handleRetrieveRequestMsg(ctx context.Context, sp *Peer, req *
return return
} }
if req.SkipCheck { if req.SkipCheck {
err = sp.Deliver(ctx, chunk, s.priority) syncing := false
err = sp.Deliver(ctx, chunk, s.priority, syncing)
if err != nil { if err != nil {
log.Warn("ERROR in handleRetrieveRequestMsg", "err", err) log.Warn("ERROR in handleRetrieveRequestMsg", "err", err)
} }
@ -189,12 +190,23 @@ func (d *Delivery) handleRetrieveRequestMsg(ctx context.Context, sp *Peer, req *
return nil return nil
} }
//Chunk delivery always uses the same message type....
type ChunkDeliveryMsg struct { type ChunkDeliveryMsg struct {
Addr storage.Address Addr storage.Address
SData []byte // the stored chunk Data (incl size) SData []byte // the stored chunk Data (incl size)
peer *Peer // set in handleChunkDeliveryMsg peer *Peer // set in handleChunkDeliveryMsg
Syncing bool // if true, this is a delivery for syncing (no SWAP accounting needed)
} }
//...but swap accounting needs to disambiguate if it is a delivery for syncing or for retrieval
//as it decides based on message type if it needs to account for this message or not
//defines a chunk delivery for retrieval (with accounting)
type ChunkDeliveryMsgRetrieval ChunkDeliveryMsg
//defines a chunk delivery for syncing (without accounting)
type ChunkDeliveryMsgSyncing ChunkDeliveryMsg
// TODO: Fix context SNAFU // TODO: Fix context SNAFU
func (d *Delivery) handleChunkDeliveryMsg(ctx context.Context, sp *Peer, req *ChunkDeliveryMsg) error { func (d *Delivery) handleChunkDeliveryMsg(ctx context.Context, sp *Peer, req *ChunkDeliveryMsg) error {
var osp opentracing.Span var osp opentracing.Span

View file

@ -347,7 +347,8 @@ func (p *Peer) handleWantedHashesMsg(ctx context.Context, req *WantedHashesMsg)
return fmt.Errorf("handleWantedHashesMsg get data %x: %v", hash, err) return fmt.Errorf("handleWantedHashesMsg get data %x: %v", hash, err)
} }
chunk := storage.NewChunk(hash, data) chunk := storage.NewChunk(hash, data)
if err := p.Deliver(ctx, chunk, s.priority); err != nil { syncing := true
if err := p.Deliver(ctx, chunk, s.priority, syncing); err != nil {
return err return err
} }
} }

View file

@ -128,16 +128,31 @@ func NewPeer(peer *protocols.Peer, streamer *Registry) *Peer {
} }
// Deliver sends a storeRequestMsg protocol message to the peer // Deliver sends a storeRequestMsg protocol message to the peer
func (p *Peer) Deliver(ctx context.Context, chunk storage.Chunk, priority uint8) error { // Depending on the `syncing` parameter we send different message types
func (p *Peer) Deliver(ctx context.Context, chunk storage.Chunk, priority uint8, syncing bool) error {
var sp opentracing.Span var sp opentracing.Span
ctx, sp = spancontext.StartSpan( ctx, sp = spancontext.StartSpan(
ctx, ctx,
"send.chunk.delivery") "send.chunk.delivery")
defer sp.Finish() defer sp.Finish()
msg := &ChunkDeliveryMsg{ var msg interface{}
Addr: chunk.Address(),
SData: chunk.Data(), //we send different types of messages if delivery is for syncing or retrievals,
//even if handling and content of the message are the same,
//because swap accounting decides which messages need accounting based on the message type
if syncing {
msg = &ChunkDeliveryMsgSyncing{
Addr: chunk.Address(),
SData: chunk.Data(),
Syncing: syncing,
}
} else {
msg = &ChunkDeliveryMsgRetrieval{
Addr: chunk.Address(),
SData: chunk.Data(),
Syncing: syncing,
}
} }
return p.SendPriority(ctx, msg, priority) return p.SendPriority(ctx, msg, priority)
} }

View file

@ -478,8 +478,13 @@ func (p *Peer) HandleMsg(ctx context.Context, msg interface{}) error {
case *WantedHashesMsg: case *WantedHashesMsg:
return p.handleWantedHashesMsg(ctx, msg) return p.handleWantedHashesMsg(ctx, msg)
case *ChunkDeliveryMsg: case *ChunkDeliveryMsgRetrieval:
return p.streamer.delivery.handleChunkDeliveryMsg(ctx, p, msg) //handling chunk delivery is the same for retrieval and syncing, so let's cast the msg
return p.streamer.delivery.handleChunkDeliveryMsg(ctx, p, ((*ChunkDeliveryMsg)(msg)))
case *ChunkDeliveryMsgSyncing:
//handling chunk delivery is the same for retrieval and syncing, so let's cast the msg
return p.streamer.delivery.handleChunkDeliveryMsg(ctx, p, ((*ChunkDeliveryMsg)(msg)))
case *RetrieveRequestMsg: case *RetrieveRequestMsg:
return p.streamer.delivery.handleRetrieveRequestMsg(ctx, p, msg) return p.streamer.delivery.handleRetrieveRequestMsg(ctx, p, msg)
@ -679,10 +684,11 @@ var Spec = &protocols.Spec{
TakeoverProofMsg{}, TakeoverProofMsg{},
SubscribeMsg{}, SubscribeMsg{},
RetrieveRequestMsg{}, RetrieveRequestMsg{},
ChunkDeliveryMsg{}, ChunkDeliveryMsgRetrieval{},
SubscribeErrorMsg{}, SubscribeErrorMsg{},
RequestSubscriptionMsg{}, RequestSubscriptionMsg{},
QuitMsg{}, QuitMsg{},
ChunkDeliveryMsgSyncing{},
}, },
} }