mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-07 16:16:37 +08:00
feat(push): implement Redis caching for push events and custom channels
- Implement cached queries in GetActivePushEventByKey and GetActivePushChannelByName. - Set TTL for cached items to 24 hours via activePushEventCacheTTL and activePushChannelCacheTTL constants. - Implement GORM hooks (AfterSave and AfterDelete) on PushEvent and PushChannel to auto-evict Redis caches, guaranteeing cache consistency. - Evict Redis caches manually inside API handlers for Create/Update/Delete/Toggle event/channel endpoints. - Update events.go to query models through the new caching methods.
This commit is contained in:
@@ -107,6 +107,9 @@ func CreateChannel(c *gin.Context) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 缓存一致性:清除渠道缓存
|
||||||
|
model.DeleteActivePushChannelCache(ctx, channel.Name)
|
||||||
|
|
||||||
c.JSON(http.StatusOK, response.OK(channel))
|
c.JSON(http.StatusOK, response.OK(channel))
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -174,6 +177,9 @@ func UpdateChannel(c *gin.Context) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 缓存一致性:清除渠道缓存
|
||||||
|
model.DeleteActivePushChannelCache(ctx, channel.Name)
|
||||||
|
|
||||||
c.JSON(http.StatusOK, response.OK(channel))
|
c.JSON(http.StatusOK, response.OK(channel))
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -210,6 +216,9 @@ func DeleteChannel(c *gin.Context) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 缓存一致性:清除渠道缓存
|
||||||
|
model.DeleteActivePushChannelCache(ctx, channel.Name)
|
||||||
|
|
||||||
c.JSON(http.StatusOK, response.OKNil())
|
c.JSON(http.StatusOK, response.OKNil())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -109,9 +109,8 @@ func (t *EventTrigger) Trigger(ctx context.Context, meta EventMetadata, body map
|
|||||||
body["user"] = getSystemUser(asyncCtx)
|
body["user"] = getSystemUser(asyncCtx)
|
||||||
}
|
}
|
||||||
|
|
||||||
// 1. Check if the event is enabled in the database
|
// 1. Check if the event is enabled (try Redis cache first)
|
||||||
var event model.PushEvent
|
eventPtr, err := model.GetActivePushEventByKey(asyncCtx, meta.Key)
|
||||||
err := db.DB(asyncCtx).Where("event_key = ? AND enabled = ?", meta.Key, true).First(&event).Error
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||||
return
|
return
|
||||||
@@ -119,6 +118,7 @@ func (t *EventTrigger) Trigger(ctx context.Context, meta EventMetadata, body map
|
|||||||
logger.ErrorF(asyncCtx, "push_event_trigger: failed to get active event %s: %v", meta.Key, err)
|
logger.ErrorF(asyncCtx, "push_event_trigger: failed to get active event %s: %v", meta.Key, err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
event := *eventPtr
|
||||||
|
|
||||||
if len(event.Channels) == 0 {
|
if len(event.Channels) == 0 {
|
||||||
return
|
return
|
||||||
@@ -222,11 +222,10 @@ func (t *EventTrigger) parseDefaultTemplate(meta EventMetadata, flatBody map[str
|
|||||||
|
|
||||||
func (t *EventTrigger) enqueuePushTasks(ctx context.Context, meta EventMetadata, event *model.PushEvent, msg NotificationMessage, flatBody map[string]any) {
|
func (t *EventTrigger) enqueuePushTasks(ctx context.Context, meta EventMetadata, event *model.PushEvent, msg NotificationMessage, flatBody map[string]any) {
|
||||||
for _, channelName := range event.Channels {
|
for _, channelName := range event.Channels {
|
||||||
// 检查是不是自定义数据库渠道
|
// 检查是不是自定义数据库渠道 (使用 Redis 缓存优先)
|
||||||
var customChannel model.PushChannel
|
customChannel, err := model.GetActivePushChannelByName(ctx, channelName)
|
||||||
err := db.DB(ctx).Where("name = ? AND enabled = ?", channelName, true).First(&customChannel).Error
|
|
||||||
if err == nil {
|
if err == nil {
|
||||||
t.enqueueCustomPushChannelTasks(ctx, meta, event, &customChannel, msg, flatBody)
|
t.enqueueCustomPushChannelTasks(ctx, meta, event, customChannel, msg, flatBody)
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -231,6 +231,9 @@ func CreateEvent(c *gin.Context) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 缓存一致性:清除旧事件缓存
|
||||||
|
model.DeleteActivePushEventCache(ctx, event.EventKey)
|
||||||
|
|
||||||
c.JSON(http.StatusOK, response.OK(event))
|
c.JSON(http.StatusOK, response.OK(event))
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -267,6 +270,9 @@ func DeleteEvent(c *gin.Context) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 缓存一致性:清除事件缓存
|
||||||
|
model.DeleteActivePushEventCache(ctx, event.EventKey)
|
||||||
|
|
||||||
c.JSON(http.StatusOK, response.OKNil())
|
c.JSON(http.StatusOK, response.OKNil())
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -320,6 +326,9 @@ func UpdateEvent(c *gin.Context) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 缓存一致性:清除事件缓存
|
||||||
|
model.DeleteActivePushEventCache(c.Request.Context(), event.EventKey)
|
||||||
|
|
||||||
c.JSON(http.StatusOK, response.OKNil())
|
c.JSON(http.StatusOK, response.OKNil())
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -360,6 +369,9 @@ func ToggleEvent(c *gin.Context) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 缓存一致性:清除事件缓存
|
||||||
|
model.DeleteActivePushEventCache(c.Request.Context(), event.EventKey)
|
||||||
|
|
||||||
c.JSON(http.StatusOK, response.OK(event.Enabled))
|
c.JSON(http.StatusOK, response.OK(event.Enabled))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -12,6 +12,7 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/Rain-kl/Wavelet/internal/db"
|
"github.com/Rain-kl/Wavelet/internal/db"
|
||||||
|
"gorm.io/gorm"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
@@ -112,3 +113,47 @@ func GetPushChannelByName(ctx context.Context, name string) (*PushChannel, error
|
|||||||
}
|
}
|
||||||
return &channel, nil
|
return &channel, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const activePushChannelCacheTTL = 24 * time.Hour
|
||||||
|
|
||||||
|
// GetActivePushChannelByName 根据名称获取启用的消息通道 (优先从 Redis 缓存获取)
|
||||||
|
func GetActivePushChannelByName(ctx context.Context, name string) (*PushChannel, error) {
|
||||||
|
cacheKey := "push:channel:active:" + name
|
||||||
|
var channel PushChannel
|
||||||
|
if db.Redis != nil {
|
||||||
|
if err := db.GetJSON(ctx, cacheKey, &channel); err == nil {
|
||||||
|
return &channel, nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
err := db.DB(ctx).Where("name = ? AND enabled = ?", name, true).First(&channel).Error
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
if db.Redis != nil {
|
||||||
|
// 缓存有效时间设置为 24 小时
|
||||||
|
_ = db.SetJSON(ctx, cacheKey, channel, activePushChannelCacheTTL)
|
||||||
|
}
|
||||||
|
|
||||||
|
return &channel, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// DeleteActivePushChannelCache 清理启用消息通道的缓存
|
||||||
|
func DeleteActivePushChannelCache(ctx context.Context, name string) {
|
||||||
|
if db.Redis != nil {
|
||||||
|
_ = db.Redis.Del(ctx, db.PrefixedKey("push:channel:active:"+name)).Err()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// AfterSave GORM 保存后钩子,用于自动清理缓存
|
||||||
|
func (pc *PushChannel) AfterSave(tx *gorm.DB) error {
|
||||||
|
DeleteActivePushChannelCache(tx.Statement.Context, pc.Name)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// AfterDelete GORM 删除后钩子,用于自动清理缓存
|
||||||
|
func (pc *PushChannel) AfterDelete(tx *gorm.DB) error {
|
||||||
|
DeleteActivePushChannelCache(tx.Statement.Context, pc.Name)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/Rain-kl/Wavelet/internal/db"
|
"github.com/Rain-kl/Wavelet/internal/db"
|
||||||
|
"gorm.io/gorm"
|
||||||
)
|
)
|
||||||
|
|
||||||
// PushEvent 系统通知事件模型
|
// PushEvent 系统通知事件模型
|
||||||
@@ -52,12 +53,46 @@ func (pe *PushEvent) Validate() error {
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// GetActivePushEventByKey 获取启用的通知事件
|
const activePushEventCacheTTL = 24 * time.Hour
|
||||||
|
|
||||||
|
// GetActivePushEventByKey 获取启用的通知事件 (优先从 Redis 缓存获取)
|
||||||
func GetActivePushEventByKey(ctx context.Context, key string) (*PushEvent, error) {
|
func GetActivePushEventByKey(ctx context.Context, key string) (*PushEvent, error) {
|
||||||
|
cacheKey := "push:event:active:" + key
|
||||||
var event PushEvent
|
var event PushEvent
|
||||||
|
if db.Redis != nil {
|
||||||
|
if err := db.GetJSON(ctx, cacheKey, &event); err == nil {
|
||||||
|
return &event, nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
err := db.DB(ctx).Where("event_key = ? AND enabled = ?", key, true).First(&event).Error
|
err := db.DB(ctx).Where("event_key = ? AND enabled = ?", key, true).First(&event).Error
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if db.Redis != nil {
|
||||||
|
// 缓存有效时间设置为 24 小时
|
||||||
|
_ = db.SetJSON(ctx, cacheKey, event, activePushEventCacheTTL)
|
||||||
|
}
|
||||||
|
|
||||||
return &event, nil
|
return &event, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// DeleteActivePushEventCache 清理启用通知事件的缓存
|
||||||
|
func DeleteActivePushEventCache(ctx context.Context, key string) {
|
||||||
|
if db.Redis != nil {
|
||||||
|
_ = db.Redis.Del(ctx, db.PrefixedKey("push:event:active:"+key)).Err()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// AfterSave GORM 保存后钩子,用于自动清理缓存
|
||||||
|
func (pe *PushEvent) AfterSave(tx *gorm.DB) error {
|
||||||
|
DeleteActivePushEventCache(tx.Statement.Context, pe.EventKey)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// AfterDelete GORM 删除后钩子,用于自动清理缓存
|
||||||
|
func (pe *PushEvent) AfterDelete(tx *gorm.DB) error {
|
||||||
|
DeleteActivePushEventCache(tx.Statement.Context, pe.EventKey)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user