swarm/newtork: WIP Span request span until delivery and put

This commit is contained in:
lash 2019-02-08 17:52:49 +01:00
parent d3ccedc767
commit b7d9719340
2 changed files with 23 additions and 14 deletions

View file

@ -209,15 +209,22 @@ type ChunkDeliveryMsgSyncing ChunkDeliveryMsg
// TODO: Fix context SNAFU
func (d *Delivery) handleChunkDeliveryMsg(ctx context.Context, sp *Peer, req *ChunkDeliveryMsg) error {
var osp opentracing.Span
ctx, osp = spancontext.StartSpan(
ctx,
"chunk.delivery")
// var osp opentracing.Span
// ctx, osp = spancontext.StartSpan(
// ctx,
// "chunk.delivery")
spanId := fmt.Sprintf("request.%v.%v", sp.ID(), req.Addr)
span, spanOk := sp.spans.Load(spanId)
sp.spans.Delete(spanId)
processReceivedChunksCount.Inc(1)
go func() {
defer osp.Finish()
//defer osp.Finish()
if spanOk {
defer span.(opentracing.Span).Finish()
}
req.peer = sp
err := d.chunkStore.Put(ctx, storage.NewChunk(req.Addr, req.SData))
@ -272,7 +279,7 @@ func (d *Delivery) RequestFromPeers(ctx context.Context, req *network.Request) (
Addr: req.Addr,
SkipCheck: req.SkipCheck,
HopCount: req.HopCount,
}, Top, "request.from.peers")
}, Top, fmt.Sprintf("request.%v.%v", sp.ID(), req.Addr))
if err != nil {
return nil, nil, err
}

View file

@ -88,11 +88,11 @@ func NewPeer(peer *protocols.Peer, streamer *Registry) *Peer {
ctx, cancel := context.WithCancel(context.Background())
go p.pq.Run(ctx, func(i interface{}) {
wmsg := i.(WrappedPriorityMsg)
defer p.spans.Delete(wmsg.Context)
sp, ok := p.spans.Load(wmsg.Context)
if ok {
defer sp.(opentracing.Span).Finish()
}
// defer p.spans.Delete(wmsg.Context)
// sp, ok := p.spans.Load(wmsg.Context)
// if ok {
// defer sp.(opentracing.Span).Finish()
// }
err := p.Send(wmsg.Context, wmsg.Msg)
if err != nil {
log.Error("Message send error, dropping peer", "peer", p.ID(), "err", err)
@ -158,7 +158,8 @@ func (p *Peer) Deliver(ctx context.Context, chunk storage.Chunk, priority uint8,
spanName += ".retrieval"
}
return p.SendPriority(ctx, msg, priority, spanName)
//return p.SendPriority(ctx, msg, priority, spanName)
return p.SendPriority(ctx, msg, priority, "")
}
// SendPriority sends message to the peer using the outgoing priority queue
@ -171,7 +172,7 @@ func (p *Peer) SendPriority(ctx context.Context, msg interface{}, priority uint8
ctx,
traceId,
)
p.spans.Store(ctx, sp)
p.spans.Store(traceId, sp)
}
wmsg := WrappedPriorityMsg{
Context: ctx,
@ -190,7 +191,8 @@ func (p *Peer) SendOfferedHashes(s *server, f, t uint64) error {
var sp opentracing.Span
ctx, sp := spancontext.StartSpan(
context.TODO(),
"send.offered.hashes")
"",
)
defer sp.Finish()
hashes, from, to, proof, err := s.setNextBatch(f, t)