Compare commits

...
9 Commits
Author SHA1 Message Date
maxlerebourg 0c2668d578 add redis database selection (#100)
*  add redis database selection

* 📝 update docs

* 📝 readme
2023-05-25 17:20:14 +02:00
maxlerebourg abae7ee028 📝 update version used (#97)
* 📝 update version used

* 📝 update doc

* 📝 documentation

* 📝 documentation

* 📝 versionning
2023-04-17 09:25:57 +02:00
maxlerebourg 1fcd4f4e2f remove down at start if crowdsec unavailable (#93)
*  remove down at start if crowdsec unavailable

* 🚨 fix lint
2023-03-12 20:49:00 +01:00
mathieuHa 39fcc38980 🐛Bump simple redis to 1.0.6 to fix bug hang with password, update doc on redis (#89) 2023-03-05 14:32:43 +01:00
Max Lerebourg dd322a966a 🍱 fix lint 2023-03-04 12:11:22 +01:00
Max Lerebourg f0bb140596 🍱 fix lint 2023-03-04 12:08:03 +01:00
Max Lerebourg 46e581eca2 🍱 fix readme to add redis pass 2023-03-04 12:05:09 +01:00
maxlerebourg 50690d1ac7 handle redis password (#87)
*  handle redis password

* 🍱 fix version
2023-03-04 11:51:54 +01:00
maxlerebourg b079073ff6 handle isHealthy in the main function and log error became… (#84)
*  handle isHealthy in the main function and log error became debug

* fix: lint

* fix: lint
2023-03-01 14:18:19 +01:00
23 changed files with 212 additions and 87 deletions
+2
View File
@@ -30,4 +30,6 @@ Steps to reproduce the behavior:
3. Scroll down to '....'
4. See error
<!---
If you like the plugin, please consider starring it, so you can get updates and we get some more visibility ✨
-->
@@ -16,4 +16,6 @@ A clear and concise description of what you want to happen.
**Additional context**
Add any other context or screenshots about the feature request here.
<!---
If you like the plugin, please consider starring it, so you can get updates and we get some more visibility ✨
-->
+11
View File
@@ -103,6 +103,14 @@ make run
- string
- default: "redis:6379"
- hostname and port for the Redis service
- RedisCachePassword
- string
- default: ""
- Password for the Redis service
- RedisCacheDatabase
- string
- default: ""
- Database selection for the Redis service
- UpdateIntervalSeconds
- int64
- default: 60
@@ -134,6 +142,7 @@ experimental:
plugins:
bouncer:
moduleName: github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin
version: vX.Y.Z # To update
```
```yaml
@@ -183,6 +192,8 @@ http:
forwardedHeadersCustomName: X-Custom-Header
redisCacheEnabled: false
redisCacheHost: "redis:6379"
redisCachePassword: password
redisCacheDatabase: "5"
crowdsecLapiTLSCertificateAuthority: |-
-----BEGIN CERTIFICATE-----
MIIEBzCCAu+gAwIBAgICEAAwDQYJKoZIhvcNAQELBQAwgZQxCzAJBgNVBAYTAlVT
+30 -16
View File
@@ -141,21 +141,26 @@ func New(ctx context.Context, next http.Handler, config *configuration.Config, n
},
cacheClient: &cache.Client{},
}
bouncer.cacheClient.New(config.RedisCacheEnabled, config.RedisCacheHost)
config.RedisCachePassword, _ = configuration.GetVariable(config, "RedisCachePassword")
bouncer.cacheClient.New(
config.RedisCacheEnabled,
config.RedisCacheHost,
config.RedisCachePassword,
config.RedisCacheDatabase,
)
if (config.CrowdsecMode == configuration.StreamMode || config.CrowdsecMode == configuration.AloneMode) && ticker == nil {
if config.CrowdsecMode == configuration.AloneMode {
err = getToken(bouncer)
if err != nil {
if err := getToken(bouncer); err != nil {
logger.Error(fmt.Sprintf("New:getToken %s", err.Error()))
return nil, err
}
}
ticker = startTicker(config, func() {
handleStreamCache(bouncer)
})
handleStreamCache(bouncer)
handleStreamTicker(bouncer)
isStartup = false
ticker = startTicker(config, func() {
handleStreamTicker(bouncer)
})
}
logger.Debug(fmt.Sprintf("New initialized mode:%s", config.CrowdsecMode))
@@ -218,7 +223,7 @@ func (bouncer *Bouncer) ServeHTTP(rw http.ResponseWriter, req *http.Request) {
if isCrowdsecStreamHealthy {
bouncer.next.ServeHTTP(rw, req)
} else {
logger.Error(fmt.Sprintf("ServeHTTP isCrowdsecStreamHealthy:false ip:%s", remoteIP))
logger.Debug(fmt.Sprintf("ServeHTTP isCrowdsecStreamHealthy:false ip:%s", remoteIP))
rw.WriteHeader(http.StatusForbidden)
}
} else {
@@ -261,6 +266,15 @@ type Login struct {
Expire string `json:"expire"`
}
func handleStreamTicker(bouncer *Bouncer) {
if err := handleStreamCache(bouncer); err != nil {
isCrowdsecStreamHealthy = false
logger.Error(err.Error())
} else {
isCrowdsecStreamHealthy = true
}
}
func startTicker(config *configuration.Config, work func()) chan bool {
ticker := time.NewTicker(time.Duration(config.UpdateIntervalSeconds) * time.Second)
stop := make(chan bool, 1)
@@ -348,7 +362,7 @@ func getToken(bouncer *Bouncer) error {
return fmt.Errorf("getToken statusCode:%d", login.Code)
}
func handleStreamCache(bouncer *Bouncer) {
func handleStreamCache(bouncer *Bouncer) error {
// TODO clean properly on exit.
// Instead of blocking the goroutine interval for all the secondary node,
// if the master service is shut down, other goroutine can take the lead
@@ -356,7 +370,10 @@ func handleStreamCache(bouncer *Bouncer) {
_, err := bouncer.cacheClient.GetDecision(cacheTimeoutKey)
if err == nil {
logger.Debug("handleStreamCache:alreadyUpdated")
return
return nil
}
if err.Error() != cache.CacheMiss {
return err
}
bouncer.cacheClient.SetDecision(cacheTimeoutKey, false, bouncer.updateInterval-1)
streamRouteURL := url.URL{
@@ -367,16 +384,12 @@ func handleStreamCache(bouncer *Bouncer) {
}
body, err := crowdsecQuery(bouncer, streamRouteURL.String(), false)
if err != nil {
logger.Error(err.Error())
isCrowdsecStreamHealthy = false
return
return err
}
var stream Stream
err = json.Unmarshal(body, &stream)
if err != nil {
logger.Error(fmt.Sprintf("handleStreamCache:parsingBody %s", err.Error()))
isCrowdsecStreamHealthy = false
return
return fmt.Errorf("handleStreamCache:parsingBody %w", err)
}
for _, decision := range stream.New {
duration, err := time.ParseDuration(decision.Duration)
@@ -389,6 +402,7 @@ func handleStreamCache(bouncer *Bouncer) {
}
logger.Debug("handleStreamCache:updated")
isCrowdsecStreamHealthy = true
return nil
}
func crowdsecQuery(bouncer *Bouncer, stringURL string, isPost bool) ([]byte, error) {
+8 -3
View File
@@ -142,14 +142,19 @@ func Test_handleStreamCache(t *testing.T) {
bouncer *Bouncer
}
tests := []struct {
name string
args args
name string
args args
wantErr bool
}{
// TODO: Add test cases.
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
handleStreamCache(tt.args.bouncer)
err := handleStreamCache(tt.args.bouncer)
if (err != nil) != tt.wantErr {
t.Errorf("handleStreamCache() error = %v, wantErr %v", err, tt.wantErr)
return
}
})
}
}
+2 -2
View File
@@ -2,7 +2,7 @@ version: "3.8"
services:
traefik:
image: "traefik:v2.9.6"
image: "traefik:v2.9.10"
container_name: "traefik"
restart: unless-stopped
command:
@@ -52,7 +52,7 @@ services:
- "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.crowdseclapikey=44c36dac5c4140af9f06f397508e82c7"
crowdsec:
image: crowdsecurity/crowdsec:v1.4.1
image: crowdsecurity/crowdsec:v1.4.6
container_name: "crowdsec"
restart: unless-stopped
environment:
+3 -3
View File
@@ -2,7 +2,7 @@ version: "3.8"
services:
traefik:
image: "traefik:v2.9.6"
image: "traefik:v2.9.10"
container_name: "traefik"
restart: unless-stopped
command:
@@ -14,7 +14,7 @@ services:
- "--entrypoints.web.address=:80"
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
- "--experimental.plugins.bouncer.version=v1.1.8"
- "--experimental.plugins.bouncer.version=v1.1.11"
volumes:
- "/var/run/docker.sock:/var/run/docker.sock:ro"
- "logs:/var/log/traefik"
@@ -63,7 +63,7 @@ services:
- "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.forwardedheaderstrustedips=172.21.0.5"
crowdsec:
image: crowdsecurity/crowdsec:v1.4.1
image: crowdsecurity/crowdsec:v1.4.6
container_name: "crowdsec"
restart: unless-stopped
environment:
@@ -2,7 +2,7 @@ version: "3.8"
services:
cloudflare:
image: "traefik:v2.9.6"
image: "traefik:v2.9.10"
container_name: "cloudflare"
restart: unless-stopped
command:
@@ -12,7 +12,6 @@ services:
- "--api.insecure=true"
- "--entrypoints.web.address=:80"
- "--providers.file.filename=/cloud.yaml"
- "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
volumes:
- /var/run/docker.sock:/var/run/docker.sock:ro
- ./cloudflare-exemple.yaml:/cloud.yaml:ro
@@ -22,7 +21,7 @@ services:
- 8080:8080
traefik:
image: "traefik:v2.9.6"
image: "traefik:v2.9.10"
container_name: "traefik"
restart: unless-stopped
command:
@@ -36,7 +35,7 @@ services:
- "--entrypoints.web.forwardedheaders.trustedips=172.21.0.5"
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
- "--experimental.plugins.bouncer.version=v1.1.7"
- "--experimental.plugins.bouncer.version=v1.1.11"
volumes:
- /var/run/docker.sock:/var/run/docker.sock:ro
- logs-traefik:/var/log/traefik
@@ -46,7 +45,7 @@ services:
depends_on:
- crowdsec
whoami1:
whoami-foo:
image: traefik/whoami
container_name: "simple-service-foo"
restart: unless-stopped
@@ -66,7 +65,7 @@ services:
- "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.forwardedheaderstrustedips=172.21.0.5"
- "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.loglevel=DEBUG"
whoami2:
whoami-bar:
image: traefik/whoami
container_name: "simple-service-bar"
restart: unless-stopped
@@ -88,7 +87,7 @@ services:
crowdsec:
image: crowdsecurity/crowdsec:v1.4.3
image: crowdsecurity/crowdsec:v1.4.6
container_name: "crowdsec"
restart: unless-stopped
environment:
@@ -2,7 +2,7 @@
DEBIAN_FRONTEND=noninteractive sudo apt-get update && sudo apt-get install wget -y
# DEBIAN_FRONTEND=noninteractive sudo apt-get upgrade -y --assume-yes
wget -O traefik.tar.gz "https://github.com/traefik/traefik/releases/download/v2.9.6/traefik_v2.9.6_linux_amd64.tar.gz"
wget -O traefik.tar.gz "https://github.com/traefik/traefik/releases/download/v2.9.10/traefik_v2.9.10_linux_amd64.tar.gz"
tar -zxvf traefik.tar.gz
# inspired from https://gist.github.com/ubergesundheit/7c9d875befc2d7bfd0bf43d8b3862d85
sudo mv ./traefik /usr/local/bin/
+1 -1
View File
@@ -1,7 +1,7 @@
#!/bin/bash
DEBIAN_FRONTEND=noninteractive sudo apt-get update && apt-get install wget -y
wget -O whoami.tar.gz "https://github.com/traefik/whoami/releases/download/v1.8.7/whoami_v1.8.7_linux_amd64.tar.gz"
wget -O whoami.tar.gz "https://github.com/traefik/whoami/releases/download/v1.9.0/whoami_v1.9.0_linux_amd64.tar.gz"
tar -zxvf whoami.tar.gz
# inspired from https://gist.github.com/ubergesundheit/7c9d875befc2d7bfd0bf43d8b3862d85
sudo mv ./whoami /usr/local/bin/
+1 -1
View File
@@ -1,5 +1,5 @@
image:
tag: v1.4.4-rc1
tag: v1.4.6
agent:
acquisition:
+2 -2
View File
@@ -1,5 +1,5 @@
image:
tag: v2.9.6
tag: v2.9.10
logs:
general:
@@ -16,4 +16,4 @@ experimental:
additionalArguments:
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
- "--experimental.plugins.bouncer.version=v1.1.7-beta1"
- "--experimental.plugins.bouncer.version=v1.1.11"
+31 -16
View File
@@ -2,7 +2,7 @@ version: "3.8"
services:
traefik:
image: "traefik:v2.9.6"
image: "traefik:v2.9.10"
container_name: "traefik"
restart: unless-stopped
command:
@@ -15,7 +15,7 @@ services:
- "--entrypoints.web.address=:80"
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
- "--experimental.plugins.bouncer.version=v1.1.7"
- "--experimental.plugins.bouncer.version=v1.1.11"
# - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
volumes:
- /var/run/docker.sock:/var/run/docker.sock:ro
@@ -26,16 +26,17 @@ services:
- 8080:8080
depends_on:
- crowdsec
- redis
- redis-insecure
- redis-secure
whoami-foo:
whoami-redis-insecure:
image: traefik/whoami
container_name: "simple-service-foo"
restart: unless-stopped
labels:
- "traefik.enable=true"
# Definition of the router
- "traefik.http.routers.router-foo.rule=Path(`/foo`)"
- "traefik.http.routers.router-foo.rule=Path(`/redis-insecure`)"
- "traefik.http.routers.router-foo.entrypoints=web"
- "traefik.http.routers.router-foo.middlewares=crowdsec-foo@docker"
# Definition of the service
@@ -45,16 +46,18 @@ services:
# crowdseclapikey must be uniq to the middleware attached to the service
- "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.crowdseclapikey=40796d93c2958f9e58345514e67740e5"
- "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.rediscacheenabled=true"
# Contact redis-unsecure without a password
- "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.rediscachehost=redis-insecure:6379"
- "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.loglevel=DEBUG"
whoami-bar:
whoami-redis-secure:
image: traefik/whoami
container_name: "simple-service-bar"
restart: unless-stopped
labels:
- "traefik.enable=true"
# Definition of the router
- "traefik.http.routers.router-bar.rule=Path(`/bar`)"
- "traefik.http.routers.router-bar.rule=Path(`/redis-secure`)"
- "traefik.http.routers.router-bar.entrypoints=web"
- "traefik.http.routers.router-bar.middlewares=crowdsec-bar@docker"
# Definition of the service
@@ -64,11 +67,14 @@ services:
# crowdseclapikey must be uniq to the middleware attached to the service
- "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.crowdseclapikey=44c36dac5c4140af9f06f397508e82c7"
- "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.rediscacheenabled=true"
- "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.rediscachepassword=FIXME"
# Contact redis-secure with password
- "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.rediscachehost=redis-secure:6379"
- "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.loglevel=DEBUG"
crowdsec:
image: crowdsecurity/crowdsec:v1.4.3
image: crowdsecurity/crowdsec:v1.4.6
container_name: "crowdsec"
restart: unless-stopped
environment:
@@ -84,18 +90,27 @@ services:
labels:
- "traefik.enable=false"
redis:
image: "redis:7.0.5-alpine"
container_name: "redis"
redis-secure:
image: "redis:7.0.9-alpine"
container_name: "redis-secure"
hostname: redis-secure
restart: unless-stopped
command: "redis-server --save 60 1"
command: "redis-server --save 60 1 --loglevel debug --requirepass FIXME"
volumes:
- redis-data:/data
ports:
- 6379:6379
- redis-secure-data:/data
redis-insecure:
image: "redis:7.0.9-alpine"
container_name: "redis-insecure"
hostname: redis-unsecure
restart: unless-stopped
command: "redis-server --save 60 1 --loglevel debug"
volumes:
- redis-unsecure-data:/data
volumes:
logs-redis:
crowdsec-db-redis:
crowdsec-config-redis:
redis-data:
redis-unsecure-data:
redis-secure-data:
@@ -2,7 +2,7 @@ version: "3.8"
services:
traefik:
image: "traefik:v2.9.6"
image: "traefik:v2.9.10"
container_name: "traefik"
restart: unless-stopped
command:
@@ -15,7 +15,7 @@ services:
- "--entrypoints.web.address=:80"
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
- "--experimental.plugins.bouncer.version=v1.1.7"
- "--experimental.plugins.bouncer.version=v1.1.11"
volumes:
- /var/run/docker.sock:/var/run/docker.sock:ro
ports:
@@ -2,7 +2,7 @@ version: "3.8"
services:
traefik:
image: "traefik:v2.9.6"
image: "traefik:v2.9.10"
container_name: "traefik"
restart: unless-stopped
command:
@@ -15,7 +15,7 @@ services:
- "--entrypoints.web.address=:80"
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
- "--experimental.plugins.bouncer.version=v1.1.7"
- "--experimental.plugins.bouncer.version=v1.1.11"
# - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
volumes:
- /var/run/docker.sock:/var/run/docker.sock:ro
@@ -64,7 +64,7 @@ services:
- "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.crowdsecLapiTLSCertificateBouncerKeyFile=/etc/traefik/crowdsec-certs/bouncer-key.pem"
crowdsec:
image: crowdsecurity/crowdsec:v1.4.3
image: crowdsecurity/crowdsec:v1.4.6
container_name: "crowdsec"
restart: unless-stopped
environment:
@@ -2,7 +2,7 @@ version: "3.8"
services:
traefik:
image: "traefik:v2.9.6"
image: "traefik:v2.9.10"
container_name: "traefik"
restart: unless-stopped
command:
@@ -15,7 +15,7 @@ services:
- "--entrypoints.web.address=:80"
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
- "--experimental.plugins.bouncer.version=v1.1.7"
- "--experimental.plugins.bouncer.version=v1.1.11"
# - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
volumes:
- /var/run/docker.sock:/var/run/docker.sock:ro
@@ -69,7 +69,7 @@ services:
crowdsec:
image: crowdsecurity/crowdsec:v1.4.3
image: crowdsecurity/crowdsec:v1.4.6
container_name: "crowdsec"
restart: unless-stopped
environment:
+1 -1
View File
@@ -4,5 +4,5 @@ go 1.19
require (
github.com/leprosus/golang-ttl-map v1.1.7
github.com/maxlerebourg/simpleredis v1.0.3
github.com/maxlerebourg/simpleredis v1.0.7
)
+2 -2
View File
@@ -1,4 +1,4 @@
github.com/leprosus/golang-ttl-map v1.1.7 h1:cF4AAFDDnJTFSV+/42sKLhmMluvLdRlCGS2UaifH6UM=
github.com/leprosus/golang-ttl-map v1.1.7/go.mod h1:4QWHJPeVBbrkhOhXdhCv9IEiyj/YzkO04/iexy4vSe0=
github.com/maxlerebourg/simpleredis v1.0.3 h1:VhXq9bVytWDqD/TS/GjHKayvQb/VUeEql5F+yUbdOiI=
github.com/maxlerebourg/simpleredis v1.0.3/go.mod h1:/DH8zOK6kDskSqoX/m5CJJdNGfkIQZd/ERBJgytDDSk=
github.com/maxlerebourg/simpleredis v1.0.7 h1:d53p3GOIgQtxxWuqsWMGTJ0dZjP8UhiDX8L1q2QZ/vY=
github.com/maxlerebourg/simpleredis v1.0.7/go.mod h1:/DH8zOK6kDskSqoX/m5CJJdNGfkIQZd/ERBJgytDDSk=
+2 -2
View File
@@ -82,9 +82,9 @@ type Client struct {
}
// New Initialize cache client.
func (client *Client) New(isRedis bool, host string) {
func (client *Client) New(isRedis bool, host, pass, database string) {
if isRedis {
redis.Init(host)
redis.Init(host, pass, database)
client.cache = &redisCache{}
} else {
client.cache = &localCache{}
+10 -1
View File
@@ -55,6 +55,9 @@ type Config struct {
ClientTrustedIPs []string `json:"clientTrustedIps,omitempty"`
RedisCacheEnabled bool `json:"redisCacheEnabled,omitempty"`
RedisCacheHost string `json:"redisCacheHost,omitempty"`
RedisCachePassword string `json:"redisCachePassword,omitempty"`
RedisCachePasswordFile string `json:"redisCachePasswordFile,omitempty"`
RedisCacheDatabase string `json:"redisCacheDatabase,omitempty"`
}
func contains(source []string, target string) bool {
@@ -83,6 +86,8 @@ func New() *Config {
ClientTrustedIPs: []string{},
RedisCacheEnabled: false,
RedisCacheHost: "redis:6379",
RedisCachePassword: "",
RedisCacheDatabase: "",
}
}
@@ -115,7 +120,7 @@ func GetVariable(config *Config, key string) (string, error) {
// ValidateParams validate all the param gave by user.
//
//nolint:gocyclo
//nolint:gocyclo,gocognit
func ValidateParams(config *Config) error {
if err := validateParamsRequired(config); err != nil {
return err
@@ -128,6 +133,10 @@ func ValidateParams(config *Config) error {
return err
}
if _, err := GetVariable(config, "RedisCachePassword"); err != nil {
return err
}
if config.CrowdsecMode == AloneMode {
if _, err := GetVariable(config, "CrowdsecCapiMachineID"); err != nil {
return err
+30 -1
View File
@@ -1,2 +1,31 @@
# simpleredis
Minimal go redis with only get, set and delete operation
Minimal go redis with only `get`, `set` and `delete` operation.
It supports password authentication with redis.
With **NO** external dependencies.
## Example
```go
import simpleredis "github.com/maxlerebourg/simpleredis"
var redis simpleredis.SimpleRedis
redis.Init("redis:6379", "", "") // redisHost, redisPass, redisDatabase
err := redis.Set("test", []bytes("whatever"), 60), // Set key "test" with "whatever" for 60 seconds
if err != nil {
...
}
val, err := redis.Get("test") // get key test
if err != nil {
// err could be only redis:unreachable, redis:miss or redis:timeout available in simpleredis.RedisUnreachable
...
}
err = redis.Del("test")
if err != nil {
...
}
```
## Author
Max Lerebourg @ [Primadviz.com](https://primadviz.com)
Mathieu Hanotaux
+58 -19
View File
@@ -17,10 +17,11 @@ const (
RedisUnreachable = "redis:unreachable"
RedisMiss = "redis:miss"
RedisTimeout = "redis:timeout"
RedisNoAuth = "redis:noauth"
)
// A RedisCmd is used to communicate with redis at low level using commands.
type RedisCmd struct {
// A redisCmd is used to communicate with redis at low level using commands.
type redisCmd struct {
Command string
Name string
Data []byte
@@ -30,7 +31,9 @@ type RedisCmd struct {
// A SimpleRedis is used to communicate with redis.
type SimpleRedis struct {
redisHost string
host string
pass string
database string
}
func genRedisArray(params ...[]byte) []byte {
@@ -49,11 +52,29 @@ func send(wr *textproto.Writer, method string, data []byte) {
}
}
func askRedis(hostnamePort string, cmd RedisCmd, channel chan RedisCmd) {
func (sr *SimpleRedis) waitRedis(reader *textproto.Reader, channel chan redisCmd) {
for {
select {
case <-time.After(time.Second * 1):
channel <- redisCmd{Error: fmt.Errorf(RedisTimeout)}
return
default:
read, _ := reader.ReadLineBytes()
if string(read) != "+OK" {
channel <- redisCmd{Error: fmt.Errorf(RedisNoAuth)}
return
}
}
// breaks out of for
break
}
}
func (sr *SimpleRedis) askRedis(cmd redisCmd, channel chan redisCmd) {
dialer := net.Dialer{Timeout: 2 * time.Second}
conn, err := dialer.Dial("tcp", hostnamePort)
conn, err := dialer.Dial("tcp", sr.host)
if err != nil {
channel <- RedisCmd{Error: fmt.Errorf(RedisUnreachable)}
channel <- redisCmd{Error: fmt.Errorf(RedisUnreachable)}
return
}
defer func() {
@@ -65,6 +86,18 @@ func askRedis(hostnamePort string, cmd RedisCmd, channel chan RedisCmd) {
writer := textproto.NewWriter(bufio.NewWriter(conn))
reader := textproto.NewReader(bufio.NewReader(conn))
if sr.pass != "" {
data := genRedisArray([]byte("AUTH"), []byte(sr.pass))
send(writer, "auth", data)
sr.waitRedis(reader, channel)
}
if sr.database != "" {
data := genRedisArray([]byte("SELECT"), []byte(sr.database))
send(writer, "select", data)
sr.waitRedis(reader, channel)
}
switch cmd.Command {
case "SET":
data := genRedisArray([]byte("SET"), []byte(cmd.Name), cmd.Data, []byte("EX"), []byte(fmt.Sprintf("%d", cmd.Duration)))
@@ -78,16 +111,20 @@ func askRedis(hostnamePort string, cmd RedisCmd, channel chan RedisCmd) {
for {
select {
case <-time.After(time.Second * 1):
channel <- RedisCmd{Error: fmt.Errorf(RedisTimeout)}
channel <- redisCmd{Error: fmt.Errorf(RedisTimeout)}
return
default:
read, _ := reader.ReadLineBytes()
if string(read) != "$1" {
channel <- RedisCmd{Error: fmt.Errorf(RedisMiss)}
str := string(read)
if strings.Contains(str, "-NOAUTH") {
channel <- redisCmd{Error: fmt.Errorf(RedisNoAuth)}
return
} else if str != "$1" {
channel <- redisCmd{Error: fmt.Errorf(RedisMiss)}
return
}
read, _ = reader.ReadLineBytes()
channel <- RedisCmd{Data: read}
channel <- redisCmd{Data: read}
return
}
}
@@ -95,18 +132,20 @@ func askRedis(hostnamePort string, cmd RedisCmd, channel chan RedisCmd) {
}
// Init sets the redisHost used to connect to redis.
func (sr *SimpleRedis) Init(redisHost string) {
sr.redisHost = redisHost
func (sr *SimpleRedis) Init(host, pass, database string) {
sr.host = host
sr.pass = pass
sr.database = database
}
// Get fetches the value for key name in redis.
func (sr *SimpleRedis) Get(name string) ([]byte, error) {
redisCmd := RedisCmd{
cmd := redisCmd{
Command: "GET",
Name: name,
}
channel := make(chan RedisCmd)
go askRedis(sr.redisHost, redisCmd, channel)
channel := make(chan redisCmd)
go sr.askRedis(cmd, channel)
resp := <-channel
if resp.Error != nil {
return nil, resp.Error
@@ -116,22 +155,22 @@ func (sr *SimpleRedis) Get(name string) ([]byte, error) {
// Set updates the value for key name in redis with value data for duration.
func (sr *SimpleRedis) Set(name string, data []byte, duration int64) error {
redisCmd := RedisCmd{
cmd := redisCmd{
Command: "SET",
Name: name,
Data: data,
Duration: duration,
}
go askRedis(sr.redisHost, redisCmd, nil)
go sr.askRedis(cmd, nil)
return nil
}
// Del removes the key name in redis.
func (sr *SimpleRedis) Del(name string) error {
redisCmd := RedisCmd{
cmd := redisCmd{
Command: "DEL",
Name: name,
}
go askRedis(sr.redisHost, redisCmd, nil)
go sr.askRedis(cmd, nil)
return nil
}
+1 -1
View File
@@ -1,6 +1,6 @@
# github.com/leprosus/golang-ttl-map v1.1.7
## explicit; go 1.15
github.com/leprosus/golang-ttl-map
# github.com/maxlerebourg/simpleredis v1.0.3
# github.com/maxlerebourg/simpleredis v1.0.7
## explicit; go 1.19
github.com/maxlerebourg/simpleredis