From ffcf4356fc0bed0c69f43e869bf804f0e5464aa4 Mon Sep 17 00:00:00 2001 From: Max Lerebourg Date: Thu, 20 Oct 2022 02:42:24 +0200 Subject: [PATCH] :sparkles: redis included --- bouncer.go | 14 +- docker-compose.local.yml | 2 - exemples/redis-cache/docker-compose.redis.yml | 30 +- go.mod | 2 + go.sum | 162 +----- pkg/cache/cache.go | 38 +- pkg/redis/redis.go | 507 ++++++----------- .../tehnerd/goUtils/netutils/netutils.go | 525 ++++++++++++++++++ .../tehnerd/goUtils/netutils/tcp_listen.go | 42 ++ vendor/modules.txt | 3 + 10 files changed, 767 insertions(+), 558 deletions(-) create mode 100644 vendor/github.com/tehnerd/goUtils/netutils/netutils.go create mode 100644 vendor/github.com/tehnerd/goUtils/netutils/tcp_listen.go diff --git a/bouncer.go b/bouncer.go index 8a9a55d..e35a50f 100644 --- a/bouncer.go +++ b/bouncer.go @@ -124,15 +124,11 @@ func New(ctx context.Context, next http.Handler, config *Config, name string) (h if config.RedisCacheEnabled { cache.InitRedisClient(config.RedisCacheHost, config.RedisCachePassword) } - if config.CrowdsecMode == streamMode { - go func() { - if ticker == nil { - go handleStreamCache(bouncer) - ticker = startTicker(config, func() { - handleStreamCache(bouncer) - }) - } - }() + if config.CrowdsecMode == streamMode && ticker == nil { + ticker = startTicker(config, func() { + handleStreamCache(bouncer) + }) + go handleStreamCache(bouncer) } return bouncer, nil diff --git a/docker-compose.local.yml b/docker-compose.local.yml index fc8f834..7849d2e 100644 --- a/docker-compose.local.yml +++ b/docker-compose.local.yml @@ -35,7 +35,6 @@ services: - "traefik.http.services.service1.loadbalancer.server.port=80" - "traefik.http.middlewares.crowdsec1.plugin.bouncer.enabled=true" - "traefik.http.middlewares.crowdsec1.plugin.bouncer.crowdseclapikey=40796d93c2958f9e58345514e67740e5" - - "traefik.http.middlewares.crowdsec1.plugin.bouncer.forwardedheaderstrustedips=172.26.0.5" whoami2: image: traefik/whoami @@ -48,7 +47,6 @@ services: - "traefik.http.services.service2.loadbalancer.server.port=80" - "traefik.http.middlewares.crowdsec2.plugin.bouncer.enabled=true" - "traefik.http.middlewares.crowdsec2.plugin.bouncer.crowdseclapikey=44c36dac5c4140af9f06f397508e82c7" - - "traefik.http.middlewares.crowdsec2.plugin.bouncer.forwardedheaderstrustedips=172.26.0.5" crowdsec: image: crowdsecurity/crowdsec:v1.4.1 diff --git a/exemples/redis-cache/docker-compose.redis.yml b/exemples/redis-cache/docker-compose.redis.yml index d8ceea3..a0d70ec 100644 --- a/exemples/redis-cache/docker-compose.redis.yml +++ b/exemples/redis-cache/docker-compose.redis.yml @@ -1,4 +1,5 @@ version: "3.8" + services: traefik: image: "traefik:v2.8.8" @@ -22,9 +23,6 @@ services: ports: - 80:80 - 8080:8080 - networks: - - redis - - services depends_on: - crowdsec - redis @@ -35,19 +33,16 @@ services: labels: - "traefik.enable=true" # Definition of the router - - "traefik.http.routers.router1.rule=Path(`/foo`)" + - "traefik.http.routers.router1.rule=Host(`localhost`) && Path(`/foo`)" - "traefik.http.routers.router1.entrypoints=web" - "traefik.http.routers.router1.middlewares=crowdsec1@docker" # Definition of the service - "traefik.http.services.service1.loadbalancer.server.port=80" - # Definitin of the middleware + # Definition of the middleware - "traefik.http.middlewares.crowdsec1.plugin.bouncer.enabled=true" # crowdseclapikey must be uniq to the middleware attached to the service - "traefik.http.middlewares.crowdsec1.plugin.bouncer.crowdseclapikey=40796d93c2958f9e58345514e67740e5" - - "traefik.http.middlewares.crowdsec1.plugin.bouncer.crowdsecmode=live" - "traefik.http.middlewares.crowdsec1.plugin.bouncer.rediscacheenabled=true" - networks: - - services whoami2: image: traefik/whoami @@ -55,19 +50,16 @@ services: labels: - "traefik.enable=true" # Definition of the router - - "traefik.http.routers.router2.rule=Path(`/bar`)" + - "traefik.http.routers.router2.rule=Host(`localhost`) && Path(`/bar`)" - "traefik.http.routers.router2.entrypoints=web" - "traefik.http.routers.router2.middlewares=crowdsec1@docker" # Definition of the service - "traefik.http.services.service2.loadbalancer.server.port=80" - # Definitin of the middleware + # Definition of the middleware - "traefik.http.middlewares.crowdsec2.plugin.bouncer.enabled=true" # crowdseclapikey must be uniq to the middleware attached to the service - "traefik.http.middlewares.crowdsec2.plugin.bouncer.crowdseclapikey=44c36dac5c4140af9f06f397508e82c7" - - "traefik.http.middlewares.crowdsec2.plugin.bouncer.crowdsecmode=live" - - "traefik.http.middlewares.crowdsec2.plugin.bouncer.rediscacheenabled=true" - networks: - - services + - "traefik.http.middlewares.crowdsec2.plugin.bouncer.rediscacheenabled=false" crowdsec: @@ -87,20 +79,16 @@ services: - "traefik.enable=false" redis: - image: "redis:7.0.5" + image: "redis:7.0.5-alpine" container_name: "redis" command: "redis-server --save 60 1" volumes: - redis-data:/data - networks: - - redis + ports: + - 6379:6379 volumes: logs-redis: crowdsec-db-redis: crowdsec-config-redis: redis-data: - -networks: - redis: - services: diff --git a/go.mod b/go.mod index f5d4f6d..c3ff139 100644 --- a/go.mod +++ b/go.mod @@ -3,3 +3,5 @@ module github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin go 1.17 require github.com/leprosus/golang-ttl-map v1.1.7 + +require github.com/tehnerd/goUtils v0.0.0-20150515130609-5a2d8fb2ded8 // indirect diff --git a/go.sum b/go.sum index fe1b8fe..9a2f881 100644 --- a/go.sum +++ b/go.sum @@ -1,162 +1,4 @@ -cloud.google.com/go v0.34.0/go.mod h1:aQUYkXzVsufM+DwF1aE+0xfcU+56JwCaLick0ClmMTw= -github.com/alecthomas/template v0.0.0-20160405071501-a0175ee3bccc/go.mod h1:LOuyumcjzFXgccqObfd/Ljyb9UuFJ6TxHnclSeseNhc= -github.com/alecthomas/template v0.0.0-20190718012654-fb15b899a751/go.mod h1:LOuyumcjzFXgccqObfd/Ljyb9UuFJ6TxHnclSeseNhc= -github.com/alecthomas/units v0.0.0-20151022065526-2efee857e7cf/go.mod h1:ybxpYRFXyAe+OPACYpWeL0wqObRcbAqCMya13uyzqw0= -github.com/alecthomas/units v0.0.0-20190717042225-c3de453c63f4/go.mod h1:ybxpYRFXyAe+OPACYpWeL0wqObRcbAqCMya13uyzqw0= -github.com/alecthomas/units v0.0.0-20190924025748-f65c72e2690d/go.mod h1:rBZYJk541a8SKzHPHnH3zbiI+7dagKZ0cgpgrD7Fyho= -github.com/beorn7/perks v0.0.0-20180321164747-3a771d992973/go.mod h1:Dwedo/Wpr24TaqPxmxbtue+5NUziq4I4S80YR8gNf3Q= -github.com/beorn7/perks v1.0.0/go.mod h1:KWe93zE9D1o94FZ5RNwFwVgaQK1VOXiVxmqh+CedLV8= -github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= -github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= -github.com/cespare/xxhash/v2 v2.1.1/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= -github.com/cespare/xxhash/v2 v2.1.2 h1:YRXhKfTDauu4ajMg1TPgFO5jnlC2HCbmLXMcTG5cbYE= -github.com/cespare/xxhash/v2 v2.1.2/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= -github.com/coreos/go-semver v0.3.0 h1:wkHLiw0WNATZnSG7epLsujiMCgPAc9xhjJ4tgnAxmfM= -github.com/coreos/go-semver v0.3.0/go.mod h1:nnelYz7RCh+5ahJtPPxZlU+153eP4D4r3EedlOD2RNk= -github.com/davecgh/go-spew v1.1.0 h1:ZDRjVQ15GmhC3fiQ8ni8+OwkZQO4DARzQgrnXU1Liz8= -github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= -github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= -github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f h1:lO4WD4F/rVNCu3HqELle0jiPLLBs70cWOduZpkS1E78= -github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f/go.mod h1:cuUVRXasLTGF7a8hSLbxyZXjz+1KgoB3wDUb6vlszIc= -github.com/go-kit/kit v0.8.0/go.mod h1:xBxKIO96dXMWWy0MnWVtmwkA9/13aqxPnvrjFYMA2as= -github.com/go-kit/kit v0.9.0/go.mod h1:xBxKIO96dXMWWy0MnWVtmwkA9/13aqxPnvrjFYMA2as= -github.com/go-kit/log v0.1.0/go.mod h1:zbhenjAZHb184qTLMA9ZjW7ThYL0H2mk7Q6pNt4vbaY= -github.com/go-logfmt/logfmt v0.3.0/go.mod h1:Qt1PoO58o5twSAckw1HlFXLmHsOX5/0LbT9GBnD5lWE= -github.com/go-logfmt/logfmt v0.4.0/go.mod h1:3RMwSq7FuexP4Kalkev3ejPJsZTpXXBr9+V4qmtdjCk= -github.com/go-logfmt/logfmt v0.5.0/go.mod h1:wCYkCAKZfumFQihp8CzCvQ3paCTfi41vtzG1KdI/P7A= -github.com/go-redis/redis/v8 v8.11.5 h1:AcZZR7igkdvfVmQTPnu9WE37LRrO/YrBH5zWyjDC0oI= -github.com/go-redis/redis/v8 v8.11.5/go.mod h1:gREzHqY1hg6oD9ngVRbLStwAWKhA0FEgq8Jd4h5lpwo= -github.com/go-stack/stack v1.8.0/go.mod h1:v0f6uXyyMGvRgIKkXu+yp6POWl0qKG85gN/melR3HDY= -github.com/gogo/protobuf v1.1.1/go.mod h1:r8qH/GZQm5c6nD/R0oafs1akxWv10x8SbQlK7atdtwQ= -github.com/golang/protobuf v1.2.0/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U= -github.com/golang/protobuf v1.3.1/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U= -github.com/golang/protobuf v1.3.2/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U= -github.com/golang/protobuf v1.4.0-rc.1/go.mod h1:ceaxUfeHdC40wWswd/P6IGgMaK3YpKi5j83Wpe3EHw8= -github.com/golang/protobuf v1.4.0-rc.1.0.20200221234624-67d41d38c208/go.mod h1:xKAWHe0F5eneWXFV3EuXVDTCmh+JuBKY0li0aMyXATA= -github.com/golang/protobuf v1.4.0-rc.2/go.mod h1:LlEzMj4AhA7rCAGe4KMBDvJI+AwstrUpVNzEA03Pprs= -github.com/golang/protobuf v1.4.0-rc.4.0.20200313231945-b860323f09d0/go.mod h1:WU3c8KckQ9AFe+yFwt9sWVRKCVIyN9cPHBJSNnbL67w= -github.com/golang/protobuf v1.4.0/go.mod h1:jodUvKwWbYaEsadDk5Fwe5c77LiNKVO9IDvqG2KuDX0= -github.com/golang/protobuf v1.4.2/go.mod h1:oDoupMAO8OvCJWAcko0GGGIgR6R6ocIYbsSw735rRwI= -github.com/golang/protobuf v1.4.3/go.mod h1:oDoupMAO8OvCJWAcko0GGGIgR6R6ocIYbsSw735rRwI= -github.com/golang/protobuf v1.5.0 h1:LUVKkCeviFUMKqHa4tXIIij/lbhnMbP7Fn5wKdKkRh4= -github.com/golang/protobuf v1.5.0/go.mod h1:FsONVRAS9T7sI+LIUmWTfcYkHO4aIWwzhcaSAoJOfIk= -github.com/gomodule/redigo v1.8.9 h1:Sl3u+2BI/kk+VEatbj0scLdrFhjPmbxOc1myhDP41ws= -github.com/gomodule/redigo v1.8.9/go.mod h1:7ArFNvsTjH8GMMzB4uy1snslv2BwmginuMs06a1uzZE= -github.com/google/go-cmp v0.3.0/go.mod h1:8QqcDgzrUqlUb/G2PQTWiueGozuR1884gddMywk6iLU= -github.com/google/go-cmp v0.3.1/go.mod h1:8QqcDgzrUqlUb/G2PQTWiueGozuR1884gddMywk6iLU= -github.com/google/go-cmp v0.4.0/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= -github.com/google/go-cmp v0.5.4/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= -github.com/google/go-cmp v0.5.5/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= -github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg= -github.com/jpillora/backoff v1.0.0/go.mod h1:J/6gKK9jxlEcS3zixgDgUAsiuZ7yrSoa/FX5e0EB2j4= -github.com/json-iterator/go v1.1.6/go.mod h1:+SdeFBvtyEkXs7REEP0seUULqWtbJapLOCVDaaPEHmU= -github.com/json-iterator/go v1.1.10/go.mod h1:KdQUCv79m/52Kvf8AW2vK1V8akMuk1QjK/uOdHXbAo4= -github.com/json-iterator/go v1.1.11/go.mod h1:KdQUCv79m/52Kvf8AW2vK1V8akMuk1QjK/uOdHXbAo4= -github.com/julienschmidt/httprouter v1.2.0/go.mod h1:SYymIcj16QtmaHHD7aYtjjsJG7VTCxuUUipMqKk8s4w= -github.com/julienschmidt/httprouter v1.3.0/go.mod h1:JR6WtHb+2LUe8TCKY3cZOxFyyO8IZAc4RVcycCCAKdM= -github.com/konsorten/go-windows-terminal-sequences v1.0.1/go.mod h1:T0+1ngSBFLxvqU3pZ+m/2kptfBszLMUkC4ZK/EgS/cQ= -github.com/konsorten/go-windows-terminal-sequences v1.0.3/go.mod h1:T0+1ngSBFLxvqU3pZ+m/2kptfBszLMUkC4ZK/EgS/cQ= -github.com/kr/logfmt v0.0.0-20140226030751-b84e30acd515/go.mod h1:+0opPa2QZZtGFBFZlji/RkVcI2GknAs/DXo4wKdlNEc= -github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo= -github.com/kr/pty v1.1.1/go.mod h1:pFQYn66WHrOpPYNljwOMqo10TkYh1fy3cYio2l3bCsQ= -github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI= 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/matttproud/golang_protobuf_extensions v1.0.1 h1:4hp9jkHxhMHkqkrB3Ix0jegS5sx/RkqARlsWZ6pIwiU= -github.com/matttproud/golang_protobuf_extensions v1.0.1/go.mod h1:D8He9yQNgCq6Z5Ld7szi9bcBfOoFv/3dc6xSMkL2PC0= -github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= -github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= -github.com/modern-go/reflect2 v0.0.0-20180701023420-4b7aa43c6742/go.mod h1:bx2lNnkwVCuqBIxFjflWJWanXIb3RllmbCylyMrvgv0= -github.com/modern-go/reflect2 v1.0.1/go.mod h1:bx2lNnkwVCuqBIxFjflWJWanXIb3RllmbCylyMrvgv0= -github.com/mwitkow/go-conntrack v0.0.0-20161129095857-cc309e4a2223/go.mod h1:qRWi+5nqEBWmkhHvq77mSJWrCKwh8bxhgT7d/eI7P4U= -github.com/mwitkow/go-conntrack v0.0.0-20190716064945-2f068394615f/go.mod h1:qRWi+5nqEBWmkhHvq77mSJWrCKwh8bxhgT7d/eI7P4U= -github.com/pkg/errors v0.8.0/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= -github.com/pkg/errors v0.8.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= -github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= -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/prometheus/client_golang v0.9.1/go.mod h1:7SWBe2y4D6OKWSNQJUaRYU/AaXPKyh/dDVn+NZz0KFw= -github.com/prometheus/client_golang v1.0.0/go.mod h1:db9x61etRT2tGnBNRi70OPL5FsnadC4Ky3P0J6CfImo= -github.com/prometheus/client_golang v1.7.1/go.mod h1:PY5Wy2awLA44sXw4AOSfFBetzPP4j5+D6mVACh+pe2M= -github.com/prometheus/client_golang v1.11.0 h1:HNkLOAEQMIDv/K+04rukrLx6ch7msSRwf3/SASFAGtQ= -github.com/prometheus/client_golang v1.11.0/go.mod h1:Z6t4BnS23TR94PD6BsDNk8yVqroYurpAkEiz0P2BEV0= -github.com/prometheus/client_model v0.0.0-20180712105110-5c3871d89910/go.mod h1:MbSGuTsp3dbXC40dX6PRTWyKYBIrTGTE9sqQNg2J8bo= -github.com/prometheus/client_model v0.0.0-20190129233127-fd36f4220a90/go.mod h1:xMI15A0UPsDsEKsMN9yxemIoYk6Tm2C1GtYGdfGttqA= -github.com/prometheus/client_model v0.2.0 h1:uq5h0d+GuxiXLJLNABMgp2qUWDPiLvgCzz2dUR+/W/M= -github.com/prometheus/client_model v0.2.0/go.mod h1:xMI15A0UPsDsEKsMN9yxemIoYk6Tm2C1GtYGdfGttqA= -github.com/prometheus/common v0.4.1/go.mod h1:TNfzLD0ON7rHzMJeJkieUDPYmFC7Snx/y86RQel1bk4= -github.com/prometheus/common v0.10.0/go.mod h1:Tlit/dnDKsSWFlCLTWaA1cyBgKHSMdTB80sz/V91rCo= -github.com/prometheus/common v0.26.0 h1:iMAkS2TDoNWnKM+Kopnx/8tnEStIfpYA0ur0xQzzhMQ= -github.com/prometheus/common v0.26.0/go.mod h1:M7rCNAaPfAosfx8veZJCuw84e35h3Cfd9VFqTh1DIvc= -github.com/prometheus/procfs v0.0.0-20181005140218-185b4288413d/go.mod h1:c3At6R/oaqEKCNdg8wHV1ftS6bRYblBhIjjI8uT2IGk= -github.com/prometheus/procfs v0.0.2/go.mod h1:TjEm7ze935MbeOT/UhFTIMYKhuLP4wbCsTZCD3I8kEA= -github.com/prometheus/procfs v0.1.3/go.mod h1:lV6e/gmhEcM9IjHGsFOCxxuZ+z1YqCvr4OA4YeYWdaU= -github.com/prometheus/procfs v0.6.0 h1:mxy4L2jP6qMonqmq+aTtOx1ifVWUgG/TAmntgbh3xv4= -github.com/prometheus/procfs v0.6.0/go.mod h1:cz+aTbrPOrUb4q7XlbU9ygM+/jj0fzG6c1xBZuNvfVA= -github.com/rueian/rueidis v0.0.80 h1:HfhTWc4gk7dMWd9vl7EVE4NkgMGjYFsyyhi6B6JV+b4= -github.com/rueian/rueidis v0.0.80/go.mod h1:LiKWMM/QnILwRfDZIhSIXi4vQqZ/UZy4+/aNkSCt8XA= -github.com/sandwich-go/redisson v1.1.16 h1:lAySrzPk8l6dermqEhsm+JUBoIcpARe09cgp/niMR1Y= -github.com/sandwich-go/redisson v1.1.16/go.mod h1:/gUQACQFRnzVXje9AyEicE8TKtrd9uJI74jMqAcOv3o= -github.com/sandwich-go/rueidis v0.1.11 h1:+lmTofyyzBb0dXLyuegE5NoAaXZxrVg+9QL5gEFnKYk= -github.com/sandwich-go/rueidis v0.1.11/go.mod h1:N894IGktaGEm5gMWxD7u6XSxkpSzxCStQEOKjqOJQ/k= -github.com/sirupsen/logrus v1.2.0/go.mod h1:LxeOpSwHxABJmUn/MG1IvRgCAasNZTLOkJPxbbu5VWo= -github.com/sirupsen/logrus v1.4.2/go.mod h1:tLMulIdttU9McNUspp0xgXVQah82FyeX6MwdIuYE2rE= -github.com/sirupsen/logrus v1.6.0/go.mod h1:7uNnSEd1DgxDLC74fIahvMZmmYsHGZGEOFrfsX/uA88= -github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= -github.com/stretchr/objx v0.1.1/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= -github.com/stretchr/testify v1.2.2/go.mod h1:a8OnRcib4nhh0OaRAV+Yts87kKdq0PP7pXfy6kDkUVs= -github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= -github.com/stretchr/testify v1.4.0/go.mod h1:j7eGeouHqKxXV5pUuKE4zz7dFj8WfuZ+81PSLYec5m4= -github.com/stretchr/testify v1.7.0 h1:nwc3DEeHmmLAfoZucVR881uASk0Mfjw8xYJ99tb5CcY= -github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= -golang.org/x/crypto v0.0.0-20180904163835-0709b304e793/go.mod h1:6SG95UA2DQfeDnfUPMdvaQW0Q7yPrPDi9nlGo2tz2b4= -golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= -golang.org/x/crypto v0.0.0-20200622213623-75b288015ac9/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= -golang.org/x/net v0.0.0-20180724234803-3673e40ba225/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= -golang.org/x/net v0.0.0-20181114220301-adae6a3d119a/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= -golang.org/x/net v0.0.0-20190108225652-1e06a53dbb7e/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= -golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= -golang.org/x/net v0.0.0-20190613194153-d28f0bde5980/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= -golang.org/x/net v0.0.0-20200625001655-4c5254603344/go.mod h1:/O7V0waA8r7cgGh81Ro3o1hOxt32SMVPicZroKQ2sZA= -golang.org/x/oauth2 v0.0.0-20190226205417-e64efc72b421/go.mod h1:gOpvHmFTYa4IltrdGE7lF6nIHvwfUNPOp7c8zoXwtLw= -golang.org/x/sync v0.0.0-20181108010431-42b317875d0f/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= -golang.org/x/sync v0.0.0-20181221193216-37e7f081c4d4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= -golang.org/x/sync v0.0.0-20190911185100-cd5d95a43a6e/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= -golang.org/x/sync v0.0.0-20201207232520-09787c993a3a/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= -golang.org/x/sys v0.0.0-20180905080454-ebe1bf3edb33/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= -golang.org/x/sys v0.0.0-20181116152217-5ac8a444bdc5/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= -golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= -golang.org/x/sys v0.0.0-20190412213103-97732733099d/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= -golang.org/x/sys v0.0.0-20190422165155-953cdadca894/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= -golang.org/x/sys v0.0.0-20200106162015-b016eb3dc98e/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= -golang.org/x/sys v0.0.0-20200323222414-85ca7c5b95cd/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= -golang.org/x/sys v0.0.0-20200615200032-f1bc736245b1/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= -golang.org/x/sys v0.0.0-20200625212154-ddb9806d33ae/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= -golang.org/x/sys v0.0.0-20210124154548-22da62e12c0c/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= -golang.org/x/sys v0.0.0-20210603081109-ebe580a85c40/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.0.0-20220731174439-a90be440212d h1:Sv5ogFZatcgIMMtBSTTAgMYsicp25MXBubjXNDKwm80= -golang.org/x/sys v0.0.0-20220731174439-a90be440212d/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= -golang.org/x/text v0.3.2/go.mod h1:bEr9sfX3Q8Zfm5fL9x+3itogRgK3+ptLWKqgva+5dAk= -golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= -golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= -google.golang.org/appengine v1.4.0/go.mod h1:xpcJRLb0r/rnEns0DIKYYv+WjYCduHsrkT7/EB5XEv4= -google.golang.org/protobuf v0.0.0-20200109180630-ec00e32a8dfd/go.mod h1:DFci5gLYBciE7Vtevhsrf46CRTquxDuWsQurQQe4oz8= -google.golang.org/protobuf v0.0.0-20200221191635-4d8936d0db64/go.mod h1:kwYJMbMJ01Woi6D6+Kah6886xMZcty6N08ah7+eCXa0= -google.golang.org/protobuf v0.0.0-20200228230310-ab0ca4ff8a60/go.mod h1:cfTl7dwQJ+fmap5saPgwCLgHXTUD7jkjRqWcaiX5VyM= -google.golang.org/protobuf v1.20.1-0.20200309200217-e05f789c0967/go.mod h1:A+miEFZTKqfCUM6K7xSMQL9OKL/b6hQv+e19PK+JZNE= -google.golang.org/protobuf v1.21.0/go.mod h1:47Nbq4nVaFHyn7ilMalzfO3qCViNmqZ2kzikPIcrTAo= -google.golang.org/protobuf v1.23.0/go.mod h1:EGpADcykh3NcUnDUJcl1+ZksZNG86OlYog2l/sGQquU= -google.golang.org/protobuf v1.26.0-rc.1/go.mod h1:jlhhOSvTdKEhbULTjvd4ARK9grFBp09yW+WbY/TyQbw= -google.golang.org/protobuf v1.27.1 h1:SnqbnDw1V7RiZcXPx5MEeqPv2s79L9i7BJUlG/+RurQ= -google.golang.org/protobuf v1.27.1/go.mod h1:9q0QmTI4eRPtz6boOQmLYwt+qCgq0jsYwAQnmE0givc= -gopkg.in/alecthomas/kingpin.v2 v2.2.6/go.mod h1:FMv+mEhP44yOT+4EoQTLFTRgOQ1FBLkstjWtayDeSgw= -gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= -gopkg.in/check.v1 v1.0.0-20190902080502-41f04d3bba15/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= -gopkg.in/yaml.v2 v2.2.1/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= -gopkg.in/yaml.v2 v2.2.2/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= -gopkg.in/yaml.v2 v2.2.4/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= -gopkg.in/yaml.v2 v2.2.5/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= -gopkg.in/yaml.v2 v2.3.0/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= -gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c h1:dUUwHk2QECo/6vqA44rthZ8ie2QXMNeKRTHCNY2nXvo= -gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +github.com/tehnerd/goUtils v0.0.0-20150515130609-5a2d8fb2ded8 h1:/b777evAfRRdUJHasZLgQ/w8D/s1HtbaeDXsCVsV0B0= +github.com/tehnerd/goUtils v0.0.0-20150515130609-5a2d8fb2ded8/go.mod h1:UuuqaOb+pZOxJZtjF1mBWTo8HYa7HQCbNkwUEaG9uU0= diff --git a/pkg/cache/cache.go b/pkg/cache/cache.go index 97fa40c..860ffd0 100644 --- a/pkg/cache/cache.go +++ b/pkg/cache/cache.go @@ -7,7 +7,7 @@ import ( ttl_map "github.com/leprosus/golang-ttl-map" logger "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/logger" - redis "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/redis" + simpleredis "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/redis" ) const ( @@ -15,9 +15,8 @@ const ( cacheNoBannedValue = "f" ) -var ctx = context.Background() var cache = ttl_map.New() - +var redis simpleredis.SimpleRedis var redisEnabled = false @@ -40,18 +39,17 @@ func DeleteDecisionLocalCache(clientIP string) { } func getDecisionRedisCache(clientIP string) (bool, error) { - // banned, err := redis.Do(ctx, redis.B().Get().Key("key").Build()).ToString() - // if err != nil { - // return false, fmt.Errorf("no cache data") - // } else { - // return banned == cacheBannedValue, nil - // } - return false, nil + banned, err := redis.Do("GET", clientIP, nil) + bannedString := string(banned) + logger.Info(fmt.Sprintf("%v banned", bannedString)) + if err == nil && len(bannedString) > 0 { + return bannedString == cacheBannedValue, nil + } + return false, fmt.Errorf("no cache data") } func SetDecisionRedisCache(clientIP string, value string, duration int64) { - // ctx := context.Background() - // redis.Do(ctx, redis.B().Set().Key("key").Value("val").Nx().Build()).Error() + redis.Do("SET", clientIP, []byte(value)) } func DeleteDecisionRedisCache(clientIP string) { @@ -70,7 +68,6 @@ func DeleteDecision(clientIP string) { } else { DeleteDecisionLocalCache(clientIP) } - } // Get Decision check in the cache if the IP has the banned / not banned value. // Otherwise return with an error to add the IP in cache if we are on. @@ -97,15 +94,8 @@ func SetDecision(clientIP string, isBanned bool, duration int64) { } } -func InitRedisClient(host string, password string) error { - // redis, err := rueidis.NewClient(rueidis.ClientOption{ - // InitAddress: []string{host}, - // }) - // if err != nil { - // return fmt.Errorf("error instanciate redis: %w", err) - // } - // redisEnabled = true - writter, reader := redis.Init(host, password) - // writter.PrintfLine("[caca, true, caca]") - return nil +func InitRedisClient(host string, password string) { + logger.Debug("connect to redis") + redisEnabled = true + redis.Init(host) } diff --git a/pkg/redis/redis.go b/pkg/redis/redis.go index ed0c89e..ea80b17 100644 --- a/pkg/redis/redis.go +++ b/pkg/redis/redis.go @@ -1,375 +1,198 @@ -package redis - -// My intepretation of the RESP protocol. -// Referenced from https://redis.io/docs/reference/protocol-spec/ -// 8-7-2022 +package simpleredis import ( - "errors" "fmt" - "io" "log" - "net/textproto" "net" - "bufio" + "strconv" "strings" - "sync" + "time" + + "github.com/tehnerd/goUtils/netutils" + + logger "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/logger" ) -type RespT byte - -const ( - SimpleString RespT = '+' - Error RespT = '-' - Integer RespT = ':' - BulkString RespT = '$' - Array RespT = '*' -) - -type Message interface { - String() string +type RedisCmd struct { + Command string + Name string + Data []byte + Error error } -type MsgSimpleStr string -type MsgError string -type MsgBulkStr string -type MsgInteger int64 -type MsgArray []Message - -func (m MsgInteger) String() string { - return fmt.Sprintf("%d", m) +type SimpleRedis struct { + redisChanRead chan RedisCmd + redisChanWrite chan RedisCmd + redisHost string + redisCmd RedisCmd } -func (m MsgBulkStr) String() string { - return string(m) -} - -func (m MsgError) String() string { - return string(m) -} - -func (m MsgSimpleStr) String() string { - return string(m) -} - -// func (m MsgArray) String() string { -// return fmt.Sprintf("[%s]", strings.Join(Map(m, Message.String), ",")) -// } - -func (m RedisMessage) String() string { - return m.Choice.String() -} - -type RedisMessage struct { - RedisType RespT - Raw string - Choice Message -} - -func NewParseError(got RespT, expected ...RespT) error { - return fmt.Errorf("got %v but expected %v", got, expected) -} - -func (msg RedisMessage) Symbol() rune { - return rune(msg.RedisType) -} - -func (msg RedisMessage) AsString() (string, error) { - if conv, ok := msg.Choice.(MsgSimpleStr); ok { - return string(conv), nil - } else if conv, ok := msg.Choice.(MsgBulkStr); ok { - return string(conv), nil - } else { - return "", NewParseError(msg.RedisType, SimpleString, BulkString) +func GenRedisArray(params ...[]byte) []byte { + CRLF := "\r\n" + MSG := "" + for cntr := 0; cntr < len(params); cntr++ { + MSG = strings.Join([]string{MSG, string(params[cntr])}, " ") } + MSG = strings.Trim(MSG, " ") + MSG = strings.Join([]string{MSG, CRLF}, "") + return []byte(MSG) } -func (msg RedisMessage) AsInteger() (int64, error) { - if conv, ok := msg.Choice.(MsgInteger); ok { - return int64(conv), nil - } else { - return 0, NewParseError(msg.RedisType, Integer) - } +func RedisSet(name string, data []byte) []byte { + return GenRedisArray([]byte("SET"), []byte(name), data) } -func (msg RedisMessage) AsError() (error, error) { - if conv, ok := msg.Choice.(MsgError); ok { - return errors.New(string(conv)), nil - } else { - return nil, NewParseError(msg.RedisType, Error) - } +func RedisGet(name string) []byte { + return GenRedisArray([]byte("GET"), []byte(name)) } -// func (msg RedisMessage) AsArray() ([]Message, error) { -// if conv, ok := (msg.Choice).(MsgArray); ok { -// return conv, nil -// } else { -// return nil, NewParseError(msg.RedisType, Array) -// } -// } - -type RedisReader struct { - Rd *textproto.Reader - Out chan RedisMessage -} - -// Reset the parsing state of the reader -func (rr *RedisReader) Reset() { - -} - -func NewRespReader(tp *textproto.Reader, out chan RedisMessage) RedisReader { - return RedisReader{ - Rd: tp, - Out: out, +func ParseRedisResponse(response []byte, dataBuf []byte, Len *int) ([]byte, []byte, error) { + dataBuf = append(dataBuf, response...) + lenCRLF := 2 + if *Len != 0 { + if len(dataBuf) < *Len { + return nil, dataBuf, nil + } else { + return dataBuf[:*Len], dataBuf[*Len:], nil + } } -} - -// Fetch grabs the next incoming RESP type from the IO input -func Fetch(rr *RedisReader) *RedisMessage { - byteCode, err := rr.TryReadByte() - if err != nil { - return nil - } - switch RespT(byteCode) { - case BulkString: - return RespBulkStr(rr) - case SimpleString: - return RespSimpleStr(rr) - case Integer: - return RespInt(rr) - case Error: - return RespError(rr) - // case Array: - // return RespArray(rr) - default: - panic(fmt.Sprint("Unknown bytecode: ", byteCode)) - } -} - -// Scan is the main API to send redis messages out to a channel. -func (rr *RedisReader) Scan() { - rr.Out <- *Fetch(rr) -} - -// RespSimpleStr attempts to parse a string from the resulting line. -// -// TODO: I don't think we need the high level facilities from net/textproto anymore -// since the parsing is quite trivial. -func RespSimpleStr(rr *RedisReader) *RedisMessage { - s, err := rr.Rd.ReadLine() - if err != nil { - log.Printf("Fatal RespSimpleStr: %+v", err) - return nil - } - - return &RedisMessage{ - RedisType: RespT(SimpleString), - Raw: s, - Choice: MsgSimpleStr(s), - } -} - -func RespError(rr *RedisReader) *RedisMessage { - s, err := rr.Rd.ReadLine() - if err != nil { - log.Printf("Fatal RespError: %+v", err) - return nil - } - - return &RedisMessage{ - RedisType: RespT(Error), - Raw: s, - Choice: MsgError(s), - } -} - -// RespSimpleStr attempts to parse a string from the resulting line. -// -// It first expects a length value parsed as an integer, -// then directly copies the subsequent payload into a string -// buffer -func RespBulkStr(rr *RedisReader) *RedisMessage { - len, err := loopReadInt(rr) - if err != nil { - return nil - } - - // Read the string data into buffer - bulkStr := make([]byte, len) - io.ReadFull(rr.Rd.R, bulkStr) - - // Discard two bytes since we don't need the '\r\n' - _, err = rr.Rd.R.Discard(2) - if err != nil { - log.Printf("Fatal: %+v", err) - } - - return &RedisMessage{ - RedisType: BulkString, - Raw: string(bulkStr), - Choice: MsgBulkStr(string(bulkStr)), - } -} - -// RespInt reads an integer from the IO stream -func RespInt(rr *RedisReader) *RedisMessage { - val, err := loopReadInt(rr) - if err != nil { - return nil - } - return &RedisMessage{ - RedisType: Integer, - Raw: fmt.Sprintf("%d", val), - Choice: MsgInteger(val), - } -} - -// RespArray reads an array of RESP objects from the IO stream -// func RespArray(rr *RedisReader) *RedisMessage { -// len, err := loopReadInt(rr) -// if err != nil { -// return nil -// } -// fetched := make([]RedisMessage, 0, len) -// for i := 0; i < len; i++ { -// fetched[i] = *Fetch(rr) -// } -// return &RedisMessage{ -// RedisType: Array, -// Raw: fmt.Sprint(Map(fetched, func(rm RedisMessage) string { return rm.Raw })), -// Choice: MsgArray(Map(fetched, func(rm RedisMessage) Message { return rm.Choice })), -// } -// } - -// loopReadInt reads from IO and returns an integer, discarding \r\n. -func loopReadInt(rr *RedisReader) (int, error) { - val := 0 - for b, _ := rr.TryReadByte(); b != '\r'; b, _ = rr.TryReadByte() { - val = (val * 10) + int(b-byte('0')) - } - _, err := rr.Rd.R.Discard(1) - if err != nil { - return 0, fmt.Errorf("Fatal loopReadInt: %+v", err) - } - return val, nil -} - -// TryReadByte attempts to read a single byte from IO and panics on any error -func (rr *RedisReader) TryReadByte() (byte, error) { - b, err := rr.Rd.R.ReadByte() - if err != nil { - return byte('0'), fmt.Errorf("Fatal TryReadByte: %+v", err) - } - return b, nil -} - -// Map. Your standard good old fashioned generic map function :) -// -// A la []T -> []K -func Map[T any, K any](slice []T, fn func(T) K) []K { - out := make([]K, 0, len(slice)) - for i := range out { - out[i] = fn(slice[i]) - } - return out -} - - - - - - - - -func loopScan(reader *textproto.Reader, msg_ch chan RedisMessage) { - rr := NewRespReader(reader, msg_ch) for { - rr.Scan() + switch string(dataBuf[0]) { + case "+", "-", ":": + //simple strings, error,int. usually ther are in format (+|-|:)DATA\r\n" + if len(dataBuf) < 3 { + return nil, dataBuf, nil + } + cntr := 1 + for ; cntr < len(dataBuf); cntr++ { + if dataBuf[cntr] == '\r' { + break + } + } + if cntr == len(dataBuf) { + return nil, dataBuf, nil + } + response = dataBuf[1:cntr] + return response, dataBuf[cntr+2:], nil + case "$": + //bulk string. format $LEN\r\nDATA\r\n. up to 512MB + cntr := 1 + for ; cntr < len(dataBuf); cntr++ { + if string(dataBuf[cntr]) == "\r" { + break + } + } + if cntr == len(dataBuf) || cntr+lenCRLF > len(dataBuf) { + return nil, dataBuf, nil + } + dataLen, err := strconv.Atoi(string(dataBuf[1:cntr])) + if err != nil { + return nil, dataBuf[cntr:], nil + } + + if dataLen == -1 { + return nil, dataBuf[cntr:], fmt.Errorf("NOT FOUND") + } + if cntr+lenCRLF > len(dataBuf)-lenCRLF { + *Len = dataLen + return nil, dataBuf[cntr+lenCRLF:], nil + } + if len(dataBuf[cntr+lenCRLF:len(dataBuf)-lenCRLF]) < dataLen { + *Len = dataLen + return nil, dataBuf[cntr+lenCRLF:], nil + } else { + return dataBuf[cntr+lenCRLF : cntr+dataLen+lenCRLF], dataBuf[cntr+dataLen+lenCRLF:], nil + } + case "*": + panic("array") + default: + if len(dataBuf) > 1 { + dataBuf = dataBuf[1:] + } else { + return nil, dataBuf, nil + } + } } + return nil, dataBuf, nil } -// repl starts an read input loop on the standard input. -// -// Commands are sent out after a new line is parsed. -// -// TODO: Improve ergonomics of the repl. -func repl(input *textproto.Reader, wr *textproto.Writer) { - for { - print("> ") - input_str, err := input.ReadLine() - if len(input_str) > 0 && err == nil { - // The redis server expects arrays of bulk strings for sending commands. - commands := strings.Split(input_str, " ") - // Specify the array length and create the bulk strings - pipe := []string{fmt.Sprintf("*%d", len(commands))} - pipe = append(pipe, createBulkStrings(commands)...) - for _, out := range pipe { - wr.PrintfLine(out) +func RedisContext(hostnamePort string, redisCmdWrite, redisCmdRead chan RedisCmd) { + tcpRemoteAddress, err := net.ResolveTCPAddr("tcp", hostnamePort) + if err != nil { + panic("cant resolve remote redis address") + } + var ladr *net.TCPAddr + msgBuf := make([]byte, 65000) + initMsg := []byte("*1\r\n$4\r\nPING\r\n") + writeChan := make(chan []byte) + readChan := make(chan []byte) + flushChan := make(chan int) + go netutils.AutoRecoonectedTCP(ladr, tcpRemoteAddress, msgBuf, initMsg, writeChan, readChan, flushChan) + <-readChan + dataBuf := make([]byte, 0) + dataLen := 0 + for true { + select { + case cmd := <-redisCmdWrite: + switch cmd.Command { + case "SET": + data := RedisSet(cmd.Name, cmd.Data) + writeChan <- data + case "GET": + data := RedisGet(cmd.Name) + writeChan <- data + } + case response := <-readChan: + data, dataBuf, err := ParseRedisResponse(response, dataBuf, &dataLen) + if dataLen != 0 { + for data == nil { + response = <-readChan + data, dataBuf, err = ParseRedisResponse(response, dataBuf, &dataLen) + } + } + if err != nil { + select { + case redisCmdRead <- RedisCmd{ + Error: err, + }: + case <-time.After(time.Second * 5): + } + } + if data != nil && string(data) != "PONG" { + select { + case redisCmdRead <- RedisCmd{ + Data: data, + }: + case <-time.After(time.Second * 5): + } + dataLen = 0 + } + case <-flushChan: + dataBuf = dataBuf[:] + dataLen = 0 + select { + case redisCmdRead <- RedisCmd{}: + case <-time.After(time.Second * 5): } } } } -// createBulkStrings converts normal strings into bulk strings to send to redis. -func createBulkStrings(commands []string) []string { - cmds := make([]string, len(commands)*2) - for i, v := range commands { - cmds[i*2] = fmt.Sprintf("$%d", len(commands[i])) - cmds[(i*2)+1] = v - } - return cmds +func (sr *SimpleRedis) Init(redisHost string) { + sr.redisHost = redisHost + sr.redisChanWrite = make(chan RedisCmd) + sr.redisChanRead = make(chan RedisCmd) + go RedisContext(sr.redisHost, sr.redisChanWrite, sr.redisChanRead) } - -// cmd/client starts a simple Read-Send-Print-Loop (RSPL?) -// with a redis server. -// -// TODO: Implement command flags -func Init(host string, password string) (*textproto.Writer, *textproto.Reader) { - log.Printf("yo") - // Create the connection to the redis server - // TODO: cli, rm hardcode - conn, err := net.Dial("tcp", host) - if err != nil { - log.Printf("[net] unable to connect: %v", err) - return nil, nil +func (sr *SimpleRedis) Do(cmd, name string, data []byte) ([]byte, error) { + sr.redisCmd.Command = cmd + sr.redisCmd.Name = name + sr.redisCmd.Data = data + sr.redisChanWrite <- sr.redisCmd + resp := <- sr.redisChanRead + if resp.Error != nil { + return nil, resp.Error } - defer conn.Close() - - // Create the readers and writers we will pass to our subroutines - msg_channel := make(chan RedisMessage) - write_sock := textproto.NewWriter(bufio.NewWriter(conn)) - read_sock := textproto.NewReader(bufio.NewReader(conn)) - - go runAsWaitGroup( - // Subscribe to redis connection - func() { - loopScan(read_sock, msg_channel) - }, - // Write to output - func() { - for { - msg := <-msg_channel - log.Printf("%v", msg) - } - }, - ).Wait() - log.Printf("yolo") - return write_sock, read_sock -} - -// runAsWaitGroup runs closures within a sync.WaitGroup -// -// Calling .Wait() on the resultant WaitGroup is a blocking operation until all closures finish running. -func runAsWaitGroup(closures ...func()) *sync.WaitGroup { - wg := sync.WaitGroup{} - for _, fn := range closures { - wg.Add(1) - go func(fn func()) { - fn() - wg.Done() - }(fn) - } - return &wg + return resp.Data, nil } \ No newline at end of file diff --git a/vendor/github.com/tehnerd/goUtils/netutils/netutils.go b/vendor/github.com/tehnerd/goUtils/netutils/netutils.go new file mode 100644 index 0000000..ddfea12 --- /dev/null +++ b/vendor/github.com/tehnerd/goUtils/netutils/netutils.go @@ -0,0 +1,525 @@ +package netutils + +import ( + "math/rand" + "net" + "strconv" + "strings" + "sync" + //"sync/atomic" + "time" +) + +//Receive msg from tcp socket and send it as a []byte to readChan +func ReadFromTCP2(sock *net.TCPConn, msgBuf []byte, readChan chan []byte, + feedbackChanFromSocket chan int) { + loop := 1 + for loop == 1 { + bytes, err := sock.Read(msgBuf) + if err != nil { + feedbackChanFromSocket <- 1 + loop = 0 + continue + } + b := make([]byte, 0) + b = append(b, msgBuf[:bytes]...) + readChan <- b + } +} + +//Receive msg from tcp socket and send it as a []byte to readChan +func ReadFromTCP(sock *net.TCPConn, msgBuf []byte, readChan chan []byte, + feedbackChanFromSocket chan int) { + feedbackToSocket := make(chan bool) + feedbackFromSocket := make(chan bool) + reuseBufferChan := make(chan []byte, 1) + go readFromTCP(sock, readChan, feedbackFromSocket, + feedbackToSocket, reuseBufferChan) + go func() { + <-feedbackFromSocket + feedbackChanFromSocket <- 1 + }() +} + +/*simple write to TCP, for oneway connections only (no communitcation w/ "read" part of the socket +in terms of error propogation)*/ +func WriteToTCPw2(sock *net.TCPConn, writeChan chan []byte, + feedbackChan chan int) { + loop := 1 + for loop == 1 { + select { + case msg := <-writeChan: + _, err := sock.Write(msg) + if err != nil { + feedbackChan <- 1 + continue + } + case <-feedbackChan: + loop = 0 + } + } +} + +/*simple write to TCP, for oneway connections only (no communitcation w/ "read" part of the socket +in terms of error propogation)*/ +func WriteToTCPw(sock *net.TCPConn, writeChan chan []byte, + feedbackChan chan int) { + feedbackToSocket := make(chan bool) + feedbackFromSocket := make(chan bool) + go writeToTCP(sock, writeChan, feedbackFromSocket, feedbackToSocket) + go func() { + <-feedbackFromSocket + feedbackChan <- 1 + }() +} + +//simple write to tcp w/ erorr propagation to/from "read" part of the socket +func WriteToTCPrw2(sock *net.TCPConn, writeChan chan []byte, + feedbackChanFromSocket, feedbackChanToSocket chan int) { + loop := 1 + for loop == 1 { + select { + case msg := <-writeChan: + _, err := sock.Write(msg) + if err != nil { + select { + case feedbackChanFromSocket <- 1: + continue + case loop = <-feedbackChanToSocket: + loop = 0 + continue + } + } + case <-feedbackChanToSocket: + loop = 0 + } + } +} + +//simple write to tcp w/ erorr propagation to/from "read" part of the socket +func WriteToTCPrw(sock *net.TCPConn, writeChan chan []byte, + feedbackChanFromSocket, feedbackChanToSocket chan int) { + feedbackToSocket := make(chan bool) + feedbackFromSocket := make(chan bool) + go writeToTCP(sock, writeChan, feedbackFromSocket, feedbackToSocket) + go func() { + select { + case <-feedbackFromSocket: + feedbackChanFromSocket <- 1 + case <-feedbackChanToSocket: + feedbackToSocket <- true + + } + }() + +} + +func makeBuffer(reuseBufferChan chan []byte) []byte { + select { + case buf := <-reuseBufferChan: + return buf + default: + return make([]byte, 10000) + } +} + +//Receive msg from tcp socket and send it as a []byte to readChan,with buffer reuse +func readFromTCP(sock *net.TCPConn, readChan chan []byte, + feedbackFromSocket, feedbackToSocket chan bool, + reuseBufferChan chan []byte) { + loop := 1 + var buf []byte + for loop == 1 { + buf = makeBuffer(reuseBufferChan) + bytes, err := sock.Read(buf) + if err != nil { + select { + case feedbackFromSocket <- true: + loop = 0 + continue + case <-feedbackToSocket: + loop = 0 + continue + } + } + select { + case readChan <- buf[:bytes]: + case <-feedbackToSocket: + loop = 0 + continue + + } + } +} + +//TCP's write routine, with feedback's chans +func writeToTCP(sock *net.TCPConn, writeChan chan []byte, + feedbackFromSocket, feedbackToSocket chan bool) { + loop := 1 + for loop == 1 { + select { + case msg := <-writeChan: + _, err := sock.Write(msg) + if err != nil { + select { + case feedbackFromSocket <- true: + case <-feedbackToSocket: + } + loop = 0 + continue + } + case <-feedbackToSocket: + loop = 0 + continue + } + } +} + +//reconnecting to remote host for both read and write purpose +func ReconnectTCPRW(ladr, radr *net.TCPAddr, msgBuf []byte, writeChan chan []byte, + readChan chan []byte, feedbackChanToSocket, feedbackChanFromSocket chan int, + init_msg []byte) { + loop := 1 + for loop == 1 { + sock, err := net.DialTCP("tcp", ladr, radr) + if err != nil { + time.Sleep(time.Duration(20+rand.Intn(15)) * time.Second) + continue + } + //testing health of the new socket. GO sometimes doesnt rise the error when + // we receive RST from remote side + _, err = sock.Write(init_msg) + if err != nil { + sock.Close() + time.Sleep(time.Duration(20+rand.Intn(15)) * time.Second) + continue + } + loop = 0 + go ReadFromTCP(sock, msgBuf, readChan, feedbackChanFromSocket) + go WriteToTCPrw(sock, writeChan, feedbackChanFromSocket, feedbackChanToSocket) + } +} + +func ReconnectTCPRWReuse(ladr, radr *net.TCPAddr, + readChan, writeChan, reuseBufferChan chan []byte, + readFeedbackFrom, readFeedbackTo chan bool, + writeFeedbackFrom, writeFeedbackTo chan bool) { + loop := 1 + for loop == 1 { + sock, err := net.DialTCP("tcp", ladr, radr) + if err != nil { + time.Sleep(time.Duration(20+rand.Intn(15)) * time.Second) + continue + } + loop = 0 + go readFromTCP(sock, readChan, readFeedbackFrom, readFeedbackTo, + reuseBufferChan) + go writeToTCP(sock, writeChan, writeFeedbackFrom, writeFeedbackTo) + } +} + +func AutoRecoonectedTCP(ladr, radr *net.TCPAddr, msgBuf, initMsg []byte, + writeChan, readChan chan []byte, flushChan chan int) { + feedbackChanFromSocket := make(chan int) + feedbackChanToSocket := make(chan int) + go ReconnectTCPRW(ladr, radr, msgBuf, writeChan, readChan, feedbackChanToSocket, + feedbackChanFromSocket, initMsg) + for { + select { + case feedbackFromSocket := <-feedbackChanFromSocket: + feedbackChanToSocket <- feedbackFromSocket + flushChan <- 1 + go ReconnectTCPRW(ladr, radr, msgBuf, writeChan, + readChan, feedbackChanToSocket, + feedbackChanFromSocket, initMsg) + + } + } + +} + +func AutoRecoonectedTCPReuse(ladr, radr *net.TCPAddr, + readChan, writeChan chan []byte, + reuseChan chan []byte, + flushChan chan bool) { + readFeedbackFrom := make(chan bool) + readFeedbackTo := make(chan bool) + writeFeedbackFrom := make(chan bool) + writeFeedbackTo := make(chan bool) + go ReconnectTCPRWReuse(ladr, radr, readChan, writeChan, reuseChan, + readFeedbackFrom, readFeedbackTo, + writeFeedbackFrom, writeFeedbackTo) + for { + select { + case <-readFeedbackFrom: + writeFeedbackTo <- true + case <-writeFeedbackFrom: + readFeedbackTo <- true + } + flushChan <- true + go ReconnectTCPRWReuse(ladr, radr, readChan, writeChan, reuseChan, + readFeedbackFrom, readFeedbackTo, + writeFeedbackFrom, writeFeedbackTo) + } + +} + +//reconnecting to remote host for write only +func ReconnectTCPW(radr net.TCPAddr, writeChan chan []byte, feedbackChan chan int) { + loop := 1 + for loop == 1 { + time.Sleep(time.Duration(20+rand.Intn(15)) * time.Second) + sock, err := net.DialTCP("tcp", nil, &radr) + if err != nil { + continue + } + //testing health of the new socket. GO sometimes doesnt rise the error when + // we receive RST from remote side + _, err = sock.Write([]byte{1}) + if err != nil { + sock.Close() + continue + } + loop = 0 + go WriteToTCPw(sock, writeChan, feedbackChan) + } +} + +/* --------------------- CONNECTION MANAGER ------------------------- +Connection manager will allow send data and receive data from remote hosts. +it will have single ConnectionMsg (see below) read chan and single +ConnectionMsg write chan toward it's clients +as well as single read chan from sockets, but multiple write sockets. +it will route msgs according to Host field in connectionManager struct +(if it received from client, it will send this msg toward Host's sockets; +if recved from socket, will proxy it toward client(and client will know from which remote host +it was received) +-------------------------------------------------------------------- */ +/* +MsgType's could be: +from Api's client to ConnectionManager: +"Data" - msg with Data to Host +"Connect" - connect to new Host +... +from ConnectionManager to Api's client: +"BufferFlush" - notification that connection to remote Host not longer working. +advice to flush all the msg buffers assosiated with remote host +*/ +type ConnectionMsg struct { + Host string + Data []byte + Type string +} + +/* +Receive msg from tcp socket and send it as a ConnectionMsg to readChan +TODO: think about more generic version to be more DRYer (to work in both CM and []byte chans +*/ +func CMReadFromTCP(sock *net.TCPConn, readChan chan ConnectionMsg, + peerAddress string) { + msgBuf := make([]byte, 65000) + loop := 1 + var msg ConnectionMsg + msg.Host = peerAddress + msg.Type = "Data" + for loop == 1 { + bytes, err := sock.Read(msgBuf) + if err != nil { + msg.Type = "ReadError" + readChan <- msg + loop = 0 + continue + } + b := make([]byte, 0) + b = append(b, msgBuf[:bytes]...) + msg.Data = b + readChan <- msg + } +} + +/* +ConnectionManager write instance to tcp w/ erorr propagation to/from "read" part of the socket +TODO: think about more generic version to be more DRYer (to work in both CM and []byte chans +*/ +func CMWriteToTCP(sock *net.TCPConn, writeChan, readChan chan ConnectionMsg, + peerAddress string) { + loop := 1 + var errorMsg ConnectionMsg + errorMsg.Host = peerAddress + errorMsg.Type = "WriteError" + for loop == 1 { + select { + case msg := <-writeChan: + switch msg.Type { + case "Data": + _, err := sock.Write(msg.Data) + if err != nil { + for loop == 1 { + select { + case readChan <- errorMsg: + case errorMsg := <-writeChan: + if errorMsg.Type != "ConnectionError" { + continue + } + } + loop = 0 + } + } + case "ConnectionError": + loop = 0 + default: + continue + } + } + } +} + +func StartConnection(tcpConn *net.TCPConn, writeChan, + readChan chan ConnectionMsg, peerAddress string) { + go CMReadFromTCP(tcpConn, readChan, peerAddress) + go CMWriteToTCP(tcpConn, writeChan, readChan, peerAddress) +} + +func CMListenForConnection(mutex *sync.RWMutex, localPort int, + writeChanMap map[string]chan ConnectionMsg, + connectionStateMap map[string]int, + readChan chan ConnectionMsg) { + laddr := strings.Join([]string{":", strconv.Itoa(localPort)}, "") + tcpLaddr, err := net.ResolveTCPAddr("tcp", laddr) + if err != nil { + panic("cant resolve local address for binding") + } + tcpListener, err := net.ListenTCP("tcp", tcpLaddr) + if err != nil { + panic("cant listen on local address for binding") + } + for { + tcpConn, err := tcpListener.AcceptTCP() + if err == nil { + radr := strings.Split(tcpConn.RemoteAddr().String(), ":")[0] + // check if we already has connection to remote peer as a client + mutex.Lock() + if val, exist := connectionStateMap[radr]; exist && val == 1 { + tcpConn.Close() + mutex.Unlock() + continue + } + connectionStateMap[radr] = 1 + if writeChan, exist := writeChanMap[radr]; exist { + mutex.Unlock() + go StartConnection(tcpConn, writeChan, readChan, radr) + } else { + writeChanMap[radr] = make(chan ConnectionMsg) + mutex.Unlock() + go StartConnection(tcpConn, writeChanMap[radr], readChan, radr) + } + } + } +} + +func CMConnectToRemotePeer(mutex *sync.RWMutex, peerTcpAddr *net.TCPAddr, + radr string, + writeChan chan ConnectionMsg, + readChan chan ConnectionMsg, + connectionStateMap map[string]int) { + connectLoop := 1 + for connectLoop == 1 { + mutex.RLock() + if connectionStateMap[radr] == 1 { + connectLoop = 0 + mutex.RUnlock() + continue + } + mutex.RUnlock() + tcpConn, err := net.DialTCP("tcp", nil, peerTcpAddr) + if err != nil { + time.Sleep(time.Second * time.Duration(rand.Int63n(15))) + continue + } + mutex.Lock() + if connectionStateMap[radr] == 0 { + connectionStateMap[radr] = 1 + go StartConnection(tcpConn, writeChan, readChan, radr) + mutex.Unlock() + connectLoop = 0 + continue + } else { + mutex.Unlock() + tcpConn.Close() + connectLoop = 0 + continue + } + } +} + +func ConnectionManager(msgChan chan ConnectionMsg, localPort int) { + writeChanMap := make(map[string]chan ConnectionMsg) + connectionStateMap := make(map[string]int) + readChan := make(chan ConnectionMsg) + var connectionMutex sync.RWMutex + go CMListenForConnection(&connectionMutex, localPort, writeChanMap, + connectionStateMap, readChan) + for { + select { + case msgToPeer := <-msgChan: + switch msgToPeer.Type { + case "Data": + if state, exists := connectionStateMap[msgToPeer.Host]; exists && state == 1 { + /*FIXME/THINK: There could be deadlock if connectin closes before we will + be able to send to the chan */ + writeChan := writeChanMap[msgToPeer.Host] + writeChan <- msgToPeer + } else { + msgChan <- ConnectionMsg{Type: "ConnectionNotExist"} + } + case "Connect": + if len(strings.Split(msgToPeer.Host, ":")) > 1 { + radr := strings.Split(msgToPeer.Host, ":")[0] + connectionMutex.Lock() + if _, exist := writeChanMap[radr]; !exist { + writeChanMap[radr] = make(chan ConnectionMsg) + } + connectionMutex.Unlock() + peerTcpAddr, err := net.ResolveTCPAddr("tcp", msgToPeer.Host) + if err != nil { + //XXX: think about , mb make something less drastic + panic("cant resolve remote address") + } + go CMConnectToRemotePeer(&connectionMutex, peerTcpAddr, radr, + writeChanMap[radr], readChan, connectionStateMap) + } else { + connectionMutex.Lock() + if _, exist := writeChanMap[msgToPeer.Host]; !exist { + writeChanMap[msgToPeer.Host] = make(chan ConnectionMsg) + } + connectionMutex.Unlock() + remoteAddr := strings.Join([]string{msgToPeer.Host, strconv.Itoa(localPort)}, ":") + peerTcpAddr, err := net.ResolveTCPAddr("tcp", remoteAddr) + if err != nil { + //XXX: again panic could be overkill + panic("cant resolve remote address") + } + go CMConnectToRemotePeer(&connectionMutex, peerTcpAddr, msgToPeer.Host, + writeChanMap[msgToPeer.Host], readChan, connectionStateMap) + } + } + case msgFromPeer := <-readChan: + switch msgFromPeer.Type { + case "Data": + msgChan <- msgFromPeer + case "WriteError", "ReadError": + connectionMutex.Lock() + connectionStateMap[msgFromPeer.Host] = 0 + connectionMutex.Unlock() + if msgFromPeer.Type == "ReadError" { + writeChanMap[msgFromPeer.Host] <- ConnectionMsg{Type: "ConnectionError"} + } + var msgToApiClient ConnectionMsg + msgToApiClient.Host = msgFromPeer.Host + msgToApiClient.Type = "BufferFlush" + msgChan <- msgToApiClient + + } + } + } +} diff --git a/vendor/github.com/tehnerd/goUtils/netutils/tcp_listen.go b/vendor/github.com/tehnerd/goUtils/netutils/tcp_listen.go new file mode 100644 index 0000000..a7ad4fb --- /dev/null +++ b/vendor/github.com/tehnerd/goUtils/netutils/tcp_listen.go @@ -0,0 +1,42 @@ +package netutils + +import ( + "errors" + "net" +) + +/* + we need to provide a function, which will read/write to/from socket and read/write to/from sockets feedback chans +*/ +func ListenForConnection(port string, fn func(chan []byte, chan []byte, chan int, chan int)) error { + addr := ":" + port + tcpAddr, err := net.ResolveTCPAddr("tcp", addr) + if err != nil { + return errors.New("cant resolve local tcp address") + } + loop := 1 + servSock, err := net.ListenTCP("tcp", tcpAddr) + if err != nil { + return errors.New("cant bind to local tcp address") + } + + for loop == 1 { + sock, err := servSock.AcceptTCP() + if err != nil { + continue + } + go ServeTcpConn(sock, fn) + } + return nil +} + +func ServeTcpConn(sock *net.TCPConn, fn func(chan []byte, chan []byte, chan int, chan int)) { + readChan := make(chan []byte) + writeChan := make(chan []byte) + feedbackFrom := make(chan int, 1) + feedbackTo := make(chan int, 1) + buf := make([]byte, 65535) + go ReadFromTCP(sock, buf, readChan, feedbackFrom) + go WriteToTCPrw(sock, writeChan, feedbackFrom, feedbackTo) + fn(readChan, writeChan, feedbackFrom, feedbackTo) +} diff --git a/vendor/modules.txt b/vendor/modules.txt index 5e64cb5..bbb1f06 100644 --- a/vendor/modules.txt +++ b/vendor/modules.txt @@ -1,3 +1,6 @@ # github.com/leprosus/golang-ttl-map v1.1.7 ## explicit; go 1.15 github.com/leprosus/golang-ttl-map +# github.com/tehnerd/goUtils v0.0.0-20150515130609-5a2d8fb2ded8 +## explicit +github.com/tehnerd/goUtils/netutils