les: don't use me, this, self as receiver names

This commit is contained in:
Delweng Zheng 2018-06-03 16:01:43 +08:00
parent 0ad4915ff7
commit ec75743567
3 changed files with 101 additions and 101 deletions

View file

@ -89,16 +89,16 @@ func NewClientManager(rcTarget, maxSimReq, maxRcSum uint64) *ClientManager {
return cm
}
func (self *ClientManager) Stop() {
self.lock.Lock()
defer self.lock.Unlock()
func (cm *ClientManager) Stop() {
cm.lock.Lock()
defer cm.lock.Unlock()
// signal any waiting accept routines to return false
self.nodes = make(map[*cmNode]struct{})
close(self.resumeQueue)
cm.nodes = make(map[*cmNode]struct{})
close(cm.resumeQueue)
}
func (self *ClientManager) addNode(cnode *ClientNode) *cmNode {
func (cm *ClientManager) addNode(cnode *ClientNode) *cmNode {
time := mclock.Now()
node := &cmNode{
node: cnode,
@ -106,28 +106,28 @@ func (self *ClientManager) addNode(cnode *ClientNode) *cmNode {
finishRecharge: time,
rcWeight: 1,
}
self.lock.Lock()
defer self.lock.Unlock()
cm.lock.Lock()
defer cm.lock.Unlock()
self.nodes[node] = struct{}{}
self.update(mclock.Now())
cm.nodes[node] = struct{}{}
cm.update(mclock.Now())
return node
}
func (self *ClientManager) removeNode(node *cmNode) {
self.lock.Lock()
defer self.lock.Unlock()
func (cm *ClientManager) removeNode(node *cmNode) {
cm.lock.Lock()
defer cm.lock.Unlock()
time := mclock.Now()
self.stop(node, time)
delete(self.nodes, node)
self.update(time)
cm.stop(node, time)
delete(cm.nodes, node)
cm.update(time)
}
// recalc sumWeight
func (self *ClientManager) updateNodes(time mclock.AbsTime) (rce bool) {
func (cm *ClientManager) updateNodes(time mclock.AbsTime) (rce bool) {
var sumWeight, rcSum uint64
for node := range self.nodes {
for node := range cm.nodes {
rc := node.recharging
node.update(time)
if rc && !node.recharging {
@ -138,44 +138,44 @@ func (self *ClientManager) updateNodes(time mclock.AbsTime) (rce bool) {
}
rcSum += uint64(node.rcValue)
}
self.sumWeight = sumWeight
self.rcSumValue = rcSum
cm.sumWeight = sumWeight
cm.rcSumValue = rcSum
return
}
func (self *ClientManager) update(time mclock.AbsTime) {
func (cm *ClientManager) update(time mclock.AbsTime) {
for {
firstTime := time
for node := range self.nodes {
for node := range cm.nodes {
if node.recharging && node.finishRecharge < firstTime {
firstTime = node.finishRecharge
}
}
if self.updateNodes(firstTime) {
for node := range self.nodes {
if cm.updateNodes(firstTime) {
for node := range cm.nodes {
if node.recharging {
node.set(node.serving, self.simReqCnt, self.sumWeight)
node.set(node.serving, cm.simReqCnt, cm.sumWeight)
}
}
} else {
self.time = time
cm.time = time
return
}
}
}
func (self *ClientManager) canStartReq() bool {
return self.simReqCnt < self.maxSimReq && self.rcSumValue < self.maxRcSum
func (cm *ClientManager) canStartReq() bool {
return cm.simReqCnt < cm.maxSimReq && cm.rcSumValue < cm.maxRcSum
}
func (self *ClientManager) queueProc() {
for rc := range self.resumeQueue {
func (cm *ClientManager) queueProc() {
for rc := range cm.resumeQueue {
for {
time.Sleep(time.Millisecond * 10)
self.lock.Lock()
self.update(mclock.Now())
cs := self.canStartReq()
self.lock.Unlock()
cm.lock.Lock()
cm.update(mclock.Now())
cs := cm.canStartReq()
cm.lock.Unlock()
if cs {
break
}
@ -184,41 +184,41 @@ func (self *ClientManager) queueProc() {
}
}
func (self *ClientManager) accept(node *cmNode, time mclock.AbsTime) bool {
self.lock.Lock()
defer self.lock.Unlock()
func (cm *ClientManager) accept(node *cmNode, time mclock.AbsTime) bool {
cm.lock.Lock()
defer cm.lock.Unlock()
self.update(time)
if !self.canStartReq() {
cm.update(time)
if !cm.canStartReq() {
resume := make(chan bool)
self.lock.Unlock()
self.resumeQueue <- resume
cm.lock.Unlock()
cm.resumeQueue <- resume
<-resume
self.lock.Lock()
if _, ok := self.nodes[node]; !ok {
cm.lock.Lock()
if _, ok := cm.nodes[node]; !ok {
return false // reject if node has been removed or manager has been stopped
}
}
self.simReqCnt++
node.set(true, self.simReqCnt, self.sumWeight)
cm.simReqCnt++
node.set(true, cm.simReqCnt, cm.sumWeight)
node.startValue = node.rcValue
self.update(self.time)
cm.update(cm.time)
return true
}
func (self *ClientManager) stop(node *cmNode, time mclock.AbsTime) {
func (cm *ClientManager) stop(node *cmNode, time mclock.AbsTime) {
if node.serving {
self.update(time)
self.simReqCnt--
node.set(false, self.simReqCnt, self.sumWeight)
self.update(time)
cm.update(time)
cm.simReqCnt--
node.set(false, cm.simReqCnt, cm.sumWeight)
cm.update(time)
}
}
func (self *ClientManager) processed(node *cmNode, time mclock.AbsTime) (rcValue, rcCost uint64) {
self.lock.Lock()
defer self.lock.Unlock()
func (cm *ClientManager) processed(node *cmNode, time mclock.AbsTime) (rcValue, rcCost uint64) {
cm.lock.Lock()
defer cm.lock.Unlock()
self.stop(node, time)
cm.stop(node, time)
return uint64(node.rcValue), uint64(node.rcValue - node.startValue)
}

View file

@ -1174,15 +1174,15 @@ type NodeInfo struct {
}
// NodeInfo retrieves some protocol metadata about the running host node.
func (self *ProtocolManager) NodeInfo() *NodeInfo {
head := self.blockchain.CurrentHeader()
func (pm *ProtocolManager) NodeInfo() *NodeInfo {
head := pm.blockchain.CurrentHeader()
hash := head.Hash()
return &NodeInfo{
Network: self.networkId,
Difficulty: self.blockchain.GetTd(hash, head.Number.Uint64()),
Genesis: self.blockchain.Genesis().Hash(),
Config: self.blockchain.Config(),
Network: pm.networkId,
Difficulty: pm.blockchain.GetTd(hash, head.Number.Uint64()),
Genesis: pm.blockchain.Genesis().Hash(),
Config: pm.blockchain.Config(),
Head: hash,
}
}

View file

@ -50,47 +50,47 @@ func NewLesTxRelay(ps *peerSet, reqDist *requestDistributor) *LesTxRelay {
return r
}
func (self *LesTxRelay) registerPeer(p *peer) {
self.lock.Lock()
defer self.lock.Unlock()
func (relay *LesTxRelay) registerPeer(p *peer) {
relay.lock.Lock()
defer relay.lock.Unlock()
self.peerList = self.ps.AllPeers()
relay.peerList = relay.ps.AllPeers()
}
func (self *LesTxRelay) unregisterPeer(p *peer) {
self.lock.Lock()
defer self.lock.Unlock()
func (relay *LesTxRelay) unregisterPeer(p *peer) {
relay.lock.Lock()
defer relay.lock.Unlock()
self.peerList = self.ps.AllPeers()
relay.peerList = relay.ps.AllPeers()
}
// send sends a list of transactions to at most a given number of peers at
// once, never resending any particular transaction to the same peer twice
func (self *LesTxRelay) send(txs types.Transactions, count int) {
func (relay *LesTxRelay) send(txs types.Transactions, count int) {
sendTo := make(map[*peer]types.Transactions)
self.peerStartPos++ // rotate the starting position of the peer list
if self.peerStartPos >= len(self.peerList) {
self.peerStartPos = 0
relay.peerStartPos++ // rotate the starting position of the peer list
if relay.peerStartPos >= len(relay.peerList) {
relay.peerStartPos = 0
}
for _, tx := range txs {
hash := tx.Hash()
ltr, ok := self.txSent[hash]
ltr, ok := relay.txSent[hash]
if !ok {
ltr = &ltrInfo{
tx: tx,
sentTo: make(map[*peer]struct{}),
}
self.txSent[hash] = ltr
self.txPending[hash] = struct{}{}
relay.txSent[hash] = ltr
relay.txPending[hash] = struct{}{}
}
if len(self.peerList) > 0 {
if len(relay.peerList) > 0 {
cnt := count
pos := self.peerStartPos
pos := relay.peerStartPos
for {
peer := self.peerList[pos]
peer := relay.peerList[pos]
if _, ok := ltr.sentTo[peer]; !ok {
sendTo[peer] = append(sendTo[peer], tx)
ltr.sentTo[peer] = struct{}{}
@ -100,10 +100,10 @@ func (self *LesTxRelay) send(txs types.Transactions, count int) {
break // sent it to the desired number of peers
}
pos++
if pos == len(self.peerList) {
if pos == len(relay.peerList) {
pos = 0
}
if pos == self.peerStartPos {
if pos == relay.peerStartPos {
break // tried all available peers
}
}
@ -130,46 +130,46 @@ func (self *LesTxRelay) send(txs types.Transactions, count int) {
return func() { peer.SendTxs(reqID, cost, ll) }
},
}
self.reqDist.queue(rq)
relay.reqDist.queue(rq)
}
}
func (self *LesTxRelay) Send(txs types.Transactions) {
self.lock.Lock()
defer self.lock.Unlock()
func (relay *LesTxRelay) Send(txs types.Transactions) {
relay.lock.Lock()
defer relay.lock.Unlock()
self.send(txs, 3)
relay.send(txs, 3)
}
func (self *LesTxRelay) NewHead(head common.Hash, mined []common.Hash, rollback []common.Hash) {
self.lock.Lock()
defer self.lock.Unlock()
func (relay *LesTxRelay) NewHead(head common.Hash, mined []common.Hash, rollback []common.Hash) {
relay.lock.Lock()
defer relay.lock.Unlock()
for _, hash := range mined {
delete(self.txPending, hash)
delete(relay.txPending, hash)
}
for _, hash := range rollback {
self.txPending[hash] = struct{}{}
relay.txPending[hash] = struct{}{}
}
if len(self.txPending) > 0 {
txs := make(types.Transactions, len(self.txPending))
if len(relay.txPending) > 0 {
txs := make(types.Transactions, len(relay.txPending))
i := 0
for hash := range self.txPending {
txs[i] = self.txSent[hash].tx
for hash := range relay.txPending {
txs[i] = relay.txSent[hash].tx
i++
}
self.send(txs, 1)
relay.send(txs, 1)
}
}
func (self *LesTxRelay) Discard(hashes []common.Hash) {
self.lock.Lock()
defer self.lock.Unlock()
func (relay *LesTxRelay) Discard(hashes []common.Hash) {
relay.lock.Lock()
defer relay.lock.Unlock()
for _, hash := range hashes {
delete(self.txSent, hash)
delete(self.txPending, hash)
delete(relay.txSent, hash)
delete(relay.txPending, hash)
}
}