mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-20 10:52:25 +00:00
les: implement token sale
This commit is contained in:
parent
3354df83f9
commit
923fc54190
4 changed files with 117 additions and 22 deletions
|
|
@ -147,7 +147,7 @@ func (api *PrivateLightServerAPI) setParams(params map[string]interface{}, clien
|
||||||
setFactor(&negFactors.requestFactor)
|
setFactor(&negFactors.requestFactor)
|
||||||
case !defParams && name == "capacity":
|
case !defParams && name == "capacity":
|
||||||
if capacity, ok := value.(float64); ok && uint64(capacity) >= api.server.minCapacity {
|
if capacity, ok := value.(float64); ok && uint64(capacity) >= api.server.minCapacity {
|
||||||
_, _, err = api.server.clientPool.setCapacity(client, uint64(capacity))
|
_, _, err = api.server.clientPool.setCapacity(client.id, client.freeID, uint64(capacity), 0, true)
|
||||||
// Don't have to call factor update explicitly. It's already done
|
// Don't have to call factor update explicitly. It's already done
|
||||||
// in setCapacity function.
|
// in setCapacity function.
|
||||||
} else {
|
} else {
|
||||||
|
|
|
||||||
|
|
@ -42,6 +42,7 @@ const (
|
||||||
persistCumulativeTimeRefresh = time.Minute * 5 // refresh period of the cumulative running time persistence
|
persistCumulativeTimeRefresh = time.Minute * 5 // refresh period of the cumulative running time persistence
|
||||||
posBalanceCacheLimit = 8192 // the maximum number of cached items in positive balance queue
|
posBalanceCacheLimit = 8192 // the maximum number of cached items in positive balance queue
|
||||||
negBalanceCacheLimit = 8192 // the maximum number of cached items in negative balance queue
|
negBalanceCacheLimit = 8192 // the maximum number of cached items in negative balance queue
|
||||||
|
fullRatioTC = time.Hour
|
||||||
|
|
||||||
// connectedBias is applied to already connected clients So that
|
// connectedBias is applied to already connected clients So that
|
||||||
// already connected client won't be kicked out very soon and we
|
// already connected client won't be kicked out very soon and we
|
||||||
|
|
@ -85,6 +86,10 @@ type clientPool struct {
|
||||||
connectedMap map[enode.ID]*clientInfo
|
connectedMap map[enode.ID]*clientInfo
|
||||||
connectedQueue *prque.LazyQueue
|
connectedQueue *prque.LazyQueue
|
||||||
|
|
||||||
|
connectedBalances, disconnectedBalances uint64
|
||||||
|
lastConnectedBalanceUpdate, fullRatioLastUpdate mclock.AbsTime
|
||||||
|
fullRatio float64
|
||||||
|
|
||||||
defaultPosFactors, defaultNegFactors priceFactors
|
defaultPosFactors, defaultNegFactors priceFactors
|
||||||
|
|
||||||
connLimit int // The maximum number of connections that clientpool can support
|
connLimit int // The maximum number of connections that clientpool can support
|
||||||
|
|
@ -175,6 +180,23 @@ func newClientPool(db ethdb.Database, minCap, freeClientCap uint64, clock mclock
|
||||||
cumulativeTime: ndb.getCumulativeTime(),
|
cumulativeTime: ndb.getCumulativeTime(),
|
||||||
stopCh: make(chan struct{}),
|
stopCh: make(chan struct{}),
|
||||||
}
|
}
|
||||||
|
// calculate total token balance amount
|
||||||
|
var start enode.ID
|
||||||
|
for {
|
||||||
|
ids := pool.ndb.getPosBalanceIDs(start, enode.ID{}, 1000)
|
||||||
|
var stop bool
|
||||||
|
l := len(ids)
|
||||||
|
if l == 1000 {
|
||||||
|
stop = true
|
||||||
|
l--
|
||||||
|
}
|
||||||
|
for i := 0; i < l; i++ {
|
||||||
|
pool.disconnectedBalances += pool.ndb.getOrNewPB(ids[i]).value
|
||||||
|
}
|
||||||
|
if stop {
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
// If the negative balance of free client is even lower than 1,
|
// If the negative balance of free client is even lower than 1,
|
||||||
// delete this entry.
|
// delete this entry.
|
||||||
ndb.nbEvictCallBack = func(now mclock.AbsTime, b negBalance) bool {
|
ndb.nbEvictCallBack = func(now mclock.AbsTime, b negBalance) bool {
|
||||||
|
|
@ -208,6 +230,59 @@ func (f *clientPool) stop() {
|
||||||
f.ndb.close()
|
f.ndb.close()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (f *clientPool) updateFullRatio() {
|
||||||
|
full := float64(1)
|
||||||
|
if f.priorityConnected < f.capLimit {
|
||||||
|
freeCap := f.capLimit - f.priorityConnected
|
||||||
|
if freeCap > f.freeClientCap {
|
||||||
|
freeCapThreshold := f.capLimit / 4
|
||||||
|
if freeCap > freeCapThreshold {
|
||||||
|
full = 0
|
||||||
|
} else {
|
||||||
|
full = float64(freeCapThreshold-freeCap) / float64(freeCapThreshold-f.freeClientCap)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
now := f.clock.Now()
|
||||||
|
dt := now - f.fullRatioLastUpdate
|
||||||
|
f.fullRatioLastUpdate = now
|
||||||
|
if dt < 0 {
|
||||||
|
dt = 0
|
||||||
|
}
|
||||||
|
d := math.Exp(-float64(dt) / float64(fullRatioTC))
|
||||||
|
f.fullRatio = full - (full-f.fullRatio)*d
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *clientPool) totalTokenLimit() uint64 {
|
||||||
|
f.lock.Lock()
|
||||||
|
defer f.lock.Unlock()
|
||||||
|
|
||||||
|
f.updateFullRatio()
|
||||||
|
d := 1 - f.fullRatio
|
||||||
|
if d > 0.5 {
|
||||||
|
d = -math.Log(0.5/d) * float64(fullRatioTC)
|
||||||
|
} else {
|
||||||
|
d = 0
|
||||||
|
}
|
||||||
|
return uint64(d * float64(f.capLimit) * f.defaultPosFactors.capacityFactor)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *clientPool) totalTokenAmount() uint64 {
|
||||||
|
f.lock.Lock()
|
||||||
|
defer f.lock.Unlock()
|
||||||
|
|
||||||
|
now := f.clock.Now()
|
||||||
|
if now > f.lastConnectedBalanceUpdate+mclock.AbsTime(time.Second) {
|
||||||
|
f.connectedBalances = 0
|
||||||
|
for _, c := range f.connectedMap {
|
||||||
|
pos, _ := c.balanceTracker.getBalance(now)
|
||||||
|
f.connectedBalances += pos
|
||||||
|
}
|
||||||
|
f.lastConnectedBalanceUpdate = now
|
||||||
|
}
|
||||||
|
return f.connectedBalances + f.disconnectedBalances
|
||||||
|
}
|
||||||
|
|
||||||
// connect should be called after a successful handshake. If the connection was
|
// connect should be called after a successful handshake. If the connection was
|
||||||
// rejected, there is no need to call disconnect.
|
// rejected, there is no need to call disconnect.
|
||||||
func (f *clientPool) connect(peer clientPeer, capacity uint64) bool {
|
func (f *clientPool) connect(peer clientPeer, capacity uint64) bool {
|
||||||
|
|
@ -249,6 +324,8 @@ func (f *clientPool) connect(peer clientPeer, capacity uint64) bool {
|
||||||
}
|
}
|
||||||
f.initBalanceTracker(&e.balanceTracker, pb, nb, capacity)
|
f.initBalanceTracker(&e.balanceTracker, pb, nb, capacity)
|
||||||
// Register new client to connection queue.
|
// Register new client to connection queue.
|
||||||
|
f.disconnectedBalances -= pb.value
|
||||||
|
f.connectedBalances += pb.value
|
||||||
f.connectedMap[id] = e
|
f.connectedMap[id] = e
|
||||||
f.connectedQueue.Push(e)
|
f.connectedQueue.Push(e)
|
||||||
f.connectedCap += e.capacity
|
f.connectedCap += e.capacity
|
||||||
|
|
@ -256,6 +333,7 @@ func (f *clientPool) connect(peer clientPeer, capacity uint64) bool {
|
||||||
// If the current client is a paid client, monitor the status of client,
|
// If the current client is a paid client, monitor the status of client,
|
||||||
// downgrade it to normal client if positive balance is used up.
|
// downgrade it to normal client if positive balance is used up.
|
||||||
if e.priority {
|
if e.priority {
|
||||||
|
f.updateFullRatio()
|
||||||
f.priorityConnected += capacity
|
f.priorityConnected += capacity
|
||||||
e.balanceTracker.addCallback(balanceCallbackZero, 0, func() { f.balanceExhausted(id) })
|
e.balanceTracker.addCallback(balanceCallbackZero, 0, func() { f.balanceExhausted(id) })
|
||||||
}
|
}
|
||||||
|
|
@ -301,13 +379,6 @@ func (f *clientPool) disconnect(p clientPeer) {
|
||||||
f.dropClient(e, f.clock.Now(), false)
|
f.dropClient(e, f.clock.Now(), false)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (f *clientPool) balanceMissing(id enode.ID, freeID string, capacity uint64, minConnTime time.Duration) (uint64, uint64) {
|
|
||||||
f.lock.Lock()
|
|
||||||
defer f.lock.Unlock()
|
|
||||||
|
|
||||||
return f.capAvailable(id, freeID, capacity, minConnTime, false)
|
|
||||||
}
|
|
||||||
|
|
||||||
// capAvailable checks whether the current priority level of the given client is enough to
|
// capAvailable checks whether the current priority level of the given client is enough to
|
||||||
// connect or change capacity to the requested level and then stay connected for at least
|
// connect or change capacity to the requested level and then stay connected for at least
|
||||||
// the specified duration. If not then the additional required amount of positive balance is returned.
|
// the specified duration. If not then the additional required amount of positive balance is returned.
|
||||||
|
|
@ -420,6 +491,7 @@ func (f *clientPool) dropClient(e *clientInfo, now mclock.AbsTime, kick bool) {
|
||||||
delete(f.connectedMap, e.id)
|
delete(f.connectedMap, e.id)
|
||||||
f.connectedCap -= e.capacity
|
f.connectedCap -= e.capacity
|
||||||
if e.priority {
|
if e.priority {
|
||||||
|
f.updateFullRatio()
|
||||||
f.priorityConnected -= e.capacity
|
f.priorityConnected -= e.capacity
|
||||||
}
|
}
|
||||||
totalConnectedGauge.Update(int64(f.connectedCap))
|
totalConnectedGauge.Update(int64(f.connectedCap))
|
||||||
|
|
@ -447,6 +519,8 @@ func (f *clientPool) capacityInfo() (uint64, uint64, uint64) {
|
||||||
func (f *clientPool) finalizeBalance(c *clientInfo, now mclock.AbsTime) {
|
func (f *clientPool) finalizeBalance(c *clientInfo, now mclock.AbsTime) {
|
||||||
c.balanceTracker.stop(now)
|
c.balanceTracker.stop(now)
|
||||||
pos, neg := c.balanceTracker.getBalance(now)
|
pos, neg := c.balanceTracker.getBalance(now)
|
||||||
|
f.disconnectedBalances += pos
|
||||||
|
f.connectedBalances -= pos
|
||||||
|
|
||||||
pb, nb := f.ndb.getOrNewPB(c.id), f.ndb.getOrNewNB(c.address)
|
pb, nb := f.ndb.getOrNewPB(c.id), f.ndb.getOrNewNB(c.address)
|
||||||
pb.value = pos
|
pb.value = pos
|
||||||
|
|
@ -472,6 +546,7 @@ func (f *clientPool) balanceExhausted(id enode.ID) {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
if c.priority {
|
if c.priority {
|
||||||
|
f.updateFullRatio()
|
||||||
f.priorityConnected -= c.capacity
|
f.priorityConnected -= c.capacity
|
||||||
}
|
}
|
||||||
c.priority = false
|
c.priority = false
|
||||||
|
|
@ -493,6 +568,7 @@ func (f *clientPool) setLimits(totalConn int, totalCap uint64) {
|
||||||
f.lock.Lock()
|
f.lock.Lock()
|
||||||
defer f.lock.Unlock()
|
defer f.lock.Unlock()
|
||||||
|
|
||||||
|
f.updateFullRatio()
|
||||||
f.connLimit = totalConn
|
f.connLimit = totalConn
|
||||||
f.capLimit = totalCap
|
f.capLimit = totalCap
|
||||||
if f.connectedCap > f.capLimit || f.connectedQueue.Size() > f.connLimit {
|
if f.connectedCap > f.capLimit || f.connectedQueue.Size() > f.connLimit {
|
||||||
|
|
@ -504,26 +580,24 @@ func (f *clientPool) setLimits(totalConn int, totalCap uint64) {
|
||||||
}
|
}
|
||||||
|
|
||||||
// setCapacity sets the assigned capacity of a connected client
|
// setCapacity sets the assigned capacity of a connected client
|
||||||
func (f *clientPool) setCapacity(c *clientInfo, capacity uint64) (uint64, uint64, error) {
|
func (f *clientPool) setCapacity(id enode.ID, freeID string, capacity uint64, minConnTime time.Duration, setCap bool) (uint64, uint64, error) {
|
||||||
if capacity == 0 {
|
c := f.connectedMap[id]
|
||||||
capacity = f.freeClientCap
|
if c != nil {
|
||||||
}
|
|
||||||
if capacity < f.minCap {
|
|
||||||
capacity = f.minCap
|
|
||||||
}
|
|
||||||
if f.connectedMap[c.id] != c {
|
|
||||||
return 0, capacity, fmt.Errorf("client %064x is not connected", c.id[:])
|
|
||||||
}
|
|
||||||
if c.capacity == capacity {
|
if c.capacity == capacity {
|
||||||
return 0, capacity, nil
|
return 0, capacity, nil
|
||||||
}
|
}
|
||||||
|
}
|
||||||
var missing uint64
|
var missing uint64
|
||||||
missing, capacity = f.capAvailable(c.id, c.freeID, capacity, 0, true)
|
missing, capacity = f.capAvailable(id, freeID, capacity, 0, setCap && c != nil)
|
||||||
if missing != 0 {
|
if missing != 0 {
|
||||||
return missing, capacity, errNoPriority
|
return missing, capacity, errNoPriority
|
||||||
}
|
}
|
||||||
|
if setCap && c == nil {
|
||||||
|
return missing, capacity, fmt.Errorf("client %064x is not connected", c.id[:])
|
||||||
|
}
|
||||||
// capacity update is possible
|
// capacity update is possible
|
||||||
f.connectedCap += capacity - c.capacity
|
f.connectedCap += capacity - c.capacity
|
||||||
|
f.updateFullRatio()
|
||||||
f.priorityConnected += capacity - c.capacity
|
f.priorityConnected += capacity - c.capacity
|
||||||
c.capacity = capacity
|
c.capacity = capacity
|
||||||
c.balanceTracker.setCapacity(capacity)
|
c.balanceTracker.setCapacity(capacity)
|
||||||
|
|
@ -534,6 +608,13 @@ func (f *clientPool) setCapacity(c *clientInfo, capacity uint64) (uint64, uint64
|
||||||
return 0, capacity, nil
|
return 0, capacity, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (f *clientPool) setCapacityLocked(id enode.ID, freeID string, capacity uint64, minConnTime time.Duration, setCap bool) (uint64, uint64, error) {
|
||||||
|
f.lock.Lock()
|
||||||
|
defer f.lock.Unlock()
|
||||||
|
|
||||||
|
return f.setCapacity(id, freeID, capacity, minConnTime, setCap)
|
||||||
|
}
|
||||||
|
|
||||||
// requestCost feeds request cost after serving a request from the given peer.
|
// requestCost feeds request cost after serving a request from the given peer.
|
||||||
func (f *clientPool) requestCost(p *peer, cost uint64) {
|
func (f *clientPool) requestCost(p *peer, cost uint64) {
|
||||||
f.lock.Lock()
|
f.lock.Lock()
|
||||||
|
|
@ -605,6 +686,7 @@ func (f *clientPool) addBalance(id enode.ID, amount int64, meta string) (uint64,
|
||||||
// The capacity should be adjusted based on the requirement,
|
// The capacity should be adjusted based on the requirement,
|
||||||
// but we have no idea about the new capacity, need a second
|
// but we have no idea about the new capacity, need a second
|
||||||
// call to udpate it.
|
// call to udpate it.
|
||||||
|
f.updateFullRatio()
|
||||||
c.priority = true
|
c.priority = true
|
||||||
f.priorityConnected += c.capacity
|
f.priorityConnected += c.capacity
|
||||||
c.balanceTracker.addCallback(balanceCallbackZero, 0, func() { f.balanceExhausted(id) })
|
c.balanceTracker.addCallback(balanceCallbackZero, 0, func() { f.balanceExhausted(id) })
|
||||||
|
|
@ -612,6 +694,9 @@ func (f *clientPool) addBalance(id enode.ID, amount int64, meta string) (uint64,
|
||||||
// if balance is set to zero then reverting to non-priority status
|
// if balance is set to zero then reverting to non-priority status
|
||||||
// is handled by the balanceExhausted callback
|
// is handled by the balanceExhausted callback
|
||||||
c.balanceMetaInfo = meta
|
c.balanceMetaInfo = meta
|
||||||
|
f.connectedBalances += pb.value - oldBalance
|
||||||
|
} else {
|
||||||
|
f.disconnectedBalances += pb.value - oldBalance
|
||||||
}
|
}
|
||||||
return oldBalance, pb.value, nil
|
return oldBalance, pb.value, nil
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -50,6 +50,7 @@ type LesServer struct {
|
||||||
defParams flowcontrol.ServerParams
|
defParams flowcontrol.ServerParams
|
||||||
servingQueue *servingQueue
|
servingQueue *servingQueue
|
||||||
clientPool *clientPool
|
clientPool *clientPool
|
||||||
|
tokenSale *tokenSale
|
||||||
|
|
||||||
minCapacity, maxCapacity, freeCapacity uint64
|
minCapacity, maxCapacity, freeCapacity uint64
|
||||||
threadsIdle int // Request serving threads count when system is idle.
|
threadsIdle int // Request serving threads count when system is idle.
|
||||||
|
|
|
||||||
|
|
@ -957,5 +957,14 @@ func (h *serverHandler) broadcastHeaders() {
|
||||||
}
|
}
|
||||||
|
|
||||||
func (h *serverHandler) talkRequestHandler(id enode.ID, addr *net.UDPAddr, payload []byte) ([]byte, bool) {
|
func (h *serverHandler) talkRequestHandler(id enode.ID, addr *net.UDPAddr, payload []byte) ([]byte, bool) {
|
||||||
return payload, true // dummy handler, just returns the same payload
|
if h.server.tokenSale == nil {
|
||||||
|
return nil, false
|
||||||
|
}
|
||||||
|
var cmds [][]byte
|
||||||
|
if err := rlp.DecodeBytes(payload, &cmds); err != nil {
|
||||||
|
return nil, false
|
||||||
|
}
|
||||||
|
results := h.server.tokenSale.runCommands(cmds, id, addr.IP.String())
|
||||||
|
res, _ := rlp.EncodeToBytes(&results)
|
||||||
|
return res, true
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue