Automatic merge from release/2.1.0 -> release/2.2.0

* commit '12b1ef939eb190af6c25790a0b8b5357aa678173':
  remove binary
  增加:镜像缓存加锁,避免批量创建时重复上传镜像
  remove binary
  使用inmemory lockman
  initial commit
This commit is contained in:
邱剑
2018-09-12 14:54:15 +08:00
19 changed files with 244 additions and 66 deletions
+2 -2
View File
@@ -20,8 +20,8 @@ func InitDB(options *DBOptions) {
}
sqlchemy.SetDB(dbConn)
// lm := lockman.NewInMemoryLockManager()
lm := lockman.NewNoopLockManager()
lm := lockman.NewInMemoryLockManager()
// lm := lockman.NewNoopLockManager()
lockman.Init(lm)
}
+2
View File
@@ -1106,7 +1106,9 @@ func (dispatcher *DBModelDispatcher) Delete(ctx context.Context, idstr string, q
return nil, httperrors.NewGeneralError(err)
}
log.Debugf("Delete %s", model.GetShortDesc())
lockman.LockObject(ctx, model)
defer lockman.ReleaseObject(ctx, model)
return deleteItem(dispatcher.modelManager, model, ctx, userCred, query, data)
}
@@ -0,0 +1,60 @@
package main
import (
"time"
"sync"
"context"
"yunion.io/x/log"
"yunion.io/x/pkg/util/stringutils"
"yunion.io/x/onecloud/pkg/cloudcommon/db/lockman"
"math/rand"
)
type FakeObject struct {
Id string
}
func (o *FakeObject) GetId() string {
return o.Id
}
func (o *FakeObject) Keyword() string {
return "fake"
}
func run(ctx context.Context, obj lockman.ILockedObject, id int, sleep time.Duration) {
log.Infof("ready to run at %d [%p]", id, ctx)
lockman.LockObject(ctx, obj)
defer lockman.ReleaseObject(ctx, obj)
log.Infof("Acquire obj at %d [%p]", id, ctx)
time.Sleep(sleep)
log.Infof("Release obj at %d [%p]", id, ctx)
}
func main() {
lockman.Init(lockman.NewInMemoryLockManager())
objId := stringutils.UUID4()
cycle := 10
var wg sync.WaitGroup
log.Infof("Start")
for id := 0; id <= 3; id += 1 {
wg.Add(1)
go func(localId int) {
log.Infof("Start %d", localId)
ctx := context.WithValue(context.Background(), "ID", localId)
for i := 0; i < cycle; i += 1 {
obj := &FakeObject{Id: objId}
run(ctx, obj, localId, time.Duration(rand.Intn(1000))*time.Millisecond)
}
wg.Done()
}(id)
}
wg.Wait()
}
+46
View File
@@ -0,0 +1,46 @@
package lockman
import (
"container/list"
"yunion.io/x/log"
)
type FIFO struct {
fifo *list.List
}
func NewFIFO() *FIFO {
return &FIFO{fifo: list.New()}
}
func (f *FIFO) Push(ele interface{}) {
f.fifo.PushBack(ele)
}
func (f *FIFO) Pop(ele interface{}) interface{} {
e := f.fifo.Front()
for e != nil && e.Value != ele {
e = e.Next()
}
if e != nil {
v := f.fifo.Remove(e)
if v != ele {
log.Fatalf("remove element not identical!!")
}
}
return nil
}
func (f *FIFO) Len() int {
return f.fifo.Len()
}
type ElementInspectFunc func(ele interface{})
func (f *FIFO) Enum(eif ElementInspectFunc) {
e := f.fifo.Front()
for e != nil {
eif(e.Value)
e = e.Next()
}
}
+19
View File
@@ -0,0 +1,19 @@
package lockman
import "testing"
func TestFIFO_Pop(t *testing.T) {
fifo := NewFIFO()
data := []int{1, 2, 3}
for i := 0; i < len(data); i += 1 {
fifo.Push(data[i])
}
t.Logf("FIFO size: %d", fifo.Len())
fifo.Enum(func(ele interface{}) {
t.Logf("FIFO ele %#v", ele)
})
for i := 0; i < len(data); i += 1 {
fifo.Pop(data[i])
}
t.Logf("FIFO size: %d", fifo.Len())
}
+79 -48
View File
@@ -2,55 +2,93 @@ package lockman
import (
"context"
"runtime/debug"
"sync"
"yunion.io/x/log"
"yunion.io/x/pkg/util/fifoutils"
"runtime/debug"
)
type SInMemoryLockOwner struct {
const (
debug_log = false
)
/*type SInMemoryLockOwner struct {
owner context.Context
ready chan bool
}
}*/
type SInMemoryLockRecord struct {
lock *sync.Mutex
holder context.Context
counter int
queue *fifoutils.FIFO
lock *sync.Mutex
cond *sync.Cond
holder context.Context
depth int
waiter *FIFO
}
func newInMemoryLockRecord(ctx context.Context) *SInMemoryLockRecord {
rec := SInMemoryLockRecord{lock: &sync.Mutex{}, queue: fifoutils.NewFIFO(), holder: ctx, counter: 0}
lock := &sync.Mutex{}
cond := sync.NewCond(lock)
rec := SInMemoryLockRecord{lock: lock, cond: cond, holder: ctx, depth: 0, waiter: NewFIFO()}
return &rec
}
func (rec *SInMemoryLockRecord) lockContext(ctx context.Context) *SInMemoryLockOwner {
func (rec *SInMemoryLockRecord) lockContext(ctx context.Context) {
rec.lock.Lock()
defer rec.lock.Unlock()
if rec.holder == nil {
rec.holder = ctx
rec.depth = 1
return
}
if debug_log {
log.Debugf("rec.hold=[%p] ctx=[%p] %v", rec.holder, ctx, rec.holder==ctx)
}
if rec.holder == ctx {
rec.counter += 1
log.Infof("lockContext: same ctx, counter: %d [%p]", rec.counter, rec.holder)
if rec.counter > 32 {
// MUST BE BUG
rec.depth += 1
if debug_log {
log.Infof("lockContext: same ctx, depth: %d [%p]", rec.depth, rec.holder)
}
if rec.depth > 32 {
// XXX MUST BE BUG ???
debug.PrintStack()
panic("Too many recursive locks!!!")
}
return nil
return
}
// check
for i := 0; i < rec.queue.Len(); i += 1 {
ele := rec.queue.ElementAt(i).(*SInMemoryLockOwner)
if ele.owner == ctx {
log.Fatalf("try to lock from a wait context????")
}
}
owner := SInMemoryLockOwner{owner: ctx, ready: make(chan bool)}
rec.queue.Push(&owner)
return &owner
// check
rec.waiter.Enum(func(ele interface{}) {
electx := ele.(context.Context)
if electx == ctx {
log.Fatalf("try to lock from a waiter context????")
}
})
rec.waiter.Push(ctx)
if debug_log {
log.Debugf("waiter size %d after push", rec.waiter.Len())
log.Debugf("Start to wait ... [%p]", ctx)
}
for rec.holder != nil {
rec.cond.Wait()
}
if debug_log {
log.Debugf("End of wait ... [%p]", ctx)
}
rec.waiter.Pop(ctx)
if debug_log {
log.Debugf("waiter size %d after pop", rec.waiter.Len())
}
rec.holder = ctx
rec.depth = 1
}
func (rec *SInMemoryLockRecord) unlockContext(ctx context.Context) (needClean bool) {
@@ -61,31 +99,27 @@ func (rec *SInMemoryLockRecord) unlockContext(ctx context.Context) (needClean bo
log.Fatalf("try to unlock a wait context???")
}
rec.counter -= 1
if debug_log {
log.Debugf("unlockContext depth %d [%p]", rec.depth, ctx)
}
if rec.counter <= 0 {
if rec.queue.Len() == 0 {
rec.depth -= 1
if rec.depth <= 0 {
if debug_log {
log.Debugf("depth 0, to release lock for context [%p]", ctx)
}
rec.holder = nil
if rec.waiter.Len() == 0 {
return true
}
newHolder := rec.queue.Pop().(*SInMemoryLockOwner)
rec.holder = newHolder.owner
rec.counter = 1
newHolder.notify()
rec.cond.Signal()
}
return false
}
func (owner *SInMemoryLockOwner) wait() {
// log.Infof("wait for notify %p", owner.owner)
<-owner.ready
}
func (owner *SInMemoryLockOwner) notify() {
// log.Infof("notify %p", owner.owner)
owner.ready <- true
}
type SInMemoryLockManager struct {
tableLock *sync.Mutex
lockTable map[string]*SInMemoryLockRecord
@@ -117,10 +151,7 @@ func (lockman *SInMemoryLockManager) getRecord(ctx context.Context, key string,
func (lockman *SInMemoryLockManager) LockKey(ctx context.Context, key string) {
record := lockman.getRecordWithLock(ctx, key)
owner := record.lockContext(ctx)
if owner != nil {
owner.wait()
}
record.lockContext(ctx)
}
func (lockman *SInMemoryLockManager) UnlockKey(ctx context.Context, key string) {
@@ -137,4 +168,4 @@ func (lockman *SInMemoryLockManager) UnlockKey(ctx context.Context, key string)
if needClean {
delete(lockman.lockTable, key)
}
}
}
+3 -3
View File
@@ -24,10 +24,10 @@ func (o *FakeObject) Keyword() string {
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 LockObject(ctx, obj)
t.Logf("Acquire obj at %s [%p]", id, ctx)
defer ReleaseObject(ctx, obj)
t.Logf("Acquire obj at %d [%p]", id, ctx)
time.Sleep(sleep)
t.Logf("Release obj at %s [%p]", id, ctx)
t.Logf("Release obj at %d [%p]", id, ctx)
}
func TestInMemoryLockManager(t *testing.T) {
+2
View File
@@ -282,8 +282,10 @@ func (model *SVirtualResourceBase) VirtualModelManager() IVirtualModelManager {
func (model *SVirtualResourceBase) CancelPendingDelete(ctx context.Context, userCred mcclient.TokenCredential) error {
ownerProjId := model.GetOwnerProjectId()
lockman.LockClass(ctx, model.GetModelManager(), ownerProjId)
defer lockman.ReleaseClass(ctx, model.GetModelManager(), ownerProjId)
_, err := model.GetModelManager().TableSpec().Update(model, func() error {
model.Name = GenerateName(model.GetModelManager(), ownerProjId, model.Name)
model.PendingDeleted = false
+11 -1
View File
@@ -10,6 +10,7 @@ import (
"yunion.io/x/onecloud/pkg/cloudprovider"
"yunion.io/x/onecloud/pkg/compute/models"
"yunion.io/x/onecloud/pkg/httperrors"
"yunion.io/x/onecloud/pkg/cloudcommon/db/lockman"
)
type SAliyunHostDriver struct {
@@ -25,7 +26,7 @@ func (self *SAliyunHostDriver) GetHostType() string {
return models.HOST_TYPE_ALIYUN
}
func (self *SAliyunHostDriver) CheckAndSetCacheImage(ctx context.Context, host *models.SHost, storageCache *models.SStoragecache, scimg *models.SStoragecachedimage, task taskman.ITask) error {
func (self *SAliyunHostDriver) CheckAndSetCacheImage(ctx context.Context, host *models.SHost, storageCache *models.SStoragecache, task taskman.ITask) error {
params := task.GetParams()
imageId, err := params.GetString("image_id")
if err != nil {
@@ -36,9 +37,16 @@ func (self *SAliyunHostDriver) CheckAndSetCacheImage(ctx context.Context, host *
osType, _ := params.GetString("os_type")
osDist, _ := params.GetString("os_distribution")
isForce := jsonutils.QueryBoolean(params, "is_force", false)
userCred := task.GetUserCred()
taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) {
lockman.LockRawObject(ctx, "cachedimages", fmt.Sprintf("%s-%s", storageCache.Id, imageId))
defer lockman.ReleaseRawObject(ctx, "cachedimages", fmt.Sprintf("%s-%s", storageCache.Id, imageId))
scimg := models.StoragecachedimageManager.Register(ctx, task.GetUserCred(), storageCache.Id, imageId)
iStorageCache, err := storageCache.GetIStorageCache()
if err != nil {
return nil, err
@@ -49,6 +57,8 @@ func (self *SAliyunHostDriver) CheckAndSetCacheImage(ctx context.Context, host *
if err != nil {
return nil, err
} else {
scimg.SetExternalId(extImgId)
ret := jsonutils.NewDict()
ret.Add(jsonutils.NewString(extImgId), "image_id")
return ret, nil
-1
View File
@@ -1 +0,0 @@
package hostdrivers
+1 -1
View File
@@ -25,7 +25,7 @@ func (self *SKVMHostDriver) GetHostType() string {
return models.HOST_TYPE_HYPERVISOR
}
func (self *SKVMHostDriver) CheckAndSetCacheImage(ctx context.Context, host *models.SHost, storageCache *models.SStoragecache, scimg *models.SStoragecachedimage, task taskman.ITask) error {
func (self *SKVMHostDriver) CheckAndSetCacheImage(ctx context.Context, host *models.SHost, storageCache *models.SStoragecache, task taskman.ITask) error {
params := task.GetParams()
imageId, err := params.GetString("image_id")
if err != nil {
+4
View File
@@ -19,6 +19,7 @@ import (
"yunion.io/x/onecloud/pkg/cloudprovider"
"yunion.io/x/onecloud/pkg/httperrors"
"yunion.io/x/onecloud/pkg/mcclient"
"yunion.io/x/onecloud/pkg/cloudcommon/db/lockman"
)
const (
@@ -485,6 +486,9 @@ func (self *SElasticip) PerformAssociate(ctx context.Context, userCred mcclient.
server := vmObj.(*SGuest)
lockman.LockObject(ctx, server)
defer lockman.ReleaseObject(ctx, server)
if server.PendingDeleted {
return nil, httperrors.NewInvalidStatusError("cannot associate pending delete server")
}
+3 -1
View File
@@ -1637,8 +1637,10 @@ func (self *SGuest) PerformSaveImage(ctx context.Context, userCred mcclient.Toke
properties.Add(jsonutils.NewString(self.OsType), "os_type")
kwargs.Add(properties, "properties")
kwargs.Add(jsonutils.NewBool(restart), "restart")
lockman.LockObject(ctx, disks.Root)
defer lockman.ReleaseObject(ctx, disks.Root)
if imageId, err := disks.Root.PrepareSaveImage(ctx, userCred, kwargs); err != nil {
return nil, err
} else {
@@ -2031,7 +2033,7 @@ func (self *SGuest) CreateDisksOnHost(ctx context.Context, userCred mcclient.Tok
func (self *SGuest) createDiskOnStorage(ctx context.Context, userCred mcclient.TokenCredential, storage *SStorage, diskConfig *SDiskConfig, pendingUsage quotas.IQuota) (*SDisk, error) {
lockman.LockObject(ctx, storage)
defer lockman.LockObject(ctx, storage)
defer lockman.ReleaseObject(ctx, storage)
lockman.LockClass(ctx, QuotaManager, self.ProjectId)
defer lockman.ReleaseClass(ctx, QuotaManager, self.ProjectId)
+1 -1
View File
@@ -11,7 +11,7 @@ import (
type IHostDriver interface {
GetHostType() string
CheckAndSetCacheImage(ctx context.Context, host *SHost, storagecache *SStoragecache, scimg *SStoragecachedimage, task taskman.ITask) error
CheckAndSetCacheImage(ctx context.Context, host *SHost, storagecache *SStoragecache, task taskman.ITask) error
RequestPrepareSaveDiskOnHost(ctx context.Context, host *SHost, disk *SDisk, imageId string, task taskman.ITask) error
RequestSaveUploadImageOnHost(ctx context.Context, host *SHost, disk *SDisk, imageId string, task taskman.ITask, data jsonutils.JSONObject) error
RequestAllocateDiskOnStorage(ctx context.Context, host *SHost, storage *SStorage, disk *SDisk, task taskman.ITask, content *jsonutils.JSONDict) error
+2
View File
@@ -270,8 +270,10 @@ func (manager *SSecurityGroupManager) DelaySync(ctx context.Context, userCred mc
log.Errorf("DelaySync secgroup failed")
} else {
needSync := false
lockman.LockObject(ctx, secgrp)
defer lockman.ReleaseObject(ctx, secgrp)
if secgrp.IsDirty {
if _, err := secgrp.GetModelManager().TableSpec().Update(secgrp, func() error {
secgrp.IsDirty = false
@@ -73,6 +73,7 @@ func (self *GuestBatchCreateTask) onSchedulerRequestFail(ctx context.Context, gu
func (self *GuestBatchCreateTask) onScheduleFail(ctx context.Context, guest *models.SGuest, msg string) {
lockman.LockObject(ctx, guest)
defer lockman.ReleaseObject(ctx, guest)
reason := "No matching resources"
if len(msg) > 0 {
reason = fmt.Sprintf("%s: %s", reason, msg)
@@ -181,8 +181,10 @@ func (self *GuestChangeConfigTask) OnCreateDisksComplete(ctx context.Context, ob
if addMem > 0 {
cancelUsage.Memory = addMem
}
lockman.LockClass(ctx, guest.GetModelManager(), guest.ProjectId)
defer lockman.ReleaseClass(ctx, guest.GetModelManager(), guest.ProjectId)
err = models.QuotaManager.CancelPendingUsage(ctx, self.UserCred, guest.ProjectId, &pendingUsage, &cancelUsage)
if err != nil {
self.SetStageFailed(ctx, err.Error())
@@ -34,7 +34,7 @@ func (self *StorageCacheImageTask) OnInit(ctx context.Context, obj db.IStandalon
self.SetStage("on_image_cache_complete", nil)
host, _ := storageCache.GetHost()
err := host.GetHostDriver().CheckAndSetCacheImage(ctx, host, storageCache, scimg, self)
err := host.GetHostDriver().CheckAndSetCacheImage(ctx, host, storageCache, self)
if err != nil {
errData := taskman.Error2TaskData(err)
self.OnImageCacheCompleteFailed(ctx, storageCache, errData)
@@ -46,8 +46,8 @@ func (self *StorageCacheImageTask) OnImageCacheComplete(ctx context.Context, obj
storageCache := obj.(*models.SStoragecache)
imageId, _ := self.Params.GetString("image_id")
scimg := models.StoragecachedimageManager.Register(ctx, self.UserCred, storageCache.Id, imageId)
extImgId, _ := data.GetString("image_id")
self.OnCacheSucc(ctx, storageCache, imageId, scimg, extImgId)
// extImgId, _ := data.GetString("image_id")
self.OnCacheSucc(ctx, storageCache, imageId, scimg)
}
func (self *StorageCacheImageTask) OnImageCacheCompleteFailed(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
@@ -70,11 +70,8 @@ func (self *StorageCacheImageTask) OnCacheFailed(ctx context.Context, cache *mod
self.SetStageFailed(ctx, err.Error())
}
func (self *StorageCacheImageTask) OnCacheSucc(ctx context.Context, cache *models.SStoragecache, imageId string, scimg *models.SStoragecachedimage, extImgId string) {
func (self *StorageCacheImageTask) OnCacheSucc(ctx context.Context, cache *models.SStoragecache, imageId string, scimg *models.SStoragecachedimage) {
scimg.SetStatus(self.UserCred, models.CACHED_IMAGE_STATUS_READY, "cached")
if len(cache.ExternalId) > 0 && len(extImgId) > 0 && scimg.ExternalId != extImgId {
scimg.SetExternalId(extImgId)
}
models.CachedimageManager.ImageAddRefCount(imageId)
db.OpsLog.LogEvent(cache, db.ACT_CACHED_IMAGE, imageId, self.UserCred)
self.SetStageComplete(ctx, nil)
+2 -1
View File
@@ -5,7 +5,6 @@ import (
"os"
"strings"
"time"
"github.com/aliyun/aliyun-oss-go-sdk/oss"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
@@ -87,12 +86,14 @@ func (self *SStoragecache) GetIImages() ([]cloudprovider.ICloudImage, error) {
}
func (self *SStoragecache) UploadImage(userCred mcclient.TokenCredential, imageId string, osArch, osType, osDist string, extId string, isForce bool) (string, error) {
if len(extId) > 0 {
status, _ := self.region.GetImageStatus(extId)
if status == ImageStatusAvailable && !isForce {
return extId, nil
}
}
return self.uploadImage(userCred, imageId, osArch, osType, osDist, isForce)
}