Skip to content
项目
群组
代码片段
帮助
当前项目
正在载入...
登录 / 注册
切换导航面板
G
go-ipfs
概览
概览
详情
活动
周期分析
版本库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
统计图
问题
0
议题
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
CI / CD
CI / CD
流水线
作业
日程表
图表
维基
Wiki
代码片段
代码片段
成员
成员
折叠边栏
关闭边栏
活动
图像
聊天
创建新问题
作业
提交
问题看板
Open sidebar
jihao
go-ipfs
Commits
d357b0ac
提交
d357b0ac
authored
1月 03, 2015
作者:
Juan Batiz-Benet
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
bitswap debug logging
上级
c100390a
隐藏空白字符变更
内嵌
并排
正在显示
4 个修改的文件
包含
26 行增加
和
25 行删除
+26
-25
bitswap.go
exchange/bitswap/bitswap.go
+16
-17
engine.go
exchange/bitswap/decision/engine.go
+9
-2
ipfs_impl.go
exchange/bitswap/network/ipfs_impl.go
+1
-0
dht_net.go
routing/dht/dht_net.go
+0
-6
没有找到文件。
exchange/bitswap/bitswap.go
浏览文件 @
d357b0ac
...
...
@@ -3,7 +3,6 @@
package
bitswap
import
(
"fmt"
"math"
"sync"
"time"
...
...
@@ -172,14 +171,14 @@ func (bs *bitswap) HasBlock(ctx context.Context, blk *blocks.Block) error {
}
func
(
bs
*
bitswap
)
sendWantlistMsgToPeer
(
ctx
context
.
Context
,
m
bsmsg
.
BitSwapMessage
,
p
peer
.
ID
)
error
{
log
d
:=
fmt
.
Sprintf
(
"%s
bitswap.sendWantlistMsgToPeer(%d, %s)"
,
bs
.
self
,
len
(
m
.
Wantlist
()),
p
)
log
:=
log
.
Prefix
(
"bitswap(%s).
bitswap.sendWantlistMsgToPeer(%d, %s)"
,
bs
.
self
,
len
(
m
.
Wantlist
()),
p
)
log
.
Debug
f
(
"%s sending wantlist"
,
logd
)
log
.
Debug
(
"sending wantlist"
)
if
err
:=
bs
.
send
(
ctx
,
p
,
m
);
err
!=
nil
{
log
.
Errorf
(
"
%s send wantlist error: %s"
,
logd
,
err
)
log
.
Errorf
(
"
send wantlist error: %s"
,
err
)
return
err
}
log
.
Debugf
(
"
%s send wantlist success"
,
logd
)
log
.
Debugf
(
"
send wantlist success"
)
return
nil
}
...
...
@@ -188,20 +187,20 @@ func (bs *bitswap) sendWantlistMsgToPeers(ctx context.Context, m bsmsg.BitSwapMe
panic
(
"Cant send wantlist to nil peerchan"
)
}
log
d
:=
fmt
.
Sprintf
(
"%s bitswap.sendWantlistMsgTo
(%d)"
,
bs
.
self
,
len
(
m
.
Wantlist
()))
log
.
Debugf
(
"
%s begin"
,
logd
)
defer
log
.
Debugf
(
"
%s end"
,
logd
)
log
:=
log
.
Prefix
(
"bitswap(%s).sendWantlistMsgToPeers
(%d)"
,
bs
.
self
,
len
(
m
.
Wantlist
()))
log
.
Debugf
(
"
begin"
)
defer
log
.
Debugf
(
"
end"
)
set
:=
pset
.
New
()
wg
:=
sync
.
WaitGroup
{}
for
peerToQuery
:=
range
peers
{
log
.
Event
(
ctx
,
"PeerToQuery"
,
peerToQuery
)
logd
:=
fmt
.
Sprintf
(
"%sto(%s)"
,
logd
,
peerToQuery
)
if
!
set
.
TryAdd
(
peerToQuery
)
{
//Do once per peer
log
.
Debugf
(
"%s skipped (already sent)"
,
logd
)
log
.
Debugf
(
"%s skipped (already sent)"
,
peerToQuery
)
continue
}
log
.
Debugf
(
"%s sending"
,
peerToQuery
)
wg
.
Add
(
1
)
go
func
(
p
peer
.
ID
)
{
...
...
@@ -223,9 +222,9 @@ func (bs *bitswap) sendWantlistToPeers(ctx context.Context, peers <-chan peer.ID
}
func
(
bs
*
bitswap
)
sendWantlistToProviders
(
ctx
context
.
Context
)
{
log
d
:=
fmt
.
Sprintf
(
"%s bitswap.sendWantlistToProviders
"
,
bs
.
self
)
log
.
Debugf
(
"
%s begin"
,
logd
)
defer
log
.
Debugf
(
"
%s end"
,
logd
)
log
:=
log
.
Prefix
(
"bitswap(%s).sendWantlistToProviders
"
,
bs
.
self
)
log
.
Debugf
(
"
begin"
)
defer
log
.
Debugf
(
"
end"
)
ctx
,
cancel
:=
context
.
WithCancel
(
ctx
)
defer
cancel
()
...
...
@@ -240,13 +239,13 @@ func (bs *bitswap) sendWantlistToProviders(ctx context.Context) {
go
func
(
k
u
.
Key
)
{
defer
wg
.
Done
()
log
d
:=
fmt
.
Sprintf
(
"%s(entry: %s)"
,
logd
,
k
)
log
.
Debug
f
(
"%s asking dht for providers"
,
logd
)
log
:=
log
.
Prefix
(
"(entry: %s) "
,
k
)
log
.
Debug
(
"asking dht for providers"
)
child
,
_
:=
context
.
WithTimeout
(
ctx
,
providerRequestTimeout
)
providers
:=
bs
.
network
.
FindProvidersAsync
(
child
,
k
,
maxProvidersPerRequest
)
for
prov
:=
range
providers
{
log
.
Debugf
(
"
%s dht returned provider %s. send wantlist"
,
logd
,
prov
)
log
.
Debugf
(
"
dht returned provider %s. send wantlist"
,
prov
)
sendToPeers
<-
prov
}
}(
e
.
Key
)
...
...
@@ -259,7 +258,7 @@ func (bs *bitswap) sendWantlistToProviders(ctx context.Context) {
err
:=
bs
.
sendWantlistToPeers
(
ctx
,
sendToPeers
)
if
err
!=
nil
{
log
.
Errorf
(
"
%s sendWantlistToPeers error: %s"
,
logd
,
err
)
log
.
Errorf
(
"
sendWantlistToPeers error: %s"
,
err
)
}
}
...
...
exchange/bitswap/decision/engine.go
浏览文件 @
d357b0ac
...
...
@@ -8,7 +8,7 @@ import (
bsmsg
"github.com/jbenet/go-ipfs/exchange/bitswap/message"
wl
"github.com/jbenet/go-ipfs/exchange/bitswap/wantlist"
peer
"github.com/jbenet/go-ipfs/p2p/peer"
u
"github.com/jbenet/go-ipfs/util
"
eventlog
"github.com/jbenet/go-ipfs/util/eventlog
"
)
// TODO consider taking responsibility for other types of requests. For
...
...
@@ -41,7 +41,7 @@ import (
// whatever it sees fit to produce desired outcomes (get wanted keys
// quickly, maintain good relationships with peers, etc).
var
log
=
u
.
Logger
(
"engine"
)
var
log
=
eventlog
.
Logger
(
"engine"
)
const
(
sizeOutboxChan
=
4
...
...
@@ -140,6 +140,10 @@ func (e *Engine) Peers() []peer.ID {
// MessageReceived performs book-keeping. Returns error if passed invalid
// arguments.
func
(
e
*
Engine
)
MessageReceived
(
p
peer
.
ID
,
m
bsmsg
.
BitSwapMessage
)
error
{
log
:=
log
.
Prefix
(
"Engine.MessageReceived(%s)"
,
p
)
log
.
Debugf
(
"enter"
)
defer
log
.
Debugf
(
"exit"
)
newWorkExists
:=
false
defer
func
()
{
if
newWorkExists
{
...
...
@@ -156,9 +160,11 @@ func (e *Engine) MessageReceived(p peer.ID, m bsmsg.BitSwapMessage) error {
}
for
_
,
entry
:=
range
m
.
Wantlist
()
{
if
entry
.
Cancel
{
log
.
Debug
(
"cancel"
,
entry
.
Key
)
l
.
CancelWant
(
entry
.
Key
)
e
.
peerRequestQueue
.
Remove
(
entry
.
Key
,
p
)
}
else
{
log
.
Debug
(
"wants"
,
entry
.
Key
,
entry
.
Priority
)
l
.
Wants
(
entry
.
Key
,
entry
.
Priority
)
if
exists
,
err
:=
e
.
bs
.
Has
(
entry
.
Key
);
err
==
nil
&&
exists
{
newWorkExists
=
true
...
...
@@ -169,6 +175,7 @@ func (e *Engine) MessageReceived(p peer.ID, m bsmsg.BitSwapMessage) error {
for
_
,
block
:=
range
m
.
Blocks
()
{
// FIXME extract blocks.NumBytes(block) or block.NumBytes() method
log
.
Debug
(
"got block %s %d bytes"
,
block
.
Key
(),
len
(
block
.
Data
))
l
.
ReceivedBytes
(
len
(
block
.
Data
))
for
_
,
l
:=
range
e
.
ledgerMap
{
if
l
.
WantListContains
(
block
.
Key
())
{
...
...
exchange/bitswap/network/ipfs_impl.go
浏览文件 @
d357b0ac
...
...
@@ -55,6 +55,7 @@ func (bsnet *impl) SendRequest(
p
peer
.
ID
,
outgoing
bsmsg
.
BitSwapMessage
)
(
bsmsg
.
BitSwapMessage
,
error
)
{
log
.
Debugf
(
"bsnet SendRequest to %s"
,
p
)
s
,
err
:=
bsnet
.
host
.
NewStream
(
ProtocolBitswap
,
p
)
if
err
!=
nil
{
return
nil
,
err
...
...
routing/dht/dht_net.go
浏览文件 @
d357b0ac
...
...
@@ -87,15 +87,11 @@ func (dht *IpfsDHT) sendRequest(ctx context.Context, p peer.ID, pmes *pb.Message
start
:=
time
.
Now
()
log
.
Debugf
(
"%s writing"
,
dht
.
self
)
if
err
:=
w
.
WriteMsg
(
pmes
);
err
!=
nil
{
return
nil
,
err
}
log
.
Event
(
ctx
,
"dhtSentMessage"
,
dht
.
self
,
p
,
pmes
)
log
.
Debugf
(
"%s reading"
,
dht
.
self
)
defer
log
.
Debugf
(
"%s done"
,
dht
.
self
)
rpmes
:=
new
(
pb
.
Message
)
if
err
:=
r
.
ReadMsg
(
rpmes
);
err
!=
nil
{
return
nil
,
err
...
...
@@ -125,12 +121,10 @@ func (dht *IpfsDHT) sendMessage(ctx context.Context, p peer.ID, pmes *pb.Message
cw
:=
ctxutil
.
NewWriter
(
ctx
,
s
)
// ok to use. we defer close stream in this func
w
:=
ggio
.
NewDelimitedWriter
(
cw
)
log
.
Debugf
(
"%s writing"
,
dht
.
self
)
if
err
:=
w
.
WriteMsg
(
pmes
);
err
!=
nil
{
return
err
}
log
.
Event
(
ctx
,
"dhtSentMessage"
,
dht
.
self
,
p
,
pmes
)
log
.
Debugf
(
"%s done"
,
dht
.
self
)
return
nil
}
...
...
编写
预览
Markdown
格式
0%
重试
或
添加新文件
添加附件
取消
您添加了
0
人
到此讨论。请谨慎行事。
请先完成此评论的编辑!
取消
请
注册
或者
登录
后发表评论