mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-19 10:22:23 +00:00
swarm/swap: own network test file for swap; more tests and bug fixes
This commit is contained in:
parent
6f27cf5ac4
commit
370b5dabbe
5 changed files with 402 additions and 149 deletions
|
|
@ -21,7 +21,6 @@ import (
|
||||||
"flag"
|
"flag"
|
||||||
"fmt"
|
"fmt"
|
||||||
"io/ioutil"
|
"io/ioutil"
|
||||||
"math/big"
|
|
||||||
"math/rand"
|
"math/rand"
|
||||||
"os"
|
"os"
|
||||||
"sync"
|
"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
|
// uploadFile, uploads a short file to the swarm instance
|
||||||
// using the api.Put method.
|
// using the api.Put method.
|
||||||
func uploadFile(swarm *Swarm) (storage.Address, string, error) {
|
func uploadFile(swarm *Swarm) (storage.Address, string, error) {
|
||||||
|
|
|
||||||
|
|
@ -187,7 +187,6 @@ func (p *SwapProtocolPeer) handleSwapMsg(ctx context.Context, msg interface{}) e
|
||||||
default:
|
default:
|
||||||
return fmt.Errorf("unknown message type: %T", msg)
|
return fmt.Errorf("unknown message type: %T", msg)
|
||||||
}
|
}
|
||||||
return nil
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (sp *SwapProtocolPeer) handleIssueChequeMsg(ctx context.Context, msg interface{}) (err error) {
|
func (sp *SwapProtocolPeer) handleIssueChequeMsg(ctx context.Context, msg interface{}) (err error) {
|
||||||
|
|
|
||||||
|
|
@ -191,16 +191,17 @@ func (sp *SwapPeer) checkAvailableFunds(ctx context.Context, msg interface{}, di
|
||||||
price := accounted.GetMsgPrice()
|
price := accounted.GetMsgPrice()
|
||||||
//local node is being credited (in its favor), so check upper limit
|
//local node is being credited (in its favor), so check upper limit
|
||||||
if direction == CreditEntry {
|
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!
|
//It should actually have been done on the client side, the debitor!
|
||||||
//creditor could theoretically go over payAt, but if well done,
|
//creditor could theoretically go over payAt, but if well done,
|
||||||
//should have been checked on the client side so this shouldn't happen?
|
//should have been checked on the client side so this shouldn't happen?
|
||||||
checkBalance := sp.balance.Add(sp.balance, price)
|
checkBalance := &big.Int{}
|
||||||
//(checkBalance *Int) Cmp(payAt)
|
checkBalance.Add(sp.balance, price)
|
||||||
// -1 if checkBalance < payAt
|
//(checkBalance *Int) CmpAbs(payAt)
|
||||||
// 0 if checkBalance == payAt
|
// -1 if |checkBalance| < |payAt|
|
||||||
// +1 if checkBalance > payAt
|
// 0 if |checkBalance| == |payAt|
|
||||||
if checkBalance.Cmp(payAt) == 1 {
|
// +1 if |checkBalance| > |payAt|
|
||||||
|
if checkBalance.CmpAbs(payAt) == 1 {
|
||||||
return nil, ErrInsufficientFunds
|
return nil, ErrInsufficientFunds
|
||||||
}
|
}
|
||||||
} else if direction == DebitEntry {
|
} else if direction == DebitEntry {
|
||||||
|
|
@ -213,7 +214,8 @@ func (sp *SwapPeer) checkAvailableFunds(ctx context.Context, msg interface{}, di
|
||||||
// -1 if checkBalance < dropAt
|
// -1 if checkBalance < dropAt
|
||||||
// 0 if checkBalance == dropAt
|
// 0 if checkBalance == dropAt
|
||||||
// +1 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 {
|
if checkBalance.Cmp(dropAt) == -1 {
|
||||||
return nil, ErrInsufficientFunds
|
return nil, ErrInsufficientFunds
|
||||||
}
|
}
|
||||||
|
|
@ -234,10 +236,10 @@ func (sp *SwapPeer) AccountMsgForPeer(ctx context.Context, msg interface{}, pric
|
||||||
if direction == CreditEntry {
|
if direction == CreditEntry {
|
||||||
//NOTE: do we need to check for sufficient funds again?
|
//NOTE: do we need to check for sufficient funds again?
|
||||||
//operations are not atomic/transactional, so balance may have changed in the meanwhile!
|
//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
|
//local node is being debited (in favor of remote peer), so its balance decreases
|
||||||
} else if direction == DebitEntry {
|
} 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?
|
//TODO: save to store here? init store?
|
||||||
sp.swapAccount.stateStore.Put(sp.storeID, sp.balance)
|
sp.swapAccount.stateStore.Put(sp.storeID, sp.balance)
|
||||||
|
|
|
||||||
|
|
@ -53,6 +53,7 @@ var testSpec = &protocols.Spec{
|
||||||
Messages: []interface{}{
|
Messages: []interface{}{
|
||||||
testExceedsPayAtMsg{},
|
testExceedsPayAtMsg{},
|
||||||
testExceedsDropAtMsg{},
|
testExceedsDropAtMsg{},
|
||||||
|
testCheapMsg{},
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -73,6 +74,7 @@ func (d *dummyRW) ReadMsg() (p2p.Msg, error) {
|
||||||
|
|
||||||
type testExceedsPayAtMsg struct{}
|
type testExceedsPayAtMsg struct{}
|
||||||
type testExceedsDropAtMsg struct{}
|
type testExceedsDropAtMsg struct{}
|
||||||
|
type testCheapMsg struct{}
|
||||||
|
|
||||||
func (tmsg *testExceedsPayAtMsg) GetMsgPrice() *big.Int {
|
func (tmsg *testExceedsPayAtMsg) GetMsgPrice() *big.Int {
|
||||||
diff := &big.Int{}
|
diff := &big.Int{}
|
||||||
|
|
@ -84,6 +86,10 @@ func (tmsg *testExceedsDropAtMsg) GetMsgPrice() *big.Int {
|
||||||
return diff.Sub(dropAt, big.NewInt(1))
|
return diff.Sub(dropAt, big.NewInt(1))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (tmsg *testCheapMsg) GetMsgPrice() *big.Int {
|
||||||
|
return big.NewInt(100)
|
||||||
|
}
|
||||||
|
|
||||||
func init() {
|
func init() {
|
||||||
flag.Parse()
|
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) {
|
func createTestSwap(t *testing.T) (*Swap, string) {
|
||||||
dir, err := ioutil.TempDir("", "swap_test_store")
|
dir, err := ioutil.TempDir("", "swap_test_store")
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|
|
||||||
334
swarm/swap_test.go
Normal file
334
swarm/swap_test.go
Normal file
|
|
@ -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
|
||||||
|
}
|
||||||
Loading…
Reference in a new issue