mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-20 10:52:25 +00:00
les: implement server API
This commit is contained in:
parent
b9107b3ba4
commit
2de75e9e1c
6 changed files with 802 additions and 393 deletions
|
|
@ -445,6 +445,11 @@ web3._extend({
|
||||||
params: 2,
|
params: 2,
|
||||||
inputFormatter:[null, null],
|
inputFormatter:[null, null],
|
||||||
}),
|
}),
|
||||||
|
new web3._extend.Method({
|
||||||
|
name: 'freezeClient',
|
||||||
|
call: 'debug_freezeClient',
|
||||||
|
params: 1,
|
||||||
|
}),
|
||||||
],
|
],
|
||||||
properties: []
|
properties: []
|
||||||
});
|
});
|
||||||
|
|
@ -772,6 +777,16 @@ web3._extend({
|
||||||
call: 'les_getCheckpoint',
|
call: 'les_getCheckpoint',
|
||||||
params: 1
|
params: 1
|
||||||
}),
|
}),
|
||||||
|
new web3._extend.Method({
|
||||||
|
name: 'clientInfo',
|
||||||
|
call: 'les_clientInfo',
|
||||||
|
params: 2
|
||||||
|
}),
|
||||||
|
new web3._extend.Method({
|
||||||
|
name: 'setClientParams',
|
||||||
|
call: 'les_setClientParams',
|
||||||
|
params: 3
|
||||||
|
}),
|
||||||
],
|
],
|
||||||
properties:
|
properties:
|
||||||
[
|
[
|
||||||
|
|
@ -783,6 +798,10 @@ web3._extend({
|
||||||
name: 'checkpointContractAddress',
|
name: 'checkpointContractAddress',
|
||||||
getter: 'les_getCheckpointContractAddress'
|
getter: 'les_getCheckpointContractAddress'
|
||||||
}),
|
}),
|
||||||
|
new web3._extend.Property({
|
||||||
|
name: 'serverInfo',
|
||||||
|
getter: 'les_serverInfo'
|
||||||
|
}),
|
||||||
]
|
]
|
||||||
});
|
});
|
||||||
`
|
`
|
||||||
|
|
|
||||||
1057
les/api.go
1057
les/api.go
File diff suppressed because it is too large
Load diff
|
|
@ -97,17 +97,14 @@ func testCapacityAPI(t *testing.T, clientCount int) {
|
||||||
t.Fatalf("Failed to obtain rpc client: %v", err)
|
t.Fatalf("Failed to obtain rpc client: %v", err)
|
||||||
}
|
}
|
||||||
headNum, headHash := getHead(ctx, t, serverRpcClient)
|
headNum, headHash := getHead(ctx, t, serverRpcClient)
|
||||||
totalCap := getTotalCap(ctx, t, serverRpcClient)
|
minCap, freeCap, totalCap := getCapacityInfo(ctx, t, serverRpcClient)
|
||||||
minCap := getMinCap(ctx, t, serverRpcClient)
|
|
||||||
testCap := totalCap * 3 / 4
|
testCap := totalCap * 3 / 4
|
||||||
fmt.Printf("Server testCap: %d minCap: %d head number: %d head hash: %064x\n", testCap, minCap, headNum, headHash)
|
fmt.Printf("Server testCap: %d minCap: %d head number: %d head hash: %064x\n", testCap, minCap, headNum, headHash)
|
||||||
reqMinCap := uint64(float64(testCap) * minRelCap / (minRelCap + float64(len(clients)-1)))
|
reqMinCap := uint64(float64(testCap) * minRelCap / (minRelCap + float64(len(clients)-1)))
|
||||||
if minCap > reqMinCap {
|
if minCap > reqMinCap {
|
||||||
t.Fatalf("Minimum client capacity (%d) bigger than required minimum for this test (%d)", minCap, reqMinCap)
|
t.Fatalf("Minimum client capacity (%d) bigger than required minimum for this test (%d)", minCap, reqMinCap)
|
||||||
}
|
}
|
||||||
|
|
||||||
freeIdx := rand.Intn(len(clients))
|
freeIdx := rand.Intn(len(clients))
|
||||||
freeCap := getFreeCap(ctx, t, serverRpcClient)
|
|
||||||
|
|
||||||
for i, client := range clients {
|
for i, client := range clients {
|
||||||
var err error
|
var err error
|
||||||
|
|
@ -147,7 +144,7 @@ func testCapacityAPI(t *testing.T, clientCount int) {
|
||||||
i, c := i, c
|
i, c := i, c
|
||||||
go func() {
|
go func() {
|
||||||
queue := make(chan struct{}, 100)
|
queue := make(chan struct{}, 100)
|
||||||
var count uint64
|
reqCount[i] = 0
|
||||||
for {
|
for {
|
||||||
select {
|
select {
|
||||||
case queue <- struct{}{}:
|
case queue <- struct{}{}:
|
||||||
|
|
@ -165,8 +162,11 @@ func testCapacityAPI(t *testing.T, clientCount int) {
|
||||||
wg.Done()
|
wg.Done()
|
||||||
<-queue
|
<-queue
|
||||||
if ok {
|
if ok {
|
||||||
count++
|
count := atomic.AddUint64(&reqCount[i], 1)
|
||||||
atomic.StoreUint64(&reqCount[i], count)
|
if count%10000 == 0 {
|
||||||
|
fmt.Println("freeze", clients[i].ID())
|
||||||
|
freezeClient(ctx, t, serverRpcClient, clients[i].ID())
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
}
|
}
|
||||||
|
|
@ -239,7 +239,7 @@ func testCapacityAPI(t *testing.T, clientCount int) {
|
||||||
default:
|
default:
|
||||||
}
|
}
|
||||||
|
|
||||||
totalCap = getTotalCap(ctx, t, serverRpcClient)
|
_, _, totalCap = getCapacityInfo(ctx, t, serverRpcClient)
|
||||||
if totalCap < testCap {
|
if totalCap < testCap {
|
||||||
fmt.Println("Total capacity underrun")
|
fmt.Println("Total capacity underrun")
|
||||||
close(stop)
|
close(stop)
|
||||||
|
|
@ -328,58 +328,61 @@ func testRequest(ctx context.Context, t *testing.T, client *rpc.Client) bool {
|
||||||
return err == nil
|
return err == nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func freezeClient(ctx context.Context, t *testing.T, server *rpc.Client, clientID enode.ID) {
|
||||||
|
if err := server.CallContext(ctx, nil, "debug_freezeClient", clientID); err != nil {
|
||||||
|
t.Fatalf("Failed to freeze client: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
}
|
||||||
|
|
||||||
func setCapacity(ctx context.Context, t *testing.T, server *rpc.Client, clientID enode.ID, cap uint64) {
|
func setCapacity(ctx context.Context, t *testing.T, server *rpc.Client, clientID enode.ID, cap uint64) {
|
||||||
if err := server.CallContext(ctx, nil, "les_setClientCapacity", clientID, cap); err != nil {
|
params := make(map[string]interface{})
|
||||||
|
params["capacity"] = cap
|
||||||
|
if err := server.CallContext(ctx, nil, "les_setClientParams", []enode.ID{clientID}, []string{}, params); err != nil {
|
||||||
t.Fatalf("Failed to set client capacity: %v", err)
|
t.Fatalf("Failed to set client capacity: %v", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func getCapacity(ctx context.Context, t *testing.T, server *rpc.Client, clientID enode.ID) uint64 {
|
func getCapacity(ctx context.Context, t *testing.T, server *rpc.Client, clientID enode.ID) uint64 {
|
||||||
var s string
|
var res map[enode.ID]map[string]interface{}
|
||||||
if err := server.CallContext(ctx, &s, "les_getClientCapacity", clientID); err != nil {
|
if err := server.CallContext(ctx, &res, "les_clientInfo", []enode.ID{clientID}, []string{}); err != nil {
|
||||||
t.Fatalf("Failed to get client capacity: %v", err)
|
t.Fatalf("Failed to get client info: %v", err)
|
||||||
}
|
}
|
||||||
cap, err := hexutil.DecodeUint64(s)
|
info, ok := res[clientID]
|
||||||
if err != nil {
|
if !ok {
|
||||||
t.Fatalf("Failed to decode client capacity: %v", err)
|
t.Fatalf("Missing client info")
|
||||||
}
|
}
|
||||||
return cap
|
v, ok := info["capacity"]
|
||||||
|
if !ok {
|
||||||
|
t.Fatalf("Missing field in client info: capacity")
|
||||||
|
}
|
||||||
|
vv, ok := v.(float64)
|
||||||
|
if !ok {
|
||||||
|
t.Fatalf("Failed to decode capacity field")
|
||||||
|
}
|
||||||
|
return uint64(vv)
|
||||||
}
|
}
|
||||||
|
|
||||||
func getTotalCap(ctx context.Context, t *testing.T, server *rpc.Client) uint64 {
|
func getCapacityInfo(ctx context.Context, t *testing.T, server *rpc.Client) (minCap, freeCap, totalCap uint64) {
|
||||||
var s string
|
var res map[string]interface{}
|
||||||
if err := server.CallContext(ctx, &s, "les_totalCapacity"); err != nil {
|
if err := server.CallContext(ctx, &res, "les_serverInfo"); err != nil {
|
||||||
t.Fatalf("Failed to query total capacity: %v", err)
|
t.Fatalf("Failed to query server info: %v", err)
|
||||||
}
|
}
|
||||||
total, err := hexutil.DecodeUint64(s)
|
decode := func(s string) uint64 {
|
||||||
if err != nil {
|
v, ok := res[s]
|
||||||
t.Fatalf("Failed to decode total capacity: %v", err)
|
if !ok {
|
||||||
|
t.Fatalf("Missing field in server info: %s", s)
|
||||||
|
}
|
||||||
|
vv, ok := v.(float64)
|
||||||
|
if !ok {
|
||||||
|
t.Fatalf("Failed to decode server info field: %s", s)
|
||||||
|
}
|
||||||
|
return uint64(vv)
|
||||||
}
|
}
|
||||||
return total
|
minCap = decode("minimumCapacity")
|
||||||
}
|
freeCap = decode("freeClientCapacity")
|
||||||
|
totalCap = decode("totalCapacity")
|
||||||
func getMinCap(ctx context.Context, t *testing.T, server *rpc.Client) uint64 {
|
return
|
||||||
var s string
|
|
||||||
if err := server.CallContext(ctx, &s, "les_minimumCapacity"); err != nil {
|
|
||||||
t.Fatalf("Failed to query minimum capacity: %v", err)
|
|
||||||
}
|
|
||||||
min, err := hexutil.DecodeUint64(s)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("Failed to decode minimum capacity: %v", err)
|
|
||||||
}
|
|
||||||
return min
|
|
||||||
}
|
|
||||||
|
|
||||||
func getFreeCap(ctx context.Context, t *testing.T, server *rpc.Client) uint64 {
|
|
||||||
var s string
|
|
||||||
if err := server.CallContext(ctx, &s, "les_freeClientCapacity"); err != nil {
|
|
||||||
t.Fatalf("Failed to query free client capacity: %v", err)
|
|
||||||
}
|
|
||||||
free, err := hexutil.DecodeUint64(s)
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("Failed to decode free client capacity: %v", err)
|
|
||||||
}
|
|
||||||
return free
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func init() {
|
func init() {
|
||||||
|
|
|
||||||
|
|
@ -408,6 +408,7 @@ func (pm *ProtocolManager) handleMsg(p *peer) error {
|
||||||
realCost = pm.server.costTracker.realCost(servingTime, msg.Size, replySize)
|
realCost = pm.server.costTracker.realCost(servingTime, msg.Size, replySize)
|
||||||
if amount != 0 {
|
if amount != 0 {
|
||||||
pm.server.costTracker.updateStats(msg.Code, amount, servingTime, realCost)
|
pm.server.costTracker.updateStats(msg.Code, amount, servingTime, realCost)
|
||||||
|
p.priceTracker.requestCost(realCost)
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
realCost = maxCost
|
realCost = maxCost
|
||||||
|
|
|
||||||
|
|
@ -103,10 +103,11 @@ type peer struct {
|
||||||
updateTime mclock.AbsTime
|
updateTime mclock.AbsTime
|
||||||
frozen uint32 // 1 if client is in frozen state
|
frozen uint32 // 1 if client is in frozen state
|
||||||
|
|
||||||
fcClient *flowcontrol.ClientNode // nil if the peer is server only
|
fcClient *flowcontrol.ClientNode // nil if the peer is server only
|
||||||
fcServer *flowcontrol.ServerNode // nil if the peer is client only
|
fcServer *flowcontrol.ServerNode // nil if the peer is client only
|
||||||
fcParams flowcontrol.ServerParams
|
fcParams flowcontrol.ServerParams
|
||||||
fcCosts requestCostTable
|
fcCosts requestCostTable
|
||||||
|
priceTracker *priceTracker
|
||||||
|
|
||||||
isTrusted bool
|
isTrusted bool
|
||||||
isOnlyAnnounce bool
|
isOnlyAnnounce bool
|
||||||
|
|
|
||||||
|
|
@ -137,13 +137,19 @@ func (s *LesServer) APIs() []rpc.API {
|
||||||
{
|
{
|
||||||
Namespace: "les",
|
Namespace: "les",
|
||||||
Version: "1.0",
|
Version: "1.0",
|
||||||
Service: NewPrivateLightServerAPI(s),
|
Service: NewPrivateLightAPI(&s.lesCommons, s.protocolManager.reg),
|
||||||
Public: false,
|
Public: false,
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
Namespace: "les",
|
Namespace: "les",
|
||||||
Version: "1.0",
|
Version: "1.0",
|
||||||
Service: NewPrivateLightAPI(&s.lesCommons, s.protocolManager.reg),
|
Service: NewPrivateLightServerAPI(s),
|
||||||
|
Public: false,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
Namespace: "debug",
|
||||||
|
Version: "1.0",
|
||||||
|
Service: NewPrivateDebugAPI(s),
|
||||||
Public: false,
|
Public: false,
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue