* feature extend readthrough for cache module (#5116) * feature 增加readthrough * feature: add write though for cache mode (#5117) * feature: add writethough for cache mode * feature add singleflight cache (#5119) * build(deps): bump go.opentelemetry.io/otel/trace from 1.8.0 to 1.11.2 Bumps [go.opentelemetry.io/otel/trace](https://github.com/open-telemetry/opentelemetry-go) from 1.8.0 to 1.11.2. - [Release notes](https://github.com/open-telemetry/opentelemetry-go/releases) - [Changelog](https://github.com/open-telemetry/opentelemetry-go/blob/main/CHANGELOG.md) - [Commits](https://github.com/open-telemetry/opentelemetry-go/compare/v1.8.0...v1.11.2) --- updated-dependencies: - dependency-name: go.opentelemetry.io/otel/trace dependency-type: direct:production update-type: version-update:semver-minor ... Signed-off-by: dependabot[bot] <support@github.com> * fix 5129: must set formatter after init the logger * remove beego.vip * build(deps): bump actions/stale from 5 to 7 Bumps [actions/stale](https://github.com/actions/stale) from 5 to 7. - [Release notes](https://github.com/actions/stale/releases) - [Changelog](https://github.com/actions/stale/blob/main/CHANGELOG.md) - [Commits](https://github.com/actions/stale/compare/v5...v7) --- updated-dependencies: - dependency-name: actions/stale dependency-type: direct:production update-type: version-update:semver-major ... Signed-off-by: dependabot[bot] <support@github.com> * fix 5079: only log msg when the channel is not closed (#5132) * optimize test * upgrade otel dependencies to v1.11.2 * format code * Bloom filter cache (#5126) * feature: add bloom filter cache * feature upload remove all temp file * bugfix Controller SaveToFile remove all temp file * rft: motify BeeLogger signalChan (#5139) * add non-block write log in asynchronous mode (#5150) * add non-block write log in asynchronous mode --------- Co-authored-by: chenhaokun <chenhaokun@itiger.com> * fix the docsite URL (#5173) * Unified gopkg.in/yaml version to v2 (#5169) * Unified gopkg.in/yaml version to v2 and go mod tidy * update CHANGELOG * bugfix: protect field access with lock to avoid possible data race (#5211) * fix some comments (#5194) Signed-off-by: cui fliter <imcusg@gmail.com> * build(deps): bump github.com/prometheus/client_golang (#5213) Bumps [github.com/prometheus/client_golang](https://github.com/prometheus/client_golang) from 1.14.0 to 1.15.1. - [Release notes](https://github.com/prometheus/client_golang/releases) - [Changelog](https://github.com/prometheus/client_golang/blob/main/CHANGELOG.md) - [Commits](https://github.com/prometheus/client_golang/compare/v1.14.0...v1.15.1) --- updated-dependencies: - dependency-name: github.com/prometheus/client_golang dependency-type: direct:production update-type: version-update:semver-minor ... Signed-off-by: dependabot[bot] <support@github.com> Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com> * build(deps): bump go.etcd.io/etcd/client/v3 from 3.5.4 to 3.5.9 (#5209) Bumps [go.etcd.io/etcd/client/v3](https://github.com/etcd-io/etcd) from 3.5.4 to 3.5.9. - [Release notes](https://github.com/etcd-io/etcd/releases) - [Commits](https://github.com/etcd-io/etcd/compare/v3.5.4...v3.5.9) --- updated-dependencies: - dependency-name: go.etcd.io/etcd/client/v3 dependency-type: direct:production update-type: version-update:semver-patch ... Signed-off-by: dependabot[bot] <support@github.com> Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com> * cache: fix typo and optimize the naming * Release 2.1.0 change log * bugfix: beegoAppConfig String and Strings function has bug * httplib: fix unstable test, do not use httplib.org * chore: pkg imported more than once * chore: fmt modify * chore: Use github.com/go-kit/log * chore: unnecessary use of fmt.Sprintf * fix: golangci-lint error * orm: refactor ORM introducing internal/models pkg * remove adapter package * build(deps): bump github.com/bits-and-blooms/bloom/v3 Bumps [github.com/bits-and-blooms/bloom/v3](https://github.com/bits-and-blooms/bloom) from 3.3.1 to 3.5.0. - [Release notes](https://github.com/bits-and-blooms/bloom/releases) - [Commits](https://github.com/bits-and-blooms/bloom/compare/v3.3.1...v3.5.0) --- updated-dependencies: - dependency-name: github.com/bits-and-blooms/bloom/v3 dependency-type: direct:production update-type: version-update:semver-minor ... Signed-off-by: dependabot[bot] <support@github.com> * feat: add write-delete cache mode * fix: unnecessary assignment to the blank identifier * fix: add change into .CHANGELOG file * build(deps): bump golang.org/x/sync from 0.1.0 to 0.3.0 Bumps [golang.org/x/sync](https://github.com/golang/sync) from 0.1.0 to 0.3.0. - [Commits](https://github.com/golang/sync/compare/v0.1.0...v0.3.0) --- updated-dependencies: - dependency-name: golang.org/x/sync dependency-type: direct:production update-type: version-update:semver-minor ... Signed-off-by: dependabot[bot] <support@github.com> * build(deps): bump golang.org/x/crypto Bumps [golang.org/x/crypto](https://github.com/golang/crypto) from 0.0.0-20220315160706-3147a52a75dd to 0.10.0. - [Commits](https://github.com/golang/crypto/commits/v0.10.0) --- updated-dependencies: - dependency-name: golang.org/x/crypto dependency-type: direct:production update-type: version-update:semver-minor ... Signed-off-by: dependabot[bot] <support@github.com> * remove golang--lint-ci * Beego web.Run() runs the server twice * fix 5255: Check the rows.Err() if rows.Next() is false * closes 5254: %COL% should be a common placeholder * build(deps): bump github.com/prometheus/client_golang Bumps [github.com/prometheus/client_golang](https://github.com/prometheus/client_golang) from 1.15.1 to 1.16.0. - [Release notes](https://github.com/prometheus/client_golang/releases) - [Changelog](https://github.com/prometheus/client_golang/blob/main/CHANGELOG.md) - [Commits](https://github.com/prometheus/client_golang/compare/v1.15.1...v1.16.0) --- updated-dependencies: - dependency-name: github.com/prometheus/client_golang dependency-type: direct:production update-type: version-update:semver-minor ... Signed-off-by: dependabot[bot] <support@github.com> * fix: use of ioutil package (#5261) * fix ioutil.NopCloser * fix ioutil.ReadAll * fix ioutil.ReadFile * fix ioutil.WriteFile * run goimports -w -format-only ./ * update CHANGELOG.md * feature: add write-double-delete cache mode (#5263) * cache/redis: support skipEmptyPrefix option (#5264) * fix: refactor InsertValue method (#5267) * fix: refactor insertValue method and add the test * fix: exec goimports and add Licence file header * fix: modify construct method of dbBase * fix: add modify record into CHANGELOG * fix: modify InsertOrUpdate method (#5269) * fix: modify InsertOrUpdate method, Remove the isMulti variable and its associated code * fix: Delete unnecessary judgment branches * fix: add modify record into CHANGELOG * cache/redis: use redisConfig to receive incoming JSON (previously using a map) (#5268) * refactor cache/redis: Use redisConfig to receive incoming JSON (previously using a map). * refactor cache/redis: Use the string type to receive JSON parameters. --------- Co-authored-by: Tan <tanqianheng@gmail.com> * fix: refactor Delete method (#5271) * fix: refactor Delete method and add test * fix: add modify record into CHANGELOG * fix: refactor update sql (#5274) * fix: refactor UpdateSQL method and add test * fix: add modify record into CHANGELOG * fix: modify url in the CHANGELOG * fix: modify pr url in the CHANGELOG * Fix setPK function for table without primary key (#5276) --------- Signed-off-by: dependabot[bot] <support@github.com> Signed-off-by: cui fliter <imcusg@gmail.com> Co-authored-by: Stone-afk <73482944+Stone-afk@users.noreply.github.com> Co-authored-by: hookokoko <hooko@tju.edu.cn> Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com> Co-authored-by: hookokoko <648646891@qq.com> Co-authored-by: Stone-afk <1711865140@qq.com> Co-authored-by: chenhaokun <chenhaokun@itiger.com> Co-authored-by: Xuing <admin@xuing.cn> Co-authored-by: cui fliter <imcusg@gmail.com> Co-authored-by: guoguangwu <guoguangwu@magic-shield.com> Co-authored-by: uzziah <uzziahlin@gmail.com> Co-authored-by: Hanjiang Yu <delacroix.yu@gmail.com> Co-authored-by: Kota <mdryzk64smsh@gmail.com> Co-authored-by: Uzziah <120019273+uzziahlin@users.noreply.github.com> Co-authored-by: Handkerchiefs-t <59816423+Handkerchiefs-t@users.noreply.github.com> Co-authored-by: Tan <tanqianheng@gmail.com> Co-authored-by: mlgd <mlgd17@gmail.com>
		
			
				
	
	
		
			351 lines
		
	
	
		
			8.3 KiB
		
	
	
	
		
			Go
		
	
	
	
	
	
			
		
		
	
	
			351 lines
		
	
	
		
			8.3 KiB
		
	
	
	
		
			Go
		
	
	
	
	
	
| // Copyright 2014 beego Author. All Rights Reserved.
 | |
| //
 | |
| // Licensed under the Apache License, Version 2.0 (the "License");
 | |
| // you may not use this file except in compliance with the License.
 | |
| // You may obtain a copy of the License at
 | |
| //
 | |
| //      http://www.apache.org/licenses/LICENSE-2.0
 | |
| //
 | |
| // Unless required by applicable law or agreed to in writing, software
 | |
| // distributed under the License is distributed on an "AS IS" BASIS,
 | |
| // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
 | |
| // See the License for the specific language governing permissions and
 | |
| // limitations under the License.
 | |
| 
 | |
| // Package redis for cache provider
 | |
| //
 | |
| // depend on github.com/gomodule/redigo/redis
 | |
| //
 | |
| // go install github.com/gomodule/redigo/redis
 | |
| //
 | |
| // Usage:
 | |
| // import(
 | |
| //
 | |
| //	_ "github.com/beego/beego/v2/client/cache/redis"
 | |
| //	"github.com/beego/beego/v2/client/cache"
 | |
| //
 | |
| // )
 | |
| //
 | |
| //	bm, err := cache.NewCache("redis", `{"conn":"127.0.0.1:11211"}`)
 | |
| package redis
 | |
| 
 | |
| import (
 | |
| 	"context"
 | |
| 	"encoding/json"
 | |
| 	"fmt"
 | |
| 	"strconv"
 | |
| 	"strings"
 | |
| 	"time"
 | |
| 
 | |
| 	"github.com/gomodule/redigo/redis"
 | |
| 
 | |
| 	"github.com/beego/beego/v2/client/cache"
 | |
| 	"github.com/beego/beego/v2/core/berror"
 | |
| )
 | |
| 
 | |
| const (
 | |
| 	// DefaultKey defines the collection name of redis for the cache adapter.
 | |
| 	DefaultKey = "beecacheRedis"
 | |
| 	// defaultMaxIdle defines the default max idle connection number.
 | |
| 	defaultMaxIdle = 3
 | |
| 	// defaultTimeout defines the default timeout .
 | |
| 	defaultTimeout = time.Second * 180
 | |
| )
 | |
| 
 | |
| // Cache is Redis cache adapter.
 | |
| type Cache struct {
 | |
| 	p        *redis.Pool // redis connection pool
 | |
| 	conninfo string
 | |
| 	dbNum    int
 | |
| 	// key actually is prefix.
 | |
| 	key      string
 | |
| 	password string
 | |
| 	maxIdle  int
 | |
| 
 | |
| 	// skipEmptyPrefix for backward compatible,
 | |
| 	// check function associate
 | |
| 	// see https://github.com/beego/beego/issues/5248
 | |
| 	skipEmptyPrefix bool
 | |
| 
 | |
| 	// Timeout value (less than the redis server's timeout value).
 | |
| 	// Timeout used for idle connection
 | |
| 	timeout time.Duration
 | |
| }
 | |
| 
 | |
| // NewRedisCache creates a new redis cache with default collection name.
 | |
| func NewRedisCache() cache.Cache {
 | |
| 	return &Cache{key: DefaultKey}
 | |
| }
 | |
| 
 | |
| // Execute the redis commands. args[0] must be the key name
 | |
| func (rc *Cache) do(commandName string, args ...interface{}) (interface{}, error) {
 | |
| 	args[0] = rc.associate(args[0])
 | |
| 	c := rc.p.Get()
 | |
| 	defer func() {
 | |
| 		_ = c.Close()
 | |
| 	}()
 | |
| 
 | |
| 	reply, err := c.Do(commandName, args...)
 | |
| 	if err != nil {
 | |
| 		return nil, berror.Wrapf(err, cache.RedisCacheCurdFailed,
 | |
| 			"could not execute this command: %s", commandName)
 | |
| 	}
 | |
| 
 | |
| 	return reply, nil
 | |
| }
 | |
| 
 | |
| // associate with config key.
 | |
| func (rc *Cache) associate(originKey interface{}) string {
 | |
| 	if rc.key == "" && rc.skipEmptyPrefix {
 | |
| 		return fmt.Sprintf("%s", originKey)
 | |
| 	}
 | |
| 	return fmt.Sprintf("%s:%s", rc.key, originKey)
 | |
| }
 | |
| 
 | |
| // Get cache from redis.
 | |
| func (rc *Cache) Get(ctx context.Context, key string) (interface{}, error) {
 | |
| 	if v, err := rc.do("GET", key); err == nil {
 | |
| 		return v, nil
 | |
| 	} else {
 | |
| 		return nil, err
 | |
| 	}
 | |
| }
 | |
| 
 | |
| // GetMulti gets cache from redis.
 | |
| func (rc *Cache) GetMulti(ctx context.Context, keys []string) ([]interface{}, error) {
 | |
| 	c := rc.p.Get()
 | |
| 	defer func() {
 | |
| 		_ = c.Close()
 | |
| 	}()
 | |
| 	var args []interface{}
 | |
| 	for _, key := range keys {
 | |
| 		args = append(args, rc.associate(key))
 | |
| 	}
 | |
| 	return redis.Values(c.Do("MGET", args...))
 | |
| }
 | |
| 
 | |
| // Put puts cache into redis.
 | |
| func (rc *Cache) Put(ctx context.Context, key string, val interface{}, timeout time.Duration) error {
 | |
| 	_, err := rc.do("SETEX", key, int64(timeout/time.Second), val)
 | |
| 	return err
 | |
| }
 | |
| 
 | |
| // Delete deletes a key's cache in redis.
 | |
| func (rc *Cache) Delete(ctx context.Context, key string) error {
 | |
| 	_, err := rc.do("DEL", key)
 | |
| 	return err
 | |
| }
 | |
| 
 | |
| // IsExist checks cache's existence in redis.
 | |
| func (rc *Cache) IsExist(ctx context.Context, key string) (bool, error) {
 | |
| 	v, err := redis.Bool(rc.do("EXISTS", key))
 | |
| 	if err != nil {
 | |
| 		return false, err
 | |
| 	}
 | |
| 	return v, nil
 | |
| }
 | |
| 
 | |
| // Incr increases a key's counter in redis.
 | |
| func (rc *Cache) Incr(ctx context.Context, key string) error {
 | |
| 	_, err := redis.Bool(rc.do("INCRBY", key, 1))
 | |
| 	return err
 | |
| }
 | |
| 
 | |
| // Decr decreases a key's counter in redis.
 | |
| func (rc *Cache) Decr(ctx context.Context, key string) error {
 | |
| 	_, err := redis.Bool(rc.do("INCRBY", key, -1))
 | |
| 	return err
 | |
| }
 | |
| 
 | |
| // ClearAll deletes all cache in the redis collection
 | |
| // Be careful about this method, because it scans all keys and the delete them one by one
 | |
| func (rc *Cache) ClearAll(context.Context) error {
 | |
| 	cachedKeys, err := rc.Scan(rc.key + ":*")
 | |
| 	if err != nil {
 | |
| 		return err
 | |
| 	}
 | |
| 	c := rc.p.Get()
 | |
| 	defer func() {
 | |
| 		_ = c.Close()
 | |
| 	}()
 | |
| 	for _, str := range cachedKeys {
 | |
| 		if _, err = c.Do("DEL", str); err != nil {
 | |
| 			return err
 | |
| 		}
 | |
| 	}
 | |
| 	return err
 | |
| }
 | |
| 
 | |
| // Scan scans all keys matching a given pattern.
 | |
| func (rc *Cache) Scan(pattern string) (keys []string, err error) {
 | |
| 	c := rc.p.Get()
 | |
| 	defer func() {
 | |
| 		_ = c.Close()
 | |
| 	}()
 | |
| 	var (
 | |
| 		cursor uint64 = 0 // start
 | |
| 		result []interface{}
 | |
| 		list   []string
 | |
| 	)
 | |
| 	for {
 | |
| 		result, err = redis.Values(c.Do("SCAN", cursor, "MATCH", pattern, "COUNT", 1024))
 | |
| 		if err != nil {
 | |
| 			return
 | |
| 		}
 | |
| 		list, err = redis.Strings(result[1], nil)
 | |
| 		if err != nil {
 | |
| 			return
 | |
| 		}
 | |
| 		keys = append(keys, list...)
 | |
| 		cursor, err = redis.Uint64(result[0], nil)
 | |
| 		if err != nil {
 | |
| 			return
 | |
| 		}
 | |
| 		if cursor == 0 { // over
 | |
| 			return
 | |
| 		}
 | |
| 	}
 | |
| }
 | |
| 
 | |
| // StartAndGC starts the redis cache adapter.
 | |
| // config: must be in this format {"key":"collection key","conn":"connection info","dbNum":"0", "skipEmptyPrefix":"true"}
 | |
| // Cached items in redis are stored forever, no garbage collection happens
 | |
| func (rc *Cache) StartAndGC(config string) error {
 | |
| 	err := rc.parseConf(config)
 | |
| 	if err != nil {
 | |
| 		return err
 | |
| 	}
 | |
| 
 | |
| 	rc.connectInit()
 | |
| 
 | |
| 	c := rc.p.Get()
 | |
| 	defer func() {
 | |
| 		_ = c.Close()
 | |
| 	}()
 | |
| 
 | |
| 	// test connection
 | |
| 	if err = c.Err(); err != nil {
 | |
| 		return berror.Wrapf(err, cache.InvalidConnection,
 | |
| 			"can not connect to remote redis server, please check the connection info and network state: %s", config)
 | |
| 	}
 | |
| 	return nil
 | |
| }
 | |
| 
 | |
| func (rc *Cache) parseConf(config string) error {
 | |
| 	var cf redisConfig
 | |
| 	err := json.Unmarshal([]byte(config), &cf)
 | |
| 	if err != nil {
 | |
| 		return berror.Wrapf(err, cache.InvalidRedisCacheCfg, "could not unmarshal the config: %s", config)
 | |
| 	}
 | |
| 
 | |
| 	err = cf.parse()
 | |
| 	if err != nil {
 | |
| 		return err
 | |
| 	}
 | |
| 
 | |
| 	rc.dbNum = cf.dbNum
 | |
| 	rc.key = cf.Key
 | |
| 	rc.conninfo = cf.Conn
 | |
| 	rc.password = cf.password
 | |
| 	rc.maxIdle = cf.maxIdle
 | |
| 	rc.timeout = cf.timeout
 | |
| 	rc.skipEmptyPrefix = cf.skipEmptyPrefix
 | |
| 
 | |
| 	return nil
 | |
| }
 | |
| 
 | |
| type redisConfig struct {
 | |
| 	DbNum           string `json:"dbNum"`
 | |
| 	SkipEmptyPrefix string `json:"skipEmptyPrefix"`
 | |
| 	Key             string `json:"key"`
 | |
| 	// Format redis://<password>@<host>:<port>
 | |
| 	Conn       string `json:"conn"`
 | |
| 	MaxIdle    string `json:"maxIdle"`
 | |
| 	TimeoutStr string `json:"timeout"`
 | |
| 
 | |
| 	dbNum           int
 | |
| 	skipEmptyPrefix bool
 | |
| 	maxIdle         int
 | |
| 	// parse from Conn
 | |
| 	password string
 | |
| 	// timeout used for idle connection, default is 180 seconds.
 | |
| 	timeout time.Duration
 | |
| }
 | |
| 
 | |
| // parse parses the config.
 | |
| // If the necessary settings have not been set, it will return an error.
 | |
| // It will fill the default values if some fields are missing.
 | |
| func (cf *redisConfig) parse() error {
 | |
| 	if cf.Conn == "" {
 | |
| 		return berror.Error(cache.InvalidRedisCacheCfg, "config missing conn field")
 | |
| 	}
 | |
| 
 | |
| 	// Format redis://<password>@<host>:<port>
 | |
| 	cf.Conn = strings.Replace(cf.Conn, "redis://", "", 1)
 | |
| 	if i := strings.Index(cf.Conn, "@"); i > -1 {
 | |
| 		cf.password = cf.Conn[0:i]
 | |
| 		cf.Conn = cf.Conn[i+1:]
 | |
| 	}
 | |
| 
 | |
| 	if cf.Key == "" {
 | |
| 		cf.Key = DefaultKey
 | |
| 	}
 | |
| 
 | |
| 	if cf.DbNum != "" {
 | |
| 		cf.dbNum, _ = strconv.Atoi(cf.DbNum)
 | |
| 	}
 | |
| 
 | |
| 	if cf.SkipEmptyPrefix != "" {
 | |
| 		cf.skipEmptyPrefix, _ = strconv.ParseBool(cf.SkipEmptyPrefix)
 | |
| 	}
 | |
| 
 | |
| 	if cf.MaxIdle == "" {
 | |
| 		cf.maxIdle = defaultMaxIdle
 | |
| 	} else {
 | |
| 		cf.maxIdle, _ = strconv.Atoi(cf.MaxIdle)
 | |
| 	}
 | |
| 
 | |
| 	if v, err := time.ParseDuration(cf.TimeoutStr); err == nil {
 | |
| 		cf.timeout = v
 | |
| 	} else {
 | |
| 		cf.timeout = defaultTimeout
 | |
| 	}
 | |
| 
 | |
| 	return nil
 | |
| }
 | |
| 
 | |
| // connect to redis.
 | |
| func (rc *Cache) connectInit() {
 | |
| 	dialFunc := func() (c redis.Conn, err error) {
 | |
| 		c, err = redis.Dial("tcp", rc.conninfo)
 | |
| 		if err != nil {
 | |
| 			return nil, berror.Wrapf(err, cache.DialFailed,
 | |
| 				"could not dial to remote server: %s ", rc.conninfo)
 | |
| 		}
 | |
| 
 | |
| 		if rc.password != "" {
 | |
| 			if _, err = c.Do("AUTH", rc.password); err != nil {
 | |
| 				_ = c.Close()
 | |
| 				return nil, err
 | |
| 			}
 | |
| 		}
 | |
| 
 | |
| 		_, selecterr := c.Do("SELECT", rc.dbNum)
 | |
| 		if selecterr != nil {
 | |
| 			_ = c.Close()
 | |
| 			return nil, selecterr
 | |
| 		}
 | |
| 		return
 | |
| 	}
 | |
| 	// initialize a new pool
 | |
| 	rc.p = &redis.Pool{
 | |
| 		MaxIdle:     rc.maxIdle,
 | |
| 		IdleTimeout: rc.timeout,
 | |
| 		Dial:        dialFunc,
 | |
| 	}
 | |
| }
 | |
| 
 | |
| func init() {
 | |
| 	cache.Register("redis", NewRedisCache)
 | |
| }
 |