redis.go 7.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361
  1. package cache
  2. import (
  3. "context"
  4. "errors"
  5. "fmt"
  6. "strings"
  7. "time"
  8. "github.com/langgenius/dify-plugin-daemon/internal/utils/parser"
  9. "github.com/redis/go-redis/v9"
  10. )
  11. var (
  12. client *redis.Client
  13. ctx = context.Background()
  14. ErrDBNotInit = errors.New("redis client not init")
  15. ErrNotFound = errors.New("key not found")
  16. )
  17. func InitRedisClient(addr, password string) error {
  18. client = redis.NewClient(&redis.Options{
  19. Addr: addr,
  20. Password: password,
  21. })
  22. if _, err := client.Ping(ctx).Result(); err != nil {
  23. return err
  24. }
  25. return nil
  26. }
  27. func Close() error {
  28. if client == nil {
  29. return ErrDBNotInit
  30. }
  31. return client.Close()
  32. }
  33. func getCmdable(context ...redis.Cmdable) redis.Cmdable {
  34. if len(context) > 0 {
  35. return context[0]
  36. }
  37. return client
  38. }
  39. func serialKey(keys ...string) string {
  40. return strings.Join(append(
  41. []string{"plugin_daemon"},
  42. keys...,
  43. ), ":")
  44. }
  45. func Store(key string, value any, time time.Duration, context ...redis.Cmdable) error {
  46. if client == nil {
  47. return ErrDBNotInit
  48. }
  49. if _, ok := value.(string); !ok {
  50. value = parser.MarshalJson(value)
  51. }
  52. return getCmdable(context...).Set(ctx, serialKey(key), value, time).Err()
  53. }
  54. func Get[T any](key string, context ...redis.Cmdable) (*T, error) {
  55. if client == nil {
  56. return nil, ErrDBNotInit
  57. }
  58. val, err := getCmdable(context...).Get(ctx, serialKey(key)).Result()
  59. if err != nil {
  60. if err == redis.Nil {
  61. return nil, ErrNotFound
  62. }
  63. return nil, err
  64. }
  65. if val == "" {
  66. return nil, ErrNotFound
  67. }
  68. result, err := parser.UnmarshalJson[T](val)
  69. return &result, err
  70. }
  71. func GetString(key string, context ...redis.Cmdable) (string, error) {
  72. if client == nil {
  73. return "", ErrDBNotInit
  74. }
  75. v, err := getCmdable(context...).Get(ctx, serialKey(key)).Result()
  76. if err != nil {
  77. if err == redis.Nil {
  78. return "", ErrNotFound
  79. }
  80. }
  81. return v, err
  82. }
  83. func Del(key string, context ...redis.Cmdable) error {
  84. if client == nil {
  85. return ErrDBNotInit
  86. }
  87. _, err := getCmdable(context...).Del(ctx, serialKey(key)).Result()
  88. if err != nil {
  89. if err == redis.Nil {
  90. return ErrNotFound
  91. }
  92. return err
  93. }
  94. return nil
  95. }
  96. func Exist(key string, context ...redis.Cmdable) (int64, error) {
  97. if client == nil {
  98. return 0, ErrDBNotInit
  99. }
  100. return getCmdable(context...).Exists(ctx, serialKey(key)).Result()
  101. }
  102. func Increase(key string, context ...redis.Cmdable) (int64, error) {
  103. if client == nil {
  104. return 0, ErrDBNotInit
  105. }
  106. num, err := getCmdable(context...).Incr(ctx, serialKey(key)).Result()
  107. if err != nil {
  108. if err == redis.Nil {
  109. return 0, ErrNotFound
  110. }
  111. return 0, err
  112. }
  113. return num, nil
  114. }
  115. func Decrease(key string, context ...redis.Cmdable) (int64, error) {
  116. if client == nil {
  117. return 0, ErrDBNotInit
  118. }
  119. return getCmdable(context...).Decr(ctx, serialKey(key)).Result()
  120. }
  121. func SetExpire(key string, time time.Duration, context ...redis.Cmdable) error {
  122. if client == nil {
  123. return ErrDBNotInit
  124. }
  125. return getCmdable(context...).Expire(ctx, serialKey(key), time).Err()
  126. }
  127. func SetMapField(key string, v map[string]any, context ...redis.Cmdable) error {
  128. if client == nil {
  129. return ErrDBNotInit
  130. }
  131. return getCmdable(context...).HMSet(ctx, serialKey(key), v).Err()
  132. }
  133. func SetMapOneField(key string, field string, value any, context ...redis.Cmdable) error {
  134. if client == nil {
  135. return ErrDBNotInit
  136. }
  137. if _, ok := value.(string); !ok {
  138. value = parser.MarshalJson(value)
  139. }
  140. return getCmdable(context...).HSet(ctx, serialKey(key), field, value).Err()
  141. }
  142. func GetMapField[T any](key string, field string, context ...redis.Cmdable) (*T, error) {
  143. if client == nil {
  144. return nil, ErrDBNotInit
  145. }
  146. val, err := getCmdable(context...).HGet(ctx, serialKey(key), field).Result()
  147. if err != nil {
  148. if err == redis.Nil {
  149. return nil, ErrNotFound
  150. }
  151. return nil, err
  152. }
  153. result, err := parser.UnmarshalJson[T](val)
  154. return &result, err
  155. }
  156. func DelMapField(key string, field string, context ...redis.Cmdable) error {
  157. if client == nil {
  158. return ErrDBNotInit
  159. }
  160. return getCmdable(context...).HDel(ctx, serialKey(key), field).Err()
  161. }
  162. func GetMap[V any](key string, context ...redis.Cmdable) (map[string]V, error) {
  163. if client == nil {
  164. return nil, ErrDBNotInit
  165. }
  166. val, err := getCmdable(context...).HGetAll(ctx, serialKey(key)).Result()
  167. if err != nil {
  168. if err == redis.Nil {
  169. return nil, ErrNotFound
  170. }
  171. return nil, err
  172. }
  173. result := make(map[string]V)
  174. for k, v := range val {
  175. value, err := parser.UnmarshalJson[V](v)
  176. if err != nil {
  177. continue
  178. }
  179. result[k] = value
  180. }
  181. return result, nil
  182. }
  183. func ScanMap[V any](key string, prefix string, context ...redis.Cmdable) (map[string]V, error) {
  184. if client == nil {
  185. return nil, ErrDBNotInit
  186. }
  187. result := make(map[string]V)
  188. ScanMapAsync[V](key, prefix, func(m map[string]V) error {
  189. for k, v := range m {
  190. result[k] = v
  191. }
  192. return nil
  193. })
  194. return result, nil
  195. }
  196. func ScanMapAsync[V any](key string, prefix string, fn func(map[string]V) error, context ...redis.Cmdable) error {
  197. if client == nil {
  198. return ErrDBNotInit
  199. }
  200. cursor := uint64(0)
  201. for {
  202. kvs, new_cursor, err := getCmdable(context...).
  203. HScan(ctx, serialKey(key), cursor, fmt.Sprintf("%s*", prefix), 32).
  204. Result()
  205. if err != nil {
  206. return err
  207. }
  208. result := make(map[string]V)
  209. for i := 0; i < len(kvs); i += 2 {
  210. value, err := parser.UnmarshalJson[V](kvs[i+1])
  211. if err != nil {
  212. continue
  213. }
  214. result[kvs[i]] = value
  215. }
  216. if err := fn(result); err != nil {
  217. return err
  218. }
  219. if new_cursor == 0 {
  220. break
  221. }
  222. cursor = new_cursor
  223. }
  224. return nil
  225. }
  226. func SetNX[T any](key string, value T, expire time.Duration, context ...redis.Cmdable) (bool, error) {
  227. if client == nil {
  228. return false, ErrDBNotInit
  229. }
  230. return getCmdable(context...).SetNX(ctx, serialKey(key), value, expire).Result()
  231. }
  232. var (
  233. ErrLockTimeout = errors.New("lock timeout")
  234. )
  235. // Lock key, expire time takes responsibility for expiration time
  236. // try_lock_timeout takes responsibility for the timeout of trying to lock
  237. func Lock(key string, expire time.Duration, try_lock_timeout time.Duration, context ...redis.Cmdable) error {
  238. if client == nil {
  239. return ErrDBNotInit
  240. }
  241. const LOCK_DURATION = 20 * time.Millisecond
  242. ticker := time.NewTicker(LOCK_DURATION)
  243. defer ticker.Stop()
  244. for range ticker.C {
  245. if _, err := getCmdable(context...).SetNX(ctx, serialKey(key), "1", expire).Result(); err == nil {
  246. return nil
  247. }
  248. try_lock_timeout -= LOCK_DURATION
  249. if try_lock_timeout <= 0 {
  250. return ErrLockTimeout
  251. }
  252. }
  253. return nil
  254. }
  255. func Unlock(key string, context ...redis.Cmdable) error {
  256. if client == nil {
  257. return ErrDBNotInit
  258. }
  259. return getCmdable(context...).Del(ctx, serialKey(key)).Err()
  260. }
  261. func Expire(key string, time time.Duration, context ...redis.Cmdable) (bool, error) {
  262. if client == nil {
  263. return false, ErrDBNotInit
  264. }
  265. return getCmdable(context...).Expire(ctx, serialKey(key), time).Result()
  266. }
  267. func Transaction(fn func(redis.Pipeliner) error) error {
  268. if client == nil {
  269. return ErrDBNotInit
  270. }
  271. return client.Watch(ctx, func(tx *redis.Tx) error {
  272. _, err := tx.TxPipelined(ctx, func(p redis.Pipeliner) error {
  273. return fn(p)
  274. })
  275. if err == redis.Nil {
  276. return nil
  277. }
  278. return err
  279. })
  280. }