From 370b5dabbe135243f79de26cc1103752f9386fb0 Mon Sep 17 00:00:00 2001 From: Fabio Barone Date: Tue, 21 Aug 2018 17:25:33 -0500 Subject: [PATCH] swarm/swap: own network test file for swap; more tests and bug fixes --- swarm/network_test.go | 138 ----------------- swarm/swap/protocol.go | 1 - swarm/swap/swap.go | 22 +-- swarm/swap/swap_test.go | 56 +++++++ swarm/swap_test.go | 334 ++++++++++++++++++++++++++++++++++++++++ 5 files changed, 402 insertions(+), 149 deletions(-) create mode 100644 swarm/swap_test.go diff --git a/swarm/network_test.go b/swarm/network_test.go index 629095c936..4771676856 100644 --- a/swarm/network_test.go +++ b/swarm/network_test.go @@ -21,7 +21,6 @@ import ( "flag" "fmt" "io/ioutil" - "math/big" "math/rand" "os" "sync" @@ -377,143 +376,6 @@ 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 - } - } - }) - - balancesMap := make(map[discover.NodeID]map[discover.NodeID]*big.Int) - - for _, node := range sim.NodeIDs() { - item, ok := sim.NodeItem(node, bucketKeySwarm) - if !ok { - log.Error("No swarm") - return - } - swarm := item.(*Swarm) - - subBalances := make(map[discover.NodeID]*big.Int) - - for _, n := range sim.NodeIDs() { - if node == n { - continue - } - balance := swarm.swap.GetPeerBalance(n) - if balance != nil { - subBalances[n] = balance - log.Debug(fmt.Sprintf("Balance of node %s to node %s: %s", node.TerminalString(), n.TerminalString(), swarm.swap.GetPeerBalance(n).String())) - } else { - log.Debug(fmt.Sprintf("Node %s has no balance with node %s", node.TerminalString(), n.TerminalString())) - } - } - balancesMap[node] = subBalances - } - - if *printStats { - for k, v := range balancesMap { - fmt.Println(fmt.Sprintf("node %s balances:", k.TerminalString())) - for kk, vv := range v { - fmt.Println(fmt.Sprintf(".........with node %s: balance %s", kk.TerminalString(), vv.String())) - } - } - } - - for k, mapForK := range balancesMap { - for n, balanceKwithN := range mapForK { - for subK, mapForSubK := range balancesMap { - if n == subK { - log.Trace(fmt.Sprintf("balance of %s with %s: %s", k.TerminalString(), n.TerminalString(), balanceKwithN)) - log.Trace(fmt.Sprintf("balance of %s with %s: %s", n.TerminalString(), k.TerminalString(), mapForSubK[k])) - if balanceKwithN.CmpAbs(mapForSubK[k]) != 0 && balanceKwithN.Cmp(big.NewInt(0)) != 0 { - log.Error("Expected balances to be |abs| = 0 AND balance1 != 0, but they are not") - } - } - } - } - } - - 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/protocol.go b/swarm/swap/protocol.go index 9863ce20f4..78517e82e4 100644 --- a/swarm/swap/protocol.go +++ b/swarm/swap/protocol.go @@ -187,7 +187,6 @@ func (p *SwapProtocolPeer) handleSwapMsg(ctx context.Context, msg interface{}) e default: return fmt.Errorf("unknown message type: %T", msg) } - return nil } func (sp *SwapProtocolPeer) handleIssueChequeMsg(ctx context.Context, msg interface{}) (err error) { diff --git a/swarm/swap/swap.go b/swarm/swap/swap.go index 3334ab9905..f0f27443e1 100644 --- a/swarm/swap/swap.go +++ b/swarm/swap/swap.go @@ -191,16 +191,17 @@ func (sp *SwapPeer) checkAvailableFunds(ctx context.Context, msg interface{}, di price := accounted.GetMsgPrice() //local node is being credited (in its favor), so check upper limit if direction == CreditEntry { - //TODO: is there are a check needed here? + //TODO: is there a check needed here? //It should actually have been done on the client side, the debitor! //creditor could theoretically go over payAt, but if well done, //should have been checked on the client side so this shouldn't happen? - checkBalance := sp.balance.Add(sp.balance, price) - //(checkBalance *Int) Cmp(payAt) - // -1 if checkBalance < payAt - // 0 if checkBalance == payAt - // +1 if checkBalance > payAt - if checkBalance.Cmp(payAt) == 1 { + checkBalance := &big.Int{} + checkBalance.Add(sp.balance, price) + //(checkBalance *Int) CmpAbs(payAt) + // -1 if |checkBalance| < |payAt| + // 0 if |checkBalance| == |payAt| + // +1 if |checkBalance| > |payAt| + if checkBalance.CmpAbs(payAt) == 1 { return nil, ErrInsufficientFunds } } else if direction == DebitEntry { @@ -213,7 +214,8 @@ func (sp *SwapPeer) checkAvailableFunds(ctx context.Context, msg interface{}, di // -1 if checkBalance < dropAt // 0 if checkBalance == dropAt // +1 if checkBalance > dropAt - checkBalance := sp.balance.Sub(sp.balance, price.Abs(price)) + checkBalance := &big.Int{} + checkBalance.Sub(sp.balance, price.Abs(price)) if checkBalance.Cmp(dropAt) == -1 { return nil, ErrInsufficientFunds } @@ -234,10 +236,10 @@ func (sp *SwapPeer) AccountMsgForPeer(ctx context.Context, msg interface{}, pric if direction == CreditEntry { //NOTE: do we need to check for sufficient funds again? //operations are not atomic/transactional, so balance may have changed in the meanwhile! - sp.balance = sp.balance.Add(sp.balance, price) + sp.balance.Add(sp.balance, price) //local node is being debited (in favor of remote peer), so its balance decreases } else if direction == DebitEntry { - sp.balance = sp.balance.Sub(sp.balance, price) + sp.balance.Sub(sp.balance, price) } //TODO: save to store here? init store? sp.swapAccount.stateStore.Put(sp.storeID, sp.balance) diff --git a/swarm/swap/swap_test.go b/swarm/swap/swap_test.go index 99408d9092..403716125f 100644 --- a/swarm/swap/swap_test.go +++ b/swarm/swap/swap_test.go @@ -53,6 +53,7 @@ var testSpec = &protocols.Spec{ Messages: []interface{}{ testExceedsPayAtMsg{}, testExceedsDropAtMsg{}, + testCheapMsg{}, }, } @@ -73,6 +74,7 @@ func (d *dummyRW) ReadMsg() (p2p.Msg, error) { type testExceedsPayAtMsg struct{} type testExceedsDropAtMsg struct{} +type testCheapMsg struct{} func (tmsg *testExceedsPayAtMsg) GetMsgPrice() *big.Int { diff := &big.Int{} @@ -84,6 +86,10 @@ func (tmsg *testExceedsDropAtMsg) GetMsgPrice() *big.Int { return diff.Sub(dropAt, big.NewInt(1)) } +func (tmsg *testCheapMsg) GetMsgPrice() *big.Int { + return big.NewInt(100) +} + func init() { flag.Parse() @@ -141,6 +147,56 @@ func TestExceedsDropAt(t *testing.T) { } } +func TestSendCheapMessage(t *testing.T) { + swap, testDir := createTestSwap(t) + defer os.RemoveAll(testDir) + + testPeer := newDummyPeer() + sp := NewSwapPeer(testPeer, swap) + testBalance := big.NewInt(1234567890) + sp.balance = testBalance + + msg := &testCheapMsg{} + ctx := context.Background() + err := sp.Send(ctx, msg) + if err != nil { + t.Fatal("Unexpected error sending message") + } + + if sp.balance.Cmp(testBalance.Sub(testBalance, msg.GetMsgPrice())) != 0 { + t.Fatal(fmt.Sprintf("Unexpected balance value after sending cheap message test. Expected balance: %s, balance is: %s", + testBalance.Sub(testBalance, msg.GetMsgPrice()).String(), sp.balance.String())) + } +} + +func TestRestoreBalanceFromStateStore(t *testing.T) { + swap, testDir := createTestSwap(t) + defer os.RemoveAll(testDir) + + testPeer := newDummyPeer() + sp := NewSwapPeer(testPeer, swap) + testBalance := big.NewInt(1234567890) + sp.balance = big.NewInt(1234567890) + + //send a message, should trigger saving to stateStore + msg := &testCheapMsg{} + ctx := context.Background() + err := sp.Send(ctx, msg) + if err != nil { + log.Error(err.Error()) + t.Fatal("Unexpected error sending message") + } + + sp2 := NewSwapPeer(testPeer, swap) + + expectedBalance := &big.Int{} + expectedBalance.Sub(testBalance, msg.GetMsgPrice()) + if sp2.balance.Cmp(expectedBalance) != 0 { + t.Fatal(fmt.Sprintf("Unexpected balance value after sending cheap message test. Expected balance: %s, balance is: %s", + expectedBalance.String(), sp2.balance.String())) + } +} + func createTestSwap(t *testing.T) (*Swap, string) { dir, err := ioutil.TempDir("", "swap_test_store") if err != nil { diff --git a/swarm/swap_test.go b/swarm/swap_test.go new file mode 100644 index 0000000000..dfca84197f --- /dev/null +++ b/swarm/swap_test.go @@ -0,0 +1,334 @@ +package swarm + +import ( + "context" + "fmt" + "io/ioutil" + "math/big" + "math/rand" + "os" + "strconv" + "sync" + "testing" + "time" + + "github.com/ethereum/go-ethereum/crypto" + "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/swarm/api" + "github.com/ethereum/go-ethereum/swarm/log" + "github.com/ethereum/go-ethereum/swarm/network/simulation" + "github.com/ethereum/go-ethereum/swarm/storage" +) + +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 + } + } + }) + + balancesMap := make(map[discover.NodeID]map[discover.NodeID]*big.Int) + + for _, node := range sim.NodeIDs() { + item, ok := sim.NodeItem(node, bucketKeySwarm) + if !ok { + log.Error("No swarm") + return + } + swarm := item.(*Swarm) + + subBalances := make(map[discover.NodeID]*big.Int) + + for _, n := range sim.NodeIDs() { + if node == n { + continue + } + balance := swarm.swap.GetPeerBalance(n) + if balance != nil { + subBalances[n] = balance + log.Debug(fmt.Sprintf("Balance of node %s to node %s: %s", node.TerminalString(), n.TerminalString(), swarm.swap.GetPeerBalance(n).String())) + } else { + log.Debug(fmt.Sprintf("Node %s has no balance with node %s", node.TerminalString(), n.TerminalString())) + } + } + balancesMap[node] = subBalances + } + + if *printStats { + for k, v := range balancesMap { + fmt.Println(fmt.Sprintf("node %s balances:", k.TerminalString())) + for kk, vv := range v { + fmt.Println(fmt.Sprintf(".........with node %s: balance %s", kk.TerminalString(), vv.String())) + } + } + } + + for k, mapForK := range balancesMap { + for n, balanceKwithN := range mapForK { + for subK, mapForSubK := range balancesMap { + if n == subK { + log.Trace(fmt.Sprintf("balance of %s with %s: %s", k.TerminalString(), n.TerminalString(), balanceKwithN)) + log.Trace(fmt.Sprintf("balance of %s with %s: %s", n.TerminalString(), k.TerminalString(), mapForSubK[k])) + if balanceKwithN.CmpAbs(mapForSubK[k]) != 0 && balanceKwithN.Cmp(big.NewInt(0)) != 0 { + log.Error("Expected balances to be |abs| = 0 AND balance1 != 0, but they are not") + } + } + } + } + } + + if result.Error != nil { + t.Fatal(result.Error) + } + log.Debug("test terminated") +} + +func TestSwapNetworkAsymmetricFileUpload(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) + } + + const maxFileSize = 1024 * 1024 * 4 //1024 bytes * 1024 * 4 = 4MB + const minfileSize = 1024 + + //pseudo random algo to define if a node will upload or not + //if a bit is 0, do not upload + pseudoRandomNum := rand.Int63() + pseudoRandomBitMask := strconv.FormatInt(pseudoRandomNum, 2) + + 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 i, id := range nodeIDs { + //if the position in random num is 0, don't upload + if string(pseudoRandomBitMask[i]) != "0" { + size := rand.Intn(maxFileSize-minfileSize) + minfileSize + key, data, err := uploadRandomFileSize(sim.Service("swarm", id).(*Swarm), size) + 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 + } + } + }) + + balancesMap := make(map[discover.NodeID]map[discover.NodeID]*big.Int) + + for _, node := range sim.NodeIDs() { + item, ok := sim.NodeItem(node, bucketKeySwarm) + if !ok { + log.Error("No swarm") + return + } + swarm := item.(*Swarm) + + subBalances := make(map[discover.NodeID]*big.Int) + + for _, n := range sim.NodeIDs() { + if node == n { + continue + } + balance := swarm.swap.GetPeerBalance(n) + if balance != nil { + subBalances[n] = balance + log.Debug(fmt.Sprintf("Balance of node %s to node %s: %s", node.TerminalString(), n.TerminalString(), swarm.swap.GetPeerBalance(n).String())) + } else { + log.Debug(fmt.Sprintf("Node %s has no balance with node %s", node.TerminalString(), n.TerminalString())) + } + } + balancesMap[node] = subBalances + } + + if *printStats { + for k, v := range balancesMap { + fmt.Println(fmt.Sprintf("node %s balances:", k.TerminalString())) + for kk, vv := range v { + fmt.Println(fmt.Sprintf(".........with node %s: balance %s", kk.TerminalString(), vv.String())) + } + } + } + + /* + TODO: Should balances in this case also be symmetric? + + for k, mapForK := range balancesMap { + for n, balanceKwithN := range mapForK { + for subK, mapForSubK := range balancesMap { + if n == subK { + log.Trace(fmt.Sprintf("balance of %s with %s: %s", k.TerminalString(), n.TerminalString(), balanceKwithN)) + log.Trace(fmt.Sprintf("balance of %s with %s: %s", n.TerminalString(), k.TerminalString(), mapForSubK[k])) + if balanceKwithN.CmpAbs(mapForSubK[k]) != 0 && balanceKwithN.Cmp(big.NewInt(0)) != 0 { + log.Error("Expected balances to be |abs| = 0 AND balance1 != 0, but they are not") + } + } + } + } + } + */ + + 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 uploadRandomFileSize(swarm *Swarm, size int) (storage.Address, string, error) { + b := make([]byte, size) + _, err := rand.Read(b) + if err != nil { + return nil, "", err + } + // uniqueness is very certain. + data := fmt.Sprintf("test content %s %x", time.Now().Round(0), b) + ctx := context.TODO() + k, wait, err := swarm.api.Put(ctx, data, "text/plain", false) + if err != nil { + return nil, "", err + } + if wait != nil { + err = wait(ctx) + } + return k, data, err +}