redis.go 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528
  1. package cache
  2. import (
  3. "context"
  4. "crypto/tls"
  5. "errors"
  6. "strings"
  7. "time"
  8. "github.com/langgenius/dify-plugin-daemon/internal/utils/log"
  9. "github.com/langgenius/dify-plugin-daemon/internal/utils/parser"
  10. "github.com/redis/go-redis/v9"
  11. )
  12. var (
  13. client *redis.Client
  14. ctx = context.Background()
  15. ErrDBNotInit = errors.New("redis client not init")
  16. ErrNotFound = errors.New("key not found")
  17. )
  18. func getRedisOptions(addr, password string, useSsl bool) *redis.Options {
  19. opts := &redis.Options{
  20. Addr: addr,
  21. Password: password,
  22. }
  23. if useSsl {
  24. opts.TLSConfig = &tls.Config{}
  25. }
  26. return opts
  27. }
  28. func InitRedisClient(addr, password string, useSsl bool) error {
  29. opts := getRedisOptions(addr, password, useSsl)
  30. client = redis.NewClient(opts)
  31. if _, err := client.Ping(ctx).Result(); err != nil {
  32. return err
  33. }
  34. return nil
  35. }
  36. // Close the redis client
  37. func Close() error {
  38. if client == nil {
  39. return ErrDBNotInit
  40. }
  41. return client.Close()
  42. }
  43. func getCmdable(context ...redis.Cmdable) redis.Cmdable {
  44. if len(context) > 0 {
  45. return context[0]
  46. }
  47. return client
  48. }
  49. func serialKey(keys ...string) string {
  50. return strings.Join(append(
  51. []string{"plugin_daemon"},
  52. keys...,
  53. ), ":")
  54. }
  55. // Store the key-value pair
  56. func Store(key string, value any, time time.Duration, context ...redis.Cmdable) error {
  57. return store(serialKey(key), value, time, context...)
  58. }
  59. // store the key-value pair, without serialKey
  60. func store(key string, value any, time time.Duration, context ...redis.Cmdable) error {
  61. if client == nil {
  62. return ErrDBNotInit
  63. }
  64. if _, ok := value.(string); !ok {
  65. var err error
  66. value, err = parser.MarshalCBOR(value)
  67. if err != nil {
  68. return err
  69. }
  70. }
  71. return getCmdable(context...).Set(ctx, key, value, time).Err()
  72. }
  73. // Get the value with key
  74. func Get[T any](key string, context ...redis.Cmdable) (*T, error) {
  75. return get[T](serialKey(key), context...)
  76. }
  77. func get[T any](key string, context ...redis.Cmdable) (*T, error) {
  78. if client == nil {
  79. return nil, ErrDBNotInit
  80. }
  81. val, err := getCmdable(context...).Get(ctx, key).Bytes()
  82. if err != nil {
  83. if err == redis.Nil {
  84. return nil, ErrNotFound
  85. }
  86. return nil, err
  87. }
  88. if len(val) == 0 {
  89. return nil, ErrNotFound
  90. }
  91. result, err := parser.UnmarshalCBOR[T](val)
  92. return &result, err
  93. }
  94. // GetString get the string with key
  95. func GetString(key string, context ...redis.Cmdable) (string, error) {
  96. if client == nil {
  97. return "", ErrDBNotInit
  98. }
  99. v, err := getCmdable(context...).Get(ctx, serialKey(key)).Result()
  100. if err != nil {
  101. if err == redis.Nil {
  102. return "", ErrNotFound
  103. }
  104. }
  105. return v, err
  106. }
  107. // Del the key
  108. func Del(key string, context ...redis.Cmdable) error {
  109. return del(serialKey(key), context...)
  110. }
  111. func del(key string, context ...redis.Cmdable) error {
  112. if client == nil {
  113. return ErrDBNotInit
  114. }
  115. _, err := getCmdable(context...).Del(ctx, key).Result()
  116. if err != nil {
  117. if err == redis.Nil {
  118. return ErrNotFound
  119. }
  120. return err
  121. }
  122. return nil
  123. }
  124. // Exist check the key exist or not
  125. func Exist(key string, context ...redis.Cmdable) (int64, error) {
  126. if client == nil {
  127. return 0, ErrDBNotInit
  128. }
  129. return getCmdable(context...).Exists(ctx, serialKey(key)).Result()
  130. }
  131. // Increase the key
  132. func Increase(key string, context ...redis.Cmdable) (int64, error) {
  133. if client == nil {
  134. return 0, ErrDBNotInit
  135. }
  136. num, err := getCmdable(context...).Incr(ctx, serialKey(key)).Result()
  137. if err != nil {
  138. if err == redis.Nil {
  139. return 0, ErrNotFound
  140. }
  141. return 0, err
  142. }
  143. return num, nil
  144. }
  145. // Decrease the key
  146. func Decrease(key string, context ...redis.Cmdable) (int64, error) {
  147. if client == nil {
  148. return 0, ErrDBNotInit
  149. }
  150. return getCmdable(context...).Decr(ctx, serialKey(key)).Result()
  151. }
  152. // SetExpire set the expire time for the key
  153. func SetExpire(key string, time time.Duration, context ...redis.Cmdable) error {
  154. if client == nil {
  155. return ErrDBNotInit
  156. }
  157. return getCmdable(context...).Expire(ctx, serialKey(key), time).Err()
  158. }
  159. // SetMapField set the map field with key
  160. func SetMapField(key string, v map[string]any, context ...redis.Cmdable) error {
  161. if client == nil {
  162. return ErrDBNotInit
  163. }
  164. return getCmdable(context...).HMSet(ctx, serialKey(key), v).Err()
  165. }
  166. // SetMapOneField set the map field with key
  167. func SetMapOneField(key string, field string, value any, context ...redis.Cmdable) error {
  168. if client == nil {
  169. return ErrDBNotInit
  170. }
  171. if _, ok := value.(string); !ok {
  172. value = parser.MarshalJson(value)
  173. }
  174. return getCmdable(context...).HSet(ctx, serialKey(key), field, value).Err()
  175. }
  176. // GetMapField get the map field with key
  177. func GetMapField[T any](key string, field string, context ...redis.Cmdable) (*T, error) {
  178. if client == nil {
  179. return nil, ErrDBNotInit
  180. }
  181. val, err := getCmdable(context...).HGet(ctx, serialKey(key), field).Result()
  182. if err != nil {
  183. if err == redis.Nil {
  184. return nil, ErrNotFound
  185. }
  186. return nil, err
  187. }
  188. result, err := parser.UnmarshalJson[T](val)
  189. return &result, err
  190. }
  191. // GetMapFieldString get the string
  192. func GetMapFieldString(key string, field string, context ...redis.Cmdable) (string, error) {
  193. if client == nil {
  194. return "", ErrDBNotInit
  195. }
  196. val, err := getCmdable(context...).HGet(ctx, serialKey(key), field).Result()
  197. if err != nil {
  198. if err == redis.Nil {
  199. return "", ErrNotFound
  200. }
  201. return "", err
  202. }
  203. return val, nil
  204. }
  205. // DelMapField delete the map field with key
  206. func DelMapField(key string, field string, context ...redis.Cmdable) error {
  207. if client == nil {
  208. return ErrDBNotInit
  209. }
  210. return getCmdable(context...).HDel(ctx, serialKey(key), field).Err()
  211. }
  212. // GetMap get the map with key
  213. func GetMap[V any](key string, context ...redis.Cmdable) (map[string]V, error) {
  214. if client == nil {
  215. return nil, ErrDBNotInit
  216. }
  217. val, err := getCmdable(context...).HGetAll(ctx, serialKey(key)).Result()
  218. if err != nil {
  219. if err == redis.Nil {
  220. return nil, ErrNotFound
  221. }
  222. return nil, err
  223. }
  224. result := make(map[string]V)
  225. for k, v := range val {
  226. value, err := parser.UnmarshalJson[V](v)
  227. if err != nil {
  228. continue
  229. }
  230. result[k] = value
  231. }
  232. return result, nil
  233. }
  234. // ScanKeys scan the keys with match pattern
  235. func ScanKeys(match string, context ...redis.Cmdable) ([]string, error) {
  236. if client == nil {
  237. return nil, ErrDBNotInit
  238. }
  239. result := make([]string, 0)
  240. if err := ScanKeysAsync(match, func(keys []string) error {
  241. result = append(result, keys...)
  242. return nil
  243. }); err != nil {
  244. return nil, err
  245. }
  246. return result, nil
  247. }
  248. // ScanKeysAsync scan the keys with match pattern, format like "key*"
  249. func ScanKeysAsync(match string, fn func([]string) error, context ...redis.Cmdable) error {
  250. if client == nil {
  251. return ErrDBNotInit
  252. }
  253. cursor := uint64(0)
  254. for {
  255. keys, newCursor, err := getCmdable(context...).Scan(ctx, cursor, match, 32).Result()
  256. if err != nil {
  257. return err
  258. }
  259. if err := fn(keys); err != nil {
  260. return err
  261. }
  262. if newCursor == 0 {
  263. break
  264. }
  265. cursor = newCursor
  266. }
  267. return nil
  268. }
  269. // ScanMap scan the map with match pattern, format like "key*"
  270. func ScanMap[V any](key string, match string, context ...redis.Cmdable) (map[string]V, error) {
  271. if client == nil {
  272. return nil, ErrDBNotInit
  273. }
  274. result := make(map[string]V)
  275. ScanMapAsync[V](key, match, func(m map[string]V) error {
  276. for k, v := range m {
  277. result[k] = v
  278. }
  279. return nil
  280. })
  281. return result, nil
  282. }
  283. // ScanMapAsync scan the map with match pattern, format like "key*"
  284. func ScanMapAsync[V any](key string, match string, fn func(map[string]V) error, context ...redis.Cmdable) error {
  285. if client == nil {
  286. return ErrDBNotInit
  287. }
  288. cursor := uint64(0)
  289. for {
  290. kvs, newCursor, err := getCmdable(context...).
  291. HScan(ctx, serialKey(key), cursor, match, 32).
  292. Result()
  293. if err != nil {
  294. return err
  295. }
  296. result := make(map[string]V)
  297. for i := 0; i < len(kvs); i += 2 {
  298. value, err := parser.UnmarshalJson[V](kvs[i+1])
  299. if err != nil {
  300. continue
  301. }
  302. result[kvs[i]] = value
  303. }
  304. if err := fn(result); err != nil {
  305. return err
  306. }
  307. if newCursor == 0 {
  308. break
  309. }
  310. cursor = newCursor
  311. }
  312. return nil
  313. }
  314. // SetNX set the key-value pair with expire time
  315. func SetNX[T any](key string, value T, expire time.Duration, context ...redis.Cmdable) (bool, error) {
  316. if client == nil {
  317. return false, ErrDBNotInit
  318. }
  319. // marshal the value
  320. bytes, err := parser.MarshalCBOR(value)
  321. if err != nil {
  322. return false, err
  323. }
  324. return getCmdable(context...).SetNX(ctx, serialKey(key), bytes, expire).Result()
  325. }
  326. var (
  327. ErrLockTimeout = errors.New("lock timeout")
  328. )
  329. // Lock key, expire time takes responsibility for expiration time
  330. // try_lock_timeout takes responsibility for the timeout of trying to lock
  331. func Lock(key string, expire time.Duration, tryLockTimeout time.Duration, context ...redis.Cmdable) error {
  332. if client == nil {
  333. return ErrDBNotInit
  334. }
  335. const LOCK_DURATION = 20 * time.Millisecond
  336. ticker := time.NewTicker(LOCK_DURATION)
  337. defer ticker.Stop()
  338. for range ticker.C {
  339. if _, err := getCmdable(context...).SetNX(ctx, serialKey(key), "1", expire).Result(); err == nil {
  340. return nil
  341. }
  342. tryLockTimeout -= LOCK_DURATION
  343. if tryLockTimeout <= 0 {
  344. return ErrLockTimeout
  345. }
  346. }
  347. return nil
  348. }
  349. func Unlock(key string, context ...redis.Cmdable) error {
  350. if client == nil {
  351. return ErrDBNotInit
  352. }
  353. return getCmdable(context...).Del(ctx, serialKey(key)).Err()
  354. }
  355. func Expire(key string, time time.Duration, context ...redis.Cmdable) (bool, error) {
  356. if client == nil {
  357. return false, ErrDBNotInit
  358. }
  359. return getCmdable(context...).Expire(ctx, serialKey(key), time).Result()
  360. }
  361. func Transaction(fn func(redis.Pipeliner) error) error {
  362. if client == nil {
  363. return ErrDBNotInit
  364. }
  365. return client.Watch(ctx, func(tx *redis.Tx) error {
  366. _, err := tx.TxPipelined(ctx, func(p redis.Pipeliner) error {
  367. return fn(p)
  368. })
  369. if err == redis.Nil {
  370. return nil
  371. }
  372. return err
  373. })
  374. }
  375. func Publish(channel string, message any, context ...redis.Cmdable) error {
  376. if client == nil {
  377. return ErrDBNotInit
  378. }
  379. if _, ok := message.(string); !ok {
  380. message = parser.MarshalJson(message)
  381. }
  382. return getCmdable(context...).Publish(ctx, channel, message).Err()
  383. }
  384. func Subscribe[T any](channel string) (<-chan T, func()) {
  385. pubsub := client.Subscribe(ctx, channel)
  386. ch := make(chan T)
  387. connectionEstablished := make(chan bool)
  388. go func() {
  389. defer close(ch)
  390. defer close(connectionEstablished)
  391. alive := true
  392. for alive {
  393. iface, err := pubsub.Receive(context.Background())
  394. if err != nil {
  395. log.Error("failed to receive message from redis: %s, will retry in 1 second", err.Error())
  396. time.Sleep(1 * time.Second)
  397. continue
  398. }
  399. switch data := iface.(type) {
  400. case *redis.Subscription:
  401. connectionEstablished <- true
  402. case *redis.Message:
  403. v, err := parser.UnmarshalJson[T](data.Payload)
  404. if err != nil {
  405. continue
  406. }
  407. ch <- v
  408. case *redis.Pong:
  409. default:
  410. alive = false
  411. }
  412. }
  413. }()
  414. // wait for the connection to be established
  415. <-connectionEstablished
  416. return ch, func() {
  417. pubsub.Close()
  418. }
  419. }