diff --git a/swarm/network/stream/common_test.go b/swarm/network/stream/common_test.go index 491dc9fd58..0559a942ea 100644 --- a/swarm/network/stream/common_test.go +++ b/swarm/network/stream/common_test.go @@ -109,7 +109,7 @@ func newStreamerTester(t *testing.T) (*p2ptest.ProtocolTester, *Registry, *stora db := storage.NewDBAPI(localStore) delivery := NewDelivery(to, db) - streamer := NewRegistry(addr, delivery, db, state.NewInmemoryStore(), nil) + streamer := NewRegistry(addr, delivery, db, state.NewInmemoryStore(), nil, nil) teardown := func() { streamer.Close() removeDataDir() diff --git a/swarm/network/stream/delivery_test.go b/swarm/network/stream/delivery_test.go index 972cc859a3..7b2c1f9dcf 100644 --- a/swarm/network/stream/delivery_test.go +++ b/swarm/network/stream/delivery_test.go @@ -328,7 +328,7 @@ func testDeliveryFromNodes(t *testing.T, nodes, conns, chunkCount int, skipCheck kad := network.NewKademlia(addr.Over(), network.NewKadParams()) delivery := NewDelivery(kad, db) - r := NewRegistry(addr, delivery, db, state.NewInmemoryStore(), &RegistryOptions{ + r := NewRegistry(addr, delivery, db, state.NewInmemoryStore(), nil, &RegistryOptions{ SkipCheck: skipCheck, }) bucket.Store(bucketKeyRegistry, r) @@ -515,7 +515,7 @@ func benchmarkDeliveryFromNodes(b *testing.B, nodes, conns, chunkCount int, skip kad := network.NewKademlia(addr.Over(), network.NewKadParams()) delivery := NewDelivery(kad, db) - r := NewRegistry(addr, delivery, db, state.NewInmemoryStore(), &RegistryOptions{ + r := NewRegistry(addr, delivery, db, state.NewInmemoryStore(), nil, &RegistryOptions{ SkipCheck: skipCheck, DoSync: true, SyncUpdateDelay: 0, diff --git a/swarm/network/stream/intervals_test.go b/swarm/network/stream/intervals_test.go index f4294134b9..02626ab65b 100644 --- a/swarm/network/stream/intervals_test.go +++ b/swarm/network/stream/intervals_test.go @@ -74,7 +74,7 @@ func testIntervals(t *testing.T, live bool, history *Range, skipCheck bool) { kad := network.NewKademlia(addr.Over(), network.NewKadParams()) delivery := NewDelivery(kad, db) - r := NewRegistry(addr, delivery, db, state.NewInmemoryStore(), &RegistryOptions{ + r := NewRegistry(addr, delivery, db, state.NewInmemoryStore(), nil, &RegistryOptions{ SkipCheck: skipCheck, }) bucket.Store(bucketKeyRegistry, r) diff --git a/swarm/network/stream/snapshot_retrieval_test.go b/swarm/network/stream/snapshot_retrieval_test.go index 4ff947b215..32c1f63c23 100644 --- a/swarm/network/stream/snapshot_retrieval_test.go +++ b/swarm/network/stream/snapshot_retrieval_test.go @@ -133,7 +133,7 @@ func runFileRetrievalTest(nodeCount int) error { kad := network.NewKademlia(addr.Over(), network.NewKadParams()) delivery := NewDelivery(kad, db) - r := NewRegistry(addr, delivery, db, state.NewInmemoryStore(), &RegistryOptions{ + r := NewRegistry(addr, delivery, db, state.NewInmemoryStore(), nil, &RegistryOptions{ DoSync: true, SyncUpdateDelay: 3 * time.Second, }) @@ -276,7 +276,7 @@ func runRetrievalTest(chunkCount int, nodeCount int) error { kad := network.NewKademlia(addr.Over(), network.NewKadParams()) delivery := NewDelivery(kad, db) - r := NewRegistry(addr, delivery, db, state.NewInmemoryStore(), &RegistryOptions{ + r := NewRegistry(addr, delivery, db, state.NewInmemoryStore(), nil, &RegistryOptions{ DoSync: true, SyncUpdateDelay: 0, }) diff --git a/swarm/network/stream/snapshot_sync_test.go b/swarm/network/stream/snapshot_sync_test.go index 6acab50af4..1b264011e7 100644 --- a/swarm/network/stream/snapshot_sync_test.go +++ b/swarm/network/stream/snapshot_sync_test.go @@ -139,7 +139,7 @@ func testSyncingViaGlobalSync(t *testing.T, chunkCount int, nodeCount int) { kad := network.NewKademlia(addr.Over(), network.NewKadParams()) delivery := NewDelivery(kad, db) - r := NewRegistry(addr, delivery, db, state.NewInmemoryStore(), &RegistryOptions{ + r := NewRegistry(addr, delivery, db, state.NewInmemoryStore(), nil, &RegistryOptions{ DoSync: true, SyncUpdateDelay: 3 * time.Second, }) @@ -297,7 +297,7 @@ func testSyncingViaDirectSubscribe(chunkCount int, nodeCount int) error { kad := network.NewKademlia(addr.Over(), network.NewKadParams()) delivery := NewDelivery(kad, db) - r := NewRegistry(addr, delivery, db, state.NewInmemoryStore(), nil) + r := NewRegistry(addr, delivery, db, state.NewInmemoryStore(), nil, nil) bucket.Store(bucketKeyRegistry, r) fileStore := storage.NewFileStore(storage.NewNetStore(localStore, nil), storage.NewFileStoreParams()) diff --git a/swarm/network/stream/stream.go b/swarm/network/stream/stream.go index d88303b419..70cb238974 100644 --- a/swarm/network/stream/stream.go +++ b/swarm/network/stream/stream.go @@ -75,7 +75,7 @@ type RegistryOptions struct { } // NewRegistry is Streamer constructor -func NewRegistry(addr *network.BzzAddr, delivery *Delivery, db *storage.DBAPI, intervalsStore state.Store, options *RegistryOptions) *Registry { +func NewRegistry(addr *network.BzzAddr, delivery *Delivery, db *storage.DBAPI, intervalsStore state.Store, swap *swap.Swap, options *RegistryOptions) *Registry { if options == nil { options = &RegistryOptions{} } @@ -103,11 +103,7 @@ func NewRegistry(addr *network.BzzAddr, delivery *Delivery, db *storage.DBAPI, i RegisterSwarmSyncerServer(streamer, db) RegisterSwarmSyncerClient(streamer, db) - var err error - streamer.swap, err = swap.NewSwap(swap.NewDefaultSwapParams().Params) - if err != nil { - log.Error(err.Error()) - } + streamer.swap = swap if options.DoSync { // latestIntC function ensures that diff --git a/swarm/network/stream/syncer_test.go b/swarm/network/stream/syncer_test.go index f72aa34441..e9ee4df458 100644 --- a/swarm/network/stream/syncer_test.go +++ b/swarm/network/stream/syncer_test.go @@ -108,7 +108,7 @@ func testSyncBetweenNodes(t *testing.T, nodes, conns, chunkCount int, skipCheck delivery := NewDelivery(kad, db) bucket.Store(bucketKeyDelivery, delivery) - r := NewRegistry(addr, delivery, db, state.NewInmemoryStore(), &RegistryOptions{ + r := NewRegistry(addr, delivery, db, state.NewInmemoryStore(), nil, &RegistryOptions{ SkipCheck: skipCheck, }) diff --git a/swarm/network_test.go b/swarm/network_test.go index 176c635d82..6e3c806773 100644 --- a/swarm/network_test.go +++ b/swarm/network_test.go @@ -41,9 +41,11 @@ import ( ) var ( - loglevel = flag.Int("loglevel", 2, "verbosity of logs") - longrunning = flag.Bool("longrunning", false, "do run long-running tests") - waitKademlia = flag.Bool("waitkademlia", false, "wait for healthy kademlia before checking files availability") + loglevel = flag.Int("loglevel", 2, "verbosity of logs") + longrunning = flag.Bool("longrunning", false, "do run long-running tests") + waitKademlia = flag.Bool("waitkademlia", false, "wait for healthy kademlia before checking files availability") + bucketKeySwap = simulation.BucketKey("swap") + bucketKeySwarm = simulation.BucketKey("swarm") ) func init() { @@ -373,6 +375,113 @@ func testSwarmNetwork(t *testing.T, o *testSwarmNetworkOptions, steps ...testSwa } } +func TestSwapNetworkSymmetricFileUpload(t *testing.T) { + nodeCount := 16 + + sim := simulation.New(map[string]simulation.ServiceFunc{ + "swarm": func(ctx *adapters.ServiceContext, bucket *sync.Map) (s node.Service, cleanup func(), err error) { + config := api.NewConfig() + + dir, err := ioutil.TempDir("", "swap-network-test-node") + if err != nil { + return nil, nil, err + } + cleanup = func() { + err := os.RemoveAll(dir) + if err != nil { + log.Error("cleaning up swarm temp dir", "err", err) + } + } + + config.Path = dir + + privkey, err := crypto.GenerateKey() + if err != nil { + return nil, cleanup, err + } + + config.Init(privkey) + + swarm, err := NewSwarm(config, nil) + if err != nil { + return nil, cleanup, err + } + bucket.Store(bucketKeySwarm, swarm) + log.Info("new swarm", "bzzKey", config.BzzKey, "baseAddr", fmt.Sprintf("%x", swarm.bzz.BaseAddr())) + return swarm, cleanup, nil + }, + }) + defer sim.Close() + + ctx := context.Background() + files := make([]file, 0) + + var checkStatusM sync.Map + var nodeStatusM sync.Map + var totalFoundCount uint64 + + _, err := sim.AddNodesAndConnectChain(nodeCount) + if err != nil { + t.Fatal(err) + } + + result := sim.Run(ctx, func(ctx context.Context, sim *simulation.Simulation) error { + nodeIDs := sim.UpNodeIDs() + shuffle(len(nodeIDs), func(i, j int) { + nodeIDs[i], nodeIDs[j] = nodeIDs[j], nodeIDs[i] + }) + for _, id := range nodeIDs { + key, data, err := uploadFile(sim.Service("swarm", id).(*Swarm)) + if err != nil { + return err + } + log.Trace("file uploaded", "node", id, "key", key.String()) + files = append(files, file{ + addr: key, + data: data, + nodeID: id, + }) + } + + if _, err := sim.WaitTillHealthy(ctx, 2); err != nil { + return err + } + + // File retrieval check is repeated until all uploaded files are retrieved from all nodes + // or until the timeout is reached. + for { + if retrieve(sim, files, &checkStatusM, &nodeStatusM, &totalFoundCount) == 0 { + return nil + } + } + }) + + for _, node := range sim.NodeIDs() { + item, ok := sim.NodeItem(node, bucketKeySwarm) + if !ok { + log.Error("No swarm") + return + } + swarm := item.(*Swarm) + + for _, n := range sim.NodeIDs() { + if node == n { + continue + } + if swarm.swap.GetPeerBalance(n) != nil { + log.Error(fmt.Sprintf("Balance of node %s to node %s: %s", node.TerminalString(), n.TerminalString(), swarm.swap.GetPeerBalance(n).String())) + } else { + log.Error(fmt.Sprintf("Node %s has no balance with node %s", node.TerminalString(), n.TerminalString())) + } + } + } + + if result.Error != nil { + t.Fatal(result.Error) + } + log.Debug("test terminated") +} + // uploadFile, uploads a short file to the swarm instance // using the api.Put method. func uploadFile(swarm *Swarm) (storage.Address, string, error) { diff --git a/swarm/swap/api.go b/swarm/swap/api.go index 632a62433f..c731cb972b 100644 --- a/swarm/swap/api.go +++ b/swarm/swap/api.go @@ -18,9 +18,11 @@ package swap import ( "context" + "math/big" "github.com/ethereum/go-ethereum/common" "github.com/ethereum/go-ethereum/common/hexutil" + "github.com/ethereum/go-ethereum/p2p/discover" ) // Wrapper for receiving pss messages when using the pss API @@ -32,6 +34,7 @@ type APIMsg struct { // Additional public methods accessible through API for pss type API struct { *SwapProtocol + *Swap } type SwapMetrics struct { @@ -44,8 +47,8 @@ func NewAPI(swap *SwapProtocol) *API { return &API{SwapProtocol: swap} } -func (swapapi *API) Balance(ctx context.Context) (balance int, err error) { - balance = 0 +func (swapapi *API) Balance(ctx context.Context, peer discover.NodeID) (balance *big.Int, err error) { + balance = big.NewInt(0) err = nil return } diff --git a/swarm/swap/swap.go b/swarm/swap/swap.go index b67c53768e..62b8b2bdd0 100644 --- a/swarm/swap/swap.go +++ b/swarm/swap/swap.go @@ -34,6 +34,7 @@ import ( "github.com/ethereum/go-ethereum/crypto" "github.com/ethereum/go-ethereum/p2p/discover" "github.com/ethereum/go-ethereum/swarm/log" + "github.com/ethereum/go-ethereum/swarm/state" whisper "github.com/ethereum/go-ethereum/whisper/whisperv5" ) @@ -64,9 +65,10 @@ const ( // Swift Automatic Payments // a peer to peer micropayment system type Swap struct { - lock sync.RWMutex - peers map[discover.NodeID]*swapPeer - local *Params // local peer's swap parameters + stateStore state.Store + lock sync.RWMutex + peers map[discover.NodeID]*swapPeer + local *Params // local peer's swap parameters } type EntryDirection bool @@ -83,11 +85,14 @@ type SwapAccountedMsgType interface { func (swap *Swap) AccountForMsg(ctx context.Context, msg interface{}, peer discover.NodeID) error { if accounted, ok := msg.(SwapAccountedMsgType); ok { if _, exists := swap.peers[peer]; !exists { + balance := big.NewInt(0) + swap.stateStore.Get(peer.String()[:24]+"-swap", &balance) swap.lock.Lock() swap.peers[peer] = &swapPeer{ peer: peer, swapAccount: swap, - balance: big.NewInt(0), + balance: balance, + storeID: peer.String()[:24] + "-swap", } swap.lock.Unlock() } @@ -98,6 +103,13 @@ func (swap *Swap) AccountForMsg(ctx context.Context, msg interface{}, peer disco return nil } +func (swap *Swap) GetPeerBalance(peer discover.NodeID) *big.Int { + if p, ok := swap.peers[peer]; ok { + return p.balance + } + return nil +} + // Profile - public swap profile // public parameters for SWAP, serializable config struct passed in handshake type Profile struct { @@ -161,6 +173,7 @@ type swapPeer struct { peer discover.NodeID swapAccount *Swap balance *big.Int + storeID string } func (sp *swapPeer) AccountMsgForPeer(price *big.Int, direction EntryDirection) { @@ -173,21 +186,24 @@ func (sp *swapPeer) AccountMsgForPeer(price *big.Int, direction EntryDirection) } else if direction == DebitEntry { sp.balance = sp.balance.Sub(sp.balance, price) } + //TODO: save to store here? init store? + sp.swapAccount.stateStore.Put(sp.storeID, sp.balance) if sp.balance.Cmp(payAt) > -1 { //TODO: Issue Cheque } if sp.balance.Cmp(dropAt) < 0 { //TODO: Drop peer } - log.Error(fmt.Sprintf("balance for peer %s: %s", sp.peer, sp.balance.String())) + log.Debug(fmt.Sprintf("balance for peer %s: %s", sp.peer, sp.balance.String())) } // New - swap constructor -func NewSwap(local *Params) (swap *Swap, err error) { +func NewSwap(local *Params, stateStore state.Store) (swap *Swap, err error) { swap = &Swap{ - local: local, - peers: make(map[discover.NodeID]*swapPeer), + local: local, + stateStore: stateStore, + peers: make(map[discover.NodeID]*swapPeer), } //swap.SetParams(local) diff --git a/swarm/swap/swap_test.go b/swarm/swap/swap_test.go index 2be0c5013f..060161b897 100644 --- a/swarm/swap/swap_test.go +++ b/swarm/swap/swap_test.go @@ -17,16 +17,28 @@ package swap import ( + "context" + "crypto/rand" "flag" "fmt" + "io/ioutil" "os" "path/filepath" "sync" + "sync/atomic" "testing" + "time" + "github.com/ethereum/go-ethereum/crypto" "github.com/ethereum/go-ethereum/log" "github.com/ethereum/go-ethereum/node" + "github.com/ethereum/go-ethereum/p2p/discover" + "github.com/ethereum/go-ethereum/p2p/simulations/adapters" "github.com/ethereum/go-ethereum/rpc" + "github.com/ethereum/go-ethereum/swarm" + "github.com/ethereum/go-ethereum/swarm/api" + "github.com/ethereum/go-ethereum/swarm/network/simulation" + "github.com/ethereum/go-ethereum/swarm/storage" colorable "github.com/mattn/go-colorable" ) @@ -91,6 +103,12 @@ func TestSwapProtocol(t *testing.T) { }, nil } + streamersvc := func(ctx *node.ServiceContext) (node.Service, erro) { + return &stream.API{ + streamer: NewRegistry, + }, nil + } + // register adds the service to the services the servicenode starts when started err = stack_one.Register(swapsvc) if err != nil { diff --git a/swarm/swarm.go b/swarm/swarm.go index 95f657bafb..30db6e1528 100644 --- a/swarm/swarm.go +++ b/swarm/swarm.go @@ -180,7 +180,12 @@ func NewSwarm(config *api.Config, mockStore *mock.NodeStore) (self *Swarm, err e ) delivery := stream.NewDelivery(to, db) - self.streamer = stream.NewRegistry(addr, delivery, db, stateStore, &stream.RegistryOptions{ + self.swap, err = swap.NewSwap(swap.NewDefaultSwapParams().Params, stateStore) + if err != nil { + return nil, err + } + + self.streamer = stream.NewRegistry(addr, delivery, db, stateStore, self.swap, &stream.RegistryOptions{ SkipCheck: config.DeliverySkipCheck, DoSync: config.SyncEnabled, DoRetrieve: true, @@ -208,12 +213,6 @@ func NewSwarm(config *api.Config, mockStore *mock.NodeStore) (self *Swarm, err e self.bzz = network.NewBzz(bzzconfig, to, stateStore, stream.Spec, self.streamer.Run) - /* - self.swap, err = swap.NewSwap(config.Swap.Params) - if err != nil { - return nil, err - } - */ // Pss = postal service over swarm (devp2p over bzz) self.ps, err = pss.NewPss(to, config.Pss) if err != nil {