mirror of
https://github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin.git
synced 2026-07-21 11:38:59 +02:00
✨ redis included
This commit is contained in:
+2
-6
@@ -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)
|
||||
if config.CrowdsecMode == streamMode && ticker == nil {
|
||||
ticker = startTicker(config, func() {
|
||||
handleStreamCache(bouncer)
|
||||
})
|
||||
}
|
||||
}()
|
||||
go handleStreamCache(bouncer)
|
||||
}
|
||||
|
||||
return bouncer, nil
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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=
|
||||
|
||||
Vendored
+14
-24
@@ -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)
|
||||
}
|
||||
|
||||
+148
-325
@@ -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 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 (m MsgError) String() string {
|
||||
return string(m)
|
||||
func RedisSet(name string, data []byte) []byte {
|
||||
return GenRedisArray([]byte("SET"), []byte(name), data)
|
||||
}
|
||||
|
||||
func (m MsgSimpleStr) String() string {
|
||||
return string(m)
|
||||
func RedisGet(name string) []byte {
|
||||
return GenRedisArray([]byte("GET"), []byte(name))
|
||||
}
|
||||
|
||||
// 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
|
||||
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 "", NewParseError(msg.RedisType, SimpleString, BulkString)
|
||||
return dataBuf[:*Len], dataBuf[*Len:], nil
|
||||
}
|
||||
}
|
||||
|
||||
func (msg RedisMessage) AsInteger() (int64, error) {
|
||||
if conv, ok := msg.Choice.(MsgInteger); ok {
|
||||
return int64(conv), nil
|
||||
} else {
|
||||
return 0, NewParseError(msg.RedisType, Integer)
|
||||
for {
|
||||
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
|
||||
}
|
||||
}
|
||||
|
||||
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)
|
||||
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
|
||||
}
|
||||
}
|
||||
|
||||
// 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
|
||||
if cntr == len(dataBuf) || cntr+lenCRLF > len(dataBuf) {
|
||||
return nil, dataBuf, nil
|
||||
}
|
||||
|
||||
// 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,
|
||||
}
|
||||
}
|
||||
|
||||
// Fetch grabs the next incoming RESP type from the IO input
|
||||
func Fetch(rr *RedisReader) *RedisMessage {
|
||||
byteCode, err := rr.TryReadByte()
|
||||
dataLen, err := strconv.Atoi(string(dataBuf[1:cntr]))
|
||||
if err != nil {
|
||||
return nil
|
||||
return nil, dataBuf[cntr:], 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)
|
||||
|
||||
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:
|
||||
panic(fmt.Sprint("Unknown bytecode: ", byteCode))
|
||||
if len(dataBuf) > 1 {
|
||||
dataBuf = dataBuf[1:]
|
||||
} else {
|
||||
return nil, dataBuf, nil
|
||||
}
|
||||
}
|
||||
}
|
||||
return nil, dataBuf, nil
|
||||
}
|
||||
|
||||
// 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()
|
||||
func RedisContext(hostnamePort string, redisCmdWrite, redisCmdRead chan RedisCmd) {
|
||||
tcpRemoteAddress, err := net.ResolveTCPAddr("tcp", hostnamePort)
|
||||
if err != nil {
|
||||
log.Printf("Fatal RespSimpleStr: %+v", err)
|
||||
return nil
|
||||
panic("cant resolve remote redis address")
|
||||
}
|
||||
|
||||
return &RedisMessage{
|
||||
RedisType: RespT(SimpleString),
|
||||
Raw: s,
|
||||
Choice: MsgSimpleStr(s),
|
||||
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)
|
||||
}
|
||||
}
|
||||
|
||||
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),
|
||||
select {
|
||||
case redisCmdRead <- RedisCmd{
|
||||
Error: err,
|
||||
}:
|
||||
case <-time.After(time.Second * 5):
|
||||
}
|
||||
}
|
||||
|
||||
// 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
|
||||
if data != nil && string(data) != "PONG" {
|
||||
select {
|
||||
case redisCmdRead <- RedisCmd{
|
||||
Data: data,
|
||||
}:
|
||||
case <-time.After(time.Second * 5):
|
||||
}
|
||||
|
||||
// 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)
|
||||
dataLen = 0
|
||||
}
|
||||
|
||||
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()
|
||||
}
|
||||
}
|
||||
|
||||
// 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)
|
||||
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
|
||||
}
|
||||
+525
@@ -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
|
||||
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
+42
@@ -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)
|
||||
}
|
||||
Vendored
+3
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user