From f8001a6891def53b09537360a48a74ebef997eaf Mon Sep 17 00:00:00 2001 From: lash Date: Fri, 8 Feb 2019 17:52:49 +0100 Subject: [PATCH] swarm/newtork: WIP Span request span until delivery and put --- swarm/network/stream/delivery.go | 19 +++++++++++++------ swarm/network/stream/peer.go | 18 ++++++++++-------- 2 files changed, 23 insertions(+), 14 deletions(-) diff --git a/swarm/network/stream/delivery.go b/swarm/network/stream/delivery.go index fae6994f0c..d0417a39ac 100644 --- a/swarm/network/stream/delivery.go +++ b/swarm/network/stream/delivery.go @@ -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 } diff --git a/swarm/network/stream/peer.go b/swarm/network/stream/peer.go index 68da8f44ad..690e84f7ff 100644 --- a/swarm/network/stream/peer.go +++ b/swarm/network/stream/peer.go @@ -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)