mirror of
https://github.com/openimsdk/open-im-server.git
synced 2026-09-04 14:48:15 +08:00
feat(redis): implement Redis-based locking mechanism for cron tasks
This commit is contained in:
parent
bfe8bcd9ad
commit
5882597ee1
@ -3,8 +3,12 @@ package cron
|
||||
import (
|
||||
"context"
|
||||
|
||||
"github.com/robfig/cron/v3"
|
||||
"google.golang.org/grpc"
|
||||
|
||||
"github.com/openimsdk/open-im-server/v3/pkg/common/config"
|
||||
disetcd "github.com/openimsdk/open-im-server/v3/pkg/common/discovery/etcd"
|
||||
"github.com/openimsdk/open-im-server/v3/pkg/dbbuild"
|
||||
pbconversation "github.com/openimsdk/protocol/conversation"
|
||||
"github.com/openimsdk/protocol/msg"
|
||||
"github.com/openimsdk/protocol/third"
|
||||
@ -14,14 +18,13 @@ import (
|
||||
"github.com/openimsdk/tools/log"
|
||||
"github.com/openimsdk/tools/mcontext"
|
||||
"github.com/openimsdk/tools/utils/runtimeenv"
|
||||
"github.com/robfig/cron/v3"
|
||||
"google.golang.org/grpc"
|
||||
)
|
||||
|
||||
type Config struct {
|
||||
CronTask config.CronTask
|
||||
Share config.Share
|
||||
Discovery config.Discovery
|
||||
CronTask config.CronTask
|
||||
Share config.Share
|
||||
Discovery config.Discovery
|
||||
RedisConfig config.Redis
|
||||
}
|
||||
|
||||
func Start(ctx context.Context, conf *Config, client discovery.SvcDiscoveryRegistry, service grpc.ServiceRegistrar) error {
|
||||
@ -60,10 +63,13 @@ func Start(ctx context.Context, conf *Config, client discovery.SvcDiscoveryRegis
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
if locker == nil {
|
||||
locker = emptyLocker{}
|
||||
} else {
|
||||
builder := dbbuild.NewBuilder(nil, &conf.RedisConfig)
|
||||
rdb, err := builder.Redis(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
locker = NewRedisLocker(rdb)
|
||||
}
|
||||
|
||||
srv := &cronServer{
|
||||
@ -95,16 +101,6 @@ func Start(ctx context.Context, conf *Config, client discovery.SvcDiscoveryRegis
|
||||
return nil
|
||||
}
|
||||
|
||||
type Locker interface {
|
||||
ExecuteWithLock(ctx context.Context, taskName string, task func())
|
||||
}
|
||||
|
||||
type emptyLocker struct{}
|
||||
|
||||
func (emptyLocker) ExecuteWithLock(ctx context.Context, taskName string, task func()) {
|
||||
task()
|
||||
}
|
||||
|
||||
type cronServer struct {
|
||||
ctx context.Context
|
||||
config *Config
|
||||
|
||||
@ -6,13 +6,10 @@ import (
|
||||
"os"
|
||||
"time"
|
||||
|
||||
"github.com/openimsdk/tools/log"
|
||||
clientv3 "go.etcd.io/etcd/client/v3"
|
||||
"go.etcd.io/etcd/client/v3/concurrency"
|
||||
)
|
||||
|
||||
const (
|
||||
lockLeaseTTL = 300
|
||||
"github.com/openimsdk/tools/log"
|
||||
)
|
||||
|
||||
type EtcdLocker struct {
|
||||
@ -35,7 +32,7 @@ func NewEtcdLocker(client *clientv3.Client) (*EtcdLocker, error) {
|
||||
}
|
||||
|
||||
func (e *EtcdLocker) ExecuteWithLock(ctx context.Context, taskName string, task func()) {
|
||||
session, err := concurrency.NewSession(e.client, concurrency.WithTTL(lockLeaseTTL))
|
||||
session, err := concurrency.NewSession(e.client, concurrency.WithTTL(int(lockLeaseTTL/time.Second)))
|
||||
if err != nil {
|
||||
log.ZWarn(ctx, "Failed to create etcd session", err,
|
||||
"taskName", taskName,
|
||||
14
internal/tools/cron/locker.go
Normal file
14
internal/tools/cron/locker.go
Normal file
@ -0,0 +1,14 @@
|
||||
package cron
|
||||
|
||||
import (
|
||||
"context"
|
||||
"time"
|
||||
)
|
||||
|
||||
const (
|
||||
lockLeaseTTL = time.Second * 300
|
||||
)
|
||||
|
||||
type Locker interface {
|
||||
ExecuteWithLock(ctx context.Context, taskName string, task func())
|
||||
}
|
||||
68
internal/tools/cron/redis_locker.go
Normal file
68
internal/tools/cron/redis_locker.go
Normal file
@ -0,0 +1,68 @@
|
||||
package cron
|
||||
|
||||
import (
|
||||
"context"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/redis/go-redis/v9"
|
||||
|
||||
"github.com/openimsdk/tools/log"
|
||||
)
|
||||
|
||||
func NewRedisLocker(client redis.UniversalClient) *RedisLocker {
|
||||
return &RedisLocker{
|
||||
client: client,
|
||||
script: redis.NewScript(strings.TrimSpace(`
|
||||
if redis.call("get", KEYS[1]) == ARGV[1] then
|
||||
return redis.call("del", KEYS[1])
|
||||
else
|
||||
return 0
|
||||
end
|
||||
`)),
|
||||
}
|
||||
}
|
||||
|
||||
type RedisLocker struct {
|
||||
client redis.UniversalClient
|
||||
script *redis.Script
|
||||
}
|
||||
|
||||
func (e *RedisLocker) getKey(name string) string {
|
||||
return "CRON_LOCKED:" + name
|
||||
}
|
||||
|
||||
func (e *RedisLocker) lock(ctx context.Context, name string, owner string) (bool, error) {
|
||||
ctx, cancel := context.WithTimeout(ctx, time.Second)
|
||||
defer cancel()
|
||||
return e.client.SetNX(ctx, e.getKey(name), owner, lockLeaseTTL).Result()
|
||||
}
|
||||
|
||||
func (e *RedisLocker) unlock(ctx context.Context, name string, owner string) error {
|
||||
ctx, cancel := context.WithTimeout(ctx, time.Second)
|
||||
defer cancel()
|
||||
return e.script.Run(ctx, e.client, []string{e.getKey(name)}, owner).Err()
|
||||
}
|
||||
|
||||
func (e *RedisLocker) ExecuteWithLock(ctx context.Context, taskName string, task func()) {
|
||||
owner := uuid.New().String()
|
||||
ok, err := e.lock(ctx, taskName, owner)
|
||||
if err != nil {
|
||||
log.ZWarn(ctx, "cron lock get lock", err, "taskName", taskName)
|
||||
return
|
||||
}
|
||||
log.ZDebug(ctx, "cron lock get lock", "taskName", taskName, "ok", ok, "owner", owner)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
defer func() {
|
||||
err := e.unlock(ctx, taskName, owner)
|
||||
if err == nil {
|
||||
log.ZDebug(ctx, "cron lock unlock", "taskName", taskName, "owner", owner)
|
||||
} else {
|
||||
log.ZWarn(ctx, "cron lock unlock", err, "taskName", taskName, "owner", owner)
|
||||
}
|
||||
}()
|
||||
task()
|
||||
}
|
||||
@ -17,12 +17,13 @@ package cmd
|
||||
import (
|
||||
"context"
|
||||
|
||||
"github.com/spf13/cobra"
|
||||
|
||||
"github.com/openimsdk/open-im-server/v3/internal/tools/cron"
|
||||
"github.com/openimsdk/open-im-server/v3/pkg/common/config"
|
||||
"github.com/openimsdk/open-im-server/v3/pkg/common/startrpc"
|
||||
"github.com/openimsdk/open-im-server/v3/version"
|
||||
"github.com/openimsdk/tools/system/program"
|
||||
"github.com/spf13/cobra"
|
||||
)
|
||||
|
||||
type CronTaskCmd struct {
|
||||
@ -39,6 +40,7 @@ func NewCronTaskCmd() *CronTaskCmd {
|
||||
config.OpenIMCronTaskCfgFileName: &cronTaskConfig.CronTask,
|
||||
config.ShareFileName: &cronTaskConfig.Share,
|
||||
config.DiscoveryConfigFilename: &cronTaskConfig.Discovery,
|
||||
config.RedisConfigFileName: &cronTaskConfig.RedisConfig,
|
||||
}
|
||||
ret.RootCmd = NewRootCmd(program.GetProcessName(), WithConfigMap(ret.configMap))
|
||||
ret.ctx = context.WithValue(context.Background(), "version", version.Version)
|
||||
|
||||
@ -47,6 +47,6 @@ func TestStandaloneGatewayRedisGetGatewayAddrs(t *testing.T) {
|
||||
|
||||
addrs, err := cache.GetGatewayAddrs(context.Background())
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, []string{"127.0.0.1:10001", "127.0.0.1:10003"}, addrs)
|
||||
assert.ElementsMatch(t, []string{"127.0.0.1:10001", "127.0.0.1:10003"}, addrs)
|
||||
assert.NoError(t, mock.ExpectationsWereMet())
|
||||
}
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user