diff --git a/go.mod b/go.mod index 511f686..30ebdda 100644 --- a/go.mod +++ b/go.mod @@ -5,11 +5,13 @@ go 1.26.2 require ( github.com/go-sql-driver/mysql v1.10.0 github.com/jmoiron/sqlx v1.4.0 + github.com/redis/go-redis/v9 v9.21.0 github.com/spf13/viper v1.21.0 ) require ( filippo.io/edwards25519 v1.2.0 // indirect + github.com/cespare/xxhash/v2 v2.3.0 // indirect github.com/fsnotify/fsnotify v1.9.0 // indirect github.com/go-viper/mapstructure/v2 v2.4.0 // indirect github.com/pelletier/go-toml/v2 v2.2.4 // indirect @@ -19,7 +21,8 @@ require ( github.com/spf13/cast v1.10.0 // indirect github.com/spf13/pflag v1.0.10 // indirect github.com/subosito/gotenv v1.6.0 // indirect + go.uber.org/atomic v1.11.0 // indirect go.yaml.in/yaml/v3 v3.0.4 // indirect - golang.org/x/sys v0.29.0 // indirect + golang.org/x/sys v0.30.0 // indirect golang.org/x/text v0.28.0 // indirect ) diff --git a/go.sum b/go.sum index e800caf..e8a976e 100644 --- a/go.sum +++ b/go.sum @@ -1,6 +1,12 @@ filippo.io/edwards25519 v1.1.0/go.mod h1:BxyFTGdWcka3PhytdK4V28tE5sGfRvvvRV7EaN4VDT4= filippo.io/edwards25519 v1.2.0 h1:crnVqOiS4jqYleHd9vaKZ+HKtHfllngJIiOpNpoJsjo= filippo.io/edwards25519 v1.2.0/go.mod h1:xzAOLCNug/yB62zG1bQ8uziwrIqIuxhctzJT18Q77mc= +github.com/bsm/ginkgo/v2 v2.12.0 h1:Ny8MWAHyOepLGlLKYmXG4IEkioBysk6GpaRTLC8zwWs= +github.com/bsm/ginkgo/v2 v2.12.0/go.mod h1:SwYbGRRDovPVboqFv0tPTcG1sN61LM1Z4ARdbAV9g4c= +github.com/bsm/gomega v1.27.10 h1:yeMWxP2pV2fG3FgAODIY8EiRE3dy0aeFYt4l7wh6yKA= +github.com/bsm/gomega v1.27.10/go.mod h1:JyEr/xRbxbtgWNi8tIEVPUYZ5Dzef52k01W3YH0H+O0= +github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= +github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/frankban/quicktest v1.14.6 h1:7Xjx+VpznH+oBnejlPUj8oUpdxnVs4f8XU8WnHkI4W8= @@ -16,6 +22,8 @@ github.com/google/go-cmp v0.6.0 h1:ofyhxvXcZhMsU5ulbFiLKl/XBFqE1GSq7atu8tAmTRI= github.com/google/go-cmp v0.6.0/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= github.com/jmoiron/sqlx v1.4.0 h1:1PLqN7S1UYp5t4SrVVnt4nUVNemrDAtxlulVe+Qgm3o= github.com/jmoiron/sqlx v1.4.0/go.mod h1:ZrZ7UsYB/weZdl2Bxg6jCRO9c3YHl8r3ahlKmRT4JLY= +github.com/klauspost/cpuid/v2 v2.2.10 h1:tBs3QSyvjDyFTq3uoc/9xFpCuOsJQFNPiAhYdw2skhE= +github.com/klauspost/cpuid/v2 v2.2.10/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0= github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= @@ -28,6 +36,8 @@ github.com/pelletier/go-toml/v2 v2.2.4 h1:mye9XuhQ6gvn5h28+VilKrrPoQVanw5PMw/TB0 github.com/pelletier/go-toml/v2 v2.2.4/go.mod h1:2gIqNv+qfxSVS7cM2xJQKtLSTLUE9V8t9Stt+h56mCY= github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/redis/go-redis/v9 v9.21.0 h1:FPBE4hhbAke+TLmcY3WkpbDffJEomdqPn3HYiqAtL9E= +github.com/redis/go-redis/v9 v9.21.0/go.mod h1:v/M13XI1PVCDcm01VtPFOADfZtHf8YW3baQf57KlIkA= github.com/rogpeppe/go-internal v1.9.0 h1:73kH8U+JUqXU8lRuOHeVHaa/SZPifC7BkcraZVejAe8= github.com/rogpeppe/go-internal v1.9.0/go.mod h1:WtVeX8xhTBvf0smdhujwtBcq4Qrzq/fJaraNFVN+nFs= github.com/sagikazarmark/locafero v0.11.0 h1:1iurJgmM9G3PA/I+wWYIOw/5SyBtxapeHDcg+AAIFXc= @@ -46,10 +56,14 @@ github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= github.com/subosito/gotenv v1.6.0 h1:9NlTDc1FTs4qu0DDq7AEtTPNw6SVm7uBMsUCUjABIf8= github.com/subosito/gotenv v1.6.0/go.mod h1:Dk4QP5c2W3ibzajGcXpNraDfq2IrhjMIvMSWPKKo0FU= +github.com/zeebo/xxh3 v1.1.0 h1:s7DLGDK45Dyfg7++yxI0khrfwq9661w9EN78eP/UZVs= +github.com/zeebo/xxh3 v1.1.0/go.mod h1:IisAie1LELR4xhVinxWS5+zf1lA4p0MW4T+w+W07F5s= +go.uber.org/atomic v1.11.0 h1:ZvwS0R+56ePWxUNi+Atn9dWONBPp/AUETXlHW0DxSjE= +go.uber.org/atomic v1.11.0/go.mod h1:LUxbIzbOniOlMKjJjyPfpl4v+PKK2cNJn91OQbhoJI0= go.yaml.in/yaml/v3 v3.0.4 h1:tfq32ie2Jv2UxXFdLJdh3jXuOzWiL1fo0bu/FbuKpbc= go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg= -golang.org/x/sys v0.29.0 h1:TPYlXGxvx1MGTn2GiZDhnjPA9wZzZeGKHHmKhHYvgaU= -golang.org/x/sys v0.29.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= +golang.org/x/sys v0.30.0 h1:QjkSwP/36a20jFYWkSue1YwXzLmsV5Gfq7Eiy72C1uc= +golang.org/x/sys v0.30.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= golang.org/x/text v0.28.0 h1:rhazDwis8INMIwQ4tpjLDzUhx6RlXqZNPEM0huQojng= golang.org/x/text v0.28.0/go.mod h1:U8nCwOR8jO/marOQ0QbDiOngZVEBB7MAiitBuMjXiNU= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= diff --git a/internal/cache/redis.go b/internal/cache/redis.go new file mode 100644 index 0000000..da76220 --- /dev/null +++ b/internal/cache/redis.go @@ -0,0 +1,150 @@ +package cache + +import ( + "context" + "encoding/json" + "fmt" + "strconv" + "time" + + "aiaa-notification-service/internal/config" + + "github.com/redis/go-redis/v9" +) + +type Cache struct { + rdb *redis.Client +} + +// CachedRule holds the minimal info needed after rule matching. +type CachedRule struct { + RuleID int `json:"rule_id"` + TemplateID int `json:"template_id"` + Content string `json:"content"` + Conditions string `json:"conditions"` // JSON string, empty if null +} + +type CachedChannel struct { + ID int `json:"id"` + Type string `json:"type"` + Config *json.RawMessage `json:"config"` +} + +func NewCache(cfg config.RedisConfig) (*Cache, error) { + rdb := redis.NewClient(&redis.Options{ + Addr: cfg.Addr(), + Password: cfg.Password, + DB: cfg.DB, + }) + ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second) + defer cancel() + if err := rdb.Ping(ctx).Err(); err != nil { + return nil, fmt.Errorf("redis ping: %w", err) + } + return &Cache{rdb: rdb}, nil +} + +// --- Rule cache --- + +func ruleKey(sourceID int, event string) string { + return fmt.Sprintf("notify:rule:%d:%s", sourceID, event) +} + +func (c *Cache) GetRule(ctx context.Context, sourceID int, event string) (*CachedRule, error) { + data, err := c.rdb.Get(ctx, ruleKey(sourceID, event)).Bytes() + if err != nil { + return nil, err + } + var cr CachedRule + if err := json.Unmarshal(data, &cr); err != nil { + return nil, err + } + return &cr, nil +} + +func (c *Cache) SetRule(ctx context.Context, sourceID int, event string, cr *CachedRule) error { + data, err := json.Marshal(cr) + if err != nil { + return err + } + return c.rdb.Set(ctx, ruleKey(sourceID, event), data, 5*time.Minute).Err() +} + +func (c *Cache) InvalidateRule(ctx context.Context, sourceID int, event string) error { + return c.rdb.Del(ctx, ruleKey(sourceID, event)).Err() +} + +// --- Channel cache --- + +func channelsKey(ruleID int) string { + return fmt.Sprintf("notify:channels:%d", ruleID) +} + +func (c *Cache) GetChannels(ctx context.Context, ruleID int) ([]CachedChannel, error) { + data, err := c.rdb.Get(ctx, channelsKey(ruleID)).Bytes() + if err != nil { + return nil, err + } + var channels []CachedChannel + if err := json.Unmarshal(data, &channels); err != nil { + return nil, err + } + return channels, nil +} + +func (c *Cache) SetChannels(ctx context.Context, ruleID int, channels []CachedChannel) error { + data, err := json.Marshal(channels) + if err != nil { + return err + } + return c.rdb.Set(ctx, channelsKey(ruleID), data, 5*time.Minute).Err() +} + +func (c *Cache) InvalidateChannels(ctx context.Context, ruleID int) error { + return c.rdb.Del(ctx, channelsKey(ruleID)).Err() +} + +// --- Bulk invalidation --- + +func (c *Cache) InvalidateBySource(ctx context.Context, sourceID int) error { + pattern := fmt.Sprintf("notify:rule:%d:*", sourceID) + iter := c.rdb.Scan(ctx, 0, pattern, 0).Iterator() + for iter.Next(ctx) { + c.rdb.Del(ctx, iter.Val()) + } + return iter.Err() +} + +func (c *Cache) InvalidateByTemplate(ctx context.Context, templateID int) error { + // Template changes mean all cached rule info is stale. Simplest: scan all rule keys. + // In production, you'd maintain a template to rule reverse index. For v1, scan is acceptable. + iter := c.rdb.Scan(ctx, 0, "notify:rule:*", 0).Iterator() + for iter.Next(ctx) { + c.rdb.Del(ctx, iter.Val()) + } + return iter.Err() +} + +// --- Rate limiter --- + +func (c *Cache) CheckRateLimit(ctx context.Context, sourceID int, limitPerSec int) (bool, int, error) { + key := fmt.Sprintf("ratelimit:%d", sourceID) + now := time.Now().Unix() + pipe := c.rdb.Pipeline() + pipe.ZRemRangeByScore(ctx, key, "0", strconv.FormatInt(now-1, 10)) + countCmd := pipe.ZCard(ctx, key) + pipe.Exec(ctx) + + count := countCmd.Val() + if count >= int64(limitPerSec) { + return false, 1, nil // rate limited, retry after 1s + } + + c.rdb.ZAdd(ctx, key, redis.Z{Score: float64(now), Member: strconv.FormatInt(now*1000, 10)}) + c.rdb.Expire(ctx, key, 2*time.Second) + return true, 0, nil +} + +func (c *Cache) Close() error { + return c.rdb.Close() +}