Merge pull request #3523 from yousong/feature/yousong-etcd-lock

lockman: distributed lockman with etcd
This commit is contained in:
yunion-ci-robot
2019-11-13 00:44:18 +08:00
committed by GitHub
400 changed files with 31895 additions and 28096 deletions
+24 -2
View File
@@ -52,9 +52,31 @@ func InitDB(options *common_options.DBOptions) {
}
sqlchemy.SetDB(dbConn)
lm := lockman.NewInMemoryLockManager()
switch options.LockmanMethod {
case "inmemory", "":
log.Infof("using inmemory lockman")
lm := lockman.NewInMemoryLockManager()
lockman.Init(lm)
case "etcd":
log.Infof("using etcd lockman")
tlsCfg, err := options.GetEtcdTLSConfig()
if err != nil {
log.Fatalln(err.Error())
}
lm, err := lockman.NewEtcdLockManager(&lockman.SEtcdLockManagerConfig{
Endpoints: options.EtcdEndpoints,
Username: options.EtcdUsername,
Password: options.EtcdPassword,
LockTTL: options.EtcdLockTTL,
LockPrefix: options.EtcdLockPrefix,
TLS: tlsCfg,
})
if err != nil {
log.Fatalf("etcd lockman: %v", err)
}
lockman.Init(lm)
}
// lm := lockman.NewNoopLockManager()
lockman.Init(lm)
}
func CloseDB() {
+27
View File
@@ -0,0 +1,27 @@
// Copyright 2019 Yunion
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package lockman
import (
"yunion.io/x/log"
)
type bugFunc func(fmtStr string, fmtArgs ...interface{})
var defaultBugFunc bugFunc = func(fmtStr string, fmtArgs ...interface{}) {
log.Errorf("BUG: "+fmtStr, fmtArgs...)
}
var bug bugFunc = defaultBugFunc
+267
View File
@@ -0,0 +1,267 @@
// Copyright 2019 Yunion
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package lockman
import (
"context"
"crypto/tls"
"fmt"
"runtime/debug"
"strings"
"sync"
"time"
"go.etcd.io/etcd/clientv3"
"go.etcd.io/etcd/clientv3/concurrency"
"google.golang.org/grpc"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
"yunion.io/x/onecloud/pkg/util/atexit"
)
type SEtcdLockRecord struct {
m *sync.Mutex
depth int
pfx string
sess *concurrency.Session
em *concurrency.Mutex
}
func newEtcdLockRecord(ctx context.Context, lockman *SEtcdLockManager, key string) *SEtcdLockRecord {
pfx := fmt.Sprintf("%s/%s", lockman.config.LockPrefix, key)
sess, err := concurrency.NewSession(lockman.cli, concurrency.WithTTL(lockman.config.LockTTL))
if err != nil {
panic(fmt.Sprintf("%s: create etcd session: %v", pfx, err))
}
em := concurrency.NewMutex(sess, pfx)
rec := SEtcdLockRecord{
m: &sync.Mutex{},
depth: 0,
pfx: pfx,
sess: sess,
em: em,
}
return &rec
}
func (rec *SEtcdLockRecord) lockContext(ctx context.Context) {
rec.m.Lock()
defer rec.m.Unlock()
rec.depth += 1
if rec.depth > 32 {
// NOTE callers are responsible for ensuring unlock got called
bug("%s: depth > 32", rec.pfx)
panic(debug.Stack())
}
if err := rec.em.Lock(ctx); err != nil {
msg := fmt.Sprintf("%s: etcd lock: %v", rec.pfx, err)
panic(msg)
}
if debug_log {
log.Infof("%s: lock depth %d\n%s", rec.pfx, rec.depth, debug.Stack())
}
}
func (rec *SEtcdLockRecord) unlockContext(ctx context.Context) (needClean bool) {
rec.m.Lock()
defer rec.m.Unlock()
if debug_log {
log.Infof("%s: unlock depth %d\n%s", rec.pfx, rec.depth, debug.Stack())
}
rec.depth -= 1
if rec.depth <= 0 {
if rec.depth < 0 {
bug("%s: overly unlocked", rec.pfx)
}
// Other players can make progress once lock record got removed
// after this
//
// There is no need to unlock rec.em. Revoke the session lease
// will trigger etcd to do that. Should the remote revoke call
// fail, we expect the lease auto-refresh to stop and the lease
// will expire in known time limit.
closed := false
for i := 0; i < 3; i++ {
if err := rec.sess.Close(); err != nil {
log.Errorf("%s: session close: %s", rec.pfx, err)
if i < 2 {
time.Sleep(time.Second)
}
continue
}
closed = true
break
}
if !closed {
log.Errorf("%s: session close failure", rec.pfx)
}
return true
}
return false
}
type SEtcdLockManagerConfig struct {
Endpoints []string
Username string
Password string
TLS *tls.Config
LockTTL int
LockPrefix string
dialOptions []grpc.DialOption
dialTimeout time.Duration
}
func (config *SEtcdLockManagerConfig) validate() error {
if len(config.Endpoints) == 0 {
return fmt.Errorf("no etcd endpoint configured")
}
if config.dialTimeout <= 0 {
config.dialTimeout = 3 * time.Second
}
if len(config.dialOptions) == 0 {
// let it fail right away
config.dialOptions = []grpc.DialOption{
grpc.WithBlock(),
grpc.WithTimeout(500 * time.Millisecond),
}
}
if config.LockTTL <= 0 {
config.LockTTL = 10
}
config.LockPrefix = strings.TrimSpace(config.LockPrefix)
config.LockPrefix = strings.TrimRight(config.LockPrefix, "/")
if config.LockPrefix == "" {
config.LockPrefix = "/"
}
return nil
}
type SLockTableIndex struct {
key string
holder context.Context
}
type SEtcdLockManager struct {
tableLock *sync.Mutex
lockTable map[SLockTableIndex]*SEtcdLockRecord
config *SEtcdLockManagerConfig
cli *clientv3.Client
}
func NewEtcdLockManager(config *SEtcdLockManagerConfig) (ILockManager, error) {
if err := config.validate(); err != nil {
return nil, err
}
cli, err := clientv3.New(clientv3.Config{
Endpoints: config.Endpoints,
Username: config.Username,
Password: config.Password,
TLS: config.TLS,
DialOptions: config.dialOptions,
DialTimeout: config.dialTimeout,
})
if err != nil {
return nil, errors.Wrap(err, "new etcd client")
}
lockman := SEtcdLockManager{
tableLock: &sync.Mutex{},
lockTable: map[SLockTableIndex]*SEtcdLockRecord{},
cli: cli,
config: config,
}
atexit.Register(atexit.ExitHandler{
Prio: atexit.PRIO_LOG_OTHER,
Reason: "etcd-lockman",
Value: lockman,
Func: atexit.ExitHandlerFunc(lockman.destroyAtExit),
})
return &lockman, nil
}
func (lockman *SEtcdLockManager) destroyAtExit(eh atexit.ExitHandler) {
log.Infof("closing etcd lockman session")
for _, rec := range lockman.lockTable {
if err := rec.sess.Close(); err != nil {
log.Errorf("%s: session close: %v", rec.pfx, err)
}
}
if err := lockman.cli.Close(); err != nil {
log.Errorf("etcd lockman client close: %v", err)
}
}
func (lockman *SEtcdLockManager) getRecordWithLock(ctx context.Context, key string) *SEtcdLockRecord {
lockman.tableLock.Lock()
defer lockman.tableLock.Unlock()
return lockman.getRecord(ctx, key, true)
}
func (lockman *SEtcdLockManager) getRecord(ctx context.Context, key string, alloc bool) *SEtcdLockRecord {
idx := SLockTableIndex{
key: key,
holder: ctx,
}
_, ok := lockman.lockTable[idx]
if !ok {
if !alloc {
return nil
}
lockman.lockTable[idx] = newEtcdLockRecord(ctx, lockman, key)
}
return lockman.lockTable[idx]
}
func (lockman *SEtcdLockManager) LockKey(ctx context.Context, key string) {
record := lockman.getRecordWithLock(ctx, key)
record.lockContext(ctx)
}
func (lockman *SEtcdLockManager) UnlockKey(ctx context.Context, key string) {
lockman.tableLock.Lock()
defer lockman.tableLock.Unlock()
record := lockman.getRecord(ctx, key, false)
if record == nil {
bug("%s: unlock a non-existent lock\n%s", key, debug.Stack())
return
}
needClean := record.unlockContext(ctx)
if needClean {
idx := SLockTableIndex{
key: key,
holder: ctx,
}
delete(lockman.lockTable, idx)
}
}
+52
View File
@@ -0,0 +1,52 @@
// Copyright 2019 Yunion
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package lockman
import (
"sync"
"testing"
"yunion.io/x/onecloud/pkg/util/atexit"
)
func TestEctdLockManager(t *testing.T) {
cfgs := []*testLockManagerConfig{}
shared := newSharedObject()
for i := 0; i < 4; i++ {
lockman, err := NewEtcdLockManager(&SEtcdLockManagerConfig{
Endpoints: []string{"localhost:2379"},
})
if err != nil {
t.Skipf("new etcd lockman: %v", err)
}
cfg := &testLockManagerConfig{
players: 3,
cycles: 3,
lockman: lockman,
shared: shared,
}
cfgs = append(cfgs, cfg)
}
defer atexit.Handle()
wg := &sync.WaitGroup{}
wg.Add(len(cfgs))
for _, cfg := range cfgs {
go func(cfg *testLockManagerConfig) {
testLockManager(t, cfg)
wg.Done()
}(cfg)
}
wg.Wait()
}
+1 -1
View File
@@ -174,7 +174,7 @@ func (lockman *SInMemoryLockManager) UnlockKey(ctx context.Context, key string)
record := lockman.getRecord(ctx, key, false)
if record == nil {
log.Warningf("unlock an none exist lock????")
log.Errorf("BUG: unlock an non-existent lock\n%s", debug.Stack())
return
}
+8 -47
View File
@@ -15,56 +15,17 @@
package lockman
import (
"context"
"sync"
"testing"
"time"
"yunion.io/x/pkg/util/stringutils"
)
type FakeObject struct {
Id string
}
func (o *FakeObject) GetId() string {
return o.Id
}
func (o *FakeObject) Keyword() string {
return "fake"
}
func run(t *testing.T, ctx context.Context, obj ILockedObject, id int, sleep time.Duration) {
t.Logf("ready to run at %d [%p]", id, ctx)
LockObject(ctx, obj)
defer ReleaseObject(ctx, obj)
t.Logf("Acquire obj at %d [%p]", id, ctx)
time.Sleep(sleep)
t.Logf("Release obj at %d [%p]", id, ctx)
}
func TestInMemoryLockManager(t *testing.T) {
Init(NewInMemoryLockManager())
objId := stringutils.UUID4()
cycle := 1
var wg sync.WaitGroup
t.Log("Start")
for id := 0; id <= 3; id += 1 {
wg.Add(1)
go func(localId int) {
t.Logf("Start %d", localId)
ctx := context.WithValue(context.Background(), "ID", localId)
for i := 0; i < cycle; i += 1 {
obj := &FakeObject{Id: objId}
run(t, ctx, obj, localId, time.Duration(localId)*time.Second)
}
wg.Done()
}(id)
lockman := NewInMemoryLockManager()
shared := newSharedObject()
cfg := &testLockManagerConfig{
players: 3,
cycles: 3,
lockman: lockman,
shared: shared,
}
wg.Wait()
testLockManager(t, cfg)
}
+126
View File
@@ -0,0 +1,126 @@
// Copyright 2019 Yunion
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package lockman
import (
"context"
"fmt"
"sync"
"testing"
"time"
"yunion.io/x/pkg/util/stringutils"
)
type FakeObject struct {
Id string
playerId int
playerIdCount int
}
func (o *FakeObject) GetId() string {
return o.Id
}
func (o *FakeObject) Keyword() string {
return "fake"
}
func (o *FakeObject) push(playerId int) {
if o.playerId >= 0 {
if o.playerId != playerId {
panic(fmt.Sprintf("obj locked by %d, locked again by %d", o.playerId, playerId))
}
} else {
if o.playerIdCount != 0 {
panic(fmt.Sprintf("obj unlocked but player id count %d", o.playerIdCount))
}
o.playerId = playerId
}
o.playerIdCount += 1
}
func (o *FakeObject) pop(playerId int) {
if o.playerId != playerId {
panic(fmt.Sprintf("obj previously locked by %d, now unlocked by %d", o.playerId, playerId))
}
o.playerIdCount -= 1
if o.playerIdCount < 0 {
panic(fmt.Sprintf("obj overly unlocked"))
} else if o.playerIdCount == 0 {
o.playerId = -1
}
}
func newSharedObject() *FakeObject {
obj := &FakeObject{
Id: stringutils.UUID4(),
playerId: -1,
playerIdCount: 0,
}
return obj
}
type testLockManagerConfig struct {
players int
cycles int
shared *FakeObject
lockman ILockManager
}
func testLockManager(t *testing.T, cfg *testLockManagerConfig) {
players := cfg.players
cycles := cfg.cycles
shared := cfg.shared
lockman := cfg.lockman
bug = func(fmtStr string, fmtArgs ...interface{}) {
t.Fatalf(fmtStr, fmtArgs...)
}
run := func(ctx context.Context, id int, cycle int, sleep time.Duration) {
key := getObjectKey(shared)
indents := []string{" ", " "}
logpref := fmt.Sprintf("%p: %02d: %p", lockman, id, ctx)
func() {
fmt.Printf("%s: %s >acquire\n", logpref, indents[0])
lockman.LockKey(ctx, key)
shared.push(id)
fmt.Printf("%s: %s >>acquired\n", logpref, indents[1])
}()
defer func() {
fmt.Printf("%s: %s <<release\n", logpref, indents[1])
shared.pop(id)
lockman.UnlockKey(ctx, key)
fmt.Printf("%s: %s <released\n", logpref, indents[0])
}()
time.Sleep(sleep)
}
var wg sync.WaitGroup
for id := 0; id < players; id += 1 {
wg.Add(1)
go func(localId int) {
ctx := context.WithValue(context.Background(), "ID", localId)
for i := 0; i < cycles; i += 1 {
run(ctx, localId, i, time.Duration(localId)*time.Second)
}
wg.Done()
}(id)
}
wg.Wait()
}
+69
View File
@@ -15,6 +15,9 @@
package options
import (
"crypto/tls"
"crypto/x509"
"encoding/pem"
"fmt"
"io/ioutil"
"os"
@@ -23,6 +26,7 @@ import (
"yunion.io/x/log"
"yunion.io/x/log/hooks"
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/util/reflectutils"
"yunion.io/x/pkg/util/version"
"yunion.io/x/pkg/utils"
@@ -101,6 +105,71 @@ type DBOptions struct {
DebugSqlchemy bool `default:"false" help:"Print SQL executed by sqlchemy"`
QueryOffsetOptimization bool `help:"apply query offset optimization"`
LockmanMethod string `help:"method for lock synchronization" choices:"inmemory|etcd" default:"inmemory"`
EtcdLockPrefix string `help:"prefix of etcd lock records" default:"/locks"`
EtcdLockTTL int `help:"ttl of etcd lock records"`
EtcdEndpoints []string `help:"endpoints of etcd cluster"`
EtcdUsername string `help:"username of etcd cluster"`
EtcdPassword string `help:"password of etcd cluster"`
EtcdUseTLS bool `help:"use tls transport to connect etcd cluster" default:"false"`
EtcdSkipTLSVerify bool `help:"skip tls verification" default:"false"`
EtcdCacert string `help:"path to cacert for connecting to etcd cluster"`
EtcdCert string `help:"path to cert file for connecting to etcd cluster"`
EtcdKey string `help:"path to key file for connecting to etcd cluster"`
}
func (this *DBOptions) GetEtcdTLSConfig() (*tls.Config, error) {
var (
cert tls.Certificate
certLoaded bool
capool *x509.CertPool
)
if this.EtcdCert != "" && this.EtcdKey != "" {
var err error
cert, err = tls.LoadX509KeyPair(this.EtcdCert, this.EtcdKey)
if err != nil {
return nil, errors.Wrap(err, "load etcd cert and key")
}
certLoaded = true
this.EtcdUseTLS = true
}
if this.EtcdCacert != "" {
data, err := ioutil.ReadFile(this.EtcdCacert)
if err != nil {
return nil, errors.Wrap(err, "read cacert file")
}
capool = x509.NewCertPool()
for {
var block *pem.Block
block, data = pem.Decode(data)
if block == nil {
break
}
cacert, err := x509.ParseCertificate(block.Bytes)
if err != nil {
return nil, errors.Wrap(err, "parse cacert file")
}
capool.AddCert(cacert)
}
this.EtcdUseTLS = true
}
if this.EtcdSkipTLSVerify { // it's false by default, true means user intends to use tls
this.EtcdUseTLS = true
}
if this.EtcdUseTLS {
cfg := &tls.Config{
RootCAs: capool,
InsecureSkipVerify: this.EtcdSkipTLSVerify,
}
if certLoaded {
cfg.Certificates = []tls.Certificate{cert}
}
return cfg, nil
}
return nil, nil
}
func (this *DBOptions) GetDBConnection() (dialect, connstr string, err error) {