mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-17 17:33:47 +00:00
Merge branch 'swarm-network-rewrite-syncer' into swarm-network-rewrite-syncer-intervals
This commit is contained in:
commit
366381c4b6
61 changed files with 455 additions and 431 deletions
36
.travis.yml
36
.travis.yml
|
|
@ -3,17 +3,6 @@ go_import_path: github.com/ethereum/go-ethereum
|
|||
sudo: false
|
||||
matrix:
|
||||
include:
|
||||
- os: linux
|
||||
dist: trusty
|
||||
sudo: required
|
||||
go: 1.7.x
|
||||
script:
|
||||
- sudo modprobe fuse
|
||||
- sudo chmod 666 /dev/fuse
|
||||
- sudo chown root:$USER /etc/fuse.conf
|
||||
- go run build/ci.go install
|
||||
- go run build/ci.go test -coverage
|
||||
|
||||
- os: linux
|
||||
dist: trusty
|
||||
sudo: required
|
||||
|
|
@ -25,7 +14,6 @@ matrix:
|
|||
- go run build/ci.go install
|
||||
- go run build/ci.go test -coverage
|
||||
|
||||
# These are the latest Go versions.
|
||||
- os: linux
|
||||
dist: trusty
|
||||
sudo: required
|
||||
|
|
@ -47,6 +35,28 @@ matrix:
|
|||
- go run build/ci.go install
|
||||
- go run build/ci.go test -coverage
|
||||
|
||||
# These are the latest Go versions.
|
||||
- os: linux
|
||||
dist: trusty
|
||||
sudo: required
|
||||
go: "1.10"
|
||||
script:
|
||||
- sudo modprobe fuse
|
||||
- sudo chmod 666 /dev/fuse
|
||||
- sudo chown root:$USER /etc/fuse.conf
|
||||
- go run build/ci.go install
|
||||
- go run build/ci.go test -coverage
|
||||
|
||||
- os: osx
|
||||
go: "1.10"
|
||||
script:
|
||||
- unset -f cd # workaround for https://github.com/travis-ci/travis-ci/issues/8703
|
||||
- brew update
|
||||
- brew install caskroom/cask/brew-cask
|
||||
- brew cask install osxfuse
|
||||
- go run build/ci.go install
|
||||
- go run build/ci.go test -coverage
|
||||
|
||||
# This builder only tests code linters on latest version of Go
|
||||
- os: linux
|
||||
dist: trusty
|
||||
|
|
@ -185,6 +195,8 @@ matrix:
|
|||
- xctool -version
|
||||
- xcrun simctl list
|
||||
|
||||
# Workaround for https://github.com/golang/go/issues/23749
|
||||
- export CGO_CFLAGS_ALLOW='-fmodules|-fblocks|-fobjc-arc'
|
||||
- go run build/ci.go xcode -signer IOS_SIGNING_KEY -deploy trunk -upload gethstore/builds
|
||||
|
||||
# This builder does the Azure archive purges to avoid accumulating junk
|
||||
|
|
|
|||
|
|
@ -281,8 +281,8 @@ func TestDeposit(t *testing.T) {
|
|||
t.Fatalf("expected balance %v, got %v", exp, chbook.Balance())
|
||||
}
|
||||
|
||||
// autodeposit every 30ms if new cheque issued
|
||||
interval := 30 * time.Millisecond
|
||||
// autodeposit every 200ms if new cheque issued
|
||||
interval := 200 * time.Millisecond
|
||||
chbook.AutoDeposit(interval, common.Big1, balance)
|
||||
_, err = chbook.Issue(addr1, amount)
|
||||
if err != nil {
|
||||
|
|
|
|||
|
|
@ -2307,7 +2307,7 @@ var toChecksumAddress = function (address) {
|
|||
};
|
||||
|
||||
/**
|
||||
* Transforms given string to valid 20 bytes-length addres with 0x prefix
|
||||
* Transforms given string to valid 20 bytes-length address with 0x prefix
|
||||
*
|
||||
* @method toAddress
|
||||
* @param {String} address
|
||||
|
|
|
|||
|
|
@ -126,13 +126,13 @@ type logger struct {
|
|||
h *swapHandler
|
||||
}
|
||||
|
||||
func (l *logger) write(msg string, lvl Lvl, ctx []interface{}) {
|
||||
func (l *logger) write(msg string, lvl Lvl, ctx []interface{}, skip int) {
|
||||
l.h.Log(&Record{
|
||||
Time: time.Now(),
|
||||
Lvl: lvl,
|
||||
Msg: msg,
|
||||
Ctx: newContext(l.ctx, ctx),
|
||||
Call: stack.Caller(2),
|
||||
Call: stack.Caller(skip),
|
||||
KeyNames: RecordKeyNames{
|
||||
Time: timeKey,
|
||||
Msg: msgKey,
|
||||
|
|
@ -156,27 +156,27 @@ func newContext(prefix []interface{}, suffix []interface{}) []interface{} {
|
|||
}
|
||||
|
||||
func (l *logger) Trace(msg string, ctx ...interface{}) {
|
||||
l.write(msg, LvlTrace, ctx)
|
||||
l.write(msg, LvlTrace, ctx, 2)
|
||||
}
|
||||
|
||||
func (l *logger) Debug(msg string, ctx ...interface{}) {
|
||||
l.write(msg, LvlDebug, ctx)
|
||||
l.write(msg, LvlDebug, ctx, 2)
|
||||
}
|
||||
|
||||
func (l *logger) Info(msg string, ctx ...interface{}) {
|
||||
l.write(msg, LvlInfo, ctx)
|
||||
l.write(msg, LvlInfo, ctx, 2)
|
||||
}
|
||||
|
||||
func (l *logger) Warn(msg string, ctx ...interface{}) {
|
||||
l.write(msg, LvlWarn, ctx)
|
||||
l.write(msg, LvlWarn, ctx, 2)
|
||||
}
|
||||
|
||||
func (l *logger) Error(msg string, ctx ...interface{}) {
|
||||
l.write(msg, LvlError, ctx)
|
||||
l.write(msg, LvlError, ctx, 2)
|
||||
}
|
||||
|
||||
func (l *logger) Crit(msg string, ctx ...interface{}) {
|
||||
l.write(msg, LvlCrit, ctx)
|
||||
l.write(msg, LvlCrit, ctx, 2)
|
||||
os.Exit(1)
|
||||
}
|
||||
|
||||
|
|
|
|||
17
log/root.go
17
log/root.go
|
|
@ -31,31 +31,36 @@ func Root() Logger {
|
|||
|
||||
// Trace is a convenient alias for Root().Trace
|
||||
func Trace(msg string, ctx ...interface{}) {
|
||||
root.write(msg, LvlTrace, ctx)
|
||||
root.write(msg, LvlTrace, ctx, 2)
|
||||
}
|
||||
|
||||
// Debug is a convenient alias for Root().Debug
|
||||
func Debug(msg string, ctx ...interface{}) {
|
||||
root.write(msg, LvlDebug, ctx)
|
||||
root.write(msg, LvlDebug, ctx, 2)
|
||||
}
|
||||
|
||||
// Info is a convenient alias for Root().Info
|
||||
func Info(msg string, ctx ...interface{}) {
|
||||
root.write(msg, LvlInfo, ctx)
|
||||
root.write(msg, LvlInfo, ctx, 2)
|
||||
}
|
||||
|
||||
// Warn is a convenient alias for Root().Warn
|
||||
func Warn(msg string, ctx ...interface{}) {
|
||||
root.write(msg, LvlWarn, ctx)
|
||||
root.write(msg, LvlWarn, ctx, 2)
|
||||
}
|
||||
|
||||
// Error is a convenient alias for Root().Error
|
||||
func Error(msg string, ctx ...interface{}) {
|
||||
root.write(msg, LvlError, ctx)
|
||||
root.write(msg, LvlError, ctx, 2)
|
||||
}
|
||||
|
||||
// Crit is a convenient alias for Root().Crit
|
||||
func Crit(msg string, ctx ...interface{}) {
|
||||
root.write(msg, LvlCrit, ctx)
|
||||
root.write(msg, LvlCrit, ctx, 2)
|
||||
os.Exit(1)
|
||||
}
|
||||
|
||||
// Output is a convenient alias for write
|
||||
func Output(msg string, lvl Lvl, skip int, ctx ...interface{}) {
|
||||
root.write(msg, lvl, ctx, skip)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -183,7 +183,7 @@ type Peer struct {
|
|||
|
||||
// NewPeer constructs a new peer
|
||||
// this constructor is called by the p2p.Protocol#Run function
|
||||
// the first two arguments are comming the arguments passed to p2p.Protocol.Run function
|
||||
// the first two arguments are coming the arguments passed to p2p.Protocol.Run function
|
||||
// the third argument is the CodeMap describing the protocol messages and options
|
||||
func NewPeer(p *p2p.Peer, rw p2p.MsgReadWriter, spec *Spec) *Peer {
|
||||
return &Peer{
|
||||
|
|
|
|||
|
|
@ -154,18 +154,18 @@ func protocolTester(t *testing.T, pp *p2ptest.TestPeerPool) *p2ptest.ProtocolTes
|
|||
func protoHandshakeExchange(id discover.NodeID, proto *protoHandshake) []p2ptest.Exchange {
|
||||
|
||||
return []p2ptest.Exchange{
|
||||
p2ptest.Exchange{
|
||||
{
|
||||
Expects: []p2ptest.Expect{
|
||||
p2ptest.Expect{
|
||||
{
|
||||
Code: 0,
|
||||
Msg: &protoHandshake{42, "420"},
|
||||
Peer: id,
|
||||
},
|
||||
},
|
||||
},
|
||||
p2ptest.Exchange{
|
||||
{
|
||||
Triggers: []p2ptest.Trigger{
|
||||
p2ptest.Trigger{
|
||||
{
|
||||
Code: 0,
|
||||
Msg: proto,
|
||||
Peer: id,
|
||||
|
|
@ -207,18 +207,18 @@ func TestProtoHandshakeSuccess(t *testing.T) {
|
|||
func moduleHandshakeExchange(id discover.NodeID, resp uint) []p2ptest.Exchange {
|
||||
|
||||
return []p2ptest.Exchange{
|
||||
p2ptest.Exchange{
|
||||
{
|
||||
Expects: []p2ptest.Expect{
|
||||
p2ptest.Expect{
|
||||
{
|
||||
Code: 1,
|
||||
Msg: &hs0{42},
|
||||
Peer: id,
|
||||
},
|
||||
},
|
||||
},
|
||||
p2ptest.Exchange{
|
||||
{
|
||||
Triggers: []p2ptest.Trigger{
|
||||
p2ptest.Trigger{
|
||||
{
|
||||
Code: 1,
|
||||
Msg: &hs0{resp},
|
||||
Peer: id,
|
||||
|
|
@ -255,42 +255,42 @@ func TestModuleHandshakeSuccess(t *testing.T) {
|
|||
func testMultiPeerSetup(a, b discover.NodeID) []p2ptest.Exchange {
|
||||
|
||||
return []p2ptest.Exchange{
|
||||
p2ptest.Exchange{
|
||||
{
|
||||
Label: "primary handshake",
|
||||
Expects: []p2ptest.Expect{
|
||||
p2ptest.Expect{
|
||||
{
|
||||
Code: 0,
|
||||
Msg: &protoHandshake{42, "420"},
|
||||
Peer: a,
|
||||
},
|
||||
p2ptest.Expect{
|
||||
{
|
||||
Code: 0,
|
||||
Msg: &protoHandshake{42, "420"},
|
||||
Peer: b,
|
||||
},
|
||||
},
|
||||
},
|
||||
p2ptest.Exchange{
|
||||
{
|
||||
Label: "module handshake",
|
||||
Triggers: []p2ptest.Trigger{
|
||||
p2ptest.Trigger{
|
||||
{
|
||||
Code: 0,
|
||||
Msg: &protoHandshake{42, "420"},
|
||||
Peer: a,
|
||||
},
|
||||
p2ptest.Trigger{
|
||||
{
|
||||
Code: 0,
|
||||
Msg: &protoHandshake{42, "420"},
|
||||
Peer: b,
|
||||
},
|
||||
},
|
||||
Expects: []p2ptest.Expect{
|
||||
p2ptest.Expect{
|
||||
{
|
||||
Code: 1,
|
||||
Msg: &hs0{42},
|
||||
Peer: a,
|
||||
},
|
||||
p2ptest.Expect{
|
||||
{
|
||||
Code: 1,
|
||||
Msg: &hs0{42},
|
||||
Peer: b,
|
||||
|
|
@ -298,10 +298,10 @@ func testMultiPeerSetup(a, b discover.NodeID) []p2ptest.Exchange {
|
|||
},
|
||||
},
|
||||
|
||||
p2ptest.Exchange{Label: "alternative module handshake", Triggers: []p2ptest.Trigger{p2ptest.Trigger{Code: 1, Msg: &hs0{41}, Peer: a},
|
||||
p2ptest.Trigger{Code: 1, Msg: &hs0{41}, Peer: b}}},
|
||||
p2ptest.Exchange{Label: "repeated module handshake", Triggers: []p2ptest.Trigger{p2ptest.Trigger{Code: 1, Msg: &hs0{1}, Peer: a}}},
|
||||
p2ptest.Exchange{Label: "receiving repeated module handshake", Expects: []p2ptest.Expect{p2ptest.Expect{Code: 1, Msg: &hs0{43}, Peer: a}}}}
|
||||
{Label: "alternative module handshake", Triggers: []p2ptest.Trigger{{Code: 1, Msg: &hs0{41}, Peer: a},
|
||||
{Code: 1, Msg: &hs0{41}, Peer: b}}},
|
||||
{Label: "repeated module handshake", Triggers: []p2ptest.Trigger{{Code: 1, Msg: &hs0{1}, Peer: a}}},
|
||||
{Label: "receiving repeated module handshake", Expects: []p2ptest.Expect{{Code: 1, Msg: &hs0{43}, Peer: a}}}}
|
||||
}
|
||||
|
||||
func runMultiplePeers(t *testing.T, peer int, errs ...error) {
|
||||
|
|
@ -327,7 +327,7 @@ func runMultiplePeers(t *testing.T, peer int, errs ...error) {
|
|||
// peer 0 sends kill request for peer with index <peer>
|
||||
s.TestExchanges(p2ptest.Exchange{
|
||||
Triggers: []p2ptest.Trigger{
|
||||
p2ptest.Trigger{
|
||||
{
|
||||
Code: 2,
|
||||
Msg: &kill{s.IDs[peer]},
|
||||
Peer: s.IDs[0],
|
||||
|
|
@ -338,7 +338,7 @@ func runMultiplePeers(t *testing.T, peer int, errs ...error) {
|
|||
// the peer not killed sends a drop request
|
||||
s.TestExchanges(p2ptest.Exchange{
|
||||
Triggers: []p2ptest.Trigger{
|
||||
p2ptest.Trigger{
|
||||
{
|
||||
Code: 3,
|
||||
Msg: &drop{},
|
||||
Peer: s.IDs[(peer+1)%2],
|
||||
|
|
@ -360,14 +360,14 @@ func runMultiplePeers(t *testing.T, peer int, errs ...error) {
|
|||
|
||||
}
|
||||
|
||||
func TestMultiplePeersDropSelf(t *testing.T) {
|
||||
func XTestMultiplePeersDropSelf(t *testing.T) {
|
||||
runMultiplePeers(t, 0,
|
||||
fmt.Errorf("subprotocol error"),
|
||||
fmt.Errorf("Message handler error: (msg code 3): dropped"),
|
||||
)
|
||||
}
|
||||
|
||||
func TestMultiplePeersDropOther(t *testing.T) {
|
||||
func XTestMultiplePeersDropOther(t *testing.T) {
|
||||
runMultiplePeers(t, 1,
|
||||
fmt.Errorf("Message handler error: (msg code 3): dropped"),
|
||||
fmt.Errorf("subprotocol error"),
|
||||
|
|
|
|||
|
|
@ -105,9 +105,9 @@ func (e *ExecAdapter) NewNode(config *NodeConfig) (Node, error) {
|
|||
conf.Stack.P2P.NAT = nil
|
||||
conf.Stack.NoUSB = true
|
||||
|
||||
// listen on a random localhost port (we'll get the actual port after
|
||||
// starting the node through the RPC admin.nodeInfo method)
|
||||
conf.Stack.P2P.ListenAddr = "127.0.0.1:0"
|
||||
// listen on a localhost port, which we set when we
|
||||
// initialise NodeConfig (usually a random port)
|
||||
conf.Stack.P2P.ListenAddr = fmt.Sprintf("127.0.0.1:%d", config.Port)
|
||||
|
||||
node := &ExecNode{
|
||||
ID: config.ID,
|
||||
|
|
|
|||
|
|
@ -34,11 +34,6 @@ import (
|
|||
"github.com/ethereum/go-ethereum/rpc"
|
||||
)
|
||||
|
||||
const (
|
||||
socketReadBuffer = 5000 * 1024
|
||||
socketWriteBuffer = 5000 * 1024
|
||||
)
|
||||
|
||||
// SimAdapter is a NodeAdapter which creates in-memory simulation nodes and
|
||||
// connects them using net.Pipe or OS socket connections
|
||||
type SimAdapter struct {
|
||||
|
|
@ -112,7 +107,7 @@ func (s *SimAdapter) NewNode(config *NodeConfig) (Node, error) {
|
|||
MaxPeers: math.MaxInt32,
|
||||
NoDiscovery: true,
|
||||
Dialer: s,
|
||||
EnableMsgEvents: false,
|
||||
EnableMsgEvents: config.EnableMsgEvents,
|
||||
},
|
||||
NoUSB: true,
|
||||
Logger: log.New("node.id", id.String()),
|
||||
|
|
@ -378,20 +373,10 @@ func socketPipe() (net.Conn, net.Conn, error) {
|
|||
return nil, nil, err
|
||||
}
|
||||
|
||||
err = setSocketBuffer(pipe1)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
err = setSocketBuffer(pipe2)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
return pipe1, pipe2, nil
|
||||
}
|
||||
|
||||
func setSocketBuffer(conn net.Conn) error {
|
||||
func setSocketBuffer(conn net.Conn, socketReadBuffer int, socketWriteBuffer int) error {
|
||||
switch v := conn.(type) {
|
||||
case *net.UnixConn:
|
||||
err := v.SetReadBuffer(socketReadBuffer)
|
||||
|
|
|
|||
|
|
@ -25,22 +25,29 @@ import (
|
|||
)
|
||||
|
||||
func TestSocketPipe(t *testing.T) {
|
||||
c1, c2, _ := socketPipe()
|
||||
c1, c2, err := socketPipe()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
done := make(chan struct{})
|
||||
|
||||
go func() {
|
||||
msgs := 20
|
||||
size := 8
|
||||
for i := 0; i < msgs; i++ {
|
||||
msg := make([]byte, size)
|
||||
_ = binary.PutUvarint(msg, uint64(i))
|
||||
|
||||
_, err := c1.Write(msg)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
// OS socket pipe is blocking (depending on buffer size on OS), so writes are emitted asynchronously
|
||||
go func() {
|
||||
for i := 0; i < msgs; i++ {
|
||||
msg := make([]byte, size)
|
||||
_ = binary.PutUvarint(msg, uint64(i))
|
||||
|
||||
_, err := c1.Write(msg)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
for i := 0; i < msgs; i++ {
|
||||
msg := make([]byte, size)
|
||||
|
|
@ -52,7 +59,7 @@ func TestSocketPipe(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
}
|
||||
|
||||
if bytes.Compare(msg, out) != 0 {
|
||||
if !bytes.Equal(msg, out) {
|
||||
t.Fatalf("expected %#v, got %#v", msg, out)
|
||||
}
|
||||
}
|
||||
|
|
@ -61,27 +68,34 @@ func TestSocketPipe(t *testing.T) {
|
|||
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(1 * time.Second):
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("test timeout")
|
||||
}
|
||||
}
|
||||
|
||||
func TestSocketPipeBidirections(t *testing.T) {
|
||||
c1, c2, _ := socketPipe()
|
||||
c1, c2, err := socketPipe()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
done := make(chan struct{})
|
||||
|
||||
go func() {
|
||||
msgs := 100
|
||||
size := 4
|
||||
for i := 0; i < msgs; i++ {
|
||||
msg := []byte(`ping`)
|
||||
|
||||
_, err := c1.Write(msg)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
// OS socket pipe is blocking (depending on buffer size on OS), so writes are emitted asynchronously
|
||||
go func() {
|
||||
for i := 0; i < msgs; i++ {
|
||||
msg := []byte(`ping`)
|
||||
|
||||
_, err := c1.Write(msg)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
for i := 0; i < msgs; i++ {
|
||||
out := make([]byte, size)
|
||||
|
|
@ -90,7 +104,7 @@ func TestSocketPipeBidirections(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
}
|
||||
|
||||
if bytes.Compare(out, []byte(`ping`)) == 0 {
|
||||
if bytes.Equal(out, []byte(`ping`)) {
|
||||
msg := []byte(`pong`)
|
||||
_, err := c2.Write(msg)
|
||||
if err != nil {
|
||||
|
|
@ -108,7 +122,7 @@ func TestSocketPipeBidirections(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
}
|
||||
|
||||
if bytes.Compare(out, expected) != 0 {
|
||||
if !bytes.Equal(out, expected) {
|
||||
t.Fatalf("expected %#v, got %#v", expected, out)
|
||||
}
|
||||
}
|
||||
|
|
@ -118,13 +132,16 @@ func TestSocketPipeBidirections(t *testing.T) {
|
|||
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(1 * time.Second):
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("test timeout")
|
||||
}
|
||||
}
|
||||
|
||||
func TestTcpPipe(t *testing.T) {
|
||||
c1, c2, _ := tcpPipe()
|
||||
c1, c2, err := tcpPipe()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
done := make(chan struct{})
|
||||
|
||||
|
|
@ -151,7 +168,7 @@ func TestTcpPipe(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
}
|
||||
|
||||
if bytes.Compare(msg, out) != 0 {
|
||||
if !bytes.Equal(msg, out) {
|
||||
t.Fatalf("expected %#v, got %#v", msg, out)
|
||||
}
|
||||
}
|
||||
|
|
@ -160,13 +177,16 @@ func TestTcpPipe(t *testing.T) {
|
|||
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(1 * time.Second):
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("test timeout")
|
||||
}
|
||||
}
|
||||
|
||||
func TestTcpPipeBidirections(t *testing.T) {
|
||||
c1, c2, _ := tcpPipe()
|
||||
c1, c2, err := tcpPipe()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
done := make(chan struct{})
|
||||
|
||||
|
|
@ -191,7 +211,7 @@ func TestTcpPipeBidirections(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
}
|
||||
|
||||
if bytes.Compare(expected, out) != 0 {
|
||||
if !bytes.Equal(expected, out) {
|
||||
t.Fatalf("expected %#v, got %#v", out, expected)
|
||||
} else {
|
||||
msg := []byte(fmt.Sprintf("pong %02d", i))
|
||||
|
|
@ -211,7 +231,7 @@ func TestTcpPipeBidirections(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
}
|
||||
|
||||
if bytes.Compare(expected, out) != 0 {
|
||||
if !bytes.Equal(expected, out) {
|
||||
t.Fatalf("expected %#v, got %#v", out, expected)
|
||||
}
|
||||
}
|
||||
|
|
@ -220,13 +240,16 @@ func TestTcpPipeBidirections(t *testing.T) {
|
|||
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(1 * time.Second):
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("test timeout")
|
||||
}
|
||||
}
|
||||
|
||||
func TestNetPipe(t *testing.T) {
|
||||
c1, c2, _ := netPipe()
|
||||
c1, c2, err := netPipe()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
done := make(chan struct{})
|
||||
|
||||
|
|
@ -256,7 +279,7 @@ func TestNetPipe(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
}
|
||||
|
||||
if bytes.Compare(msg, out) != 0 {
|
||||
if !bytes.Equal(msg, out) {
|
||||
t.Fatalf("expected %#v, got %#v", msg, out)
|
||||
}
|
||||
}
|
||||
|
|
@ -266,13 +289,16 @@ func TestNetPipe(t *testing.T) {
|
|||
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(1 * time.Second):
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("test timeout")
|
||||
}
|
||||
}
|
||||
|
||||
func TestNetPipeBidirections(t *testing.T) {
|
||||
c1, c2, _ := netPipe()
|
||||
c1, c2, err := netPipe()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
done := make(chan struct{})
|
||||
|
||||
|
|
@ -305,7 +331,7 @@ func TestNetPipeBidirections(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
}
|
||||
|
||||
if bytes.Compare(expected, out) != 0 {
|
||||
if !bytes.Equal(expected, out) {
|
||||
t.Fatalf("expected %#v, got %#v", expected, out)
|
||||
}
|
||||
}
|
||||
|
|
@ -323,7 +349,7 @@ func TestNetPipeBidirections(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
}
|
||||
|
||||
if bytes.Compare(expected, out) != 0 {
|
||||
if !bytes.Equal(expected, out) {
|
||||
t.Fatalf("expected %#v, got %#v", expected, out)
|
||||
} else {
|
||||
msg := []byte(fmt.Sprintf(pongTemplate, i))
|
||||
|
|
@ -338,7 +364,7 @@ func TestNetPipeBidirections(t *testing.T) {
|
|||
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(1 * time.Second):
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("test timeout")
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -23,6 +23,7 @@ import (
|
|||
"fmt"
|
||||
"net"
|
||||
"os"
|
||||
"strconv"
|
||||
|
||||
"github.com/docker/docker/pkg/reexec"
|
||||
"github.com/ethereum/go-ethereum/crypto"
|
||||
|
|
@ -97,24 +98,30 @@ type NodeConfig struct {
|
|||
|
||||
// function to sanction or prevent suggesting a peer
|
||||
Reachable func(id discover.NodeID) bool
|
||||
|
||||
Port uint16
|
||||
}
|
||||
|
||||
// nodeConfigJSON is used to encode and decode NodeConfig as JSON by encoding
|
||||
// all fields as strings
|
||||
type nodeConfigJSON struct {
|
||||
ID string `json:"id"`
|
||||
PrivateKey string `json:"private_key"`
|
||||
Name string `json:"name"`
|
||||
Services []string `json:"services"`
|
||||
ID string `json:"id"`
|
||||
PrivateKey string `json:"private_key"`
|
||||
Name string `json:"name"`
|
||||
Services []string `json:"services"`
|
||||
EnableMsgEvents bool `json:"enable_msg_events"`
|
||||
Port uint16 `json:"port"`
|
||||
}
|
||||
|
||||
// MarshalJSON implements the json.Marshaler interface by encoding the config
|
||||
// fields as strings
|
||||
func (n *NodeConfig) MarshalJSON() ([]byte, error) {
|
||||
confJSON := nodeConfigJSON{
|
||||
ID: n.ID.String(),
|
||||
Name: n.Name,
|
||||
Services: n.Services,
|
||||
ID: n.ID.String(),
|
||||
Name: n.Name,
|
||||
Services: n.Services,
|
||||
Port: n.Port,
|
||||
EnableMsgEvents: n.EnableMsgEvents,
|
||||
}
|
||||
if n.PrivateKey != nil {
|
||||
confJSON.PrivateKey = hex.EncodeToString(crypto.FromECDSA(n.PrivateKey))
|
||||
|
|
@ -152,6 +159,8 @@ func (n *NodeConfig) UnmarshalJSON(data []byte) error {
|
|||
|
||||
n.Name = confJSON.Name
|
||||
n.Services = confJSON.Services
|
||||
n.Port = confJSON.Port
|
||||
n.EnableMsgEvents = confJSON.EnableMsgEvents
|
||||
|
||||
return nil
|
||||
}
|
||||
|
|
@ -165,10 +174,34 @@ func RandomNodeConfig() *NodeConfig {
|
|||
}
|
||||
|
||||
id := discover.PubkeyID(&key.PublicKey)
|
||||
return &NodeConfig{
|
||||
ID: id,
|
||||
PrivateKey: key,
|
||||
port, err := assignTCPPort()
|
||||
if err != nil {
|
||||
panic("unable to assign tcp port")
|
||||
}
|
||||
return &NodeConfig{
|
||||
ID: id,
|
||||
Name: fmt.Sprintf("node_%s", id.String()),
|
||||
PrivateKey: key,
|
||||
Port: port,
|
||||
EnableMsgEvents: true,
|
||||
}
|
||||
}
|
||||
|
||||
func assignTCPPort() (uint16, error) {
|
||||
l, err := net.Listen("tcp", "127.0.0.1:0")
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
l.Close()
|
||||
_, port, err := net.SplitHostPort(l.Addr().String())
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
p, err := strconv.ParseInt(port, 10, 32)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
return uint16(p), nil
|
||||
}
|
||||
|
||||
// ServiceContext is a collection of options and methods which can be utilised
|
||||
|
|
|
|||
|
|
@ -561,7 +561,8 @@ func (s *Server) LoadSnapshot(w http.ResponseWriter, req *http.Request) {
|
|||
|
||||
// CreateNode creates a node in the network using the given configuration
|
||||
func (s *Server) CreateNode(w http.ResponseWriter, req *http.Request) {
|
||||
config := adapters.RandomNodeConfig()
|
||||
config := &adapters.NodeConfig{}
|
||||
|
||||
err := json.NewDecoder(req.Body).Decode(config)
|
||||
if err != nil && err != io.EOF {
|
||||
http.Error(w, err.Error(), http.StatusBadRequest)
|
||||
|
|
|
|||
|
|
@ -348,7 +348,8 @@ func startTestNetwork(t *testing.T, client *Client) []string {
|
|||
nodeCount := 2
|
||||
nodeIDs := make([]string, nodeCount)
|
||||
for i := 0; i < nodeCount; i++ {
|
||||
node, err := client.CreateNode(nil)
|
||||
config := adapters.RandomNodeConfig()
|
||||
node, err := client.CreateNode(config)
|
||||
if err != nil {
|
||||
t.Fatalf("error creating node: %s", err)
|
||||
}
|
||||
|
|
@ -527,7 +528,9 @@ func TestHTTPNodeRPC(t *testing.T) {
|
|||
|
||||
// start a node in the network
|
||||
client := NewClient(s.URL)
|
||||
node, err := client.CreateNode(nil)
|
||||
|
||||
config := adapters.RandomNodeConfig()
|
||||
node, err := client.CreateNode(config)
|
||||
if err != nil {
|
||||
t.Fatalf("error creating node: %s", err)
|
||||
}
|
||||
|
|
@ -589,7 +592,8 @@ func TestHTTPSnapshot(t *testing.T) {
|
|||
nodeCount := 2
|
||||
nodes := make([]*p2p.NodeInfo, nodeCount)
|
||||
for i := 0; i < nodeCount; i++ {
|
||||
node, err := client.CreateNode(nil)
|
||||
config := adapters.RandomNodeConfig()
|
||||
node, err := client.CreateNode(config)
|
||||
if err != nil {
|
||||
t.Fatalf("error creating node: %s", err)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -26,6 +26,7 @@ import (
|
|||
|
||||
"github.com/ethereum/go-ethereum/log"
|
||||
"github.com/ethereum/go-ethereum/p2p/discover"
|
||||
"github.com/ethereum/go-ethereum/p2p/simulations/adapters"
|
||||
)
|
||||
|
||||
//a map of mocker names to its function
|
||||
|
|
@ -165,7 +166,8 @@ func probabilistic(net *Network, quit chan struct{}, nodeCount int) {
|
|||
func connectNodesInRing(net *Network, nodeCount int) ([]discover.NodeID, error) {
|
||||
ids := make([]discover.NodeID, nodeCount)
|
||||
for i := 0; i < nodeCount; i++ {
|
||||
node, err := net.NewNode()
|
||||
conf := adapters.RandomNodeConfig()
|
||||
node, err := net.NewNodeWithConfig(conf)
|
||||
if err != nil {
|
||||
log.Error("Error creating a node! %s", err)
|
||||
return nil, err
|
||||
|
|
|
|||
|
|
@ -78,26 +78,12 @@ func (self *Network) Events() *event.Feed {
|
|||
return &self.events
|
||||
}
|
||||
|
||||
// NewNode adds a new node to the network with a random ID
|
||||
func (self *Network) NewNode() (*Node, error) {
|
||||
conf := adapters.RandomNodeConfig()
|
||||
conf.Services = []string{self.DefaultService}
|
||||
return self.NewNodeWithConfig(conf)
|
||||
}
|
||||
|
||||
// NewNodeWithConfig adds a new node to the network with the given config,
|
||||
// returning an error if a node with the same ID or name already exists
|
||||
func (self *Network) NewNodeWithConfig(conf *adapters.NodeConfig) (*Node, error) {
|
||||
self.lock.Lock()
|
||||
defer self.lock.Unlock()
|
||||
|
||||
// create a random ID and PrivateKey if not set
|
||||
if conf.ID == (discover.NodeID{}) {
|
||||
c := adapters.RandomNodeConfig()
|
||||
conf.ID = c.ID
|
||||
conf.PrivateKey = c.PrivateKey
|
||||
}
|
||||
id := conf.ID
|
||||
if conf.Reachable == nil {
|
||||
conf.Reachable = func(otherID discover.NodeID) bool {
|
||||
_, err := self.InitConn(conf.ID, otherID)
|
||||
|
|
@ -105,14 +91,9 @@ func (self *Network) NewNodeWithConfig(conf *adapters.NodeConfig) (*Node, error)
|
|||
}
|
||||
}
|
||||
|
||||
// assign a name to the node if not set
|
||||
if conf.Name == "" {
|
||||
conf.Name = fmt.Sprintf("node%02d", len(self.Nodes)+1)
|
||||
}
|
||||
|
||||
// check the node doesn't already exist
|
||||
if node := self.getNode(id); node != nil {
|
||||
return nil, fmt.Errorf("node with ID %q already exists", id)
|
||||
if node := self.getNode(conf.ID); node != nil {
|
||||
return nil, fmt.Errorf("node with ID %q already exists", conf.ID)
|
||||
}
|
||||
if node := self.getNodeByName(conf.Name); node != nil {
|
||||
return nil, fmt.Errorf("node with name %q already exists", conf.Name)
|
||||
|
|
@ -132,8 +113,8 @@ func (self *Network) NewNodeWithConfig(conf *adapters.NodeConfig) (*Node, error)
|
|||
Node: adapterNode,
|
||||
Config: conf,
|
||||
}
|
||||
log.Trace(fmt.Sprintf("node %v created", id))
|
||||
self.nodeMap[id] = len(self.Nodes)
|
||||
log.Trace(fmt.Sprintf("node %v created", conf.ID))
|
||||
self.nodeMap[conf.ID] = len(self.Nodes)
|
||||
self.Nodes = append(self.Nodes, node)
|
||||
|
||||
// emit a "control" event
|
||||
|
|
|
|||
|
|
@ -41,7 +41,8 @@ func TestNetworkSimulation(t *testing.T) {
|
|||
nodeCount := 20
|
||||
ids := make([]discover.NodeID, nodeCount)
|
||||
for i := 0; i < nodeCount; i++ {
|
||||
node, err := network.NewNode()
|
||||
conf := adapters.RandomNodeConfig()
|
||||
node, err := network.NewNodeWithConfig(conf)
|
||||
if err != nil {
|
||||
t.Fatalf("error creating node: %s", err)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -111,7 +111,7 @@ func posProximity(one, other Address, pos int) (ret int, eq bool) {
|
|||
start = pos % 8
|
||||
}
|
||||
for j := start; j < 8; j++ {
|
||||
if (uint8(oxo)>>uint8(7-j))&0x01 != 0 {
|
||||
if (oxo>>uint8(7-j))&0x01 != 0 {
|
||||
return i*8 + j, false
|
||||
}
|
||||
}
|
||||
|
|
@ -173,13 +173,13 @@ func RandomAddress() Address {
|
|||
func NewAddressFromString(s string) []byte {
|
||||
ha := [32]byte{}
|
||||
|
||||
t := s + string(zerosBin)[:len(zerosBin)-len(s)]
|
||||
t := s + zerosBin[:len(zerosBin)-len(s)]
|
||||
for i := 0; i < 4; i++ {
|
||||
n, err := strconv.ParseUint(t[i*64:(i+1)*64], 2, 64)
|
||||
if err != nil {
|
||||
panic("wrong format: " + err.Error())
|
||||
}
|
||||
binary.BigEndian.PutUint64(ha[i*8:(i+1)*8], uint64(n))
|
||||
binary.BigEndian.PutUint64(ha[i*8:(i+1)*8], n)
|
||||
}
|
||||
return ha[:]
|
||||
}
|
||||
|
|
@ -229,7 +229,7 @@ func proximityOrder(one, other []byte, pos int) (int, bool) {
|
|||
start = pos % 8
|
||||
}
|
||||
for j := start; j < 8; j++ {
|
||||
if (uint8(oxo)>>uint8(7-j))&0x01 != 0 {
|
||||
if (oxo>>uint8(7-j))&0x01 != 0 {
|
||||
return i*8 + j, false
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -48,8 +48,8 @@ concurrent routines,
|
|||
Pot
|
||||
* retrieval, insertion and deletion by key involves log(n) pointer lookups
|
||||
* for any item retrieval (defined as common prefix on the binary key)
|
||||
* provide syncronous iterators respecting proximity ordering wrt any item
|
||||
* provide asyncronous iterator (for parallel execution of operations) over n items
|
||||
* provide synchronous iterators respecting proximity ordering wrt any item
|
||||
* provide asynchronous iterator (for parallel execution of operations) over n items
|
||||
* allows cheap iteration over ranges
|
||||
* asymmetric concurrent merge (union)
|
||||
|
||||
|
|
|
|||
|
|
@ -559,7 +559,7 @@ func (t *Pot) eachBin(val Val, pof Pof, po int, f func(int, int, func(func(val V
|
|||
|
||||
}
|
||||
|
||||
// EachNeighbour is a syncronous iterator over neighbours of any target val
|
||||
// EachNeighbour is a synchronous iterator over neighbours of any target val
|
||||
// the order of elements retrieved reflect proximity order to the target
|
||||
// TODO: add maximum proxbin to start range of iteration
|
||||
func (t *Pot) EachNeighbour(val Val, pof Pof, f func(Val, int) bool) bool {
|
||||
|
|
@ -615,7 +615,7 @@ func (t *Pot) eachNeighbour(val Val, pof Pof, f func(Val, int) bool) bool {
|
|||
return true
|
||||
}
|
||||
|
||||
// EachNeighbourAsync called on (val, max, maxPos, f, wait) is an asyncronous iterator
|
||||
// EachNeighbourAsync called on (val, max, maxPos, f, wait) is an asynchronous iterator
|
||||
// over elements not closer than maxPos wrt val.
|
||||
// val does not need to be match an element of the Pot, but if it does, and
|
||||
// maxPos is keylength than it is included in the iteration
|
||||
|
|
@ -762,7 +762,7 @@ func (t *Pot) eachNeighbourAsync(val Val, pof Pof, max int, maxPos int, f func(V
|
|||
|
||||
// getPos called on (n) returns the forking node at PO n and its index if it exists
|
||||
// otherwise nil
|
||||
// caller is suppoed to hold the lock
|
||||
// caller is supposed to hold the lock
|
||||
func (t *Pot) getPos(po int) (n *Pot, i int) {
|
||||
for i, n = range t.bins {
|
||||
if po > n.po {
|
||||
|
|
|
|||
|
|
@ -271,10 +271,7 @@ func testPotEachNeighbour(n *Pot, pof Pof, val Val, expCount int, fs ...func(Val
|
|||
}
|
||||
}
|
||||
count++
|
||||
if count == expCount {
|
||||
return false
|
||||
}
|
||||
return true
|
||||
return count != expCount
|
||||
})
|
||||
if err == nil && count < expCount {
|
||||
return fmt.Errorf("not enough neighbours returned, expected %v, got %v", expCount, count)
|
||||
|
|
@ -558,10 +555,7 @@ func benchmarkEachNeighbourSync(t *testing.B, max, count int, d time.Duration) {
|
|||
n.EachNeighbour(val, pof, func(v Val, po int) bool {
|
||||
time.Sleep(d)
|
||||
m++
|
||||
if m == count {
|
||||
return false
|
||||
}
|
||||
return true
|
||||
return m != count
|
||||
})
|
||||
}
|
||||
t.StopTimer()
|
||||
|
|
|
|||
|
|
@ -109,10 +109,11 @@ func TestApiPut(t *testing.T) {
|
|||
content := "hello"
|
||||
exp := expResponse(content, "text/plain", 0)
|
||||
// exp := expResponse([]byte(content), "text/plain", 0)
|
||||
key, _, err := api.Put(content, exp.MimeType)
|
||||
key, wait, err := api.Put(content, exp.MimeType)
|
||||
if err != nil {
|
||||
t.Fatalf("unexpected error: %v", err)
|
||||
}
|
||||
wait()
|
||||
resp := testGet(t, api, key.Hex(), "")
|
||||
checkResponse(t, resp, exp)
|
||||
})
|
||||
|
|
|
|||
|
|
@ -110,7 +110,8 @@ func ShowMultipleChoices(w http.ResponseWriter, r *http.Request, list api.Manife
|
|||
//(and return the correct HTTP status code)
|
||||
func ShowError(w http.ResponseWriter, r *http.Request, msg string, code int) {
|
||||
if code == http.StatusInternalServerError {
|
||||
log.Error(msg)
|
||||
//log.Error(msg)
|
||||
log.Output(msg, log.LvlError, 3)
|
||||
}
|
||||
respond(w, r, &ErrorParams{
|
||||
Code: code,
|
||||
|
|
|
|||
|
|
@ -90,21 +90,21 @@ type Request struct {
|
|||
// body in swarm and returns the resulting storage key as a text/plain response
|
||||
func (s *Server) HandlePostRaw(w http.ResponseWriter, r *Request) {
|
||||
if r.uri.Path != "" {
|
||||
s.BadRequest(w, r, "raw POST request cannot contain a path")
|
||||
ShowError(w, &r.Request, fmt.Sprintf("Bad request %s %s: %s", r.Method, r.uri, "raw POST request cannot contain a path"), http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
|
||||
if r.Header.Get("Content-Length") == "" {
|
||||
s.BadRequest(w, r, "missing Content-Length header in request")
|
||||
ShowError(w, &r.Request, fmt.Sprintf("Bad request %s %s: %s", r.Method, r.uri, "missing Content-Length header in request"), http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
|
||||
key, _, err := s.api.Store(r.Body, r.ContentLength)
|
||||
if err != nil {
|
||||
s.Error(w, r, err)
|
||||
ShowError(w, &r.Request, fmt.Sprintf("Error serving %s %s: %s", r.Method, r.uri, err), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
s.logDebug("content for %s stored", key.Log())
|
||||
log.Debug(fmt.Sprintf("content for %s stored", key.Log()))
|
||||
|
||||
w.Header().Set("Content-Type", "text/plain")
|
||||
w.WriteHeader(http.StatusOK)
|
||||
|
|
@ -119,7 +119,7 @@ func (s *Server) HandlePostRaw(w http.ResponseWriter, r *Request) {
|
|||
func (s *Server) HandlePostFiles(w http.ResponseWriter, r *Request) {
|
||||
contentType, params, err := mime.ParseMediaType(r.Header.Get("Content-Type"))
|
||||
if err != nil {
|
||||
s.BadRequest(w, r, err.Error())
|
||||
ShowError(w, &r.Request, fmt.Sprintf("Bad request %s %s: %s", r.Method, r.uri, err), http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
|
||||
|
|
@ -127,13 +127,13 @@ func (s *Server) HandlePostFiles(w http.ResponseWriter, r *Request) {
|
|||
if r.uri.Addr != "" {
|
||||
key, err = s.api.Resolve(r.uri)
|
||||
if err != nil {
|
||||
s.Error(w, r, fmt.Errorf("error resolving %s: %s", r.uri.Addr, err))
|
||||
ShowError(w, &r.Request, fmt.Sprintf("Error serving %s %s: %s", r.Method, r.uri, fmt.Errorf("error resolving %s: %s", r.uri.Addr, err)), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
} else {
|
||||
key, err = s.api.NewManifest()
|
||||
if err != nil {
|
||||
s.Error(w, r, err)
|
||||
ShowError(w, &r.Request, fmt.Sprintf("Error serving %s %s: %s", r.Method, r.uri, err), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
}
|
||||
|
|
@ -152,7 +152,7 @@ func (s *Server) HandlePostFiles(w http.ResponseWriter, r *Request) {
|
|||
}
|
||||
})
|
||||
if err != nil {
|
||||
s.Error(w, r, fmt.Errorf("error creating manifest: %s", err))
|
||||
ShowError(w, &r.Request, fmt.Sprintf("Error serving %s %s: %s", r.Method, r.uri, fmt.Errorf("error creating manifest: %s", err)), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
||||
|
|
@ -185,12 +185,12 @@ func (s *Server) handleTarUpload(req *Request, mw *api.ManifestWriter) error {
|
|||
Size: hdr.Size,
|
||||
ModTime: hdr.ModTime,
|
||||
}
|
||||
s.logDebug("adding %s (%d bytes) to new manifest", entry.Path, entry.Size)
|
||||
log.Debug(fmt.Sprintf("adding %s (%d bytes) to new manifest", entry.Path, entry.Size))
|
||||
contentKey, err := mw.AddEntry(tr, entry)
|
||||
if err != nil {
|
||||
return fmt.Errorf("error adding manifest entry from tar stream: %s", err)
|
||||
}
|
||||
s.logDebug("content for %s stored", contentKey.Log())
|
||||
log.Debug(fmt.Sprintf("content for %s stored", contentKey.Log()))
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -242,12 +242,12 @@ func (s *Server) handleMultipartUpload(req *Request, boundary string, mw *api.Ma
|
|||
Size: size,
|
||||
ModTime: time.Now(),
|
||||
}
|
||||
s.logDebug("adding %s (%d bytes) to new manifest", entry.Path, entry.Size)
|
||||
log.Debug(fmt.Sprintf("adding %s (%d bytes) to new manifest", entry.Path, entry.Size))
|
||||
contentKey, err := mw.AddEntry(reader, entry)
|
||||
if err != nil {
|
||||
return fmt.Errorf("error adding manifest entry from multipart form: %s", err)
|
||||
}
|
||||
s.logDebug("content for %s stored", contentKey.Log())
|
||||
log.Debug(fmt.Sprintf("content for %s stored", contentKey.Log()))
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -262,7 +262,7 @@ func (s *Server) handleDirectUpload(req *Request, mw *api.ManifestWriter) error
|
|||
if err != nil {
|
||||
return err
|
||||
}
|
||||
s.logDebug("content for %s stored", key.Log())
|
||||
log.Debug(fmt.Sprintf("content for %s stored", key.Log()))
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -272,16 +272,16 @@ func (s *Server) handleDirectUpload(req *Request, mw *api.ManifestWriter) error
|
|||
func (s *Server) HandleDelete(w http.ResponseWriter, r *Request) {
|
||||
key, err := s.api.Resolve(r.uri)
|
||||
if err != nil {
|
||||
s.Error(w, r, fmt.Errorf("error resolving %s: %s", r.uri.Addr, err))
|
||||
ShowError(w, &r.Request, fmt.Sprintf("Error serving %s %s: %s", r.Method, r.uri, fmt.Errorf("error resolving %s: %s", r.uri.Addr, err)), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
||||
newKey, err := s.updateManifest(key, func(mw *api.ManifestWriter) error {
|
||||
s.logDebug("removing %s from manifest %s", r.uri.Path, key.Log())
|
||||
log.Debug(fmt.Sprintf("removing %s from manifest %s", r.uri.Path, key.Log()))
|
||||
return mw.RemoveEntry(r.uri.Path)
|
||||
})
|
||||
if err != nil {
|
||||
s.Error(w, r, fmt.Errorf("error updating manifest: %s", err))
|
||||
ShowError(w, &r.Request, fmt.Sprintf("Error serving %s %s: %s", r.Method, r.uri, fmt.Errorf("error updating manifest: %s", err)), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
||||
|
|
@ -298,7 +298,7 @@ func (s *Server) HandleDelete(w http.ResponseWriter, r *Request) {
|
|||
func (s *Server) HandleGet(w http.ResponseWriter, r *Request) {
|
||||
key, err := s.api.Resolve(r.uri)
|
||||
if err != nil {
|
||||
s.Error(w, r, fmt.Errorf("error resolving %s: %s", r.uri.Addr, err))
|
||||
ShowError(w, &r.Request, fmt.Sprintf("Error serving %s %s: %s", r.Method, r.uri, fmt.Errorf("error resolving %s: %s", r.uri.Addr, err)), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
||||
|
|
@ -307,7 +307,7 @@ func (s *Server) HandleGet(w http.ResponseWriter, r *Request) {
|
|||
if r.uri.Path != "" {
|
||||
walker, err := s.api.NewManifestWalker(key, nil)
|
||||
if err != nil {
|
||||
s.BadRequest(w, r, fmt.Sprintf("%s is not a manifest", key))
|
||||
ShowError(w, &r.Request, fmt.Sprintf("Bad request %s %s: %s", r.Method, r.uri, fmt.Sprintf("%s is not a manifest", key)), http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
var entry *api.ManifestEntry
|
||||
|
|
@ -335,7 +335,7 @@ func (s *Server) HandleGet(w http.ResponseWriter, r *Request) {
|
|||
return api.SkipManifest
|
||||
})
|
||||
if entry == nil {
|
||||
s.NotFound(w, r, fmt.Errorf("Manifest entry could not be loaded"))
|
||||
ShowError(w, &r.Request, fmt.Sprintf("NOT FOUND error serving %s %s: %s", r.Method, r.uri, fmt.Errorf("Manifest entry could not be loaded")), http.StatusNotFound)
|
||||
return
|
||||
}
|
||||
key = storage.Key(common.Hex2Bytes(entry.Hash))
|
||||
|
|
@ -344,7 +344,7 @@ func (s *Server) HandleGet(w http.ResponseWriter, r *Request) {
|
|||
// check the root chunk exists by retrieving the file's size
|
||||
reader := s.api.Retrieve(key)
|
||||
if _, err := reader.Size(nil); err != nil {
|
||||
s.NotFound(w, r, fmt.Errorf("Root chunk not found %s: %s", key, err))
|
||||
ShowError(w, &r.Request, fmt.Sprintf("NOT FOUND error serving %s %s: %s", r.Method, r.uri, fmt.Errorf("Root chunk not found %s: %s", key, err)), http.StatusNotFound)
|
||||
return
|
||||
}
|
||||
|
||||
|
|
@ -371,19 +371,19 @@ func (s *Server) HandleGet(w http.ResponseWriter, r *Request) {
|
|||
// contained in the manifest
|
||||
func (s *Server) HandleGetFiles(w http.ResponseWriter, r *Request) {
|
||||
if r.uri.Path != "" {
|
||||
s.BadRequest(w, r, "files request cannot contain a path")
|
||||
ShowError(w, &r.Request, fmt.Sprintf("Bad request %s %s: %s", r.Method, r.uri, "files request cannot contain a path"), http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
|
||||
key, err := s.api.Resolve(r.uri)
|
||||
if err != nil {
|
||||
s.Error(w, r, fmt.Errorf("error resolving %s: %s", r.uri.Addr, err))
|
||||
ShowError(w, &r.Request, fmt.Sprintf("Error serving %s %s: %s", r.Method, r.uri, fmt.Errorf("error resolving %s: %s", r.uri.Addr, err)), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
||||
walker, err := s.api.NewManifestWalker(key, nil)
|
||||
if err != nil {
|
||||
s.Error(w, r, err)
|
||||
ShowError(w, &r.Request, fmt.Sprintf("Error serving %s %s: %s", r.Method, r.uri, err), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
||||
|
|
@ -430,7 +430,7 @@ func (s *Server) HandleGetFiles(w http.ResponseWriter, r *Request) {
|
|||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
s.logError("error generating tar stream: %s", err)
|
||||
log.Error(fmt.Sprintf("error generating tar stream: %s", err))
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -446,14 +446,14 @@ func (s *Server) HandleGetList(w http.ResponseWriter, r *Request) {
|
|||
|
||||
key, err := s.api.Resolve(r.uri)
|
||||
if err != nil {
|
||||
s.Error(w, r, fmt.Errorf("error resolving %s: %s", r.uri.Addr, err))
|
||||
ShowError(w, &r.Request, fmt.Sprintf("Error serving %s %s: %s", r.Method, r.uri, fmt.Errorf("error resolving %s: %s", r.uri.Addr, err)), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
||||
list, err := s.getManifestList(key, r.uri.Path)
|
||||
|
||||
if err != nil {
|
||||
s.Error(w, r, err)
|
||||
ShowError(w, &r.Request, fmt.Sprintf("Error serving %s %s: %s", r.Method, r.uri, err), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
||||
|
|
@ -470,7 +470,7 @@ func (s *Server) HandleGetList(w http.ResponseWriter, r *Request) {
|
|||
List: &list,
|
||||
})
|
||||
if err != nil {
|
||||
s.logError("error rendering list HTML: %s", err)
|
||||
log.Error(fmt.Sprintf("error rendering list HTML: %s", err))
|
||||
}
|
||||
return
|
||||
}
|
||||
|
|
@ -546,7 +546,7 @@ func (s *Server) HandleGetFile(w http.ResponseWriter, r *Request) {
|
|||
|
||||
key, err := s.api.Resolve(r.uri)
|
||||
if err != nil {
|
||||
s.Error(w, r, fmt.Errorf("error resolving %s: %s", r.uri.Addr, err))
|
||||
ShowError(w, &r.Request, fmt.Sprintf("Error serving %s %s: %s", r.Method, r.uri, fmt.Errorf("error resolving %s: %s", r.uri.Addr, err)), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
||||
|
|
@ -554,9 +554,9 @@ func (s *Server) HandleGetFile(w http.ResponseWriter, r *Request) {
|
|||
if err != nil {
|
||||
switch status {
|
||||
case http.StatusNotFound:
|
||||
s.NotFound(w, r, err)
|
||||
ShowError(w, &r.Request, fmt.Sprintf("NOT FOUND error serving %s %s: %s", r.Method, r.uri, err), http.StatusNotFound)
|
||||
default:
|
||||
s.Error(w, r, err)
|
||||
ShowError(w, &r.Request, fmt.Sprintf("Error serving %s %s: %s", r.Method, r.uri, err), http.StatusInternalServerError)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
|
@ -567,11 +567,11 @@ func (s *Server) HandleGetFile(w http.ResponseWriter, r *Request) {
|
|||
list, err := s.getManifestList(key, r.uri.Path)
|
||||
|
||||
if err != nil {
|
||||
s.Error(w, r, err)
|
||||
ShowError(w, &r.Request, fmt.Sprintf("Error serving %s %s: %s", r.Method, r.uri, err), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
||||
s.logDebug(fmt.Sprintf("Multiple choices! --> %v", list))
|
||||
log.Debug(fmt.Sprintf("Multiple choices! --> %v", list))
|
||||
//show a nice page links to available entries
|
||||
ShowMultipleChoices(w, &r.Request, list)
|
||||
return
|
||||
|
|
@ -579,7 +579,7 @@ func (s *Server) HandleGetFile(w http.ResponseWriter, r *Request) {
|
|||
|
||||
// check the root chunk exists by retrieving the file's size
|
||||
if _, err := reader.Size(nil); err != nil {
|
||||
s.NotFound(w, r, fmt.Errorf("File not found %s: %s", r.uri, err))
|
||||
ShowError(w, &r.Request, fmt.Sprintf("NOT FOUND error serving %s %s: %s", r.Method, r.uri, fmt.Errorf("File not found %s: %s", r.uri, err)), http.StatusNotFound)
|
||||
return
|
||||
}
|
||||
|
||||
|
|
@ -589,16 +589,16 @@ func (s *Server) HandleGetFile(w http.ResponseWriter, r *Request) {
|
|||
}
|
||||
|
||||
func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
s.logDebug("HTTP %s request URL: '%s', Host: '%s', Path: '%s', Referer: '%s', Accept: '%s'", r.Method, r.RequestURI, r.URL.Host, r.URL.Path, r.Referer(), r.Header.Get("Accept"))
|
||||
log.Debug(fmt.Sprintf("HTTP %s request URL: '%s', Host: '%s', Path: '%s', Referer: '%s', Accept: '%s'", r.Method, r.RequestURI, r.URL.Host, r.URL.Path, r.Referer(), r.Header.Get("Accept")))
|
||||
|
||||
uri, err := api.Parse(strings.TrimLeft(r.URL.Path, "/"))
|
||||
req := &Request{Request: *r, uri: uri}
|
||||
if err != nil {
|
||||
s.logError("Invalid URI %q: %s", r.URL.Path, err)
|
||||
s.BadRequest(w, req, fmt.Sprintf("Invalid URI %q: %s", r.URL.Path, err))
|
||||
log.Error(fmt.Sprintf("Invalid URI %q: %s", r.URL.Path, err))
|
||||
ShowError(w, r, fmt.Sprintf("Bad request %s %s: %s", r.Method, uri, fmt.Sprintf("Invalid URI %q: %s", r.URL.Path, err)), http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
s.logDebug("%s request received for %s", r.Method, uri)
|
||||
log.Debug(fmt.Sprintf("%s request received for %s", r.Method, uri))
|
||||
|
||||
switch r.Method {
|
||||
case "POST":
|
||||
|
|
@ -666,26 +666,6 @@ func (s *Server) updateManifest(key storage.Key, update func(mw *api.ManifestWri
|
|||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
s.logDebug("generated manifest %s", key)
|
||||
log.Debug(fmt.Sprintf("generated manifest %s", key))
|
||||
return key, nil
|
||||
}
|
||||
|
||||
func (s *Server) logDebug(format string, v ...interface{}) {
|
||||
log.Debug(fmt.Sprintf("[BZZ] HTTP: "+format, v...))
|
||||
}
|
||||
|
||||
func (s *Server) logError(format string, v ...interface{}) {
|
||||
log.Error(fmt.Sprintf("[BZZ] HTTP: "+format, v...))
|
||||
}
|
||||
|
||||
func (s *Server) BadRequest(w http.ResponseWriter, r *Request, reason string) {
|
||||
ShowError(w, &r.Request, fmt.Sprintf("Bad request %s %s: %s", r.Method, r.uri, reason), http.StatusBadRequest)
|
||||
}
|
||||
|
||||
func (s *Server) Error(w http.ResponseWriter, r *Request, err error) {
|
||||
ShowError(w, &r.Request, fmt.Sprintf("Error serving %s %s: %s", r.Method, r.uri, err), http.StatusInternalServerError)
|
||||
}
|
||||
|
||||
func (s *Server) NotFound(w http.ResponseWriter, r *Request, err error) {
|
||||
ShowError(w, &r.Request, fmt.Sprintf("NOT FOUND error serving %s %s: %s", r.Method, r.uri, err), http.StatusNotFound)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -64,7 +64,8 @@ func (a *Api) NewManifest() (storage.Key, error) {
|
|||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
key, _, err := a.Store(bytes.NewReader(data), int64(len(data)))
|
||||
key, wait, err := a.Store(bytes.NewReader(data), int64(len(data)))
|
||||
wait()
|
||||
return key, err
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -16,7 +16,11 @@
|
|||
|
||||
package api
|
||||
|
||||
import "path"
|
||||
import (
|
||||
"path"
|
||||
|
||||
"github.com/ethereum/go-ethereum/swarm/storage"
|
||||
)
|
||||
|
||||
type Response struct {
|
||||
MimeType string
|
||||
|
|
@ -41,12 +45,8 @@ func NewStorage(api *Api) *Storage {
|
|||
// its content type
|
||||
//
|
||||
// DEPRECATED: Use the HTTP API instead
|
||||
func (self *Storage) Put(content, contentType string) (string, error) {
|
||||
key, _, err := self.api.Put(content, contentType)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
return key.Hex(), err
|
||||
func (self *Storage) Put(content, contentType string) (storage.Key, func(), error) {
|
||||
return self.api.Put(content, contentType)
|
||||
}
|
||||
|
||||
// Get retrieves the content from bzzpath and reads the response in full
|
||||
|
|
|
|||
|
|
@ -31,10 +31,12 @@ func TestStoragePutGet(t *testing.T) {
|
|||
content := "hello"
|
||||
exp := expResponse(content, "text/plain", 0)
|
||||
// exp := expResponse([]byte(content), "text/plain", 0)
|
||||
bzzhash, err := api.Put(content, exp.MimeType)
|
||||
bzzkey, wait, err := api.Put(content, exp.MimeType)
|
||||
if err != nil {
|
||||
t.Fatalf("unexpected error: %v", err)
|
||||
}
|
||||
wait()
|
||||
bzzhash := bzzkey.Hex()
|
||||
// to check put against the Api#Get
|
||||
resp0 := testGet(t, api.api, bzzhash, "")
|
||||
checkResponse(t, resp0, exp)
|
||||
|
|
|
|||
|
|
@ -30,7 +30,7 @@ func NewFromBytes(b []byte, l int) (bv *BitVector, err error) {
|
|||
|
||||
func (bv *BitVector) Get(i int) bool {
|
||||
bi := i / 8
|
||||
return uint8(bv.b[bi])&(0x1<<uint(i%8)) != 0
|
||||
return bv.b[bi]&(0x1<<uint(i%8)) != 0
|
||||
}
|
||||
|
||||
func (bv *BitVector) Set(i int, v bool) {
|
||||
|
|
|
|||
|
|
@ -58,11 +58,11 @@ func TestBitvectorGetSet(t *testing.T) {
|
|||
bv.Set(i, true)
|
||||
for j := 0; j < length; j++ {
|
||||
if j == i {
|
||||
if bv.Get(j) != true {
|
||||
if !bv.Get(j) {
|
||||
t.Errorf("element on index %v is not set to true", i)
|
||||
}
|
||||
} else {
|
||||
if bv.Get(j) != false {
|
||||
if bv.Get(j) {
|
||||
t.Errorf("element on index %v is not false", i)
|
||||
}
|
||||
}
|
||||
|
|
@ -70,7 +70,7 @@ func TestBitvectorGetSet(t *testing.T) {
|
|||
|
||||
bv.Set(i, false)
|
||||
|
||||
if bv.Get(i) != false {
|
||||
if bv.Get(i) {
|
||||
t.Errorf("element on index %v is not set to false", i)
|
||||
}
|
||||
}
|
||||
|
|
@ -82,7 +82,7 @@ func TestBitvectorNewFromBytesGet(t *testing.T) {
|
|||
if err != nil {
|
||||
t.Error(err)
|
||||
}
|
||||
if bv.Get(3) != true {
|
||||
if !bv.Get(3) {
|
||||
t.Fatalf("element 3 is not set to true: state %08b", bv.b[0])
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -47,7 +47,7 @@ func TestDiscovery(t *testing.T) {
|
|||
s.TestExchanges(p2ptest.Exchange{
|
||||
Label: "outgoing SubPeersMsg",
|
||||
Expects: []p2ptest.Expect{
|
||||
p2ptest.Expect{
|
||||
{
|
||||
Code: 3,
|
||||
Msg: &subPeersMsg{Depth: 0},
|
||||
Peer: s.ProtocolTester.IDs[0],
|
||||
|
|
|
|||
|
|
@ -54,14 +54,14 @@ var pof = pot.DefaultPof(256)
|
|||
// KadParams holds the config params for Kademlia
|
||||
type KadParams struct {
|
||||
// adjustable parameters
|
||||
MaxProxDisplay int // number of rows the table shows
|
||||
MinProxBinSize int // nearest neighbour core minimum cardinality
|
||||
MinBinSize int // minimum number of peers in a row
|
||||
MaxBinSize int // maximum number of peers in a row before pruning
|
||||
RetryInterval int // initial interval before a peer is first redialed
|
||||
RetryExponent int // exponent to multiply retry intervals with
|
||||
MaxRetries int // maximum number of redial attempts
|
||||
PruneInterval int // interval between peer pruning cycles
|
||||
MaxProxDisplay int // number of rows the table shows
|
||||
MinProxBinSize int // nearest neighbour core minimum cardinality
|
||||
MinBinSize int // minimum number of peers in a row
|
||||
MaxBinSize int // maximum number of peers in a row before pruning
|
||||
RetryInterval int64 // initial interval before a peer is first redialed
|
||||
RetryExponent int // exponent to multiply retry intervals with
|
||||
MaxRetries int // maximum number of redial attempts
|
||||
PruneInterval int // interval between peer pruning cycles
|
||||
// function to sanction or prevent suggesting a peer
|
||||
Reachable func(OverlayAddr) bool
|
||||
}
|
||||
|
|
@ -399,9 +399,9 @@ func (k *Kademlia) callable(val pot.Val) OverlayAddr {
|
|||
return nil
|
||||
}
|
||||
// calculate the allowed number of retries based on time lapsed since last seen
|
||||
timeAgo := int(time.Since(e.seenAt))
|
||||
div := k.RetryExponent
|
||||
div += (150000 - rand.Intn(300000)) * div / 1000000
|
||||
timeAgo := int64(time.Since(e.seenAt))
|
||||
div := int64(k.RetryExponent)
|
||||
div += (150000 - rand.Int63n(300000)) * div / 1000000
|
||||
var retries int
|
||||
for delta := timeAgo; delta > k.RetryInterval; delta /= div {
|
||||
retries++
|
||||
|
|
@ -424,7 +424,7 @@ func (k *Kademlia) callable(val pot.Val) OverlayAddr {
|
|||
return e.addr()
|
||||
}
|
||||
|
||||
// BaseAddr return the kademlia base addres
|
||||
// BaseAddr return the kademlia base address
|
||||
func (k *Kademlia) BaseAddr() []byte {
|
||||
return k.base
|
||||
}
|
||||
|
|
|
|||
|
|
@ -283,16 +283,15 @@ func TestSuggestPeerFindPeers(t *testing.T) {
|
|||
func TestSuggestPeerRetries(t *testing.T) {
|
||||
// 2 row gap, unsaturated proxbin, no callables -> want PO 0
|
||||
k := newTestKademlia("00000000")
|
||||
cycle := time.Second
|
||||
k.RetryInterval = int(cycle)
|
||||
k.RetryInterval = int64(time.Second) // cycle
|
||||
k.MaxRetries = 50
|
||||
k.RetryExponent = 2
|
||||
sleep := func(n int) {
|
||||
t := k.RetryInterval
|
||||
ts := k.RetryInterval
|
||||
for i := 1; i < n; i++ {
|
||||
t *= k.RetryExponent
|
||||
ts *= int64(k.RetryExponent)
|
||||
}
|
||||
time.Sleep(time.Duration(t))
|
||||
time.Sleep(time.Duration(ts))
|
||||
}
|
||||
|
||||
k.Register("01000000")
|
||||
|
|
|
|||
|
|
@ -88,7 +88,7 @@ func (r *RemoteSectionReader) Read(b []byte) (n int64, err error) {
|
|||
end = true
|
||||
}
|
||||
copy(b[n:], chunk.SData[:m])
|
||||
n += int64(m)
|
||||
n += m
|
||||
}
|
||||
|
||||
for {
|
||||
|
|
|
|||
|
|
@ -207,10 +207,11 @@ func (b *Bzz) RunProtocol(spec *protocols.Spec, run func(*BzzPeer) error) func(*
|
|||
// performHandshake implements the negotiation of the bzz handshake
|
||||
// shared among swarm subprotocols
|
||||
func performHandshake(p *protocols.Peer, handshake *HandshakeMsg) error {
|
||||
ctx, _ := context.WithTimeout(context.Background(), bzzHandshakeTimeout)
|
||||
// defer cancel()
|
||||
// ctx, cancel := context.WithTimeout(context.Background(), bzzHandshakeTimeout)
|
||||
defer close(handshake.done)
|
||||
ctx, cancel := context.WithTimeout(context.Background(), bzzHandshakeTimeout)
|
||||
defer func() {
|
||||
close(handshake.done)
|
||||
cancel()
|
||||
}()
|
||||
rsh, err := p.Handshake(ctx, handshake, checkHandshake)
|
||||
if err != nil {
|
||||
handshake.err = err
|
||||
|
|
@ -261,7 +262,7 @@ func NewBzzTestPeer(p *protocols.Peer, addr *BzzAddr) *BzzPeer {
|
|||
}
|
||||
}
|
||||
|
||||
// Off returns the overlay peer record for offline persistance
|
||||
// Off returns the overlay peer record for offline persistence
|
||||
func (p *BzzPeer) Off() OverlayAddr {
|
||||
return p.BzzAddr
|
||||
}
|
||||
|
|
@ -402,6 +403,15 @@ func NewAddrFromNodeID(id discover.NodeID) *BzzAddr {
|
|||
}
|
||||
}
|
||||
|
||||
// NewAddrFromNodeIDAndPort constucts a BzzAddr from a discover.NodeID and port uint16
|
||||
// the overlay address is derived as the hash of the nodeID
|
||||
func NewAddrFromNodeIDAndPort(id discover.NodeID, port uint16) *BzzAddr {
|
||||
return &BzzAddr{
|
||||
OAddr: ToOverlayAddr(id.Bytes()),
|
||||
UAddr: []byte(discover.NewNode(id, net.IP{127, 0, 0, 1}, port, port).String()),
|
||||
}
|
||||
}
|
||||
|
||||
// ToOverlayAddr creates an overlayaddress from a byte slice
|
||||
func ToOverlayAddr(id []byte) []byte {
|
||||
return crypto.Keccak256(id)
|
||||
|
|
|
|||
|
|
@ -70,18 +70,18 @@ func (t *testStore) Save(key string, v []byte) error {
|
|||
func HandshakeMsgExchange(lhs, rhs *HandshakeMsg, id discover.NodeID) []p2ptest.Exchange {
|
||||
|
||||
return []p2ptest.Exchange{
|
||||
p2ptest.Exchange{
|
||||
{
|
||||
Expects: []p2ptest.Expect{
|
||||
p2ptest.Expect{
|
||||
{
|
||||
Code: 0,
|
||||
Msg: lhs,
|
||||
Peer: id,
|
||||
},
|
||||
},
|
||||
},
|
||||
p2ptest.Exchange{
|
||||
{
|
||||
Triggers: []p2ptest.Trigger{
|
||||
p2ptest.Trigger{
|
||||
{
|
||||
Code: 0,
|
||||
Msg: rhs,
|
||||
Peer: id,
|
||||
|
|
|
|||
|
|
@ -70,7 +70,7 @@ func BenchmarkDiscovery_64_4(b *testing.B) { benchmarkDiscovery(b, 64, 4) }
|
|||
func BenchmarkDiscovery_128_4(b *testing.B) { benchmarkDiscovery(b, 128, 4) }
|
||||
func BenchmarkDiscovery_256_4(b *testing.B) { benchmarkDiscovery(b, 256, 4) }
|
||||
|
||||
func TestDiscoverySimulationDockerAdapter(t *testing.T) {
|
||||
func XTestDiscoverySimulationDockerAdapter(t *testing.T) {
|
||||
testDiscoverySimulationDockerAdapter(t, *nodeCount, *initCount)
|
||||
}
|
||||
|
||||
|
|
@ -99,13 +99,20 @@ func testDiscoverySimulationExecAdapter(t *testing.T, nodes, conns int) {
|
|||
testDiscoverySimulation(t, nodes, conns, adapters.NewExecAdapter(baseDir))
|
||||
}
|
||||
|
||||
func TestDiscoverySimulationSocketAdapter(t *testing.T) {
|
||||
testDiscoverySimulationSocketAdapter(t, *nodeCount, *initCount)
|
||||
}
|
||||
|
||||
func TestDiscoverySimulationSimAdapter(t *testing.T) {
|
||||
testDiscoverySimulationSimAdapter(t, *nodeCount, *initCount)
|
||||
}
|
||||
|
||||
func testDiscoverySimulationSimAdapter(t *testing.T, nodes, conns int) {
|
||||
testDiscoverySimulation(t, nodes, conns, adapters.NewSimAdapter(services))
|
||||
}
|
||||
|
||||
func testDiscoverySimulationSocketAdapter(t *testing.T, nodes, conns int) {
|
||||
testDiscoverySimulation(t, nodes, conns, adapters.NewSocketAdapter(services))
|
||||
// testDiscoverySimulation(t, nodes, conns, adapters.NewSimAdapter(services))
|
||||
}
|
||||
|
||||
func testDiscoverySimulation(t *testing.T, nodes, conns int, adapter adapters.NodeAdapter) {
|
||||
|
|
@ -139,7 +146,7 @@ func benchmarkDiscovery(b *testing.B, nodes, conns int) {
|
|||
for i := 0; i < b.N; i++ {
|
||||
result, err := discoverySimulation(nodes, conns, adapters.NewSimAdapter(services))
|
||||
if err != nil {
|
||||
b.Fatalf("setting up simulation failed", result)
|
||||
b.Fatalf("setting up simulation failed: %s", err)
|
||||
}
|
||||
if result.Error != nil {
|
||||
b.Logf("simulation failed: %s", result.Error)
|
||||
|
|
@ -157,7 +164,8 @@ func discoverySimulation(nodes, conns int, adapter adapters.NodeAdapter) (*simul
|
|||
trigger := make(chan discover.NodeID)
|
||||
ids := make([]discover.NodeID, nodes)
|
||||
for i := 0; i < nodes; i++ {
|
||||
node, err := net.NewNode()
|
||||
conf := adapters.RandomNodeConfig()
|
||||
node, err := net.NewNodeWithConfig(conf)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("error starting node: %s", err)
|
||||
}
|
||||
|
|
@ -297,7 +305,7 @@ func triggerChecks(trigger chan discover.NodeID, net *simulations.Network, id di
|
|||
}
|
||||
|
||||
func newService(ctx *adapters.ServiceContext) (node.Service, error) {
|
||||
addr := network.NewAddrFromNodeID(ctx.Config.ID)
|
||||
addr := network.NewAddrFromNodeIDAndPort(ctx.Config.ID, ctx.Config.Port)
|
||||
|
||||
kp := network.NewKadParams()
|
||||
kp.MinProxBinSize = testMinProxBinSize
|
||||
|
|
|
|||
|
|
@ -191,7 +191,7 @@ R:
|
|||
// this should be has locally
|
||||
chunk, err := d.db.Get(req.Key)
|
||||
if !bytes.Equal(chunk.Key, req.Key) {
|
||||
panic(fmt.Errorf("processReceivedChunks: chunk key %s != req key %s (peer %s)", chunk.Key.Hex(), storage.Key(req.Key).Hex(), req.peer.ID()))
|
||||
panic(fmt.Errorf("processReceivedChunks: chunk key %s != req key %s (peer %s)", chunk.Key.Hex(), req.Key.Hex(), req.peer.ID()))
|
||||
}
|
||||
if err == nil {
|
||||
continue R
|
||||
|
|
|
|||
|
|
@ -57,7 +57,7 @@ func TestStreamerRetrieveRequest(t *testing.T) {
|
|||
err = tester.TestExchanges(p2ptest.Exchange{
|
||||
Label: "RetrieveRequestMsg",
|
||||
Expects: []p2ptest.Expect{
|
||||
p2ptest.Expect{
|
||||
{
|
||||
Code: 5,
|
||||
Msg: &RetrieveRequestMsg{
|
||||
Key: hash0[:],
|
||||
|
|
@ -98,7 +98,7 @@ func TestStreamerUpstreamRetrieveRequestMsgExchangeWithoutStore(t *testing.T) {
|
|||
err = tester.TestExchanges(p2ptest.Exchange{
|
||||
Label: "RetrieveRequestMsg",
|
||||
Triggers: []p2ptest.Trigger{
|
||||
p2ptest.Trigger{
|
||||
{
|
||||
Code: 5,
|
||||
Msg: &RetrieveRequestMsg{
|
||||
Key: chunk.Key[:],
|
||||
|
|
@ -107,7 +107,7 @@ func TestStreamerUpstreamRetrieveRequestMsgExchangeWithoutStore(t *testing.T) {
|
|||
},
|
||||
},
|
||||
Expects: []p2ptest.Expect{
|
||||
p2ptest.Expect{
|
||||
{
|
||||
Code: 1,
|
||||
Msg: &OfferedHashesMsg{
|
||||
HandoverProof: nil,
|
||||
|
|
@ -156,7 +156,7 @@ func TestStreamerUpstreamRetrieveRequestMsgExchange(t *testing.T) {
|
|||
err = tester.TestExchanges(p2ptest.Exchange{
|
||||
Label: "RetrieveRequestMsg",
|
||||
Triggers: []p2ptest.Trigger{
|
||||
p2ptest.Trigger{
|
||||
{
|
||||
Code: 5,
|
||||
Msg: &RetrieveRequestMsg{
|
||||
Key: hash,
|
||||
|
|
@ -165,7 +165,7 @@ func TestStreamerUpstreamRetrieveRequestMsgExchange(t *testing.T) {
|
|||
},
|
||||
},
|
||||
Expects: []p2ptest.Expect{
|
||||
p2ptest.Expect{
|
||||
{
|
||||
Code: 1,
|
||||
Msg: &OfferedHashesMsg{
|
||||
HandoverProof: &HandoverProof{
|
||||
|
|
@ -195,7 +195,7 @@ func TestStreamerUpstreamRetrieveRequestMsgExchange(t *testing.T) {
|
|||
err = tester.TestExchanges(p2ptest.Exchange{
|
||||
Label: "RetrieveRequestMsg",
|
||||
Triggers: []p2ptest.Trigger{
|
||||
p2ptest.Trigger{
|
||||
{
|
||||
Code: 5,
|
||||
Msg: &RetrieveRequestMsg{
|
||||
Key: hash,
|
||||
|
|
@ -205,7 +205,7 @@ func TestStreamerUpstreamRetrieveRequestMsgExchange(t *testing.T) {
|
|||
},
|
||||
},
|
||||
Expects: []p2ptest.Expect{
|
||||
p2ptest.Expect{
|
||||
{
|
||||
Code: 6,
|
||||
Msg: &ChunkDeliveryMsg{
|
||||
Key: hash,
|
||||
|
|
@ -258,7 +258,7 @@ func TestStreamerDownstreamChunkDeliveryMsgExchange(t *testing.T) {
|
|||
err = tester.TestExchanges(p2ptest.Exchange{
|
||||
Label: "Subscribe message",
|
||||
Expects: []p2ptest.Expect{
|
||||
p2ptest.Expect{
|
||||
{
|
||||
Code: 4,
|
||||
Msg: &SubscribeMsg{
|
||||
Stream: stream,
|
||||
|
|
@ -275,7 +275,7 @@ func TestStreamerDownstreamChunkDeliveryMsgExchange(t *testing.T) {
|
|||
p2ptest.Exchange{
|
||||
Label: "ChunkDeliveryRequest message",
|
||||
Triggers: []p2ptest.Trigger{
|
||||
p2ptest.Trigger{
|
||||
{
|
||||
Code: 6,
|
||||
Msg: &ChunkDeliveryMsg{
|
||||
Key: chunkKey,
|
||||
|
|
@ -324,11 +324,12 @@ func testDeliveryFromNodes(t *testing.T, nodes, conns, chunkCount int, skipCheck
|
|||
defaultSkipCheck = skipCheck
|
||||
toAddr = network.NewAddrFromNodeID
|
||||
conf := &streamTesting.RunConfig{
|
||||
Adapter: *adapter,
|
||||
NodeCount: nodes,
|
||||
ConnLevel: conns,
|
||||
ToAddr: toAddr,
|
||||
Services: services,
|
||||
Adapter: *adapter,
|
||||
NodeCount: nodes,
|
||||
ConnLevel: conns,
|
||||
ToAddr: toAddr,
|
||||
Services: services,
|
||||
EnableMsgEvents: false,
|
||||
}
|
||||
|
||||
sim, teardown, err := streamTesting.NewSimulation(conf)
|
||||
|
|
@ -498,11 +499,12 @@ func benchmarkDeliveryFromNodes(b *testing.B, nodes, conns, chunkCount int, skip
|
|||
defer cancel()
|
||||
|
||||
conf := &streamTesting.RunConfig{
|
||||
Adapter: *adapter,
|
||||
NodeCount: nodes,
|
||||
ConnLevel: conns,
|
||||
ToAddr: toAddr,
|
||||
Services: services,
|
||||
Adapter: *adapter,
|
||||
NodeCount: nodes,
|
||||
ConnLevel: conns,
|
||||
ToAddr: toAddr,
|
||||
Services: services,
|
||||
EnableMsgEvents: false,
|
||||
}
|
||||
sim, teardown, err := streamTesting.NewSimulation(conf)
|
||||
defer teardown()
|
||||
|
|
|
|||
|
|
@ -319,9 +319,6 @@ func (m TakeoverProofMsg) String() string {
|
|||
|
||||
func (p *Peer) handleTakeoverProofMsg(req *TakeoverProofMsg) error {
|
||||
_, err := p.getServer(req.Stream)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
// store the strongest takeoverproof for the stream in streamer
|
||||
return nil
|
||||
return err
|
||||
}
|
||||
|
|
|
|||
|
|
@ -38,7 +38,7 @@ var (
|
|||
errClientParamsNotFound = errors.New("client params not found")
|
||||
)
|
||||
|
||||
// Peer is the Peer extention for the streaming protocol
|
||||
// Peer is the Peer extension for the streaming protocol
|
||||
type Peer struct {
|
||||
*protocols.Peer
|
||||
streamer *Registry
|
||||
|
|
|
|||
|
|
@ -188,7 +188,7 @@ func TestStreamerDownstreamSubscribeUnsubscribeMsgExchange(t *testing.T) {
|
|||
err = tester.TestExchanges(p2ptest.Exchange{
|
||||
Label: "Unsubscribe message",
|
||||
Expects: []p2ptest.Expect{
|
||||
p2ptest.Expect{
|
||||
{
|
||||
Code: 0,
|
||||
Msg: &UnsubscribeMsg{
|
||||
Stream: stream,
|
||||
|
|
@ -221,7 +221,7 @@ func TestStreamerUpstreamSubscribeUnsubscribeMsgExchange(t *testing.T) {
|
|||
err = tester.TestExchanges(p2ptest.Exchange{
|
||||
Label: "Subscribe message",
|
||||
Triggers: []p2ptest.Trigger{
|
||||
p2ptest.Trigger{
|
||||
{
|
||||
Code: 4,
|
||||
Msg: &SubscribeMsg{
|
||||
Stream: stream,
|
||||
|
|
@ -235,7 +235,7 @@ func TestStreamerUpstreamSubscribeUnsubscribeMsgExchange(t *testing.T) {
|
|||
},
|
||||
},
|
||||
Expects: []p2ptest.Expect{
|
||||
p2ptest.Expect{
|
||||
{
|
||||
Code: 1,
|
||||
Msg: &OfferedHashesMsg{
|
||||
Stream: stream,
|
||||
|
|
@ -258,7 +258,7 @@ func TestStreamerUpstreamSubscribeUnsubscribeMsgExchange(t *testing.T) {
|
|||
err = tester.TestExchanges(p2ptest.Exchange{
|
||||
Label: "unsubscribe message",
|
||||
Triggers: []p2ptest.Trigger{
|
||||
p2ptest.Trigger{
|
||||
{
|
||||
Code: 0,
|
||||
Msg: &UnsubscribeMsg{
|
||||
Stream: stream,
|
||||
|
|
@ -357,7 +357,7 @@ func TestStreamerUpstreamSubscribeErrorMsgExchange(t *testing.T) {
|
|||
err = tester.TestExchanges(p2ptest.Exchange{
|
||||
Label: "Subscribe message",
|
||||
Triggers: []p2ptest.Trigger{
|
||||
p2ptest.Trigger{
|
||||
{
|
||||
Code: 4,
|
||||
Msg: &SubscribeMsg{
|
||||
Stream: stream,
|
||||
|
|
@ -371,7 +371,7 @@ func TestStreamerUpstreamSubscribeErrorMsgExchange(t *testing.T) {
|
|||
},
|
||||
},
|
||||
Expects: []p2ptest.Expect{
|
||||
p2ptest.Expect{
|
||||
{
|
||||
Code: 7,
|
||||
Msg: &SubscribeErrorMsg{
|
||||
Error: "stream bar not registered",
|
||||
|
|
@ -482,7 +482,7 @@ func TestStreamerDownstreamOfferedHashesMsgExchange(t *testing.T) {
|
|||
err = tester.TestExchanges(p2ptest.Exchange{
|
||||
Label: "Subscribe message",
|
||||
Expects: []p2ptest.Expect{
|
||||
p2ptest.Expect{
|
||||
{
|
||||
Code: 4,
|
||||
Msg: &SubscribeMsg{
|
||||
Stream: stream,
|
||||
|
|
@ -499,7 +499,7 @@ func TestStreamerDownstreamOfferedHashesMsgExchange(t *testing.T) {
|
|||
p2ptest.Exchange{
|
||||
Label: "WantedHashes message",
|
||||
Triggers: []p2ptest.Trigger{
|
||||
p2ptest.Trigger{
|
||||
{
|
||||
Code: 1,
|
||||
Msg: &OfferedHashesMsg{
|
||||
HandoverProof: &HandoverProof{
|
||||
|
|
@ -514,7 +514,7 @@ func TestStreamerDownstreamOfferedHashesMsgExchange(t *testing.T) {
|
|||
},
|
||||
},
|
||||
Expects: []p2ptest.Expect{
|
||||
p2ptest.Expect{
|
||||
{
|
||||
Code: 2,
|
||||
Msg: &WantedHashesMsg{
|
||||
Stream: stream,
|
||||
|
|
|
|||
|
|
@ -65,7 +65,7 @@ const maxPO = 32
|
|||
|
||||
func RegisterSwarmSyncerServer(streamer *Registry, db *storage.DBAPI) {
|
||||
streamer.RegisterServerFunc("SYNC", func(p *Peer, t []byte, live bool) (Server, error) {
|
||||
po := uint8(t[0])
|
||||
po := t[0]
|
||||
return NewSwarmSyncerServer(live, po, db)
|
||||
})
|
||||
// streamer.RegisterServerFunc(stream, func(p *Peer) (Server, error) {
|
||||
|
|
|
|||
|
|
@ -51,11 +51,12 @@ func testSyncBetweenNodes(t *testing.T, nodes, conns, chunkCount int, skipCheck
|
|||
return addr
|
||||
}
|
||||
conf := &streamTesting.RunConfig{
|
||||
Adapter: *adapter,
|
||||
NodeCount: nodes,
|
||||
ConnLevel: conns,
|
||||
ToAddr: toAddr,
|
||||
Services: services,
|
||||
Adapter: *adapter,
|
||||
NodeCount: nodes,
|
||||
ConnLevel: conns,
|
||||
ToAddr: toAddr,
|
||||
Services: services,
|
||||
EnableMsgEvents: false,
|
||||
}
|
||||
// create context for simulation run
|
||||
timeout := 30 * time.Second
|
||||
|
|
|
|||
|
|
@ -117,12 +117,13 @@ func CheckResult(t *testing.T, result *simulations.StepResult, startedAt, finish
|
|||
}
|
||||
|
||||
type RunConfig struct {
|
||||
Adapter string
|
||||
Step *simulations.Step
|
||||
NodeCount int
|
||||
ConnLevel int
|
||||
ToAddr func(discover.NodeID) *network.BzzAddr
|
||||
Services adapters.Services
|
||||
Adapter string
|
||||
Step *simulations.Step
|
||||
NodeCount int
|
||||
ConnLevel int
|
||||
ToAddr func(discover.NodeID) *network.BzzAddr
|
||||
Services adapters.Services
|
||||
EnableMsgEvents bool
|
||||
}
|
||||
|
||||
func NewSimulation(conf *RunConfig) (*Simulation, func(), error) {
|
||||
|
|
@ -144,7 +145,9 @@ func NewSimulation(conf *RunConfig) (*Simulation, func(), error) {
|
|||
addrs := make([]network.Addr, nodes)
|
||||
// start nodes
|
||||
for i := 0; i < nodes; i++ {
|
||||
node, err := net.NewNode()
|
||||
nodeconf := adapters.RandomNodeConfig()
|
||||
nodeconf.EnableMsgEvents = conf.EnableMsgEvents
|
||||
node, err := net.NewNodeWithConfig(nodeconf)
|
||||
if err != nil {
|
||||
return nil, teardown, fmt.Errorf("error creating node: %s", err)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -96,7 +96,7 @@ func (pssapi *API) BaseAddr() (PssAddress, error) {
|
|||
func (pssapi *API) GetPublicKey() (keybytes hexutil.Bytes) {
|
||||
key := pssapi.Pss.PublicKey()
|
||||
keybytes = crypto.FromECDSAPub(key)
|
||||
return hexutil.Bytes(keybytes)
|
||||
return keybytes
|
||||
}
|
||||
|
||||
// Set Public key to associate with a particular Pss peer
|
||||
|
|
|
|||
|
|
@ -104,7 +104,8 @@ func TestClientHandshake(t *testing.T) {
|
|||
lproto := pss.NewPingProtocol(lpssping)
|
||||
rproto := pss.NewPingProtocol(rpssping)
|
||||
|
||||
ctx, _ := context.WithTimeout(context.Background(), time.Second*10)
|
||||
ctx, cancel := context.WithTimeout(context.Background(), time.Second*10)
|
||||
defer cancel()
|
||||
err = lpsc.RunProtocol(ctx, lproto)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
|
|
@ -179,9 +180,9 @@ func setupNetwork(numnodes int) (clients []*rpc.Client, err error) {
|
|||
DefaultService: "bzz",
|
||||
})
|
||||
for i := 0; i < numnodes; i++ {
|
||||
nodes[i], err = net.NewNodeWithConfig(&adapters.NodeConfig{
|
||||
Services: []string{"bzz", "pss"},
|
||||
})
|
||||
nodeconf := adapters.RandomNodeConfig()
|
||||
nodeconf.Services = []string{"bzz", "pss"}
|
||||
nodes[i], err = net.NewNodeWithConfig(nodeconf)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("error creating node 1: %v", err)
|
||||
}
|
||||
|
|
@ -231,13 +232,14 @@ func newServices() adapters.Services {
|
|||
"pss": func(ctx *adapters.ServiceContext) (node.Service, error) {
|
||||
cachedir, err := ioutil.TempDir("", "pss-cache")
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("create pss cache tmpdir failed", "error", err)
|
||||
return nil, fmt.Errorf("create pss cache tmpdir failed: %s", err)
|
||||
}
|
||||
dpa, err := storage.NewLocalDPA(cachedir, make([]byte, 32))
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("local dpa creation failed", "error", err)
|
||||
return nil, fmt.Errorf("local dpa creation failed: %s", err)
|
||||
}
|
||||
ctxlocal, _ := context.WithTimeout(context.Background(), time.Second)
|
||||
ctxlocal, cancel := context.WithTimeout(context.Background(), time.Second)
|
||||
defer cancel()
|
||||
keys, err := wapi.NewKeyPair(ctxlocal)
|
||||
privkey, err := w.GetPrivateKey(keys)
|
||||
psparams := pss.NewPssParams(privkey)
|
||||
|
|
|
|||
|
|
@ -254,7 +254,7 @@ func (self *HandshakeController) cleanHandshake(pubkeyid string, topic *Topic, i
|
|||
func (self *HandshakeController) clean() {
|
||||
peerpubkeys := self.handshakes
|
||||
for pubkeyid, peertopics := range peerpubkeys {
|
||||
for topic, _ := range peertopics {
|
||||
for topic := range peertopics {
|
||||
self.cleanHandshake(pubkeyid, &topic, true, true)
|
||||
}
|
||||
}
|
||||
|
|
@ -268,7 +268,7 @@ func (self *HandshakeController) handler(msg []byte, p *p2p.Peer, asymmetric boo
|
|||
if !asymmetric {
|
||||
if self.symKeyIndex[symkeyid] != nil {
|
||||
if self.symKeyIndex[symkeyid].count >= self.symKeyIndex[symkeyid].limit {
|
||||
return fmt.Errorf("discarding message using expired key", "symkeyid", symkeyid)
|
||||
return fmt.Errorf("discarding message using expired key: %s", symkeyid)
|
||||
}
|
||||
self.symKeyIndex[symkeyid].count++
|
||||
log.Trace("increment symkey recv use", "symsymkeyid", symkeyid, "count", self.symKeyIndex[symkeyid].count, "limit", self.symKeyIndex[symkeyid].limit, "receiver", common.ToHex(crypto.FromECDSAPub(self.pss.PublicKey())))
|
||||
|
|
@ -444,7 +444,7 @@ func (self *HandshakeAPI) Handshake(pubkeyid string, topic Topic, sync bool, flu
|
|||
keycount = self.ctrl.symKeyCapacity
|
||||
} else {
|
||||
validkeys := self.ctrl.validKeys(pubkeyid, &topic, false)
|
||||
keycount = uint8(self.ctrl.symKeyCapacity - uint8(len(validkeys)))
|
||||
keycount = self.ctrl.symKeyCapacity - uint8(len(validkeys))
|
||||
}
|
||||
if keycount == 0 {
|
||||
return keys, errors.New("Incoming symmetric key store is already full")
|
||||
|
|
@ -457,7 +457,8 @@ func (self *HandshakeAPI) Handshake(pubkeyid string, topic Topic, sync bool, flu
|
|||
return keys, err
|
||||
}
|
||||
if sync {
|
||||
ctx, _ := context.WithTimeout(context.Background(), self.ctrl.symKeyRequestTimeout)
|
||||
ctx, cancel := context.WithTimeout(context.Background(), self.ctrl.symKeyRequestTimeout)
|
||||
defer cancel()
|
||||
select {
|
||||
case keys = <-hsc:
|
||||
log.Trace("sync handshake response receive", "key", keys)
|
||||
|
|
@ -474,7 +475,7 @@ func (self *HandshakeAPI) AddHandshake(topic Topic) error {
|
|||
return nil
|
||||
}
|
||||
|
||||
// Deactivate handshake functionalty on a topic
|
||||
// Deactivate handshake functionality on a topic
|
||||
func (self *HandshakeAPI) RemoveHandshake(topic *Topic) error {
|
||||
if _, ok := self.ctrl.deregisterFuncs[*topic]; ok {
|
||||
self.ctrl.deregisterFuncs[*topic]()
|
||||
|
|
|
|||
|
|
@ -227,7 +227,7 @@ func (self *Protocol) AddPeer(p *p2p.Peer, run func(*p2p.Peer, p2p.MsgReadWriter
|
|||
}
|
||||
go func() {
|
||||
err := run(p, rw)
|
||||
log.Warn(fmt.Sprintf("pss vprotocol quit on addr %v topic %v: %v", topic, err))
|
||||
log.Warn(fmt.Sprintf("pss vprotocol quit topic %v: %v", topic, err))
|
||||
}()
|
||||
return rw, nil
|
||||
}
|
||||
|
|
|
|||
|
|
@ -73,11 +73,13 @@ func testProtocol(t *testing.T) {
|
|||
time.Sleep(time.Millisecond * 1000) // replace with hive healthy code
|
||||
|
||||
lmsgC := make(chan APIMsg)
|
||||
lctx, _ := context.WithTimeout(context.Background(), time.Second*10)
|
||||
lctx, cancel := context.WithTimeout(context.Background(), time.Second*10)
|
||||
defer cancel()
|
||||
lsub, err := clients[0].Subscribe(lctx, "pss", lmsgC, "receive", topic)
|
||||
defer lsub.Unsubscribe()
|
||||
rmsgC := make(chan APIMsg)
|
||||
rctx, _ := context.WithTimeout(context.Background(), time.Second*10)
|
||||
rctx, cancel := context.WithTimeout(context.Background(), time.Second*10)
|
||||
defer cancel()
|
||||
rsub, err := clients[1].Subscribe(rctx, "pss", rmsgC, "receive", topic)
|
||||
defer rsub.Unsubscribe()
|
||||
|
||||
|
|
|
|||
|
|
@ -190,7 +190,7 @@ var pssSpec = &protocols.Spec{
|
|||
|
||||
func (self *Pss) Protocols() []p2p.Protocol {
|
||||
return []p2p.Protocol{
|
||||
p2p.Protocol{
|
||||
{
|
||||
Name: pssSpec.Name,
|
||||
Version: pssSpec.Version,
|
||||
Length: pssSpec.Length(),
|
||||
|
|
@ -209,16 +209,14 @@ func (self *Pss) Run(p *p2p.Peer, rw p2p.MsgReadWriter) error {
|
|||
|
||||
func (self *Pss) APIs() []rpc.API {
|
||||
apis := []rpc.API{
|
||||
rpc.API{
|
||||
{
|
||||
Namespace: "pss",
|
||||
Version: "1.0",
|
||||
Service: NewAPI(self),
|
||||
Public: true,
|
||||
},
|
||||
}
|
||||
for _, auxapi := range self.auxAPIs {
|
||||
apis = append(apis, auxapi)
|
||||
}
|
||||
apis = append(apis, self.auxAPIs...)
|
||||
return apis
|
||||
}
|
||||
|
||||
|
|
@ -389,7 +387,7 @@ func (self *Pss) SetPeerPublicKey(pubkey *ecdsa.PublicKey, topic Topic, address
|
|||
address: address,
|
||||
}
|
||||
self.pubKeyPoolMu.Lock()
|
||||
if _, ok := self.pubKeyPool[pubkeyid]; ok == false {
|
||||
if _, ok := self.pubKeyPool[pubkeyid]; !ok {
|
||||
self.pubKeyPool[pubkeyid] = make(map[Topic]*pssPeer)
|
||||
}
|
||||
self.pubKeyPool[pubkeyid][topic] = psp
|
||||
|
|
@ -418,7 +416,7 @@ func (self *Pss) generateSymmetricKey(topic Topic, address *PssAddress, addToCac
|
|||
// If addtocache is set to true, the key will be added to the cache of keys
|
||||
// used to attempt symmetric decryption of incoming messages.
|
||||
//
|
||||
// Returns a string id that can be used to retreive the key bytes
|
||||
// Returns a string id that can be used to retrieve the key bytes
|
||||
// from the whisper backend (see pss.GetSymmetricKey())
|
||||
func (self *Pss) SetSymmetricKey(key []byte, topic Topic, address *PssAddress, addtocache bool) (string, error) {
|
||||
keyid, err := self.w.AddSymKeyDirect(key)
|
||||
|
|
@ -501,7 +499,7 @@ func (self *Pss) processSym(envelope *whisper.Envelope) (*whisper.ReceivedMessag
|
|||
func (self *Pss) processAsym(envelope *whisper.Envelope) (*whisper.ReceivedMessage, string, *PssAddress, error) {
|
||||
recvmsg, err := envelope.OpenAsymmetric(self.privateKey)
|
||||
if err != nil {
|
||||
return nil, "", nil, fmt.Errorf("could not decrypt message: %v", "err", err)
|
||||
return nil, "", nil, fmt.Errorf("could not decrypt message: %s", err)
|
||||
}
|
||||
// check signature (if signed), strip padding
|
||||
if !recvmsg.Validate() {
|
||||
|
|
@ -538,7 +536,7 @@ func (self *Pss) cleanKeys() (count int) {
|
|||
match = true
|
||||
}
|
||||
}
|
||||
if match == false {
|
||||
if !match {
|
||||
expiredtopics = append(expiredtopics, topic)
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -137,7 +137,8 @@ func TestTopic(t *testing.T) {
|
|||
func TestCache(t *testing.T) {
|
||||
var err error
|
||||
to, _ := hex.DecodeString("08090a0b0c0d0e0f1011121314150001020304050607161718191a1b1c1d1e1f")
|
||||
ctx, _ := context.WithTimeout(context.Background(), time.Second)
|
||||
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
|
||||
defer cancel()
|
||||
keys, err := wapi.NewKeyPair(ctx)
|
||||
privkey, err := w.GetPrivateKey(keys)
|
||||
if err != nil {
|
||||
|
|
@ -211,7 +212,8 @@ func TestAddressMatch(t *testing.T) {
|
|||
remoteaddr := []byte("feedbeef")
|
||||
kadparams := network.NewKadParams()
|
||||
kad := network.NewKademlia(localaddr, kadparams)
|
||||
ctx, _ := context.WithTimeout(context.Background(), time.Second)
|
||||
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
|
||||
defer cancel()
|
||||
keys, err := wapi.NewKeyPair(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("Could not generate private key: %v", err)
|
||||
|
|
@ -255,12 +257,14 @@ func TestAddressMatch(t *testing.T) {
|
|||
// set and generate pubkeys and symkeys
|
||||
func TestKeys(t *testing.T) {
|
||||
// make our key and init pss with it
|
||||
ctx, _ := context.WithTimeout(context.Background(), time.Second)
|
||||
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
|
||||
defer cancel()
|
||||
ourkeys, err := wapi.NewKeyPair(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("create 'our' key fail")
|
||||
}
|
||||
ctx, _ = context.WithTimeout(context.Background(), time.Second)
|
||||
ctx, cancel2 := context.WithTimeout(context.Background(), time.Second)
|
||||
defer cancel2()
|
||||
theirkeys, err := wapi.NewKeyPair(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("create 'their' key fail")
|
||||
|
|
@ -449,12 +453,14 @@ func testSymSend(t *testing.T) {
|
|||
// at this point we've verified that symkeys are saved and match on each peer
|
||||
// now try sending symmetrically encrypted message, both directions
|
||||
lmsgC := make(chan APIMsg)
|
||||
lctx, _ := context.WithTimeout(context.Background(), time.Second*10)
|
||||
lctx, lcancel := context.WithTimeout(context.Background(), time.Second*10)
|
||||
defer lcancel()
|
||||
lsub, err := clients[0].Subscribe(lctx, "pss", lmsgC, "receive", topic)
|
||||
log.Trace("lsub", "id", lsub)
|
||||
defer lsub.Unsubscribe()
|
||||
rmsgC := make(chan APIMsg)
|
||||
rctx, _ := context.WithTimeout(context.Background(), time.Second*10)
|
||||
rctx, rcancel := context.WithTimeout(context.Background(), time.Second*10)
|
||||
defer rcancel()
|
||||
rsub, err := clients[1].Subscribe(rctx, "pss", rmsgC, "receive", topic)
|
||||
log.Trace("rsub", "id", rsub)
|
||||
defer rsub.Unsubscribe()
|
||||
|
|
@ -562,12 +568,14 @@ func testAsymSend(t *testing.T) {
|
|||
time.Sleep(time.Millisecond * 500) // replace with hive healthy code
|
||||
|
||||
lmsgC := make(chan APIMsg)
|
||||
lctx, _ := context.WithTimeout(context.Background(), time.Second*10)
|
||||
lctx, lcancel := context.WithTimeout(context.Background(), time.Second*10)
|
||||
defer lcancel()
|
||||
lsub, err := clients[0].Subscribe(lctx, "pss", lmsgC, "receive", topic)
|
||||
log.Trace("lsub", "id", lsub)
|
||||
defer lsub.Unsubscribe()
|
||||
rmsgC := make(chan APIMsg)
|
||||
rctx, _ := context.WithTimeout(context.Background(), time.Second*10)
|
||||
rctx, rcancel := context.WithTimeout(context.Background(), time.Second*10)
|
||||
defer rcancel()
|
||||
rsub, err := clients[1].Subscribe(rctx, "pss", rmsgC, "receive", topic)
|
||||
log.Trace("rsub", "id", rsub)
|
||||
defer rsub.Unsubscribe()
|
||||
|
|
@ -626,7 +634,7 @@ func worker(id int, jobs <-chan Job, rpcs map[discover.NodeID]*rpc.Client, pubke
|
|||
// params in run name:
|
||||
// nodes/msgs/addrbytes/adaptertype
|
||||
// if adaptertype is exec uses execadapter, simadapter otherwise
|
||||
func TestNetwork(t *testing.T) {
|
||||
func XTestNetwork(t *testing.T) {
|
||||
t.Run("3/2000/4/sock", testNetwork)
|
||||
t.Run("4/2000/4/sock", testNetwork)
|
||||
t.Run("8/2000/4/sock", testNetwork)
|
||||
|
|
@ -711,7 +719,7 @@ func testNetwork(t *testing.T) {
|
|||
select {
|
||||
case recvmsg := <-msgC:
|
||||
idx, _ := binary.Uvarint(recvmsg.Msg)
|
||||
if recvmsgs[idx] == false {
|
||||
if !recvmsgs[idx] {
|
||||
log.Debug("msg recv", "idx", idx, "id", id)
|
||||
recvmsgs[idx] = true
|
||||
trigger <- id
|
||||
|
|
@ -834,7 +842,8 @@ func benchmarkSymKeySend(b *testing.B) {
|
|||
if err != nil {
|
||||
b.Fatalf("benchmark called with invalid msgsize param '%s': %v", msgsizestring[1], err)
|
||||
}
|
||||
ctx, _ := context.WithTimeout(context.Background(), time.Second)
|
||||
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
|
||||
defer cancel()
|
||||
keys, err := wapi.NewKeyPair(ctx)
|
||||
privkey, err := w.GetPrivateKey(keys)
|
||||
ps := newTestPss(privkey, nil, nil)
|
||||
|
|
@ -849,7 +858,7 @@ func benchmarkSymKeySend(b *testing.B) {
|
|||
}
|
||||
symkey, err := ps.w.GetSymKey(symkeyid)
|
||||
if err != nil {
|
||||
b.Fatalf("could not retreive symkey: %v", err)
|
||||
b.Fatalf("could not retrieve symkey: %v", err)
|
||||
}
|
||||
ps.SetSymmetricKey(symkey, topic, &to, false)
|
||||
|
||||
|
|
@ -877,7 +886,8 @@ func benchmarkAsymKeySend(b *testing.B) {
|
|||
if err != nil {
|
||||
b.Fatalf("benchmark called with invalid msgsize param '%s': %v", msgsizestring[1], err)
|
||||
}
|
||||
ctx, _ := context.WithTimeout(context.Background(), time.Second)
|
||||
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
|
||||
defer cancel()
|
||||
keys, err := wapi.NewKeyPair(ctx)
|
||||
privkey, err := w.GetPrivateKey(keys)
|
||||
ps := newTestPss(privkey, nil, nil)
|
||||
|
|
@ -922,7 +932,8 @@ func benchmarkSymkeyBruteforceChangeaddr(b *testing.B) {
|
|||
}
|
||||
pssmsgs := make([]*PssMsg, 0, keycount)
|
||||
var keyid string
|
||||
ctx, _ := context.WithTimeout(context.Background(), time.Second)
|
||||
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
|
||||
defer cancel()
|
||||
keys, err := wapi.NewKeyPair(ctx)
|
||||
privkey, err := w.GetPrivateKey(keys)
|
||||
if cachesize > 0 {
|
||||
|
|
@ -940,7 +951,7 @@ func benchmarkSymkeyBruteforceChangeaddr(b *testing.B) {
|
|||
}
|
||||
symkey, err := ps.w.GetSymKey(keyid)
|
||||
if err != nil {
|
||||
b.Fatalf("could not retreive symkey %s: %v", keyid, err)
|
||||
b.Fatalf("could not retrieve symkey %s: %v", keyid, err)
|
||||
}
|
||||
wparams := &whisper.MessageParams{
|
||||
TTL: defaultWhisperTTL,
|
||||
|
|
@ -1004,7 +1015,8 @@ func benchmarkSymkeyBruteforceSameaddr(b *testing.B) {
|
|||
}
|
||||
}
|
||||
addr := make([]PssAddress, keycount)
|
||||
ctx, _ := context.WithTimeout(context.Background(), time.Second)
|
||||
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
|
||||
defer cancel()
|
||||
keys, err := wapi.NewKeyPair(ctx)
|
||||
privkey, err := w.GetPrivateKey(keys)
|
||||
if cachesize > 0 {
|
||||
|
|
@ -1023,7 +1035,7 @@ func benchmarkSymkeyBruteforceSameaddr(b *testing.B) {
|
|||
}
|
||||
symkey, err := ps.w.GetSymKey(keyid)
|
||||
if err != nil {
|
||||
b.Fatalf("could not retreive symkey %s: %v", keyid, err)
|
||||
b.Fatalf("could not retrieve symkey %s: %v", keyid, err)
|
||||
}
|
||||
wparams := &whisper.MessageParams{
|
||||
TTL: defaultWhisperTTL,
|
||||
|
|
@ -1069,9 +1081,9 @@ func setupNetwork(numnodes int) (clients []*rpc.Client, err error) {
|
|||
DefaultService: "bzz",
|
||||
})
|
||||
for i := 0; i < numnodes; i++ {
|
||||
nodes[i], err = net.NewNodeWithConfig(&adapters.NodeConfig{
|
||||
Services: []string{"bzz", pssProtocolName},
|
||||
})
|
||||
nodeconf := adapters.RandomNodeConfig()
|
||||
nodeconf.Services = []string{"bzz", pssProtocolName}
|
||||
nodes[i], err = net.NewNodeWithConfig(nodeconf)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("error creating node 1: %v", err)
|
||||
}
|
||||
|
|
@ -1121,17 +1133,18 @@ func newServices() adapters.Services {
|
|||
pssProtocolName: func(ctx *adapters.ServiceContext) (node.Service, error) {
|
||||
cachedir, err := ioutil.TempDir("", "pss-cache")
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("create pss cache tmpdir failed", "error", err)
|
||||
return nil, fmt.Errorf("create pss cache tmpdir failed: %s", err)
|
||||
}
|
||||
dpa, err := storage.NewLocalDPA(cachedir, network.NewAddrFromNodeID(ctx.Config.ID).Over())
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("local dpa creation failed", "error", err)
|
||||
return nil, fmt.Errorf("local dpa creation failed: %s", err)
|
||||
}
|
||||
|
||||
// execadapter does not exec init()
|
||||
initTest()
|
||||
|
||||
ctxlocal, _ := context.WithTimeout(context.Background(), time.Second)
|
||||
ctxlocal, cancel := context.WithTimeout(context.Background(), time.Second)
|
||||
defer cancel()
|
||||
keys, err := wapi.NewKeyPair(ctxlocal)
|
||||
privkey, err := w.GetPrivateKey(keys)
|
||||
pssp := NewPssParams(privkey)
|
||||
|
|
|
|||
|
|
@ -378,7 +378,7 @@ func (self *LazyChunkReader) ReadAt(b []byte, off int64) (read int, err error) {
|
|||
return 0, err
|
||||
}
|
||||
if off+int64(len(b)) >= size {
|
||||
return int(size - int64(off)), io.EOF
|
||||
return int(size - off), io.EOF
|
||||
}
|
||||
return len(b), nil
|
||||
}
|
||||
|
|
|
|||
|
|
@ -98,7 +98,7 @@ type DbStore struct {
|
|||
|
||||
// TODO: Instead of passing the distance function, just pass the address from which distances are calculated
|
||||
// to avoid the appearance of a pluggable distance metric and opportunities of bugs associated with providing
|
||||
// a function diferent from the one that is actually used.
|
||||
// a function different from the one that is actually used.
|
||||
func NewDbStore(path string, hash SwarmHasher, capacity uint64, po func(Key) uint8) (s *DbStore, err error) {
|
||||
s = new(DbStore)
|
||||
s.hashfunc = hash
|
||||
|
|
@ -126,7 +126,7 @@ func NewDbStore(path string, hash SwarmHasher, capacity uint64, po func(Key) uin
|
|||
for i := 0; i < 0x100; i++ {
|
||||
k := make([]byte, 2)
|
||||
k[0] = keyDistanceCnt
|
||||
k[1] = byte(uint8(i))
|
||||
k[1] = uint8(i)
|
||||
cnt, _ := s.db.Get(k)
|
||||
s.bucketCnt[i] = BytesToU64(cnt)
|
||||
s.bucketCnt[i]++
|
||||
|
|
@ -211,7 +211,7 @@ func getOldDataKey(idx uint64) []byte {
|
|||
func getDataKey(idx uint64, po uint8) []byte {
|
||||
key := make([]byte, 10)
|
||||
key[0] = keyData
|
||||
key[1] = byte(po)
|
||||
key[1] = po
|
||||
binary.BigEndian.PutUint64(key[2:], idx)
|
||||
|
||||
return key
|
||||
|
|
@ -483,9 +483,9 @@ func (s *DbStore) ReIndex() {
|
|||
oldCntKey[0] = keyDistanceCnt
|
||||
newCntKey[0] = keyDistanceCnt
|
||||
key[0] = keyData
|
||||
key[1] = byte(s.po(Key(key[1:])))
|
||||
key[1] = s.po(Key(key[1:]))
|
||||
oldCntKey[1] = key[1]
|
||||
newCntKey[1] = byte(s.po(Key(newKey[1:])))
|
||||
newCntKey[1] = s.po(Key(newKey[1:]))
|
||||
copy(newKey[2:], key[1:])
|
||||
newValue := append(hash, data...)
|
||||
|
||||
|
|
@ -733,8 +733,7 @@ func (s *DbStore) setCapacity(c uint64) {
|
|||
s.capacity = c
|
||||
|
||||
if s.entryCnt > c {
|
||||
var ratio float32
|
||||
ratio = float32(1.01) - float32(c)/float32(s.entryCnt)
|
||||
ratio := float32(1.01) - float32(c)/float32(s.entryCnt)
|
||||
if ratio < gcArrayFreeRatio {
|
||||
ratio = gcArrayFreeRatio
|
||||
}
|
||||
|
|
@ -760,7 +759,7 @@ func (s *DbStore) SyncIterator(since uint64, until uint64, po uint8, f func(Key,
|
|||
|
||||
for ok := it.Seek(sincekey); ok; ok = it.Next() {
|
||||
dbkey := it.Key()
|
||||
if dbkey[0] != keyData || dbkey[1] != byte(po) || bytes.Compare(untilkey, dbkey) < 0 {
|
||||
if dbkey[0] != keyData || dbkey[1] != po || bytes.Compare(untilkey, dbkey) < 0 {
|
||||
break
|
||||
}
|
||||
key := make([]byte, 32)
|
||||
|
|
|
|||
|
|
@ -206,7 +206,7 @@ func testIterator(t *testing.T, mock bool) {
|
|||
}
|
||||
|
||||
for i = 0; i < chunkcount; i++ {
|
||||
if bytes.Compare(chunkkeys[i], chunkkeys_results[i]) != 0 {
|
||||
if !bytes.Equal(chunkkeys[i], chunkkeys_results[i]) {
|
||||
t.Fatalf("Chunk put #%d key '%v' does not match iterator's key '%v'", i, chunkkeys[i], chunkkeys_results[i])
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -625,10 +625,7 @@ func (self *ResourceHandler) verifyContent(chunkdata []byte) error {
|
|||
}
|
||||
|
||||
func (self *ResourceHandler) hasUpdate(name string, period uint32) bool {
|
||||
if self.resources[name].lastPeriod == period {
|
||||
return true
|
||||
}
|
||||
return false
|
||||
return self.resources[name].lastPeriod == period
|
||||
}
|
||||
|
||||
type resourceChunkStore struct {
|
||||
|
|
|
|||
|
|
@ -88,7 +88,7 @@ func Proximity(one, other []byte) (ret int) {
|
|||
m = MaxPO % 8
|
||||
}
|
||||
for j := 0; j < m; j++ {
|
||||
if (uint8(oxo)>>uint8(7-j))&0x01 != 0 {
|
||||
if (oxo>>uint8(7-j))&0x01 != 0 {
|
||||
return i*8 + j
|
||||
}
|
||||
}
|
||||
|
|
@ -156,10 +156,7 @@ func (c KeyCollection) Len() int {
|
|||
}
|
||||
|
||||
func (c KeyCollection) Less(i, j int) bool {
|
||||
if bytes.Compare(c[i], c[j]) == -1 {
|
||||
return true
|
||||
}
|
||||
return false
|
||||
return bytes.Compare(c[i], c[j]) == -1
|
||||
}
|
||||
|
||||
func (c KeyCollection) Swap(i, j int) {
|
||||
|
|
|
|||
|
|
@ -262,20 +262,13 @@ func (self *Swarm) Stop() error {
|
|||
|
||||
// implements the node.Service interface
|
||||
func (self *Swarm) Protocols() (protos []p2p.Protocol) {
|
||||
|
||||
for _, p := range self.bzz.Protocols() {
|
||||
protos = append(protos, p)
|
||||
}
|
||||
protos = append(protos, self.bzz.Protocols()...)
|
||||
|
||||
if self.ps != nil {
|
||||
for _, p := range self.ps.Protocols() {
|
||||
protos = append(protos, p)
|
||||
}
|
||||
protos = append(protos, self.ps.Protocols()...)
|
||||
}
|
||||
if self.streamer != nil {
|
||||
for _, p := range self.streamer.Protocols() {
|
||||
protos = append(protos, p)
|
||||
}
|
||||
protos = append(protos, self.streamer.Protocols()...)
|
||||
}
|
||||
return
|
||||
}
|
||||
|
|
@ -336,14 +329,10 @@ func (self *Swarm) APIs() []rpc.API {
|
|||
// {Namespace, Version, api.NewAdmin(self), false},
|
||||
}
|
||||
|
||||
for _, api := range self.bzz.APIs() {
|
||||
apis = append(apis, api)
|
||||
}
|
||||
apis = append(apis, self.bzz.APIs()...)
|
||||
|
||||
if self.ps != nil {
|
||||
for _, api := range self.ps.APIs() {
|
||||
apis = append(apis, api)
|
||||
}
|
||||
apis = append(apis, self.ps.APIs()...)
|
||||
}
|
||||
|
||||
return apis
|
||||
|
|
|
|||
6
vendor/github.com/rjeczalik/notify/watcher_fsevents_cgo.go
generated
vendored
6
vendor/github.com/rjeczalik/notify/watcher_fsevents_cgo.go
generated
vendored
|
|
@ -48,7 +48,7 @@ var wg sync.WaitGroup // used to wait until the runloop starts
|
|||
// started and is ready via the wg. It also serves purpose of a dummy source,
|
||||
// thanks to it the runloop does not return as it also has at least one source
|
||||
// registered.
|
||||
var source = C.CFRunLoopSourceCreate(refZero, 0, &C.CFRunLoopSourceContext{
|
||||
var source = C.CFRunLoopSourceCreate(nil, 0, &C.CFRunLoopSourceContext{
|
||||
perform: (C.CFRunLoopPerformCallBack)(C.gosource),
|
||||
})
|
||||
|
||||
|
|
@ -162,8 +162,8 @@ func (s *stream) Start() error {
|
|||
return nil
|
||||
}
|
||||
wg.Wait()
|
||||
p := C.CFStringCreateWithCStringNoCopy(refZero, C.CString(s.path), C.kCFStringEncodingUTF8, refZero)
|
||||
path := C.CFArrayCreate(refZero, (*unsafe.Pointer)(unsafe.Pointer(&p)), 1, nil)
|
||||
p := C.CFStringCreateWithCStringNoCopy(nil, C.CString(s.path), C.kCFStringEncodingUTF8, nil)
|
||||
path := C.CFArrayCreate(nil, (*unsafe.Pointer)(unsafe.Pointer(&p)), 1, nil)
|
||||
ctx := C.FSEventStreamContext{}
|
||||
ref := C.EventStreamCreate(&ctx, C.uintptr_t(s.info), path, C.FSEventStreamEventId(atomic.LoadUint64(&since)), latency, flags)
|
||||
if ref == nilstream {
|
||||
|
|
|
|||
9
vendor/github.com/rjeczalik/notify/watcher_fsevents_go1.10.go
generated
vendored
9
vendor/github.com/rjeczalik/notify/watcher_fsevents_go1.10.go
generated
vendored
|
|
@ -1,9 +0,0 @@
|
|||
// Copyright (c) 2017 The Notify Authors. All rights reserved.
|
||||
// Use of this source code is governed by the MIT license that can be
|
||||
// found in the LICENSE file.
|
||||
|
||||
// +build darwin,!kqueue,go1.10
|
||||
|
||||
package notify
|
||||
|
||||
const refZero = 0
|
||||
14
vendor/github.com/rjeczalik/notify/watcher_fsevents_go1.9.go
generated
vendored
14
vendor/github.com/rjeczalik/notify/watcher_fsevents_go1.9.go
generated
vendored
|
|
@ -1,14 +0,0 @@
|
|||
// Copyright (c) 2017 The Notify Authors. All rights reserved.
|
||||
// Use of this source code is governed by the MIT license that can be
|
||||
// found in the LICENSE file.
|
||||
|
||||
// +build darwin,!kqueue,cgo,!go1.10
|
||||
|
||||
package notify
|
||||
|
||||
/*
|
||||
#include <CoreServices/CoreServices.h>
|
||||
*/
|
||||
import "C"
|
||||
|
||||
var refZero = (*C.struct___CFAllocator)(nil)
|
||||
6
vendor/vendor.json
vendored
6
vendor/vendor.json
vendored
|
|
@ -286,10 +286,10 @@
|
|||
"revisionTime": "2016-11-28T21:05:44Z"
|
||||
},
|
||||
{
|
||||
"checksumSHA1": "1ESHllhZOIBg7MnlGHUdhz047bI=",
|
||||
"checksumSHA1": "28UVHMmHx0iqO0XiJsjx+fwILyI=",
|
||||
"path": "github.com/rjeczalik/notify",
|
||||
"revision": "27b537f07230b3f917421af6dcf044038dbe57e2",
|
||||
"revisionTime": "2018-01-03T13:19:05Z"
|
||||
"revision": "c31e5f2cb22b3e4ef3f882f413847669bf2652b9",
|
||||
"revisionTime": "2018-02-03T14:01:15Z"
|
||||
},
|
||||
{
|
||||
"checksumSHA1": "5uqO4ITTDMklKi3uNaE/D9LQ5nM=",
|
||||
|
|
|
|||
|
|
@ -92,7 +92,7 @@ var masterBloomFilter []byte
|
|||
var masterPow = 0.00000001
|
||||
var round int = 1
|
||||
|
||||
func TestSimulation(t *testing.T) {
|
||||
func XTestSimulation(t *testing.T) {
|
||||
// create a chain of whisper nodes,
|
||||
// installs the filters with shared (predefined) parameters
|
||||
initialize(t)
|
||||
|
|
|
|||
Loading…
Reference in a new issue