| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528 | package cacheimport (	"context"	"crypto/tls"	"errors"	"strings"	"time"	"github.com/langgenius/dify-plugin-daemon/internal/utils/log"	"github.com/langgenius/dify-plugin-daemon/internal/utils/parser"	"github.com/redis/go-redis/v9")var (	client *redis.Client	ctx    = context.Background()	ErrDBNotInit = errors.New("redis client not init")	ErrNotFound  = errors.New("key not found"))func getRedisOptions(addr, password string, useSsl bool) *redis.Options {	opts := &redis.Options{		Addr:     addr,		Password: password,	}	if useSsl {		opts.TLSConfig = &tls.Config{}	}	return opts}func InitRedisClient(addr, password string, useSsl bool) error {	opts := getRedisOptions(addr, password, useSsl)	client = redis.NewClient(opts)	if _, err := client.Ping(ctx).Result(); err != nil {		return err	}	return nil}// Close the redis clientfunc Close() error {	if client == nil {		return ErrDBNotInit	}	return client.Close()}func getCmdable(context ...redis.Cmdable) redis.Cmdable {	if len(context) > 0 {		return context[0]	}	return client}func serialKey(keys ...string) string {	return strings.Join(append(		[]string{"plugin_daemon"},		keys...,	), ":")}// Store the key-value pairfunc Store(key string, value any, time time.Duration, context ...redis.Cmdable) error {	return store(serialKey(key), value, time, context...)}// store the key-value pair, without serialKeyfunc store(key string, value any, time time.Duration, context ...redis.Cmdable) error {	if client == nil {		return ErrDBNotInit	}	if _, ok := value.(string); !ok {		var err error		value, err = parser.MarshalCBOR(value)		if err != nil {			return err		}	}	return getCmdable(context...).Set(ctx, key, value, time).Err()}// Get the value with keyfunc Get[T any](key string, context ...redis.Cmdable) (*T, error) {	return get[T](serialKey(key), context...)}func get[T any](key string, context ...redis.Cmdable) (*T, error) {	if client == nil {		return nil, ErrDBNotInit	}	val, err := getCmdable(context...).Get(ctx, key).Bytes()	if err != nil {		if err == redis.Nil {			return nil, ErrNotFound		}		return nil, err	}	if len(val) == 0 {		return nil, ErrNotFound	}	result, err := parser.UnmarshalCBOR[T](val)	return &result, err}// GetString get the string with keyfunc GetString(key string, context ...redis.Cmdable) (string, error) {	if client == nil {		return "", ErrDBNotInit	}	v, err := getCmdable(context...).Get(ctx, serialKey(key)).Result()	if err != nil {		if err == redis.Nil {			return "", ErrNotFound		}	}	return v, err}// Del the keyfunc Del(key string, context ...redis.Cmdable) error {	return del(serialKey(key), context...)}func del(key string, context ...redis.Cmdable) error {	if client == nil {		return ErrDBNotInit	}	_, err := getCmdable(context...).Del(ctx, key).Result()	if err != nil {		if err == redis.Nil {			return ErrNotFound		}		return err	}	return nil}// Exist check the key exist or notfunc Exist(key string, context ...redis.Cmdable) (int64, error) {	if client == nil {		return 0, ErrDBNotInit	}	return getCmdable(context...).Exists(ctx, serialKey(key)).Result()}// Increase the keyfunc Increase(key string, context ...redis.Cmdable) (int64, error) {	if client == nil {		return 0, ErrDBNotInit	}	num, err := getCmdable(context...).Incr(ctx, serialKey(key)).Result()	if err != nil {		if err == redis.Nil {			return 0, ErrNotFound		}		return 0, err	}	return num, nil}// Decrease the keyfunc Decrease(key string, context ...redis.Cmdable) (int64, error) {	if client == nil {		return 0, ErrDBNotInit	}	return getCmdable(context...).Decr(ctx, serialKey(key)).Result()}// SetExpire set the expire time for the keyfunc SetExpire(key string, time time.Duration, context ...redis.Cmdable) error {	if client == nil {		return ErrDBNotInit	}	return getCmdable(context...).Expire(ctx, serialKey(key), time).Err()}// SetMapField set the map field with keyfunc SetMapField(key string, v map[string]any, context ...redis.Cmdable) error {	if client == nil {		return ErrDBNotInit	}	return getCmdable(context...).HMSet(ctx, serialKey(key), v).Err()}// SetMapOneField set the map field with keyfunc SetMapOneField(key string, field string, value any, context ...redis.Cmdable) error {	if client == nil {		return ErrDBNotInit	}	if _, ok := value.(string); !ok {		value = parser.MarshalJson(value)	}	return getCmdable(context...).HSet(ctx, serialKey(key), field, value).Err()}// GetMapField get the map field with keyfunc GetMapField[T any](key string, field string, context ...redis.Cmdable) (*T, error) {	if client == nil {		return nil, ErrDBNotInit	}	val, err := getCmdable(context...).HGet(ctx, serialKey(key), field).Result()	if err != nil {		if err == redis.Nil {			return nil, ErrNotFound		}		return nil, err	}	result, err := parser.UnmarshalJson[T](val)	return &result, err}// GetMapFieldString get the stringfunc GetMapFieldString(key string, field string, context ...redis.Cmdable) (string, error) {	if client == nil {		return "", ErrDBNotInit	}	val, err := getCmdable(context...).HGet(ctx, serialKey(key), field).Result()	if err != nil {		if err == redis.Nil {			return "", ErrNotFound		}		return "", err	}	return val, nil}// DelMapField delete the map field with keyfunc DelMapField(key string, field string, context ...redis.Cmdable) error {	if client == nil {		return ErrDBNotInit	}	return getCmdable(context...).HDel(ctx, serialKey(key), field).Err()}// GetMap get the map with keyfunc GetMap[V any](key string, context ...redis.Cmdable) (map[string]V, error) {	if client == nil {		return nil, ErrDBNotInit	}	val, err := getCmdable(context...).HGetAll(ctx, serialKey(key)).Result()	if err != nil {		if err == redis.Nil {			return nil, ErrNotFound		}		return nil, err	}	result := make(map[string]V)	for k, v := range val {		value, err := parser.UnmarshalJson[V](v)		if err != nil {			continue		}		result[k] = value	}	return result, nil}// ScanKeys scan the keys with match patternfunc ScanKeys(match string, context ...redis.Cmdable) ([]string, error) {	if client == nil {		return nil, ErrDBNotInit	}	result := make([]string, 0)	if err := ScanKeysAsync(match, func(keys []string) error {		result = append(result, keys...)		return nil	}); err != nil {		return nil, err	}	return result, nil}// ScanKeysAsync scan the keys with match pattern, format like "key*"func ScanKeysAsync(match string, fn func([]string) error, context ...redis.Cmdable) error {	if client == nil {		return ErrDBNotInit	}	cursor := uint64(0)	for {		keys, newCursor, err := getCmdable(context...).Scan(ctx, cursor, match, 32).Result()		if err != nil {			return err		}		if err := fn(keys); err != nil {			return err		}		if newCursor == 0 {			break		}		cursor = newCursor	}	return nil}// ScanMap scan the map with match pattern, format like "key*"func ScanMap[V any](key string, match string, context ...redis.Cmdable) (map[string]V, error) {	if client == nil {		return nil, ErrDBNotInit	}	result := make(map[string]V)	ScanMapAsync[V](key, match, func(m map[string]V) error {		for k, v := range m {			result[k] = v		}		return nil	})	return result, nil}// ScanMapAsync scan the map with match pattern, format like "key*"func ScanMapAsync[V any](key string, match string, fn func(map[string]V) error, context ...redis.Cmdable) error {	if client == nil {		return ErrDBNotInit	}	cursor := uint64(0)	for {		kvs, newCursor, err := getCmdable(context...).			HScan(ctx, serialKey(key), cursor, match, 32).			Result()		if err != nil {			return err		}		result := make(map[string]V)		for i := 0; i < len(kvs); i += 2 {			value, err := parser.UnmarshalJson[V](kvs[i+1])			if err != nil {				continue			}			result[kvs[i]] = value		}		if err := fn(result); err != nil {			return err		}		if newCursor == 0 {			break		}		cursor = newCursor	}	return nil}// SetNX set the key-value pair with expire timefunc SetNX[T any](key string, value T, expire time.Duration, context ...redis.Cmdable) (bool, error) {	if client == nil {		return false, ErrDBNotInit	}	// marshal the value	bytes, err := parser.MarshalCBOR(value)	if err != nil {		return false, err	}	return getCmdable(context...).SetNX(ctx, serialKey(key), bytes, expire).Result()}var (	ErrLockTimeout = errors.New("lock timeout"))// Lock key, expire time takes responsibility for expiration time// try_lock_timeout takes responsibility for the timeout of trying to lockfunc Lock(key string, expire time.Duration, tryLockTimeout time.Duration, context ...redis.Cmdable) error {	if client == nil {		return ErrDBNotInit	}	const LOCK_DURATION = 20 * time.Millisecond	ticker := time.NewTicker(LOCK_DURATION)	defer ticker.Stop()	for range ticker.C {		if _, err := getCmdable(context...).SetNX(ctx, serialKey(key), "1", expire).Result(); err == nil {			return nil		}		tryLockTimeout -= LOCK_DURATION		if tryLockTimeout <= 0 {			return ErrLockTimeout		}	}	return nil}func Unlock(key string, context ...redis.Cmdable) error {	if client == nil {		return ErrDBNotInit	}	return getCmdable(context...).Del(ctx, serialKey(key)).Err()}func Expire(key string, time time.Duration, context ...redis.Cmdable) (bool, error) {	if client == nil {		return false, ErrDBNotInit	}	return getCmdable(context...).Expire(ctx, serialKey(key), time).Result()}func Transaction(fn func(redis.Pipeliner) error) error {	if client == nil {		return ErrDBNotInit	}	return client.Watch(ctx, func(tx *redis.Tx) error {		_, err := tx.TxPipelined(ctx, func(p redis.Pipeliner) error {			return fn(p)		})		if err == redis.Nil {			return nil		}		return err	})}func Publish(channel string, message any, context ...redis.Cmdable) error {	if client == nil {		return ErrDBNotInit	}	if _, ok := message.(string); !ok {		message = parser.MarshalJson(message)	}	return getCmdable(context...).Publish(ctx, channel, message).Err()}func Subscribe[T any](channel string) (<-chan T, func()) {	pubsub := client.Subscribe(ctx, channel)	ch := make(chan T)	connectionEstablished := make(chan bool)	go func() {		defer close(ch)		defer close(connectionEstablished)		alive := true		for alive {			iface, err := pubsub.Receive(context.Background())			if err != nil {				log.Error("failed to receive message from redis: %s, will retry in 1 second", err.Error())				time.Sleep(1 * time.Second)				continue			}			switch data := iface.(type) {			case *redis.Subscription:				connectionEstablished <- true			case *redis.Message:				v, err := parser.UnmarshalJson[T](data.Payload)				if err != nil {					continue				}				ch <- v			case *redis.Pong:			default:				alive = false			}		}	}()	// wait for the connection to be established	<-connectionEstablished	return ch, func() {		pubsub.Close()	}}
 |