1
0
镜像自地址 https://github.com/tuna/tunasync.git 已同步 2025-12-06 14:36:47 +00:00

Add redis backend for db

这个提交包含在:
jiegec
2020-10-13 14:50:19 +08:00
父节点 a2887da2dd
当前提交 a137f0676a
共有 5 个文件被更改,包括 210 次插入5 次删除

查看文件

@@ -5,6 +5,7 @@ import (
"time"
"github.com/boltdb/bolt"
"github.com/go-redis/redis/v8"
. "github.com/tuna/tunasync/internal"
)
@@ -24,6 +25,11 @@ type dbAdapter interface {
Close() error
}
const (
_workerBucketKey = "workers"
_statusBucketKey = "mirror_status"
)
func makeDBAdapter(dbType string, dbFile string) (dbAdapter, error) {
if dbType == "bolt" {
innerDB, err := bolt.Open(dbFile, 0600, &bolt.Options{
@@ -38,6 +44,15 @@ func makeDBAdapter(dbType string, dbFile string) (dbAdapter, error) {
}
err = db.Init()
return &db, err
} else if dbType == "redis" {
innerDB := redis.NewClient(&redis.Options{
Addr: dbFile,
})
db := redisAdapter{
db: innerDB,
}
err := db.Init()
return &db, err
}
// unsupported db-type
return nil, fmt.Errorf("unsupported db-type: %s", dbType)

查看文件

@@ -11,11 +11,6 @@ import (
. "github.com/tuna/tunasync/internal"
)
const (
_workerBucketKey = "workers"
_statusBucketKey = "mirror_status"
)
type boltAdapter struct {
db *bolt.DB
dbFile string

159
manager/db_redis.go 普通文件
查看文件

@@ -0,0 +1,159 @@
package manager
import (
"context"
"encoding/json"
"fmt"
"strings"
"time"
"github.com/go-redis/redis/v8"
. "github.com/tuna/tunasync/internal"
)
type redisAdapter struct {
db *redis.Client
}
var ctx = context.Background()
func (b *redisAdapter) Init() (err error) {
return nil
}
func (b *redisAdapter) ListWorkers() (ws []WorkerStatus, err error) {
var val map[string]string
val, err = b.db.HGetAll(ctx, _workerBucketKey).Result()
if err == nil {
var w WorkerStatus
for _, v := range val {
jsonErr := json.Unmarshal([]byte(v), &w)
if jsonErr != nil {
err = fmt.Errorf("%s; %s", err.Error(), jsonErr)
continue
}
ws = append(ws, w)
}
}
return
}
func (b *redisAdapter) GetWorker(workerID string) (w WorkerStatus, err error) {
var val string
val, err = b.db.HGet(ctx, _workerBucketKey, workerID).Result()
if err == nil {
err = json.Unmarshal([]byte(val), &w)
} else {
err = fmt.Errorf("invalid workerID %s", workerID)
}
return
}
func (b *redisAdapter) DeleteWorker(workerID string) (err error) {
_, err = b.db.HDel(ctx, _workerBucketKey, workerID).Result()
if err != nil {
err = fmt.Errorf("invalid workerID %s", workerID)
}
return
}
func (b *redisAdapter) CreateWorker(w WorkerStatus) (WorkerStatus, error) {
var v []byte
v, err := json.Marshal(w)
if err == nil {
_, err = b.db.HSet(ctx, _workerBucketKey, w.ID, string(v)).Result()
}
return w, err
}
func (b *redisAdapter) RefreshWorker(workerID string) (w WorkerStatus, err error) {
w, err = b.GetWorker(workerID)
if err == nil {
w.LastOnline = time.Now()
w, err = b.CreateWorker(w)
}
return w, err
}
func (b *redisAdapter) UpdateMirrorStatus(workerID, mirrorID string, status MirrorStatus) (MirrorStatus, error) {
id := mirrorID + "/" + workerID
v, err := json.Marshal(status)
if err == nil {
_, err = b.db.HSet(ctx, _statusBucketKey, id, string(v)).Result()
}
return status, err
}
func (b *redisAdapter) GetMirrorStatus(workerID, mirrorID string) (m MirrorStatus, err error) {
id := mirrorID + "/" + workerID
var val string
val, err = b.db.HGet(ctx, _statusBucketKey, id).Result()
if err == nil {
err = json.Unmarshal([]byte(val), &m)
} else {
err = fmt.Errorf("no mirror '%s' exists in worker '%s'", mirrorID, workerID)
}
return
}
func (b *redisAdapter) ListMirrorStatus(workerID string) (ms []MirrorStatus, err error) {
var val map[string]string
val, err = b.db.HGetAll(ctx, _statusBucketKey).Result()
if err == nil {
var m MirrorStatus
for k, v := range val {
if wID := strings.Split(string(k), "/")[1]; wID == workerID {
jsonErr := json.Unmarshal([]byte(v), &m)
if jsonErr != nil {
err = fmt.Errorf("%s; %s", err.Error(), jsonErr)
continue
}
ms = append(ms, m)
}
}
}
return
}
func (b *redisAdapter) ListAllMirrorStatus() (ms []MirrorStatus, err error) {
var val map[string]string
val, err = b.db.HGetAll(ctx, _statusBucketKey).Result()
if err == nil {
var m MirrorStatus
for _, v := range val {
jsonErr := json.Unmarshal([]byte(v), &m)
if jsonErr != nil {
err = fmt.Errorf("%s; %s", err.Error(), jsonErr)
continue
}
ms = append(ms, m)
}
}
return
}
func (b *redisAdapter) FlushDisabledJobs() (err error) {
var val map[string]string
val, err = b.db.HGetAll(ctx, _statusBucketKey).Result()
if err == nil {
var m MirrorStatus
for k, v := range val {
jsonErr := json.Unmarshal([]byte(v), &m)
if jsonErr != nil {
err = fmt.Errorf("%s; %s", err.Error(), jsonErr)
continue
}
if m.Status == Disabled || len(m.Name) == 0 {
_, err = b.db.HDel(ctx, _statusBucketKey, k).Result()
}
}
}
return
}
func (b *redisAdapter) Close() error {
if b.db != nil {
return b.db.Close()
}
return nil
}