diff --git a/swarm/network/request_test.go b/swarm/network/request_test.go index cd014a4509..12e820b2ad 100644 --- a/swarm/network/request_test.go +++ b/swarm/network/request_test.go @@ -20,15 +20,8 @@ import ( "bytes" "context" crand "crypto/rand" - "errors" - "flag" "fmt" "io" - "io/ioutil" - "math/rand" - "os" - "sync" - "sync/atomic" "testing" "time" @@ -40,7 +33,6 @@ import ( "github.com/ethereum/go-ethereum/p2p/simulations" "github.com/ethereum/go-ethereum/p2p/simulations/adapters" p2ptest "github.com/ethereum/go-ethereum/p2p/testing" - "github.com/ethereum/go-ethereum/rpc" "github.com/ethereum/go-ethereum/swarm/storage" ) @@ -167,9 +159,11 @@ func TestStreamerUpstreamRetrieveRequestMsgExchange(t *testing.T) { p2ptest.Expect{ Code: 1, Msg: &OfferedHashesMsg{ - HandoverProof: nil, - Hashes: hash, - From: 0, + HandoverProof: &HandoverProof{ + Handover: &Handover{}, + }, + Hashes: hash, + From: 0, // TODO: why is this 32??? To: 32, Key: []byte{}, @@ -305,102 +299,6 @@ func TestStreamerDownstreamChunkDeliveryMsgExchange(t *testing.T) { } -var services = adapters.Services{ - "delivery": newDeliveryService, - "syncer": newSyncerService, -} - -var ( - adapter = flag.String("adapter", "sim", "type of simulation: sim|socket|exec|docker") - loglevel = flag.Int("loglevel", 5, "verbosity of logs") -) - -type roundRobinStore struct { - index uint32 - stores []storage.ChunkStore -} - -func newRoundRobinStore(stores ...storage.ChunkStore) *roundRobinStore { - return &roundRobinStore{ - stores: stores, - } -} - -func (rrs *roundRobinStore) Get(key storage.Key) (*storage.Chunk, error) { - return nil, errors.New("get not well defined on round robin store") -} - -func (rrs *roundRobinStore) Put(chunk *storage.Chunk) { - log.Warn("chunksize", "size", chunk.Size, "sdata", len(chunk.SData)) - i := atomic.AddUint32(&rrs.index, 1) - idx := int(i) % len(rrs.stores) - log.Trace(fmt.Sprintf("put %v into localstore %v", chunk.Key, idx)) - rrs.stores[idx].Put(chunk) -} - -func (rrs *roundRobinStore) Close() { - for _, store := range rrs.stores { - store.Close() - } -} - -func init() { - flag.Parse() - // register the Delivery service which will run as a devp2p - // protocol when using the exec adapter - adapters.RegisterServices(services) - - log.Root().SetHandler(log.LvlFilterHandler(log.Lvl(*loglevel), log.StreamHandler(os.Stderr, log.TerminalFormat(false)))) -} - -func testSimulation(t *testing.T, simf func(adapters.NodeAdapter) (*simulations.StepResult, error)) { - var err error - var result *simulations.StepResult - startedAt := time.Now() - - switch *adapter { - case "sim": - t.Logf("simadapter") - result, err = simf(adapters.NewSimAdapter(services)) - case "socket": - result, err = simf(adapters.NewSocketAdapter(services)) - case "exec": - baseDir, err0 := ioutil.TempDir("", "swarm-test") - if err0 != nil { - t.Fatal(err0) - } - defer os.RemoveAll(baseDir) - result, err = simf(adapters.NewExecAdapter(baseDir)) - case "docker": - adapter, err0 := adapters.NewDockerAdapter() - if err0 != nil { - t.Fatal(err0) - } - result, err = simf(adapter) - default: - t.Fatal("adapter needs to be one of sim, socket, exec, docker") - } - if err != nil { - t.Fatal(err) - } - t.Logf("Simulation with %d nodes passed in %s", len(result.Passes), result.FinishedAt.Sub(result.StartedAt)) - var min, max time.Duration - var sum int - for _, pass := range result.Passes { - duration := pass.Sub(result.StartedAt) - if sum == 0 || duration < min { - min = duration - } - if duration > max { - max = duration - } - sum += int(duration.Nanoseconds()) - } - t.Logf("Min: %s, Max: %s, Average: %s", min, max, time.Duration(sum/len(result.Passes))*time.Nanosecond) - finishedAt := time.Now() - t.Logf("Setup: %s, shutdown: %s", result.StartedAt.Sub(startedAt), finishedAt.Sub(result.FinishedAt)) -} - func TestDeliveryFromNodes(t *testing.T) { testSimulation(t, testDeliveryFromNodes(2, 1, 8100, true)) testSimulation(t, testDeliveryFromNodes(2, 1, 8100, false)) @@ -408,57 +306,6 @@ func TestDeliveryFromNodes(t *testing.T) { testSimulation(t, testDeliveryFromNodes(3, 1, 8100, false)) } -var ( - delivery *Delivery - localStores []storage.ChunkStore - fileHash storage.Key - nodeCount int -) - -func setLocalStores(n int) (func(), error) { - var datadirs []string - localStores = make([]storage.ChunkStore, n) - var err error - for i := 0; i < n; i++ { - // TODO: remove temp datadir after test - var datadir string - datadir, err = ioutil.TempDir("", "streamer") - if err != nil { - break - } - var localStore *storage.LocalStore - localStore, err = storage.NewTestLocalStore(datadir) - if err != nil { - break - } - datadirs = append(datadirs, datadir) - localStores[i] = localStore - } - teardown := func() { - for _, datadir := range datadirs { - os.RemoveAll(datadir) - } - } - return teardown, err -} - -func mustReadAll(dpa *storage.DPA, hash storage.Key) (int, error) { - r := dpa.Retrieve(fileHash) - buf := make([]byte, 1024) - var n, total int - var err error - for (total == 0 || n > 0) && err == nil { - log.Warn(fmt.Sprintf("reading %v bytes at offset %v", len(buf), total)) - n, err = r.ReadAt(buf, int64(total)) - total += n - } - log.Warn(fmt.Sprintf("read %v bytes at offset %v error %v", len(buf), total, err)) - if err != nil && err != io.EOF { - return total, err - } - return total, nil -} - func testDeliveryFromNodes(nodes, conns, size int, skipCheck bool) func(adapter adapters.NodeAdapter) (*simulations.StepResult, error) { return func(adapter adapters.NodeAdapter) (*simulations.StepResult, error) { trigger := func(net *simulations.Network) chan discover.NodeID { @@ -466,12 +313,7 @@ func testDeliveryFromNodes(nodes, conns, size int, skipCheck bool) func(adapter ticker := time.NewTicker(500 * time.Millisecond) go func() { defer ticker.Stop() - // we are only testing the pivot node (net.Nodes[0]) but simulation needs - // all nodes to pass the check so we trigger each and the check function - // will trivially return true - for i := 1; i < nodes; i++ { - triggerC <- net.Nodes[i].ID() - } + // we are only testing the pivot node (net.Nodes[0]) for range ticker.C { triggerC <- net.Nodes[0].ID() } @@ -523,10 +365,9 @@ func testDeliveryFromNodes(nodes, conns, size int, skipCheck bool) func(adapter default: } // try to locally retrieve the file to check if retrieve requests have been successful - log.Warn(fmt.Sprintf("try to locally retrieve %v", fileHash)) total, err := mustReadAll(dpa, fileHash) + log.Debug(fmt.Sprintf("check if %08x is available locally: number of bytes read %v/%v (error: %v)", fileHash, total, size, err)) if err != nil || total != size { - log.Warn(fmt.Sprintf("number of bytes read %v/%v (error: %v)", total, size, err)) return false, nil } return true, nil @@ -547,7 +388,7 @@ func testDeliveryFromNodes(nodes, conns, size int, skipCheck bool) func(adapter } } - result, err := runSimulation(nodes, conns, "delivery", action, trigger, check, adapter) + result, err := runSimulation(nodes, conns, "delivery", NewAddrFromNodeID, action, trigger, check, adapter) if err != nil { return nil, fmt.Errorf("Setting up simulation failed: %v", err) } @@ -558,83 +399,6 @@ func testDeliveryFromNodes(nodes, conns, size int, skipCheck bool) func(adapter } } -func runSimulation(nodes, conns int, serviceName string, action func(*simulations.Network) func(context.Context) error, trigger func(*simulations.Network) chan discover.NodeID, check func(*simulations.Network, *storage.DPA) func(context.Context, discover.NodeID) (bool, error), adapter adapters.NodeAdapter) (*simulations.StepResult, error) { - // create network - net := simulations.NewNetwork(adapter, &simulations.NetworkConfig{ - ID: "0", - DefaultService: serviceName, - }) - defer net.Shutdown() - // set nodes number of localstores globally available - teardown, err := setLocalStores(nodes) - defer teardown() - if err != nil { - return nil, err - } - ids := make([]discover.NodeID, nodes) - nodeCount = 0 - // start nodes - for i := 0; i < nodes; i++ { - node, err := net.NewNode() - if err != nil { - return nil, fmt.Errorf("error starting node: %s", err) - } - if err := net.Start(node.ID()); err != nil { - return nil, fmt.Errorf("error starting node %s: %s", node.ID().TerminalString(), err) - } - ids[i] = node.ID() - } - - // run a simulation which connects the 10 nodes in a chain - var addrs [][]byte - wg := sync.WaitGroup{} - log.Warn("runSimulation 1") - for i := range ids { - log.Warn("runSimulation 2") - // collect the overlay addresses, to - addrs = append(addrs, ToOverlayAddr(ids[i].Bytes())) - for j := 0; j < conns; j++ { - log.Warn("runSimulation 3") - var k int - if j == 0 { - k = i - 1 - } else { - k = rand.Intn(len(ids)) - } - if i > 0 { - log.Warn("runSimulation 4") - wg.Add(1) - go func(i, k int) { - defer wg.Done() - log.Warn("net.Connect") - net.Connect(ids[i], ids[k]) - }(i, k) - } - } - } - wg.Wait() - - log.Debug(fmt.Sprintf("nodes: %v", len(addrs))) - - // create an only locally retrieving dpa for the pivot node to test - // if retriee requests have arrived - dpa := storage.NewDPA(localStores[0], storage.NewChunkerParams()) - dpa.Start() - defer dpa.Stop() - timeout := 300 * time.Second - ctx, cancel := context.WithTimeout(context.Background(), timeout) - defer cancel() - result := simulations.NewSimulation(net).Run(ctx, &simulations.Step{ - Action: action(net), - Trigger: trigger(net), - Expect: &simulations.Expectation{ - Nodes: ids, - Check: check(net, dpa), - }, - }) - return result, nil -} - // newDeliveryService func newDeliveryService(ctx *adapters.ServiceContext) (node.Service, error) { id := ctx.Config.ID @@ -649,49 +413,15 @@ func newDeliveryService(ctx *adapters.ServiceContext) (node.Service, error) { // swarm enabled dpa delivery = streamer.delivery } - nodeCount++ - - log.Warn("new service created") self := &testStreamerService{ addr: addr, streamer: streamer, } self.run = self.runDelivery + nodeCount++ return self, nil } -type testStreamerService struct { - addr *BzzAddr - streamer *Streamer - run func(p *p2p.Peer, rw p2p.MsgReadWriter) error -} - -func (tds *testStreamerService) Protocols() []p2p.Protocol { - log.Warn("Protocols function", "run", tds.run) - return []p2p.Protocol{ - { - Name: StreamerSpec.Name, - Version: StreamerSpec.Version, - Length: StreamerSpec.Length(), - Run: tds.run, - // NodeInfo: , - // PeerInfo: , - }, - } -} - -func (b *testStreamerService) APIs() []rpc.API { - return []rpc.API{} -} - -func (b *testStreamerService) Start(server *p2p.Server) error { - return nil -} - -func (b *testStreamerService) Stop() error { - return nil -} - func (b *testStreamerService) runDelivery(p *p2p.Peer, rw p2p.MsgReadWriter) error { bzzPeer := &bzzPeer{ Peer: protocols.NewPeer(p, rw, StreamerSpec), diff --git a/swarm/network/requests.go b/swarm/network/requests.go index 7cc0a47946..3495cfceb9 100644 --- a/swarm/network/requests.go +++ b/swarm/network/requests.go @@ -168,8 +168,12 @@ func (self *Delivery) processReceivedChunks() { continue } chunk.SData = req.SData - self.dbAccess.put(chunk) - close(chunk.ReqC) + select { + case <-chunk.ReqC: + default: + self.dbAccess.put(chunk) + close(chunk.ReqC) + } } } diff --git a/swarm/network/streamer.go b/swarm/network/streamer.go index 3c367f7b6f..66abf82f3a 100644 --- a/swarm/network/streamer.go +++ b/swarm/network/streamer.go @@ -20,6 +20,7 @@ import ( "context" "errors" "fmt" + "math" "sync" "github.com/ethereum/go-ethereum/log" @@ -109,6 +110,14 @@ func (self WantedHashesMsg) String() string { return fmt.Sprintf("Stream '%v', Want: %x, Next: [%v-%v]", self.Stream, self.Want, self.From, self.To) } +func keyToString(key []byte) string { + l := len(key) + if l == 0 { + return "" + } + return fmt.Sprintf("%s-%d", string(key[:l-1]), uint8(key[l-1])) +} + // Streamer registry for outgoing and incoming streamer constructors type Streamer struct { incomingLock sync.RWMutex @@ -187,6 +196,7 @@ type outgoingStreamer struct { priority uint8 currentBatch []byte stream string + key []byte } // OutgoingStreamer interface for outgoing peer Streamer @@ -200,6 +210,8 @@ type incomingStreamer struct { priority uint8 sessionAt uint64 live bool + stream string + key []byte quit chan struct{} next chan struct{} } @@ -280,26 +292,30 @@ func (self *StreamerPeer) getIncomingStreamer(s string) (*incomingStreamer, erro return streamer, nil } -func (self *StreamerPeer) setOutgoingStreamer(s string, o OutgoingStreamer, priority uint8) (*outgoingStreamer, error) { +func (self *StreamerPeer) setOutgoingStreamer(s string, key []byte, o OutgoingStreamer, priority uint8) (*outgoingStreamer, error) { self.outgoingLock.Lock() defer self.outgoingLock.Unlock() - if self.outgoing[s] != nil { - return nil, fmt.Errorf("stream %v already registered", s) + sk := s + keyToString(key) + if self.outgoing[sk] != nil { + return nil, fmt.Errorf("stream %v already registered", sk) } os := &outgoingStreamer{ OutgoingStreamer: o, priority: priority, stream: s, + key: key, } - self.outgoing[s] = os + self.outgoing[sk] = os return os, nil } -func (self *StreamerPeer) setIncomingStreamer(s string, i IncomingStreamer, priority uint8, live bool) error { +func (self *StreamerPeer) setIncomingStreamer(s string, key []byte, i IncomingStreamer, priority uint8, live bool) error { self.incomingLock.Lock() defer self.incomingLock.Unlock() - if self.incoming[s] != nil { - return fmt.Errorf("stream %v already registered", s) + + sk := s + keyToString(key) + if self.incoming[sk] != nil { + return fmt.Errorf("stream %v already registered", sk) } next := make(chan struct{}, 1) // var intervals *Intervals @@ -307,12 +323,14 @@ func (self *StreamerPeer) setIncomingStreamer(s string, i IncomingStreamer, prio // key := s + self.ID().String() // intervals = NewIntervals(key, self.streamer) // } - self.incoming[s] = &incomingStreamer{ + self.incoming[sk] = &incomingStreamer{ IncomingStreamer: i, // intervals: intervals, live: live, priority: priority, next: next, + stream: s, + key: key, } next <- struct{}{} // this is to allow wantedKeysMsg before first batch arrives return nil @@ -330,6 +348,8 @@ func (self *incomingStreamer) nextBatch(from uint64) (nextFrom uint64, nextTo ui nextFrom = from } else if from >= self.sessionAt { // history sync complete intervals = nil + nextFrom = from + nextTo = math.MaxUint64 } else if len(intervals) > 2 && from >= intervals[2] { // filled a gap in the intervals intervals = append(intervals[:1], intervals[3:]...) nextFrom = intervals[1] @@ -349,7 +369,6 @@ func (self *incomingStreamer) nextBatch(from uint64) (nextFrom uint64, nextTo ui // Subscribe initiates the streamer func (self *Streamer) Subscribe(peerId discover.NodeID, s string, t []byte, from, to uint64, priority uint8, live bool) error { - log.Warn("!!!!!! Subscribe ", "peer", peerId) f, err := self.GetIncomingStreamer(s) if err != nil { return err @@ -364,18 +383,21 @@ func (self *Streamer) Subscribe(peerId discover.NodeID, s string, t []byte, from if err != nil { return err } - err = peer.setIncomingStreamer(s+string(t), is, priority, live) + err = peer.setIncomingStreamer(s, t, is, priority, live) if err != nil { return err } msg := &SubscribeMsg{ - Stream: s, - Key: t, + Stream: s, + Key: t, + // Live: live, From: from, To: to, Priority: priority, } + log.Debug("Subscribe ", "peer", peerId, "stream", s, "key", t, "from", from, "to", to) + peer.SendPriority(msg, priority) return nil } @@ -389,11 +411,11 @@ func (self *StreamerPeer) handleSubscribeMsg(req *SubscribeMsg) error { if err != nil { return err } - key := req.Stream + string(req.Key) - os, err := self.setOutgoingStreamer(key, s, req.Priority) + os, err := self.setOutgoingStreamer(req.Stream, req.Key, s, req.Priority) if err != nil { return nil } + log.Debug("received subscription", "stream", req.Stream, "Key", req.Key, "from", req.From, "to", req.To) go self.SendOfferedHashes(os, req.From, req.To) return nil } @@ -401,7 +423,9 @@ func (self *StreamerPeer) handleSubscribeMsg(req *SubscribeMsg) error { // handleOfferedHashesMsg protocol msg handler calls the incoming streamer interface // Filter method func (self *StreamerPeer) handleOfferedHashesMsg(req *OfferedHashesMsg) error { - s, err := self.getIncomingStreamer(req.Stream) + sk := req.Stream + sk += keyToString(req.Key) + s, err := self.getIncomingStreamer(sk) if err != nil { return err } @@ -440,11 +464,14 @@ func (self *StreamerPeer) handleOfferedHashesMsg(req *OfferedHashesMsg) error { s.sessionAt = req.From } from, to := s.nextBatch(req.To) + log.Debug("received batch", "stream", req.Stream, "Key", req.Key, "from", req.From, "to", req.To) if from == to { return nil } + msg := &WantedHashesMsg{ Stream: req.Stream, + Key: req.Key, Want: want.Bytes(), From: from, To: to, @@ -455,6 +482,7 @@ func (self *StreamerPeer) handleOfferedHashesMsg(req *OfferedHashesMsg) error { case <-s.quit: return } + log.Debug("want batch", "stream", msg.Stream, "Key", msg.Key, "from", msg.From, "to", msg.To) self.SendPriority(msg, s.priority) }() return nil @@ -464,8 +492,10 @@ func (self *StreamerPeer) handleOfferedHashesMsg(req *OfferedHashesMsg) error { // * sends the next batch of unsynced keys // * sends the actual data chunks as per WantedHashesMsg func (self *StreamerPeer) handleWantedHashesMsg(req *WantedHashesMsg) error { - s, err := self.getOutgoingStreamer(req.Stream) + log.Debug("received wanted batch", "stream", req.Stream, "Key", req.Key, "from", req.From, "to", req.To) + s, err := self.getOutgoingStreamer(req.Stream + keyToString(req.Key)) if err != nil { + log.Debug(err.Error()) return err } hashes := s.currentBatch @@ -534,9 +564,9 @@ func (self *StreamerPeer) SendOfferedHashes(s *outgoingStreamer, f, t uint64) er From: from, To: to, Stream: s.stream, - // TODO: use real key here - Key: []byte{}, + Key: s.key, } + log.Debug("Swarm syncer offer batch", "stream", s.stream, "key", s.key, "len", len(hashes), "from", from, "to", to) return self.SendPriority(msg, s.priority) } diff --git a/swarm/network/streamer_common_test.go b/swarm/network/streamer_common_test.go new file mode 100644 index 0000000000..7690ed75c8 --- /dev/null +++ b/swarm/network/streamer_common_test.go @@ -0,0 +1,349 @@ +// Copyright 2016 The go-ethereum Authors +// This file is part of the go-ethereum library. +// +// The go-ethereum library is free software: you can redistribute it and/or modify +// it under the terms of the GNU Lesser General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// The go-ethereum library is distributed in the hope that it will be useful, +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Lesser General Public License for more details. +// +// You should have received a copy of the GNU Lesser General Public License +// along with the go-ethereum library. If not, see . + +package network + +import ( + "context" + "errors" + "flag" + "fmt" + "io" + "io/ioutil" + "math/rand" + "os" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/ethereum/go-ethereum/log" + "github.com/ethereum/go-ethereum/p2p" + "github.com/ethereum/go-ethereum/p2p/discover" + "github.com/ethereum/go-ethereum/p2p/protocols" + "github.com/ethereum/go-ethereum/p2p/simulations" + "github.com/ethereum/go-ethereum/p2p/simulations/adapters" + p2ptest "github.com/ethereum/go-ethereum/p2p/testing" + "github.com/ethereum/go-ethereum/rpc" + "github.com/ethereum/go-ethereum/swarm/storage" +) + +var services = adapters.Services{ + "delivery": newDeliveryService, + "syncer": newSyncerService, +} + +var ( + adapter = flag.String("adapter", "sim", "type of simulation: sim|socket|exec|docker") + loglevel = flag.Int("loglevel", 2, "verbosity of logs") +) + +func init() { + flag.Parse() + // register the Delivery service which will run as a devp2p + // protocol when using the exec adapter + adapters.RegisterServices(services) + + log.Root().SetHandler(log.LvlFilterHandler(log.Lvl(*loglevel), log.StreamHandler(os.Stderr, log.TerminalFormat(false)))) +} + +var ( + delivery *Delivery + localStores []storage.ChunkStore + addrs []Addr + fileHash storage.Key + nodeCount int +) + +func setLocalStores(addrs ...Addr) (func(), error) { + var datadirs []string + localStores = make([]storage.ChunkStore, len(addrs)) + var err error + for i, addr := range addrs { + // TODO: remove temp datadir after test + var datadir string + datadir, err = ioutil.TempDir("", "streamer") + if err != nil { + break + } + var localStore *storage.LocalStore + localStore, err = storage.NewTestLocalStoreForAddr(datadir, addr.Over()) + if err != nil { + break + } + datadirs = append(datadirs, datadir) + localStores[i] = localStore + } + teardown := func() { + for _, datadir := range datadirs { + os.RemoveAll(datadir) + } + } + return teardown, err +} + +func mustReadAll(dpa *storage.DPA, hash storage.Key) (int, error) { + r := dpa.Retrieve(fileHash) + buf := make([]byte, 1024) + var n, total int + var err error + for (total == 0 || n > 0) && err == nil { + n, err = r.ReadAt(buf, int64(total)) + total += n + } + if err != nil && err != io.EOF { + return total, err + } + return total, nil +} + +func testSimulation(t *testing.T, simf func(adapters.NodeAdapter) (*simulations.StepResult, error)) { + var err error + var result *simulations.StepResult + startedAt := time.Now() + + switch *adapter { + case "sim": + t.Logf("simadapter") + result, err = simf(adapters.NewSimAdapter(services)) + case "socket": + result, err = simf(adapters.NewSocketAdapter(services)) + case "exec": + baseDir, err0 := ioutil.TempDir("", "swarm-test") + if err0 != nil { + t.Fatal(err0) + } + defer os.RemoveAll(baseDir) + result, err = simf(adapters.NewExecAdapter(baseDir)) + case "docker": + adapter, err0 := adapters.NewDockerAdapter() + if err0 != nil { + t.Fatal(err0) + } + result, err = simf(adapter) + default: + t.Fatal("adapter needs to be one of sim, socket, exec, docker") + } + if err != nil { + t.Fatal(err) + } + t.Logf("Simulation with %d nodes passed in %s", len(result.Passes), result.FinishedAt.Sub(result.StartedAt)) + var min, max time.Duration + var sum int + for _, pass := range result.Passes { + duration := pass.Sub(result.StartedAt) + if sum == 0 || duration < min { + min = duration + } + if duration > max { + max = duration + } + sum += int(duration.Nanoseconds()) + } + t.Logf("Min: %s, Max: %s, Average: %s", min, max, time.Duration(sum/len(result.Passes))*time.Nanosecond) + finishedAt := time.Now() + t.Logf("Setup: %s, shutdown: %s", result.StartedAt.Sub(startedAt), finishedAt.Sub(result.FinishedAt)) +} + +func runSimulation(nodes, conns int, serviceName string, toAddr func(discover.NodeID) *BzzAddr, action func(*simulations.Network) func(context.Context) error, trigger func(*simulations.Network) chan discover.NodeID, check func(*simulations.Network, *storage.DPA) func(context.Context, discover.NodeID) (bool, error), adapter adapters.NodeAdapter) (*simulations.StepResult, error) { + // create network + net := simulations.NewNetwork(adapter, &simulations.NetworkConfig{ + ID: "0", + DefaultService: serviceName, + }) + defer net.Shutdown() + ids := make([]discover.NodeID, nodes) + nodeCount = 0 + addrs = make([]Addr, nodes) + // start nodes + for i := 0; i < nodes; i++ { + node, err := net.NewNode() + if err != nil { + return nil, fmt.Errorf("error creating node: %s", err) + } + ids[i] = node.ID() + addrs[i] = toAddr(ids[i]) + } + // set nodes number of localstores globally available + teardown, err := setLocalStores(addrs...) + defer teardown() + if err != nil { + return nil, err + } + + for i := 0; i < nodes; i++ { + if err := net.Start(ids[i]); err != nil { + return nil, fmt.Errorf("error starting node %s: %s", ids[i].TerminalString(), err) + } + } + + // run a simulation which connects the 10 nodes in a chain + wg := sync.WaitGroup{} + for i := range ids { + // collect the overlay addresses, to + for j := 0; j < conns; j++ { + var k int + if j == 0 { + k = i - 1 + } else { + k = rand.Intn(len(ids)) + } + if i > 0 { + wg.Add(1) + go func(i, k int) { + defer wg.Done() + net.Connect(ids[i], ids[k]) + }(i, k) + } + } + } + wg.Wait() + + log.Debug(fmt.Sprintf("nodes: %v", len(addrs))) + + // create an only locally retrieving dpa for the pivot node to test + // if retriee requests have arrived + dpa := storage.NewDPA(localStores[0], storage.NewChunkerParams()) + dpa.Start() + defer dpa.Stop() + timeout := 300 * time.Second + ctx, cancel := context.WithTimeout(context.Background(), timeout) + defer cancel() + result := simulations.NewSimulation(net).Run(ctx, &simulations.Step{ + Action: action(net), + Trigger: trigger(net), + Expect: &simulations.Expectation{ + Nodes: ids[0:1], + Check: check(net, dpa), + }, + }) + return result, nil +} + +func newStreamerTester(t *testing.T) (*p2ptest.ProtocolTester, *Streamer, *storage.LocalStore, func(), error) { + // setup + addr := RandomAddr() // tested peers peer address + to := NewKademlia(addr.OAddr, NewKadParams()) + + // temp datadir + datadir, err := ioutil.TempDir("", "streamer") + if err != nil { + return nil, nil, nil, func() {}, err + } + teardown := func() { + os.RemoveAll(datadir) + } + + localStore, err := storage.NewTestLocalStoreForAddr(datadir, addr.Over()) + if err != nil { + return nil, nil, nil, teardown, err + } + + dbAccess := NewDbAccess(localStore) + delivery := NewDelivery(to, dbAccess) + streamer := NewStreamer(delivery) + run := func(p *p2p.Peer, rw p2p.MsgReadWriter) error { + bzzPeer := &bzzPeer{ + Peer: protocols.NewPeer(p, rw, StreamerSpec), + localAddr: addr, + BzzAddr: NewAddrFromNodeID(p.ID()), + } + to.On(bzzPeer) + return streamer.Run(bzzPeer) + } + protocolTester := p2ptest.NewProtocolTester(t, NewNodeIDFromAddr(addr), 1, run) + + err = waitForPeers(streamer, 1*time.Second) + if err != nil { + return nil, nil, nil, nil, errors.New("timeout: peer is not created") + } + + return protocolTester, streamer, localStore, teardown, nil +} + +type roundRobinStore struct { + index uint32 + stores []storage.ChunkStore +} + +func newRoundRobinStore(stores ...storage.ChunkStore) *roundRobinStore { + return &roundRobinStore{ + stores: stores, + } +} + +func (rrs *roundRobinStore) Get(key storage.Key) (*storage.Chunk, error) { + return nil, errors.New("get not well defined on round robin store") +} + +func (rrs *roundRobinStore) Put(chunk *storage.Chunk) { + i := atomic.AddUint32(&rrs.index, 1) + idx := int(i) % len(rrs.stores) + rrs.stores[idx].Put(chunk) +} + +func (rrs *roundRobinStore) Close() { + for _, store := range rrs.stores { + store.Close() + } +} + +func waitForPeers(streamer *Streamer, timeout time.Duration) error { + ticker := time.NewTicker(10 * time.Millisecond) + timeoutTimer := time.NewTimer(timeout) + for { + select { + case <-ticker.C: + if len(streamer.peers) > 0 { + return nil + } + case <-timeoutTimer.C: + return errors.New("timeout") + } + } +} + +type testStreamerService struct { + index int + addr *BzzAddr + streamer *Streamer + run func(p *p2p.Peer, rw p2p.MsgReadWriter) error +} + +func (tds *testStreamerService) Protocols() []p2p.Protocol { + return []p2p.Protocol{ + { + Name: StreamerSpec.Name, + Version: StreamerSpec.Version, + Length: StreamerSpec.Length(), + Run: tds.run, + // NodeInfo: , + // PeerInfo: , + }, + } +} + +func (b *testStreamerService) APIs() []rpc.API { + return []rpc.API{} +} + +func (b *testStreamerService) Start(server *p2p.Server) error { + return nil +} + +func (b *testStreamerService) Stop() error { + return nil +} diff --git a/swarm/network/streamer_test.go b/swarm/network/streamer_test.go index 86ba965071..447713a33d 100644 --- a/swarm/network/streamer_test.go +++ b/swarm/network/streamer_test.go @@ -18,65 +18,13 @@ package network import ( "bytes" - "errors" - "io/ioutil" - "os" "testing" "time" "github.com/ethereum/go-ethereum/crypto/sha3" - "github.com/ethereum/go-ethereum/p2p" - "github.com/ethereum/go-ethereum/p2p/protocols" p2ptest "github.com/ethereum/go-ethereum/p2p/testing" - "github.com/ethereum/go-ethereum/swarm/storage" ) -// -// func init() { -// log.Root().SetHandler(log.CallerFileHandler(log.LvlFilterHandler(log.LvlWarn, log.StreamHandler(os.Stderr, log.TerminalFormat(true))))) -// } - -func newStreamerTester(t *testing.T) (*p2ptest.ProtocolTester, *Streamer, *storage.LocalStore, func(), error) { - // setup - addr := RandomAddr() // tested peers peer address - to := NewKademlia(addr.OAddr, NewKadParams()) - - // temp datadir - datadir, err := ioutil.TempDir("", "streamer") - if err != nil { - return nil, nil, nil, func() {}, err - } - teardown := func() { - os.RemoveAll(datadir) - } - - localStore, err := storage.NewTestLocalStore(datadir) - if err != nil { - return nil, nil, nil, teardown, err - } - - dbAccess := NewDbAccess(localStore) - delivery := NewDelivery(to, dbAccess) - streamer := NewStreamer(delivery) - run := func(p *p2p.Peer, rw p2p.MsgReadWriter) error { - bzzPeer := &bzzPeer{ - Peer: protocols.NewPeer(p, rw, StreamerSpec), - localAddr: addr, - BzzAddr: NewAddrFromNodeID(p.ID()), - } - to.On(bzzPeer) - return streamer.Run(bzzPeer) - } - protocolTester := p2ptest.NewProtocolTester(t, NewNodeIDFromAddr(addr), 1, run) - - err = waitForPeers(streamer, 1*time.Second) - if err != nil { - return nil, nil, nil, nil, errors.New("timeout: peer is not created") - } - - return protocolTester, streamer, localStore, teardown, nil -} - func TestStreamerSubscribe(t *testing.T) { tester, streamer, _, teardown, err := newStreamerTester(t) defer teardown() @@ -214,6 +162,7 @@ func TestStreamerUpstreamSubscribeMsgExchange(t *testing.T) { Code: 1, Msg: &OfferedHashesMsg{ Stream: "foo", + Key: []byte{}, HandoverProof: &HandoverProof{ Handover: &Handover{}, }, @@ -329,18 +278,3 @@ func TestStreamerDownstreamOfferedHashesMsgExchange(t *testing.T) { } } - -func waitForPeers(streamer *Streamer, timeout time.Duration) error { - ticker := time.NewTicker(10 * time.Millisecond) - timeoutTimer := time.NewTimer(timeout) - for { - select { - case <-ticker.C: - if len(streamer.peers) > 0 { - return nil - } - case <-timeoutTimer.C: - return errors.New("timeout") - } - } -} diff --git a/swarm/network/syncer.go b/swarm/network/syncer.go index db2e551826..ec8a42808e 100644 --- a/swarm/network/syncer.go +++ b/swarm/network/syncer.go @@ -25,12 +25,12 @@ import ( "time" "github.com/ethereum/go-ethereum/log" - "github.com/ethereum/go-ethereum/p2p/discover" "github.com/ethereum/go-ethereum/swarm/storage" ) const ( - batchSize = 128 + batchSize = 2 + // batchSize = 128 ) // wrapper of db-s to provide mockable custom local chunk store access to syncer @@ -60,7 +60,6 @@ func (self *DbAccess) iterator(from uint64, to uint64, po uint8, f func(storage. // to obtain the chunks from key or request db entry only func (self *DbAccess) getOrCreateRequest(key storage.Key) (*storage.Chunk, bool) { - log.Warn("getOrCreateRequest", "self", self) return self.loc.GetOrCreateRequest(key) } @@ -100,17 +99,9 @@ const maxPO = 32 func RegisterOutgoingSyncer(streamer *Streamer, db *DbAccess) { streamer.RegisterOutgoingStreamer("SYNC", func(p *StreamerPeer, t []byte) (OutgoingStreamer, error) { - syncType, po := parseSyncLabel(t) + po := uint8(t[0]) // TODO: make this work for HISTORY too - syncType = "LIVE" - switch syncType { - case "LIVE": - return NewOutgoingSwarmSyncer(true, po, db) - case "HISTORY": - return NewOutgoingSwarmSyncer(false, po, db) - default: - return nil, errors.New("invalid sync type") - } + return NewOutgoingSwarmSyncer(false, po, db) }) // streamer.RegisterOutgoingStreamer(stream, func(p *StreamerPeer) (OutgoingStreamer, error) { // return NewOutgoingProvableSwarmSyncer(po, db) @@ -133,10 +124,9 @@ func (self *OutgoingSwarmSyncer) SetNextBatch(from, to uint64) ([]byte, uint64, if from == 0 { from = self.start } - if to <= from { + if to <= from || from >= self.sessionAt { to = math.MaxUint64 } - log.Warn("!!!!!!!!!!!!! setNextBatch", "from", from, "to", to, "currentStoreCount", self.db.currentBucketStorageIndex(1)) ticker := time.NewTicker(10 * time.Millisecond) defer ticker.Stop() for range ticker.C { @@ -154,7 +144,7 @@ func (self *OutgoingSwarmSyncer) SetNextBatch(from, to uint64) ([]byte, uint64, } } - log.Debug("Swarm batch", "po", self.po, "len", i, "from", from, "to", to) + log.Debug("Swarm syncer offer batch", "po", self.po, "len", i, "from", from, "to", to, "current store count", self.db.currentBucketStorageIndex(self.po)) return batch, from, to + 1, nil, nil } @@ -204,40 +194,24 @@ func NewIncomingSwarmSyncer(p Peer, dbAccess *DbAccess, chunker storage.Chunker) // return self // } -func newSyncLabel(typ string, po uint8) []byte { - t := []byte(typ) - t = append(t, byte(po)) - return t -} - -func parseSyncLabel(t []byte) (string, uint8) { - l := len(t) - 1 - return string(t[:l]), uint8(t[l]) -} - -// StartSyncing is called on the StreamerPeer to start the syncing process -// the idea is that it is called only after kademlia is close to healthy -func StartSyncing(s *Streamer, peerId discover.NodeID, po uint8, nn bool) { - lastPO := po - if nn { - lastPO = maxPO - } - - for i := po; i <= lastPO; i++ { - s.Subscribe(peerId, "SYNC", newSyncLabel("LIVE", po), 0, 0, High, true) - s.Subscribe(peerId, "SYNC", newSyncLabel("HISTORY", po), 0, 0, Mid, false) - } -} +// // StartSyncing is called on the StreamerPeer to start the syncing process +// // the idea is that it is called only after kademlia is close to healthy +// func StartSyncing(s *Streamer, peerId discover.NodeID, po uint8, nn bool) { +// lastPO := po +// if nn { +// lastPO = maxPO +// } +// +// for i := po; i <= lastPO; i++ { +// s.Subscribe(peerId, "SYNC", newSyncLabel("LIVE", po), 0, 0, High, true) +// s.Subscribe(peerId, "SYNC", newSyncLabel("HISTORY", po), 0, 0, Mid, false) +// } +// } func RegisterIncomingSyncer(streamer *Streamer, db *DbAccess) { streamer.RegisterIncomingStreamer("SYNC", func(p *StreamerPeer, t []byte) (IncomingStreamer, error) { return NewIncomingSwarmSyncer(p, db, nil) }) - // stream = fmt.Sprintf("SYNC-%02d-delete", po) - // streamer.RegisterIncomingStreamer(stream, func(p *StreamerPeer) (OutgoingStreamer, error) { - // intervals := loadIntervals(p, po, true) - // return NewIncomingSwarmSyncer(po, Mid, sessionAt, intervals, p) - // }) } // NeedData @@ -287,9 +261,10 @@ func (self *IncomingSwarmSyncer) TakeoverProof(s string, from uint64, hashes []b self.end += uint64(len(hashes)) / HashSize takeover := &Takeover{ Stream: s, - Start: self.start, - End: self.end, - Root: root, + // Key: self.Key, + Start: self.start, + End: self.end, + Root: root, } // serialise and sign return &TakeoverProof{ diff --git a/swarm/network/syncer_test.go b/swarm/network/syncer_test.go index 5962f79583..893be367f5 100644 --- a/swarm/network/syncer_test.go +++ b/swarm/network/syncer_test.go @@ -22,7 +22,6 @@ import ( "fmt" "io" "math" - "net" "testing" "time" @@ -36,26 +35,19 @@ import ( "github.com/ethereum/go-ethereum/swarm/storage" ) -var nodeAddrById map[discover.NodeID]*BzzAddr - func TestSyncerSimulation(t *testing.T) { testSimulation(t, testSyncBetweenNodes(2, 1, 81000, true, 1)) + testSimulation(t, testSyncBetweenNodes(3, 1, 81000, true, 1)) } func testSyncBetweenNodes(nodes, conns, size int, skipCheck bool, po uint8) func(adapter adapters.NodeAdapter) (*simulations.StepResult, error) { return func(adapter adapters.NodeAdapter) (*simulations.StepResult, error) { - nodeAddrById = make(map[discover.NodeID]*BzzAddr) trigger := func(net *simulations.Network) chan discover.NodeID { triggerC := make(chan discover.NodeID) ticker := time.NewTicker(500 * time.Millisecond) go func() { defer ticker.Stop() - // we are only testing the pivot node (net.Nodes[0]) but simulation needs - // all nodes to pass the check so we trigger each and the check function - // will trivially return true - for i := 1; i < nodes; i++ { - triggerC <- net.Nodes[i].ID() - } + // we are only testing the pivot node (net.Nodes[0]) for range ticker.C { triggerC <- net.Nodes[0].ID() } @@ -68,29 +60,15 @@ func testSyncBetweenNodes(nodes, conns, size int, skipCheck bool, po uint8) func rrdpa := storage.NewDPA(newRoundRobinStore(localStores[1:]...), storage.NewChunkerParams()) rrdpa.Start() // create a retriever dpa for the pivot node - dpacs := storage.NewDpaChunkStore(localStores[0].(*storage.LocalStore), func(chunk *storage.Chunk) error { return delivery.RequestFromPeers(chunk.Key[:], skipCheck) }) - dpa := storage.NewDPA(dpacs, storage.NewChunkerParams()) - dpa.Start() return func(context.Context) error { defer rrdpa.Stop() // upload an actual random file of size size - _, _, err := rrdpa.Store(io.LimitReader(crand.Reader, int64(size)), int64(size)) + _, wait, err := rrdpa.Store(io.LimitReader(crand.Reader, int64(size)), int64(size)) if err != nil { return err } - // // wait until all chunks stored - // wait() - // // assign the fileHash to a global so that it is available for the check function - // fileHash = hash - // go func() { - // defer dpa.Stop() - // log.Debug(fmt.Sprintf("retrieve %v", fileHash)) - // // start the retrieval on the pivot node - this will spawn retrieve requests for missing chunks - // // we must wait for the peer connections to have started before requesting - // time.Sleep(2 * time.Second) - // n, err := mustReadAll(dpa, fileHash) - // log.Debug(fmt.Sprintf("retrieved %v", fileHash), "read", n, "err", err) - // }() + // wait until all chunks stored + wait() return nil } } @@ -102,52 +80,37 @@ func testSyncBetweenNodes(nodes, conns, size int, skipCheck bool, po uint8) func dbAccesses[i] = NewDbAccess(localStores[i].(*storage.LocalStore)) } return func(ctx context.Context, id discover.NodeID) (bool, error) { - var found, total int - dbAccesses[1].iterator(0, math.MaxUint64, po, func(key storage.Key, index uint64) bool { - _, err := dbAccesses[0].get(key) - if err == nil { - found++ - } - total++ - return true - }) - - // - // if id != net.Nodes[0].ID() { - // return true, nil - // } + if id != net.Nodes[0].ID() { + return true, nil + } select { case <-ctx.Done(): return false, ctx.Err() default: } + + var found, total int + for i := 1; i < nodes; i++ { + dbAccesses[i].iterator(0, math.MaxUint64, po, func(key storage.Key, index uint64) bool { + _, err := dbAccesses[0].get(key) + if err == nil { + found++ + } + total++ + return true + }) + } + log.Debug("sync check", "bin", po, "found", found, "total", total) return found == total, nil - // // try to locally retrieve the file to check if retrieve requests have been successful - // log.Warn(fmt.Sprintf("try to locally retrieve %v", fileHash)) - // total, err := mustReadAll(dpa, fileHash) - // if err != nil || total != size { - // log.Warn(fmt.Sprintf("number of bytes read %v/%v (error: %v)", total, size, err)) - // return false, nil - // } - // return true, nil - // node := net.GetNode(id) - // if node == nil { - // return false, fmt.Errorf("unknown node: %s", id) - // } - // client, err := node.Client() - // if err != nil { - // return false, fmt.Errorf("error getting node client: %s", err) - // } - // var response int - // if err := client.Call(&response, "test_haslocal", hash); err != nil { - // return false, fmt.Errorf("error getting bzz_has response: %s", err) - // } - // log.Debug(fmt.Sprintf("node has: %v\n%v", id, response)) - // return response == 0, nil } } + toAddr := func(id discover.NodeID) *BzzAddr { + addr := NewAddrFromNodeID(id) + addr.OAddr[0] = byte(0) + return addr + } - result, err := runSimulation(nodes, conns, "syncer", action, trigger, check, adapter) + result, err := runSimulation(nodes, conns, "syncer", toAddr, action, trigger, check, adapter) if err != nil { return nil, fmt.Errorf("Setting up simulation failed: %v", err) } @@ -161,60 +124,45 @@ func testSyncBetweenNodes(nodes, conns, size int, skipCheck bool, po uint8) func func newSyncerService(ctx *adapters.ServiceContext) (node.Service, error) { id := ctx.Config.ID addr := NewAddrFromNodeID(id) + // for the test we make all peers share 8 bits so that syncing full bins make sense + addr.OAddr[0] = byte(0) kad := NewKademlia(addr.Over(), NewKadParams()) localStore := localStores[nodeCount] dbAccess := NewDbAccess(localStore.(*storage.LocalStore)) streamer := NewStreamer(NewDelivery(kad, dbAccess)) - log.Warn("!!!!!!!! Registering syncers") RegisterIncomingSyncer(streamer, dbAccess) RegisterOutgoingSyncer(streamer, dbAccess) - addrBytes := addr.Address() - if nodeCount == 0 { - // the delivery service for the pivot node is assigned globally - // so that the simulation action call can use it for the - // swarm enabled dpa - delivery = streamer.delivery - addrBytes[0] = 0x0 - } else { - addrBytes[0] = 0xF0 - } - addr = &BzzAddr{ - OAddr: addrBytes, - UAddr: []byte(discover.NewNode(id, net.IP{127, 0, 0, 1}, 30303, 30303).String()), - } - nodeAddrById[id] = addr - //else { - // RegisterOutgoingSyncer(streamer, dbAccess) - // } - nodeCount++ - - log.Warn("new service created") self := &testStreamerService{ + index: nodeCount, addr: addr, streamer: streamer, } self.run = self.runSyncer + nodeCount++ return self, nil } func (b *testStreamerService) runSyncer(p *p2p.Peer, rw p2p.MsgReadWriter) error { + addr := NewAddrFromNodeID(p.ID()) + addr.OAddr[0] = byte(0) bzzPeer := &bzzPeer{ Peer: protocols.NewPeer(p, rw, StreamerSpec), localAddr: b.addr, - BzzAddr: nodeAddrById[p.ID()], + BzzAddr: addr, } b.streamer.delivery.overlay.On(bzzPeer) defer b.streamer.delivery.overlay.Off(bzzPeer) + // if len(addr) > b.index+1 && bytes.Equal(addrs[b.index+1], addr) { go func() { // each node Subscribes to each other's retrieveRequestStream // need to wait till an aynchronous process registers the peers in streamer.peers // that is used by Subscribe time.Sleep(1 * time.Second) - err := b.streamer.Subscribe(p.ID(), "SYNC", []byte{uint8(1)}, 0, 0, Top, true) - if err != nil { + if err := b.streamer.Subscribe(p.ID(), "SYNC", []byte{uint8(1)}, 0, 0, Top, false); err != nil { log.Warn("error in subscribe", "err", err) } }() + // } return b.streamer.Run(bzzPeer) }