Compare commits

..
7 Commits
Author SHA1 Message Date
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
maxlerebourg 976cbb7d1f 81 bug stream mode stops blocking (#82)
*  fix isHealthy issue at startup

* 🍱 not added in first commit ?

* 🍱 remove unused import

* 🍱 fix lint

* fix: lint
2023-01-30 14:03:10 +01:00
Max Lerebourg 80726df450 🐛 fix alone mode 2023-01-25 20:43:07 +01:00
12 changed files with 188 additions and 56 deletions
+5
View File
@@ -103,6 +103,10 @@ make run
- string - string
- default: "redis:6379" - default: "redis:6379"
- hostname and port for the Redis service - hostname and port for the Redis service
- RedisCachePassword
- string
- default: ""
- Password for the Redis service
- UpdateIntervalSeconds - UpdateIntervalSeconds
- int64 - int64
- default: 60 - default: 60
@@ -183,6 +187,7 @@ http:
forwardedHeadersCustomName: X-Custom-Header forwardedHeadersCustomName: X-Custom-Header
redisCacheEnabled: false redisCacheEnabled: false
redisCacheHost: "redis:6379" redisCacheHost: "redis:6379"
redisCachePassword: password
crowdsecLapiTLSCertificateAuthority: |- crowdsecLapiTLSCertificateAuthority: |-
-----BEGIN CERTIFICATE----- -----BEGIN CERTIFICATE-----
MIIEBzCCAu+gAwIBAgICEAAwDQYJKoZIhvcNAQELBQAwgZQxCzAJBgNVBAYTAlVT MIIEBzCCAu+gAwIBAgICEAAwDQYJKoZIhvcNAQELBQAwgZQxCzAJBgNVBAYTAlVT
+32 -22
View File
@@ -15,8 +15,6 @@ import (
"text/template" "text/template"
"time" "time"
simpleredis "github.com/maxlerebourg/simpleredis"
cache "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/cache" cache "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/cache"
configuration "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/configuration" configuration "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/configuration"
ip "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/ip" ip "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/ip"
@@ -35,6 +33,7 @@ const (
//nolint:gochecknoglobals //nolint:gochecknoglobals
var ( var (
isStartup = true
isCrowdsecStreamHealthy = true isCrowdsecStreamHealthy = true
ticker chan bool ticker chan bool
) )
@@ -85,7 +84,7 @@ func New(ctx context.Context, next http.Handler, config *configuration.Config, n
crowdsecStreamRoute := "" crowdsecStreamRoute := ""
crowdsecHeader := "" crowdsecHeader := ""
if config.CrowdsecMode == configuration.AloneMode { if config.CrowdsecMode == configuration.AloneMode {
config.CrowdsecCapiMachineID, _ = configuration.GetVariable(config, "CrowdsecCapiMachineId") config.CrowdsecCapiMachineID, _ = configuration.GetVariable(config, "CrowdsecCapiMachineID")
config.CrowdsecCapiPassword, _ = configuration.GetVariable(config, "CrowdsecCapiPassword") config.CrowdsecCapiPassword, _ = configuration.GetVariable(config, "CrowdsecCapiPassword")
config.CrowdsecLapiHost = "api.crowdsec.net" config.CrowdsecLapiHost = "api.crowdsec.net"
config.CrowdsecLapiScheme = "https" config.CrowdsecLapiScheme = "https"
@@ -142,20 +141,29 @@ func New(ctx context.Context, next http.Handler, config *configuration.Config, n
}, },
cacheClient: &cache.Client{}, 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)
//nolint:nestif
if (config.CrowdsecMode == configuration.StreamMode || config.CrowdsecMode == configuration.AloneMode) && ticker == nil { if (config.CrowdsecMode == configuration.StreamMode || config.CrowdsecMode == configuration.AloneMode) && ticker == nil {
if config.CrowdsecMode == configuration.AloneMode { if config.CrowdsecMode == configuration.AloneMode {
err = getToken(bouncer) if err := getToken(bouncer); err != nil {
if err != nil {
logger.Error(fmt.Sprintf("New:getToken %s", err.Error())) logger.Error(fmt.Sprintf("New:getToken %s", err.Error()))
return nil, err return nil, err
} }
} }
if err := handleStreamCache(bouncer); err != nil {
return nil, err
}
isStartup = false
ticker = startTicker(config, func() { ticker = startTicker(config, func() {
handleStreamCache(bouncer) if err := handleStreamCache(bouncer); err != nil {
isCrowdsecStreamHealthy = false
logger.Error(err.Error())
} else {
isCrowdsecStreamHealthy = true
}
}) })
go handleStreamCache(bouncer)
} }
logger.Debug(fmt.Sprintf("New initialized mode:%s", config.CrowdsecMode)) logger.Debug(fmt.Sprintf("New initialized mode:%s", config.CrowdsecMode))
@@ -193,10 +201,12 @@ func (bouncer *Bouncer) ServeHTTP(rw http.ResponseWriter, req *http.Request) {
// TODO This should be simplified // TODO This should be simplified
if bouncer.crowdsecMode != configuration.NoneMode { if bouncer.crowdsecMode != configuration.NoneMode {
isBanned, erro := bouncer.cacheClient.GetDecision(remoteIP) isBanned, cacheErr := bouncer.cacheClient.GetDecision(remoteIP)
if erro != nil { if cacheErr != nil {
logger.Debug(fmt.Sprintf("ServeHTTP:getDecision ip:%s isBanned:true %s", remoteIP, erro.Error())) errString := cacheErr.Error()
if erro.Error() == simpleredis.RedisUnreachable { logger.Debug(fmt.Sprintf("ServeHTTP:getDecision ip:%s isBanned:false %s", remoteIP, errString))
if errString != cache.CacheMiss {
logger.Error(fmt.Sprintf("ServeHTTP:getDecision ip:%s %s", remoteIP, errString))
rw.WriteHeader(http.StatusForbidden) rw.WriteHeader(http.StatusForbidden)
return return
} }
@@ -216,7 +226,7 @@ func (bouncer *Bouncer) ServeHTTP(rw http.ResponseWriter, req *http.Request) {
if isCrowdsecStreamHealthy { if isCrowdsecStreamHealthy {
bouncer.next.ServeHTTP(rw, req) bouncer.next.ServeHTTP(rw, req)
} else { } 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) rw.WriteHeader(http.StatusForbidden)
} }
} else { } else {
@@ -346,7 +356,7 @@ func getToken(bouncer *Bouncer) error {
return fmt.Errorf("getToken statusCode:%d", login.Code) return fmt.Errorf("getToken statusCode:%d", login.Code)
} }
func handleStreamCache(bouncer *Bouncer) { func handleStreamCache(bouncer *Bouncer) error {
// TODO clean properly on exit. // TODO clean properly on exit.
// Instead of blocking the goroutine interval for all the secondary node, // Instead of blocking the goroutine interval for all the secondary node,
// if the master service is shut down, other goroutine can take the lead // if the master service is shut down, other goroutine can take the lead
@@ -354,27 +364,26 @@ func handleStreamCache(bouncer *Bouncer) {
_, err := bouncer.cacheClient.GetDecision(cacheTimeoutKey) _, err := bouncer.cacheClient.GetDecision(cacheTimeoutKey)
if err == nil { if err == nil {
logger.Debug("handleStreamCache:alreadyUpdated") logger.Debug("handleStreamCache:alreadyUpdated")
return return nil
}
if err.Error() != cache.CacheMiss {
return err
} }
bouncer.cacheClient.SetDecision(cacheTimeoutKey, false, bouncer.updateInterval-1) bouncer.cacheClient.SetDecision(cacheTimeoutKey, false, bouncer.updateInterval-1)
streamRouteURL := url.URL{ streamRouteURL := url.URL{
Scheme: bouncer.crowdsecScheme, Scheme: bouncer.crowdsecScheme,
Host: bouncer.crowdsecHost, Host: bouncer.crowdsecHost,
Path: bouncer.crowdsecStreamRoute, Path: bouncer.crowdsecStreamRoute,
RawQuery: fmt.Sprintf("startup=%t", !isCrowdsecStreamHealthy), RawQuery: fmt.Sprintf("startup=%t", !isCrowdsecStreamHealthy || isStartup),
} }
body, err := crowdsecQuery(bouncer, streamRouteURL.String(), false) body, err := crowdsecQuery(bouncer, streamRouteURL.String(), false)
if err != nil { if err != nil {
logger.Error(err.Error()) return err
isCrowdsecStreamHealthy = false
return
} }
var stream Stream var stream Stream
err = json.Unmarshal(body, &stream) err = json.Unmarshal(body, &stream)
if err != nil { if err != nil {
logger.Error(fmt.Sprintf("handleStreamCache:parsingBody %s", err.Error())) return fmt.Errorf("handleStreamCache:parsingBody %w", err)
isCrowdsecStreamHealthy = false
return
} }
for _, decision := range stream.New { for _, decision := range stream.New {
duration, err := time.ParseDuration(decision.Duration) duration, err := time.ParseDuration(decision.Duration)
@@ -387,6 +396,7 @@ func handleStreamCache(bouncer *Bouncer) {
} }
logger.Debug("handleStreamCache:updated") logger.Debug("handleStreamCache:updated")
isCrowdsecStreamHealthy = true isCrowdsecStreamHealthy = true
return nil
} }
func crowdsecQuery(bouncer *Bouncer, stringURL string, isPost bool) ([]byte, error) { 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 bouncer *Bouncer
} }
tests := []struct { tests := []struct {
name string name string
args args args args
wantErr bool
}{ }{
// TODO: Add test cases. // TODO: Add test cases.
} }
for _, tt := range tests { for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) { 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
}
}) })
} }
} }
+3 -2
View File
@@ -11,10 +11,11 @@ These CAPI credentials must be set in your docker-compose.yml or in your config
... ...
whoami: whoami:
labels: labels:
- "traefik.http.middlewares.crowdsec.plugin.bouncer.enabled=true"
- "traefik.http.middlewares.crowdsec.plugin.bouncer.crowdsecMode=alone"
- "traefik.http.middlewares.crowdsec.plugin.bouncer.crowdsecCapiMachineId=LOGIN" - "traefik.http.middlewares.crowdsec.plugin.bouncer.crowdsecCapiMachineId=LOGIN"
- "traefik.http.middlewares.crowdsec.plugin.bouncer.crowdsecCapiPassword=PASSWORD" - "traefik.http.middlewares.crowdsec.plugin.bouncer.crowdsecCapiPassword=PASSWORD"
- "traefik.http.middlewares.crowdsec.plugin.bouncer.crowdseccapiscenarios=crowdsecurity/http-generic-bf,crowdsecurity/http-xss-probing,..." - "traefik.http.middlewares.crowdsec.plugin.bouncer.crowdsecCapiScenarios=crowdsecurity/http-generic-bf,crowdsecurity/http-xss-probing,..."
- "traefik.http.middlewares.crowdsec.plugin.bouncer.enabled=true"
``` ```
You can then run all the containers: You can then run all the containers:
@@ -0,0 +1,45 @@
version: "3.8"
services:
traefik:
image: "traefik:v2.9.6"
container_name: "traefik"
restart: unless-stopped
command:
# - "--log.level=DEBUG"
- "--accesslog"
- "--accesslog.filepath=/var/log/traefik/access.log"
- "--api.insecure=true"
- "--providers.docker=true"
- "--providers.docker.exposedbydefault=false"
- "--entrypoints.web.address=:80"
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
- "--experimental.plugins.bouncer.version=v1.1.7"
volumes:
- /var/run/docker.sock:/var/run/docker.sock:ro
ports:
- 80:80
- 8080:8080
whoami-foo:
image: traefik/whoami
container_name: "simple-service-foo"
restart: unless-stopped
labels:
- "traefik.enable=true"
- "traefik.http.routers.router-foo.rule=Path(`/foo`)"
- "traefik.http.routers.router-foo.entrypoints=web"
- "traefik.http.routers.router-foo.middlewares=crowdsec-foo@docker"
- "traefik.http.services.service-foo.loadbalancer.server.port=80"
- "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.enabled=true"
# - "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.loglevel=DEBUG"
- "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.crowdsecmode=alone"
- "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.crowdseclapikey=40796d93c2958f9e58345514e67740e5"
- "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.CrowdsecCapiMachineId=logincacacalfkrjebfreifgzfblezgyfoerxsqxsqxsqxsr"
- "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.CrowdsecCapiPassword=Password2"
- "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.crowdseccapiscenarios=crowdsecurity/sshd,crowdsecurity/asterisk_bf,crowdsecurity/asterisk_user_enum,crowdsecurity/base-http-scenarios"
volumes:
logs-local:
+1 -1
View File
@@ -4,5 +4,5 @@ go 1.19
require ( require (
github.com/leprosus/golang-ttl-map v1.1.7 github.com/leprosus/golang-ttl-map v1.1.7
github.com/maxlerebourg/simpleredis v1.0.3 github.com/maxlerebourg/simpleredis v1.0.5
) )
+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 h1:cF4AAFDDnJTFSV+/42sKLhmMluvLdRlCGS2UaifH6UM=
github.com/leprosus/golang-ttl-map v1.1.7/go.mod h1:4QWHJPeVBbrkhOhXdhCv9IEiyj/YzkO04/iexy4vSe0= 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.5 h1:1ubyIpTgIb+dadpILivAxuhZTHOOapEHRhcMakveuYY=
github.com/maxlerebourg/simpleredis v1.0.3/go.mod h1:/DH8zOK6kDskSqoX/m5CJJdNGfkIQZd/ERBJgytDDSk= github.com/maxlerebourg/simpleredis v1.0.5/go.mod h1:/DH8zOK6kDskSqoX/m5CJJdNGfkIQZd/ERBJgytDDSk=
+9 -3
View File
@@ -16,6 +16,9 @@ const (
cacheNoBannedValue = "f" cacheNoBannedValue = "f"
) )
// CacheMiss error string when cache is miss.
const CacheMiss = "cache:miss"
//nolint:gochecknoglobals //nolint:gochecknoglobals
var ( var (
redis simpleredis.SimpleRedis redis simpleredis.SimpleRedis
@@ -30,7 +33,7 @@ func (localCache) getDecision(clientIP string) (bool, error) {
if isCached && isValid && len(bannedString) > 0 { if isCached && isValid && len(bannedString) > 0 {
return bannedString == cacheBannedValue, nil return bannedString == cacheBannedValue, nil
} }
return false, fmt.Errorf("cache:miss") return false, fmt.Errorf(CacheMiss)
} }
func (localCache) setDecision(clientIP string, value string, duration int64) { func (localCache) setDecision(clientIP string, value string, duration int64) {
@@ -49,6 +52,9 @@ func (redisCache) getDecision(clientIP string) (bool, error) {
if err == nil && len(bannedString) > 0 { if err == nil && len(bannedString) > 0 {
return bannedString == cacheBannedValue, nil return bannedString == cacheBannedValue, nil
} }
if err.Error() == simpleredis.RedisMiss {
return false, fmt.Errorf(CacheMiss)
}
return false, err return false, err
} }
@@ -76,9 +82,9 @@ type Client struct {
} }
// New Initialize cache client. // New Initialize cache client.
func (client *Client) New(isRedis bool, host string) { func (client *Client) New(isRedis bool, host string, pass string) {
if isRedis { if isRedis {
redis.Init(host) redis.Init(host, pass)
client.cache = &redisCache{} client.cache = &redisCache{}
} else { } else {
client.cache = &localCache{} client.cache = &localCache{}
+9 -2
View File
@@ -55,6 +55,8 @@ type Config struct {
ClientTrustedIPs []string `json:"clientTrustedIps,omitempty"` ClientTrustedIPs []string `json:"clientTrustedIps,omitempty"`
RedisCacheEnabled bool `json:"redisCacheEnabled,omitempty"` RedisCacheEnabled bool `json:"redisCacheEnabled,omitempty"`
RedisCacheHost string `json:"redisCacheHost,omitempty"` RedisCacheHost string `json:"redisCacheHost,omitempty"`
RedisCachePassword string `json:"redisCachePassword,omitempty"`
RedisCachePasswordFile string `json:"redisCachePasswordFile,omitempty"`
} }
func contains(source []string, target string) bool { func contains(source []string, target string) bool {
@@ -83,6 +85,7 @@ func New() *Config {
ClientTrustedIPs: []string{}, ClientTrustedIPs: []string{},
RedisCacheEnabled: false, RedisCacheEnabled: false,
RedisCacheHost: "redis:6379", RedisCacheHost: "redis:6379",
RedisCachePassword: "",
} }
} }
@@ -115,7 +118,7 @@ func GetVariable(config *Config, key string) (string, error) {
// ValidateParams validate all the param gave by user. // ValidateParams validate all the param gave by user.
// //
//nolint:gocyclo //nolint:gocyclo,gocognit
func ValidateParams(config *Config) error { func ValidateParams(config *Config) error {
if err := validateParamsRequired(config); err != nil { if err := validateParamsRequired(config); err != nil {
return err return err
@@ -128,8 +131,12 @@ func ValidateParams(config *Config) error {
return err return err
} }
if _, err := GetVariable(config, "RedisCachePassword"); err != nil {
return err
}
if config.CrowdsecMode == AloneMode { if config.CrowdsecMode == AloneMode {
if _, err := GetVariable(config, "CrowdsecCapiMachineId"); err != nil { if _, err := GetVariable(config, "CrowdsecCapiMachineID"); err != nil {
return err return err
} }
if _, err := GetVariable(config, "CrowdsecCapiPassword"); err != nil { if _, err := GetVariable(config, "CrowdsecCapiPassword"); err != nil {
+28 -1
View File
@@ -1,2 +1,29 @@
# simpleredis # simpleredis
Minimal go redis with only get, set and delete operation Minimal go redis with only `get`, `set` and `delete` operation.
With **NO** extern dependencies.
## Example
```go
import simpleredis "github.com/maxlerebourg/simpleredis"
var redis simpleredis.SimpleRedis
redis.Init("redis:6379", "") // redisHost, redisPass
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)
+45 -19
View File
@@ -17,10 +17,11 @@ const (
RedisUnreachable = "redis:unreachable" RedisUnreachable = "redis:unreachable"
RedisMiss = "redis:miss" RedisMiss = "redis:miss"
RedisTimeout = "redis:timeout" RedisTimeout = "redis:timeout"
RedisNoAuth = "redis:noauth"
) )
// A RedisCmd is used to communicate with redis at low level using commands. // A redisCmd is used to communicate with redis at low level using commands.
type RedisCmd struct { type redisCmd struct {
Command string Command string
Name string Name string
Data []byte Data []byte
@@ -30,7 +31,8 @@ type RedisCmd struct {
// A SimpleRedis is used to communicate with redis. // A SimpleRedis is used to communicate with redis.
type SimpleRedis struct { type SimpleRedis struct {
redisHost string host string
pass string
} }
func genRedisArray(params ...[]byte) []byte { func genRedisArray(params ...[]byte) []byte {
@@ -49,11 +51,11 @@ func send(wr *textproto.Writer, method string, data []byte) {
} }
} }
func askRedis(hostnamePort string, cmd RedisCmd, channel chan RedisCmd) { func askRedis(sr *SimpleRedis, cmd redisCmd, channel chan redisCmd) {
dialer := net.Dialer{Timeout: 2 * time.Second} dialer := net.Dialer{Timeout: 2 * time.Second}
conn, err := dialer.Dial("tcp", hostnamePort) conn, err := dialer.Dial("tcp", sr.host)
if err != nil { if err != nil {
channel <- RedisCmd{Error: fmt.Errorf(RedisUnreachable)} channel <- redisCmd{Error: fmt.Errorf(RedisUnreachable)}
return return
} }
defer func() { defer func() {
@@ -65,6 +67,25 @@ func askRedis(hostnamePort string, cmd RedisCmd, channel chan RedisCmd) {
writer := textproto.NewWriter(bufio.NewWriter(conn)) writer := textproto.NewWriter(bufio.NewWriter(conn))
reader := textproto.NewReader(bufio.NewReader(conn)) reader := textproto.NewReader(bufio.NewReader(conn))
if sr.pass != "" {
data := genRedisArray([]byte("AUTH"), []byte(sr.pass))
send(writer, "auth", data)
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
}
break
}
}
}
switch cmd.Command { switch cmd.Command {
case "SET": case "SET":
data := genRedisArray([]byte("SET"), []byte(cmd.Name), cmd.Data, []byte("EX"), []byte(fmt.Sprintf("%d", cmd.Duration))) data := genRedisArray([]byte("SET"), []byte(cmd.Name), cmd.Data, []byte("EX"), []byte(fmt.Sprintf("%d", cmd.Duration)))
@@ -78,16 +99,20 @@ func askRedis(hostnamePort string, cmd RedisCmd, channel chan RedisCmd) {
for { for {
select { select {
case <-time.After(time.Second * 1): case <-time.After(time.Second * 1):
channel <- RedisCmd{Error: fmt.Errorf(RedisTimeout)} channel <- redisCmd{Error: fmt.Errorf(RedisTimeout)}
return return
default: default:
read, _ := reader.ReadLineBytes() read, _ := reader.ReadLineBytes()
if string(read) != "$1" { str := string(read)
channel <- RedisCmd{Error: fmt.Errorf(RedisMiss)} if strings.Contains(str, "-NOAUTH") {
channel <- redisCmd{Error: fmt.Errorf(RedisNoAuth)}
return
} else if str != "$1" {
channel <- redisCmd{Error: fmt.Errorf(RedisMiss)}
return return
} }
read, _ = reader.ReadLineBytes() read, _ = reader.ReadLineBytes()
channel <- RedisCmd{Data: read} channel <- redisCmd{Data: read}
return return
} }
} }
@@ -95,18 +120,19 @@ func askRedis(hostnamePort string, cmd RedisCmd, channel chan RedisCmd) {
} }
// Init sets the redisHost used to connect to redis. // Init sets the redisHost used to connect to redis.
func (sr *SimpleRedis) Init(redisHost string) { func (sr *SimpleRedis) Init(host string, pass string) {
sr.redisHost = redisHost sr.host = host
sr.pass = pass
} }
// Get fetches the value for key name in redis. // Get fetches the value for key name in redis.
func (sr *SimpleRedis) Get(name string) ([]byte, error) { func (sr *SimpleRedis) Get(name string) ([]byte, error) {
redisCmd := RedisCmd{ cmd := redisCmd{
Command: "GET", Command: "GET",
Name: name, Name: name,
} }
channel := make(chan RedisCmd) channel := make(chan redisCmd)
go askRedis(sr.redisHost, redisCmd, channel) go askRedis(sr, cmd, channel)
resp := <-channel resp := <-channel
if resp.Error != nil { if resp.Error != nil {
return nil, resp.Error return nil, resp.Error
@@ -116,22 +142,22 @@ func (sr *SimpleRedis) Get(name string) ([]byte, error) {
// Set updates the value for key name in redis with value data for duration. // 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 { func (sr *SimpleRedis) Set(name string, data []byte, duration int64) error {
redisCmd := RedisCmd{ cmd := redisCmd{
Command: "SET", Command: "SET",
Name: name, Name: name,
Data: data, Data: data,
Duration: duration, Duration: duration,
} }
go askRedis(sr.redisHost, redisCmd, nil) go askRedis(sr, cmd, nil)
return nil return nil
} }
// Del removes the key name in redis. // Del removes the key name in redis.
func (sr *SimpleRedis) Del(name string) error { func (sr *SimpleRedis) Del(name string) error {
redisCmd := RedisCmd{ cmd := redisCmd{
Command: "DEL", Command: "DEL",
Name: name, Name: name,
} }
go askRedis(sr.redisHost, redisCmd, nil) go askRedis(sr, cmd, nil)
return nil return nil
} }
+1 -1
View File
@@ -1,6 +1,6 @@
# github.com/leprosus/golang-ttl-map v1.1.7 # github.com/leprosus/golang-ttl-map v1.1.7
## explicit; go 1.15 ## explicit; go 1.15
github.com/leprosus/golang-ttl-map github.com/leprosus/golang-ttl-map
# github.com/maxlerebourg/simpleredis v1.0.3 # github.com/maxlerebourg/simpleredis v1.0.5
## explicit; go 1.19 ## explicit; go 1.19
github.com/maxlerebourg/simpleredis github.com/maxlerebourg/simpleredis