Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
111 changes: 84 additions & 27 deletions blockproducer/metastate.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ package blockproducer

import (
"bytes"
"sort"
"time"

pi "github.com/CovenantSQL/CovenantSQL/blockproducer/interfaces"
Expand All @@ -44,6 +45,20 @@ type metaState struct {
dirty, readonly *metaIndex
}

// MinerInfos is MinerInfo array
type MinerInfos []*types.MinerInfo

// Len returns the length of the uints array.
func (x MinerInfos) Len() int { return len(x) }

// Less returns true if MinerInfo i is less than node j.
func (x MinerInfos) Less(i, j int) bool {
return x[i].NodeID < x[j].NodeID
}

// Swap exchanges MinerInfo i and j.
func (x MinerInfos) Swap(i, j int) { x[i], x[j] = x[j], x[i] }

func newMetaState() *metaState {
return &metaState{
dirty: newMetaIndex(),
Expand Down Expand Up @@ -641,7 +656,7 @@ func (s *metaState) matchProvidersWithUser(tx *types.CreateDatabase) (err error)
return
}

miners := make([]*types.MinerInfo, 0, minerCount)
miners := make(MinerInfos, 0, minerCount)

for _, m := range tx.ResourceMeta.TargetMiners {
if po, loaded := s.loadProviderObject(m); !loaded {
Expand Down Expand Up @@ -669,28 +684,14 @@ func (s *metaState) matchProvidersWithUser(tx *types.CreateDatabase) (err error)
err = errors.Wrapf(err, "miners match target are not enough %d:%d", len(miners), minerCount)
return
}
// try old miners first
for _, po := range s.readonly.provider {
miners, _ = filterAndAppendMiner(miners, po, tx, sender)
// if got enough, break
if uint64(len(miners)) == minerCount {
break
}
}
// try fresh miners
if uint64(len(miners)) < minerCount {
for _, po := range s.dirty.provider {
miners, _ = filterAndAppendMiner(miners, po, tx, sender)
// if got enough, break
if uint64(len(miners)) == minerCount {
break
}
}
}
if uint64(len(miners)) < minerCount {
err = ErrNoEnoughMiner
var newMiners MinerInfos
// create new merged map
newMiners, err = s.filterNMiners(tx, sender, int(minerCount)-len(miners))
if err != nil {
return
}

miners = append(miners, newMiners...)
}

// generate new sqlchain id and address
Expand Down Expand Up @@ -747,7 +748,7 @@ func (s *metaState) matchProvidersWithUser(tx *types.CreateDatabase) (err error)
Owner: sender,
Users: users,
EncodedGenesis: enc.Bytes(),
Miners: miners[:],
Miners: miners,
}

if _, loaded := s.loadSQLChainObject(dbID); loaded {
Expand All @@ -756,13 +757,52 @@ func (s *metaState) matchProvidersWithUser(tx *types.CreateDatabase) (err error)
}
s.dirty.accounts[dbAddr] = &types.Account{Address: dbAddr}
s.dirty.databases[dbID] = sp
for _, miner := range tx.ResourceMeta.TargetMiners {
s.deleteProviderObject(miner)
for _, miner := range miners {
s.deleteProviderObject(miner.Address)
}
log.Infof("success create sqlchain with database ID: %s", dbID)
return
}

func (s *metaState) filterNMiners(
tx *types.CreateDatabase,
user proto.AccountAddress,
minerCount int) (
m MinerInfos, err error,
) {
// create new merged map
allProviderMap := make(map[proto.AccountAddress]*types.ProviderProfile)
for k, v := range s.readonly.provider {
allProviderMap[k] = v
}
for k, v := range s.dirty.provider {
if v == nil {
delete(allProviderMap, k)
} else {
allProviderMap[k] = v
}
}

// delete selected target miners
for _, m := range tx.ResourceMeta.TargetMiners {
delete(allProviderMap, m)
}

// suppose 1/4 miners match
newMiners := make(MinerInfos, 0, len(allProviderMap)/4)
// filter all miners to slice and sort
for _, po := range allProviderMap {
newMiners, _ = filterAndAppendMiner(newMiners, po, tx, user)
}
if len(newMiners) < minerCount {
err = ErrNoEnoughMiner
return
}

sort.Slice(newMiners, newMiners.Less)
return newMiners[:minerCount], nil
}

func filterAndAppendMiner(
miners []*types.MinerInfo,
po *types.ProviderProfile,
Expand Down Expand Up @@ -939,7 +979,7 @@ func (s *metaState) updateBilling(tx *types.UpdateBilling) (err error) {
err = errors.Wrap(ErrDatabaseNotFound, "update billing failed")
return
}
log.Debugf("update billing addr: %s, tx: %v", tx.GetAccountAddress(), tx)
log.Debugf("update billing addr: %s, user: %d, tx: %v", tx.GetAccountAddress(), len(tx.Users), tx)

if newProfile.GasPrice == 0 {
return
Expand Down Expand Up @@ -979,15 +1019,32 @@ func (s *metaState) updateBilling(tx *types.UpdateBilling) (err error) {
miner.PendingIncome += userMap[user.Address][miner.Address] * newProfile.GasPrice
}
} else {
rate := 1 - float64(user.AdvancePayment)/float64(costMap[user.Address]*newProfile.GasPrice)
rate := float64(user.AdvancePayment) / float64(costMap[user.Address]*newProfile.GasPrice)
user.AdvancePayment = 0
user.Status = types.Arrears
for _, miner := range newProfile.Miners {
income := userMap[user.Address][miner.Address] * newProfile.GasPrice
minerIncome := uint64(float64(income) * rate)
miner.PendingIncome += minerIncome
if miner.UserArrears == nil {
miner.UserArrears = make([]*types.UserArrears, 0)
}
exist := false
for i := range miner.UserArrears {
miner.UserArrears[i].Arrears += (income - minerIncome)
if miner.UserArrears[i].User == user.Address {
exist = true
diff := income - minerIncome
miner.UserArrears[i].Arrears += diff
user.Arrears += diff
}
}
if !exist {
diff := income - minerIncome
miner.UserArrears = append(miner.UserArrears, &types.UserArrears{
User: user.Address,
Arrears: diff,
})
user.Arrears += diff
}
}
}
Expand Down
90 changes: 84 additions & 6 deletions blockproducer/metastate_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,8 @@ import (
"github.com/CovenantSQL/CovenantSQL/proto"
"github.com/CovenantSQL/CovenantSQL/route"
"github.com/CovenantSQL/CovenantSQL/types"
"github.com/CovenantSQL/CovenantSQL/utils/log"

Comment thread
auxten marked this conversation as resolved.
"github.com/pkg/errors"
. "github.com/smartystreets/goconvey/convey"
)
Expand Down Expand Up @@ -639,7 +641,7 @@ func TestMetaState(t *testing.T) {
types.NewBaseAccount(
&types.Account{
Address: addr2,
TokenBalance: [types.SupportTokenNumber]uint64{10000000, 100},
TokenBalance: [types.SupportTokenNumber]uint64{10000000000, 100},
},
),
types.NewBaseAccount(
Expand Down Expand Up @@ -755,7 +757,7 @@ func TestMetaState(t *testing.T) {
Owner: addr3,
ResourceMeta: types.ResourceMeta{
TargetMiners: []proto.AccountAddress{addr2},
Node: 10,
Node: 2,
Space: 9,
Memory: 9,
LoadAvgPerCPU: 0.1,
Expand All @@ -764,11 +766,30 @@ func TestMetaState(t *testing.T) {
},
Nonce: 1,
GasPrice: 1,
AdvancePayment: uint64(conf.GConf.QPS) * conf.GConf.BillingBlockCount * 10,
AdvancePayment: uint64(conf.GConf.QPS) * conf.GConf.BillingBlockCount * 2,
},
}
err = invalidCd7.Sign(privKey3)
So(err, ShouldBeNil)
invalidCd8 := types.CreateDatabase{
CreateDatabaseHeader: types.CreateDatabaseHeader{
Owner: addr2,
ResourceMeta: types.ResourceMeta{
TargetMiners: []proto.AccountAddress{addr2},
Node: 2,
Space: 9,
Memory: 9,
LoadAvgPerCPU: 0.1,
UseEventualConsistency: false,
ConsistencyLevel: 0,
},
Nonce: 1,
GasPrice: 1,
AdvancePayment: uint64(conf.GConf.QPS) * uint64(conf.GConf.BillingBlockCount) * 2,
},
}
err = invalidCd8.Sign(privKey2)
So(err, ShouldBeNil)

err = ms.apply(&invalidPs)
So(errors.Cause(err), ShouldEqual, ErrInsufficientBalance)
Expand Down Expand Up @@ -833,6 +854,63 @@ func TestMetaState(t *testing.T) {
}
err = ms.apply(&invalidCd7)
So(errors.Cause(err), ShouldEqual, ErrNoEnoughMiner)

ms.readonly.provider[proto.AccountAddress(hash.HashH([]byte("9")))] = &types.ProviderProfile{
TargetUser: []proto.AccountAddress{addr2},
GasPrice: 1,
LoadAvgPerCPU: 0.001,
Memory: 100,
Space: 100,
TokenType: 0,
NodeID: "0001111",
}
ms.dirty.provider[proto.AccountAddress(hash.HashH([]byte("9")))] = &types.ProviderProfile{
TargetUser: []proto.AccountAddress{addr2},
GasPrice: 1,
LoadAvgPerCPU: 0.001,
Memory: 100,
Space: 100,
TokenType: 0,
NodeID: "0002111",
}
ms.dirty.provider[proto.AccountAddress(hash.HashH([]byte("10")))] = &types.ProviderProfile{
TargetUser: []proto.AccountAddress{addr2},
GasPrice: 1,
LoadAvgPerCPU: 0.001,
Memory: 100,
Space: 100,
TokenType: 0,
NodeID: "0003111",
}
ms.dirty.provider[proto.AccountAddress(hash.HashH([]byte("11")))] = &types.ProviderProfile{
TargetUser: []proto.AccountAddress{addr2},
GasPrice: 1,
LoadAvgPerCPU: 0.001,
Memory: 100,
Space: 100,
TokenType: 0,
NodeID: "0000003",
}
ms.dirty.provider[proto.AccountAddress(hash.HashH([]byte("12")))] = &types.ProviderProfile{
TargetUser: []proto.AccountAddress{addr2},
GasPrice: 1,
LoadAvgPerCPU: 0.001,
Memory: 100,
Space: 100,
TokenType: 0,
NodeID: "0000001",
}
err = ms.apply(&invalidCd8)
So(err, ShouldBeNil)
dbID := proto.FromAccountAndNonce(addr2, uint32(invalidCd8.Nonce))

mIDs := make([]string, 0)
for _, m := range ms.dirty.databases[dbID].Miners {
mIDs = append(mIDs, string(m.NodeID))
}
log.Debugf("mIDs: %v", mIDs)
So(mIDs, ShouldContain, "0000003")
So(mIDs, ShouldContain, "0000001")
})
Convey("When SQLChain create", func() {
ps := types.ProvideService{
Expand Down Expand Up @@ -1096,7 +1174,7 @@ func TestMetaState(t *testing.T) {
sqlchain, loaded := ms.loadSQLChainObject(dbID)
So(loaded, ShouldBeTrue)
So(len(sqlchain.Miners), ShouldEqual, 1)
So(sqlchain.Miners[0].PendingIncome, ShouldEqual, 125)
So(sqlchain.Miners[0].PendingIncome, ShouldEqual, 100)
users = [3]*types.UserCost{
&types.UserCost{
User: addr1,
Expand Down Expand Up @@ -1143,8 +1221,8 @@ func TestMetaState(t *testing.T) {
sqlchain, loaded = ms.loadSQLChainObject(dbID)
So(loaded, ShouldBeTrue)
So(len(sqlchain.Miners), ShouldEqual, 1)
So(sqlchain.Miners[0].PendingIncome, ShouldEqual, 115)
So(sqlchain.Miners[0].ReceivedIncome, ShouldEqual, 125)
So(sqlchain.Miners[0].PendingIncome, ShouldEqual, 100)
So(sqlchain.Miners[0].ReceivedIncome, ShouldEqual, 100)
})
})
})
Expand Down
6 changes: 2 additions & 4 deletions blockproducer/rpc.go
Original file line number Diff line number Diff line change
Expand Up @@ -158,14 +158,13 @@ func WaitDatabaseCreation(
req = &types.QuerySQLChainProfileReq{
DBID: dbID,
}
resp = &types.QuerySQLChainProfileResp{}
)
defer ticker.Stop()
for {
select {
case <-ticker.C:
if err = rpc.RequestBP(
route.MCCQuerySQLChainProfile.String(), req, resp,
route.MCCQuerySQLChainProfile.String(), req, nil,
); err != nil {
if !strings.Contains(err.Error(), ErrDatabaseNotFound.Error()) {
// err != nil && err != ErrDatabaseNotFound (unexpected error)
Expand Down Expand Up @@ -195,14 +194,13 @@ func WaitBPChainService(ctx context.Context, period time.Duration) (err error) {
req = &types.FetchBlockReq{
Height: 0, // Genesis block
}
resp = &types.FetchTxBillingResp{}
)
defer ticker.Stop()
for {
select {
case <-ticker.C:
if err = rpc.RequestBP(
route.MCCFetchBlock.String(), req, resp,
route.MCCFetchBlock.String(), req, nil,
); err == nil || !strings.Contains(err.Error(), "can't find service") {
return
}
Expand Down
11 changes: 11 additions & 0 deletions client/clientbench_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,12 +17,15 @@
package client

import (
"context"
"database/sql"
"os"
"path/filepath"
"sync"
"testing"
"time"

"github.com/CovenantSQL/CovenantSQL/blockproducer"
"github.com/CovenantSQL/CovenantSQL/utils"
"github.com/CovenantSQL/CovenantSQL/utils/log"
)
Expand Down Expand Up @@ -52,6 +55,14 @@ func BenchmarkCovenantSQLDriver(b *testing.B) {
}
})

// wait for chain service
var ctx1, cancel1 = context.WithTimeout(context.Background(), 1*time.Minute)
defer cancel1()
err = blockproducer.WaitBPChainService(ctx1, 3*time.Second)
if err != nil {
b.Fatalf("wait for chain service failed: %v", err)
}

// create
meta := ResourceMeta{}
meta.Node = 3
Expand Down
Loading