From bef192d5faaf777468ed98b5fb0a3982dc16de12 Mon Sep 17 00:00:00 2001 From: Fabio Barone Date: Mon, 8 Oct 2018 17:25:11 -0500 Subject: [PATCH] swarm/swap: first accounting metrics --- p2p/protocols/accounting.go | 12 ++--- swarm/swap/api.go | 48 +++++++++++++++-- swarm/swap/swap.go | 102 ++++++++++++++++++++++++++++-------- swarm/swap/swap_test.go | 12 +++-- swarm/swap_test.go | 18 ++++++- 5 files changed, 154 insertions(+), 38 deletions(-) diff --git a/p2p/protocols/accounting.go b/p2p/protocols/accounting.go index a81fe5c318..3395588d62 100644 --- a/p2p/protocols/accounting.go +++ b/p2p/protocols/accounting.go @@ -27,9 +27,9 @@ type PriceOracle interface { type BalanceManager interface { //Credit is crediting the peer, charging local node - Credit(peer *Peer, amount uint64) error + Credit(peer *Peer, amount uint64, size uint32) error //Debit is crediting the local node, charging the remote peer - Debit(peer *Peer, amount uint64) error + Debit(peer *Peer, amount uint64, size uint32) error } type EntryDirection bool @@ -62,9 +62,9 @@ func (ah *AccountingHook) Send(peer *Peer, size uint32, msg interface{}) error { } direction, price := ah.PriceOracle.Price(size, msg) if direction == ChargeSender { - err = ah.BalanceManager.Debit(peer, price) + err = ah.BalanceManager.Debit(peer, price, size) } else { - err = ah.BalanceManager.Credit(peer, price) + err = ah.BalanceManager.Credit(peer, price, size) } return err } @@ -78,9 +78,9 @@ func (ah *AccountingHook) Receive(peer *Peer, size uint32, msg interface{}) erro } direction, price := ah.PriceOracle.Price(size, msg) if direction == ChargeReceiver { - err = ah.BalanceManager.Debit(peer, price) + err = ah.BalanceManager.Debit(peer, price, size) } else { - err = ah.BalanceManager.Credit(peer, price) + err = ah.BalanceManager.Credit(peer, price, size) } return err } diff --git a/swarm/swap/api.go b/swarm/swap/api.go index efdd94f092..70aecb876c 100644 --- a/swarm/swap/api.go +++ b/swarm/swap/api.go @@ -32,9 +32,21 @@ type API struct { swap *Swap } -//TODO: define metrics //Get metrics about swap for this node -type SwapMetrics struct { +//The current metrics are for accounted message types only +//(i.e. BytesTransferred is amount of bytes sent but only for a +//accounted message types) +type Metrics struct { + BalanceCredited uint64 + BalanceDebited uint64 + BytesCredited uint64 + BytesDebited uint64 + MsgCredited uint64 + MsgDebited uint64 + ChequesIssued uint64 + ChequesReceived uint64 + PeerDrops uint64 + SelfDrops uint64 } //Create a new API instance @@ -55,7 +67,7 @@ func (swapapi *API) BalanceWithPeer(ctx context.Context, peer enode.ID) (balance //Get the overall balance of the node //Iterates over all peers this node is having accounted interaction //and just adds up balances. -//It assumes that if a disfavorable balance is represented as a negative value +//It assumes that a disfavorable balance is represented as a negative value func (swapapi *API) Balance(ctx context.Context) (balance int64, err error) { balance = 0 for _, peerBalance := range swapapi.swap.balances { @@ -65,6 +77,32 @@ func (swapapi *API) Balance(ctx context.Context) (balance int64, err error) { } //Just return the Swap metrics -func (swapapi *API) GetSwapMetrics() (*SwapMetrics, error) { - return nil, nil +func (swapapi *API) GetSwapMetricsForPeer(ctx context.Context, peer enode.ID) (*Metrics, error) { + var ok bool + var metrics *Metrics + metrics, ok = swapapi.swap.metrics[peer] + if !ok { + return nil, ErrNoSuchPeerAccounting + } + return metrics, nil +} + +//Just return the Swap metrics +func (swapapi *API) GetSwapMetrics(ctx context.Context, peer enode.ID) (*Metrics, error) { + var metrics *Metrics + + for _, m := range swapapi.swap.metrics { + metrics.BalanceCredited += m.BalanceCredited + metrics.BalanceDebited += m.BalanceDebited + metrics.BytesCredited += m.BytesCredited + metrics.BytesDebited += m.BytesDebited + metrics.ChequesIssued += m.ChequesIssued + metrics.ChequesReceived += m.ChequesReceived + metrics.MsgCredited += m.MsgCredited + metrics.MsgDebited += m.MsgDebited + metrics.PeerDrops += m.PeerDrops + metrics.SelfDrops += m.SelfDrops + } + + return metrics, nil } diff --git a/swarm/swap/swap.go b/swarm/swap/swap.go index 5e96514b0c..65a9e63c51 100644 --- a/swarm/swap/swap.go +++ b/swarm/swap/swap.go @@ -40,7 +40,6 @@ var ( payAt = int64(-4096 * 10000000) // threshold that triggers payment {request} (bytes) dropAt = int64(-4096 * 12000000) // threshold that triggers disconnect (bytes) - ErrNotAccountedMsg = errors.New("Message does not need accounting") ErrInsufficientFunds = errors.New("Insufficient funds") ) @@ -49,30 +48,28 @@ var ( // A node maintains an individual balance with every peer // Only messages which have a price will be accounted for type Swap struct { - chequeManager *ChequeManager //cheque manager keeps track of issued cheques - stateStore state.Store //stateStore is needed in order to keep balances across sessions - lock sync.RWMutex //lock the balances - balances map[enode.ID]int64 //map of balances for each peer - protocol *Protocol //reference to the cheque exchange protocol + chequeManager *ChequeManager //cheque manager keeps track of issued cheques + stateStore state.Store //stateStore is needed in order to keep balances across sessions + lock sync.RWMutex //lock the balances + balances map[enode.ID]int64 //map of balances for each peer + metrics map[enode.ID]*Metrics //map of metrics for each peer + protocol *Protocol //reference to the cheque exchange protocol } //Credit us and debit remote -func (s *Swap) Credit(peer *protocols.Peer, amount uint64) (err error) { +func (s *Swap) Credit(peer *protocols.Peer, amount uint64, size uint32) (err error) { s.lock.Lock() defer s.lock.Unlock() - peerBalance := s.balances[peer.ID()] - if _, ok := s.balances[peer.ID()]; !ok { - s.stateStore.Get(peer.ID().String(), &peerBalance) - s.balances[peer.ID()] = peerBalance - } - //local node is being credited(in favor of local node), so the balance increases + s.loadState(peer) + s.balances[peer.ID()] += int64(amount) - peerBalance = s.balances[peer.ID()] + peerBalance := s.balances[peer.ID()] s.stateStore.Put(peer.ID().String(), &peerBalance) if float64(peerBalance) > math.Abs(float64(payAt)) { ctx := context.TODO() + //WARNING: Should we do this? Otherwise anyone could create a WantChequeMsg requiring us to handle special situations... err := s.wantCheque(ctx, peer.ID()) if err != nil { //TODO: special error handling, as at this point the accounting has been done @@ -81,29 +78,32 @@ func (s *Swap) Credit(peer *protocols.Peer, amount uint64) (err error) { } } if float64(peerBalance) > math.Abs(float64(dropAt)) { + s.metrics[peer.ID()].PeerDrops += 1 return ErrInsufficientFunds } + //TODO: size for metrics is currently misleading: should only account for size based messages(?) + s.updatePeerMetrics(peer, true, amount, size) + log.Debug(fmt.Sprintf("balance for peer %s: %s", peer.ID().String(), strconv.FormatInt(peerBalance, 10))) return err } //Debit us and credit remote -func (s *Swap) Debit(peer *protocols.Peer, amount uint64) (err error) { +func (s *Swap) Debit(peer *protocols.Peer, amount uint64, size uint32) (err error) { s.lock.Lock() defer s.lock.Unlock() - peerBalance := s.balances[peer.ID()] - if _, ok := s.balances[peer.ID()]; !ok { - s.stateStore.Get(peer.ID().String(), &peerBalance) - s.balances[peer.ID()] = peerBalance - } + s.loadState(peer) + //local node is being debited (in favor of remote peer), so its balance decreases s.balances[peer.ID()] -= int64(amount) - peerBalance = s.balances[peer.ID()] + peerBalance := s.balances[peer.ID()] s.stateStore.Put(peer.ID().String(), &peerBalance) + if peerBalance < payAt { ctx := context.TODO() err := s.issueCheque(ctx, peer.ID()) + s.metrics[peer.ID()].ChequesIssued += 1 if err != nil { //TODO: special error handling, as at this point the accounting has been done //but the cheque could not be sent? @@ -111,8 +111,12 @@ func (s *Swap) Debit(peer *protocols.Peer, amount uint64) (err error) { } } if peerBalance < dropAt { + s.metrics[peer.ID()].SelfDrops += 1 return ErrInsufficientFunds } + + s.updatePeerMetrics(peer, false, amount, size) + log.Debug(fmt.Sprintf("balance for peer %s: %s", peer.ID().String(), strconv.FormatInt(peerBalance, 10))) return nil } @@ -127,6 +131,61 @@ func (swap *Swap) GetPeerBalance(peer enode.ID) (int64, error) { return 0, errors.New("Peer not found") } +func (swap *Swap) GetPeerMetrics(peer enode.ID) (*Metrics, error) { + swap.lock.RLock() + defer swap.lock.RUnlock() + if p, ok := swap.metrics[peer]; ok { + return p, nil + } + return nil, errors.New("Peer not found") +} + +func (s *Swap) loadState(peer *protocols.Peer) { + var peerBalance int64 + var peerMetrics *Metrics + peerID := peer.ID() + if _, ok := s.metrics[peerID]; !ok { + s.stateStore.Get("metrics"+peerID.String(), &peerMetrics) + if peerMetrics == nil { + peerMetrics = &Metrics{ + BalanceCredited: 0, + BalanceDebited: 0, + BytesCredited: 0, + BytesDebited: 0, + MsgCredited: 0, + MsgDebited: 0, + ChequesIssued: 0, + ChequesReceived: 0, + PeerDrops: 0, + SelfDrops: 0, + } + s.metrics[peerID] = peerMetrics + } + } + if _, ok := s.balances[peerID]; !ok { + s.stateStore.Get(peerID.String(), &peerBalance) + s.balances[peerID] = peerBalance + } +} + +//local node is being credited(in favor of local node), so the balance increases +func (s *Swap) updatePeerMetrics(peer *protocols.Peer, credit bool, amount uint64, size uint32) { + + metrics := s.metrics[peer.ID()] + + if credit { + metrics.BalanceCredited += amount + metrics.BytesCredited += uint64(size) + metrics.MsgCredited += 1 + } else { + metrics.BalanceDebited += amount + metrics.BytesDebited += uint64(size) + metrics.MsgDebited += 1 + } + + s.stateStore.Put("metrics"+peer.ID().String(), metrics) +} + //Issue a cheque for the remote peer. Happens if we are indebted with the peer //and crossed the payment threshold func (s *Swap) issueCheque(ctx context.Context, id enode.ID) error { @@ -158,6 +217,7 @@ func New(stateStore state.Store) (swap *Swap) { chequeManager: NewChequeManager(stateStore), stateStore: stateStore, balances: make(map[enode.ID]int64), + metrics: make(map[enode.ID]*Metrics), protocol: NewProtocol(), } return diff --git a/swarm/swap/swap_test.go b/swarm/swap/swap_test.go index 297fbf798d..61e6955a83 100644 --- a/swarm/swap/swap_test.go +++ b/swarm/swap/swap_test.go @@ -156,10 +156,12 @@ func TestRepeatedBookings(t *testing.T) { defer os.RemoveAll(testDir) testPeer := newDummyPeer() + //size is irrelevant for this test + size := uint32(0) amount := mrand.Intn(100) cnt := 1 + mrand.Intn(10) for i := 0; i < cnt; i++ { - swap.Credit(testPeer.Peer.Peer, uint64(amount)) + swap.Credit(testPeer.Peer.Peer, uint64(amount), size) } expectedBalance := int64(cnt * amount) realBalance := swap.balances[testPeer.ID()] @@ -171,7 +173,7 @@ func TestRepeatedBookings(t *testing.T) { amount = mrand.Intn(100) cnt = 1 + mrand.Intn(10) for i := 0; i < cnt; i++ { - swap.Debit(testPeer2.Peer.Peer, uint64(amount)) + swap.Debit(testPeer2.Peer.Peer, uint64(amount), size) } expectedBalance = int64(0 - (cnt * amount)) realBalance = swap.balances[testPeer2.ID()] @@ -183,9 +185,9 @@ func TestRepeatedBookings(t *testing.T) { amount1 := mrand.Intn(100) amount2 := mrand.Intn(100) amount3 := mrand.Intn(100) - swap.Credit(testPeer2.Peer.Peer, uint64(amount1)) - swap.Credit(testPeer2.Peer.Peer, uint64(amount2)) - swap.Debit(testPeer2.Peer.Peer, uint64(amount3)) + swap.Credit(testPeer2.Peer.Peer, uint64(amount1), size) + swap.Credit(testPeer2.Peer.Peer, uint64(amount2), size) + swap.Debit(testPeer2.Peer.Peer, uint64(amount3), size) expectedBalance = expectedBalance + int64(amount1+amount2-amount3) realBalance = swap.balances[testPeer2.ID()] diff --git a/swarm/swap_test.go b/swarm/swap_test.go index 1a34626dcc..dc47ffd85e 100644 --- a/swarm/swap_test.go +++ b/swarm/swap_test.go @@ -166,13 +166,29 @@ func TestSwapNetworkSymmetricFileUpload(t *testing.T) { //get the peer's balance with this node balance, err := swarm.swap.GetPeerBalance(n) - fmt.Println(balance) if err == nil { subBalances[n] = balance log.Debug(fmt.Sprintf("Balance of node %s to node %s: %d", node.TerminalString(), n.TerminalString(), balance)) } else { log.Debug(fmt.Sprintf("Node %s has no balance with node %s", node.TerminalString(), n.TerminalString())) } + metrics, err := swarm.swap.GetPeerMetrics(n) + if err == nil && *printStats { + fmt.Println(fmt.Sprintf("********** Metrics for node %s with node %s: *************", node.TerminalString(), n.TerminalString())) + fmt.Println(fmt.Sprintf("Total units credited: %d", metrics.BalanceCredited)) + fmt.Println(fmt.Sprintf("Total units debited: %d", metrics.BalanceDebited)) + fmt.Println(fmt.Sprintf("Total bytes credited: %d", metrics.BytesCredited)) + fmt.Println(fmt.Sprintf("Total bytes debited: %d", metrics.BytesDebited)) + fmt.Println(fmt.Sprintf("Cheques issued: %d", metrics.ChequesIssued)) + fmt.Println(fmt.Sprintf("Cheques received: %d", metrics.ChequesReceived)) + fmt.Println(fmt.Sprintf("Number of messages credited: %d", metrics.MsgCredited)) + fmt.Println(fmt.Sprintf("Number of messages debited: %d", metrics.MsgDebited)) + fmt.Println(fmt.Sprintf("Peers dropped: %d", metrics.PeerDrops)) + fmt.Println(fmt.Sprintf("Number of times node dropped itself: %d", metrics.SelfDrops)) + } else { + //not all peers have metrics with every node, so probably can be ignored + log.Debug("Error getting metrics", "err", err) + } } //update the map for this node balancesMap[node] = subBalances