seaweedfs/weed/cluster/lock_client.go

178 lines
4.7 KiB
Go
Raw Normal View History

2023-06-26 08:38:34 +08:00
package cluster
import (
"context"
"fmt"
2023-06-26 11:30:20 +08:00
"github.com/seaweedfs/seaweedfs/weed/cluster/lock_manager"
2023-06-26 08:38:34 +08:00
"github.com/seaweedfs/seaweedfs/weed/glog"
"github.com/seaweedfs/seaweedfs/weed/pb"
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
"github.com/seaweedfs/seaweedfs/weed/util"
"google.golang.org/grpc"
"time"
)
type LockClient struct {
grpcDialOption grpc.DialOption
maxLockDuration time.Duration
sleepDuration time.Duration
2023-09-17 06:05:38 +08:00
seedFiler pb.ServerAddress
2023-06-26 08:38:34 +08:00
}
2023-09-17 06:05:38 +08:00
func NewLockClient(grpcDialOption grpc.DialOption, seedFiler pb.ServerAddress) *LockClient {
2023-06-26 08:38:34 +08:00
return &LockClient{
grpcDialOption: grpcDialOption,
maxLockDuration: 5 * time.Second,
2023-09-17 06:05:38 +08:00
sleepDuration: 2473 * time.Millisecond,
seedFiler: seedFiler,
2023-06-26 08:38:34 +08:00
}
}
type LiveLock struct {
key string
renewToken string
expireAtNs int64
filer pb.ServerAddress
cancelCh chan struct{}
grpcDialOption grpc.DialOption
2024-02-05 01:20:21 +08:00
isLocked bool
self string
lc *LockClient
owner string
2023-06-26 08:38:34 +08:00
}
2024-02-02 15:01:44 +08:00
// NewShortLivedLock creates a lock with a 5-second duration
func (lc *LockClient) NewShortLivedLock(key string, owner string) (lock *LiveLock) {
2023-09-17 06:05:38 +08:00
lock = &LiveLock{
key: key,
filer: lc.seedFiler,
cancelCh: make(chan struct{}),
2024-02-05 01:20:21 +08:00
expireAtNs: time.Now().Add(5 * time.Second).UnixNano(),
2023-09-17 06:05:38 +08:00
grpcDialOption: lc.grpcDialOption,
2024-02-03 07:54:57 +08:00
self: owner,
2023-09-17 06:05:38 +08:00
lc: lc,
}
2024-02-05 01:20:21 +08:00
lock.retryUntilLocked(5 * time.Second)
2023-09-17 06:05:38 +08:00
return
2023-06-26 10:31:25 +08:00
}
2024-02-03 07:54:57 +08:00
// StartLongLivedLock starts a goroutine to lock the key and returns immediately.
func (lc *LockClient) StartLongLivedLock(key string, owner string, onLockOwnerChange func(newLockOwner string)) (lock *LiveLock) {
2023-06-26 08:38:34 +08:00
lock = &LiveLock{
key: key,
2023-09-17 06:05:38 +08:00
filer: lc.seedFiler,
2023-06-26 08:38:34 +08:00
cancelCh: make(chan struct{}),
2024-02-05 01:20:21 +08:00
expireAtNs: time.Now().Add(lock_manager.LiveLockTTL).UnixNano(),
2023-06-26 08:38:34 +08:00
grpcDialOption: lc.grpcDialOption,
2024-02-03 07:54:57 +08:00
self: owner,
2023-09-17 06:05:38 +08:00
lc: lc,
2023-06-26 08:38:34 +08:00
}
2024-02-02 15:01:44 +08:00
go func() {
2024-02-03 07:54:57 +08:00
isLocked := false
lockOwner := ""
for {
if isLocked {
2024-02-05 01:20:21 +08:00
if err := lock.AttemptToLock(lock_manager.LiveLockTTL); err != nil {
2024-02-03 07:54:57 +08:00
glog.V(0).Infof("Lost lock %s: %v", key, err)
isLocked = false
}
} else {
2024-02-05 01:20:21 +08:00
if err := lock.AttemptToLock(lock_manager.LiveLockTTL); err == nil {
2024-02-03 07:54:57 +08:00
isLocked = true
}
}
if lockOwner != lock.LockOwner() && lock.LockOwner() != "" {
glog.V(0).Infof("Lock owner changed from %s to %s", lockOwner, lock.LockOwner())
onLockOwnerChange(lock.LockOwner())
lockOwner = lock.LockOwner()
}
select {
case <-lock.cancelCh:
return
default:
2024-02-05 01:20:21 +08:00
time.Sleep(lock_manager.RenewInterval)
2024-02-03 07:54:57 +08:00
}
}
2024-02-02 15:01:44 +08:00
}()
2024-01-30 14:46:23 +08:00
return
}
2024-02-02 15:01:44 +08:00
func (lock *LiveLock) retryUntilLocked(lockDuration time.Duration) {
2024-01-30 14:46:23 +08:00
util.RetryUntil("create lock:"+lock.key, func() error {
2024-02-02 15:01:44 +08:00
return lock.AttemptToLock(lockDuration)
2023-06-26 08:38:34 +08:00
}, func(err error) (shouldContinue bool) {
if err != nil {
2024-01-30 14:46:23 +08:00
glog.Warningf("create lock %s: %s", lock.key, err)
2023-06-26 08:38:34 +08:00
}
return lock.renewToken == ""
})
2024-01-30 14:46:23 +08:00
}
2023-06-26 08:38:34 +08:00
2024-02-02 15:01:44 +08:00
func (lock *LiveLock) AttemptToLock(lockDuration time.Duration) error {
2024-01-30 14:46:23 +08:00
errorMessage, err := lock.doLock(lockDuration)
if err != nil {
time.Sleep(time.Second)
return err
2023-06-26 08:38:34 +08:00
}
2024-01-30 14:46:23 +08:00
if errorMessage != "" {
time.Sleep(time.Second)
return fmt.Errorf("%v", errorMessage)
}
lock.isLocked = true
return nil
2023-06-26 08:38:34 +08:00
}
2024-02-02 15:01:44 +08:00
func (lock *LiveLock) StopShortLivedLock() error {
2023-09-17 06:05:38 +08:00
if !lock.isLocked {
return nil
}
2024-02-02 15:01:44 +08:00
defer func() {
lock.isLocked = false
}()
2023-06-26 08:38:34 +08:00
return pb.WithFilerClient(false, 0, lock.filer, lock.grpcDialOption, func(client filer_pb.SeaweedFilerClient) error {
2023-09-17 06:05:38 +08:00
_, err := client.DistributedUnlock(context.Background(), &filer_pb.UnlockRequest{
2023-06-26 08:38:34 +08:00
Name: lock.key,
RenewToken: lock.renewToken,
})
return err
})
}
func (lock *LiveLock) doLock(lockDuration time.Duration) (errorMessage string, err error) {
err = pb.WithFilerClient(false, 0, lock.filer, lock.grpcDialOption, func(client filer_pb.SeaweedFilerClient) error {
2023-09-17 06:05:38 +08:00
resp, err := client.DistributedLock(context.Background(), &filer_pb.LockRequest{
2023-06-26 08:38:34 +08:00
Name: lock.key,
SecondsToLock: int64(lockDuration.Seconds()),
RenewToken: lock.renewToken,
IsMoved: false,
2024-02-03 07:54:57 +08:00
Owner: lock.self,
2023-06-26 08:38:34 +08:00
})
2024-02-03 07:54:57 +08:00
if err == nil && resp != nil {
2023-06-26 08:38:34 +08:00
lock.renewToken = resp.RenewToken
2024-02-03 07:54:57 +08:00
} else {
2024-02-05 01:20:21 +08:00
//this can be retried. Need to remember the last valid renewToken
lock.renewToken = ""
2023-06-26 08:38:34 +08:00
}
if resp != nil {
errorMessage = resp.Error
2024-02-03 07:54:57 +08:00
if resp.LockHostMovedTo != "" {
lock.filer = pb.ServerAddress(resp.LockHostMovedTo)
2023-09-17 06:05:38 +08:00
lock.lc.seedFiler = lock.filer
2023-06-26 08:38:34 +08:00
}
2024-02-03 07:54:57 +08:00
if resp.LockOwner != "" {
lock.owner = resp.LockOwner
2024-02-05 04:47:07 +08:00
// fmt.Printf("lock %s owner: %s\n", lock.key, lock.owner)
2024-02-03 07:54:57 +08:00
} else {
2024-02-05 04:47:07 +08:00
// fmt.Printf("lock %s has no owner\n", lock.key)
2024-02-03 07:54:57 +08:00
lock.owner = ""
}
2023-06-26 08:38:34 +08:00
}
return err
})
return
}
2024-02-03 07:54:57 +08:00
func (lock *LiveLock) LockOwner() string {
return lock.owner
}