mirror of
https://github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin.git
synced 2026-09-02 20:28:50 +02:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
4a5f1eca6a | ||
|
|
cdda369fb8 | ||
|
|
f0f28fecef | ||
|
|
23620207f7 | ||
|
|
4a674da9f9 | ||
|
|
0dfd18f18e | ||
|
|
d8ee0a34eb | ||
|
|
b50074dca4 | ||
|
|
5044004ec2 | ||
|
|
41a46c0584 | ||
|
|
0a186cf9a9 | ||
|
|
ffcf4356fc | ||
|
|
f94e48aa03 | ||
|
|
a197194591 | ||
|
|
44e329cd57 | ||
|
|
4058836678 | ||
|
|
87ed9e9c4e | ||
|
|
395c80dccf | ||
|
|
59268ee33d | ||
|
|
781a83465e |
@@ -28,6 +28,9 @@ run_local:
|
|||||||
run_behindproxy:
|
run_behindproxy:
|
||||||
docker-compose -f exemples/behind-proxy/docker-compose.cloudflare.yml up -d --remove-orphans
|
docker-compose -f exemples/behind-proxy/docker-compose.cloudflare.yml up -d --remove-orphans
|
||||||
|
|
||||||
|
run_cacheredis:
|
||||||
|
docker-compose -f exemples/redis-cache/docker-compose.redis.yml up -d --remove-orphans
|
||||||
|
|
||||||
run:
|
run:
|
||||||
docker-compose -f docker-compose.yml up -d --remove-orphans
|
docker-compose -f docker-compose.yml up -d --remove-orphans
|
||||||
|
|
||||||
@@ -51,6 +54,7 @@ show_dev_logs:
|
|||||||
|
|
||||||
clean_all_docker:
|
clean_all_docker:
|
||||||
docker-compose -f exemples/behind-proxy/docker-compose.cloudflare.yml down --remove-orphans
|
docker-compose -f exemples/behind-proxy/docker-compose.cloudflare.yml down --remove-orphans
|
||||||
|
docker-compose -f exemples/behind-proxy/docker-compose.redis.yml down --remove-orphans
|
||||||
docker-compose -f docker-compose.local.yml down --remove-orphans
|
docker-compose -f docker-compose.local.yml down --remove-orphans
|
||||||
docker-compose -f docker-compose.yml down --remove-orphans
|
docker-compose -f docker-compose.yml down --remove-orphans
|
||||||
|
|
||||||
|
|||||||
@@ -26,6 +26,9 @@ There are 3 operating modes (CrowdsecMode) for this plugin:
|
|||||||
|
|
||||||
The recommanded mode for performance is the streaming mode, decisions are updated every 60 sec by default and that's the only communication between traefik and crowdsec. Every requests that happens hits the cache for quick decisions.
|
The recommanded mode for performance is the streaming mode, decisions are updated every 60 sec by default and that's the only communication between traefik and crowdsec. Every requests that happens hits the cache for quick decisions.
|
||||||
|
|
||||||
|
The cache can be local to the Traefik instance using the filesystem or use of a separated redis instance.
|
||||||
|
The redis instance is currently in beta and support Redis 7.0.X version
|
||||||
|
|
||||||
## Usage
|
## Usage
|
||||||
|
|
||||||
To get started, use the `docker-compose.yml` file.
|
To get started, use the `docker-compose.yml` file.
|
||||||
@@ -36,19 +39,16 @@ make run
|
|||||||
```
|
```
|
||||||
|
|
||||||
### Note
|
### Note
|
||||||
Each middleware in traefik has it's own data and is instanciated by service.
|
|
||||||
This means if there are 10 services protected by the bouncer in streaming alone or live mode, the cache will be duplicated to all 10 services.
|
|
||||||
This is because traefik does not allow plugins to store data locally that can be consummed.
|
|
||||||
|
|
||||||
The synchronisation with the crowdsec service will happen also 10 times in the period selected.
|
**/!\ Since Release 1.10, cache is no longer duplicated but shared by all services**
|
||||||
It should be taken into account when fixing this period so each middleware has time to sync data from crowdsec.
|
*This lowers the overhead of the cache in memory and the numbers of cache to fetch it from crowdsec in situation with many services*
|
||||||
|
|
||||||
At each start of synchronisation, the middleware will wait a random number of seconds to avoid simultaneous calls to crowdsec.
|
|
||||||
|
|
||||||
### Variables
|
### Variables
|
||||||
- Enabled
|
- Enabled
|
||||||
- bool
|
- bool
|
||||||
- enable the plugin
|
- enable the plugin
|
||||||
|
- default: false
|
||||||
- LogLevel
|
- LogLevel
|
||||||
- string
|
- string
|
||||||
- default: `INFO`, expected value are: `INFO`, `DEBUG`
|
- default: `INFO`, expected value are: `INFO`, `DEBUG`
|
||||||
@@ -61,7 +61,7 @@ At each start of synchronisation, the middleware will wait a random number of se
|
|||||||
- CrowdsecLapiHost
|
- CrowdsecLapiHost
|
||||||
- string
|
- string
|
||||||
- default: "crowdsec:8080"
|
- default: "crowdsec:8080"
|
||||||
- Crowdsec LAPI available on which host.
|
- Crowdsec LAPI available on which host and port.
|
||||||
- CrowdsecLapiKey
|
- CrowdsecLapiKey
|
||||||
- string
|
- string
|
||||||
- Crowdsec LAPI generated key for the bouncer : **must be unique by service**.
|
- Crowdsec LAPI generated key for the bouncer : **must be unique by service**.
|
||||||
@@ -79,8 +79,17 @@ At each start of synchronisation, the middleware will wait a random number of se
|
|||||||
- List of IPs of trusted Proxies that are in front of traefik (ex: Cloudflare)
|
- List of IPs of trusted Proxies that are in front of traefik (ex: Cloudflare)
|
||||||
- ForwardedHeadersCustomName
|
- ForwardedHeadersCustomName
|
||||||
- string
|
- string
|
||||||
- default: X-Forwarded-For
|
- default: "X-Forwarded-For"
|
||||||
- Name of the header where the real IP of the client should be retrieved
|
- Name of the header where the real IP of the client should be retrieved
|
||||||
|
- RedisCacheEnabled
|
||||||
|
- bool
|
||||||
|
- default: false
|
||||||
|
- enable redis cache instead of filesystem cache
|
||||||
|
- RedisCacheHost
|
||||||
|
- string
|
||||||
|
- default: "redis:6379"
|
||||||
|
- hostname and port for the redis service
|
||||||
|
|
||||||
|
|
||||||
### Configuration
|
### Configuration
|
||||||
|
|
||||||
@@ -131,6 +140,8 @@ http:
|
|||||||
- 10.0.10.23/32
|
- 10.0.10.23/32
|
||||||
- 10.0.20.0/24
|
- 10.0.20.0/24
|
||||||
forwardedHeadersCustomName: X-Custom-Header
|
forwardedHeadersCustomName: X-Custom-Header
|
||||||
|
redisCacheEnabled: false
|
||||||
|
redisCacheHost: "redis:6379"
|
||||||
```
|
```
|
||||||
These are the default values of the plugin except for LapiKey.
|
These are the default values of the plugin except for LapiKey.
|
||||||
|
|
||||||
@@ -224,9 +235,21 @@ We configure the middleware to trust as well the IP:
|
|||||||
|
|
||||||
To run the environnement run:
|
To run the environnement run:
|
||||||
```bash
|
```bash
|
||||||
make run_behind_proxy
|
make run_behindproxy
|
||||||
```
|
```
|
||||||
|
|
||||||
|
2. With Redis as an external shared cache
|
||||||
|
|
||||||
|
The plugin must be configured to connect to a redis instance
|
||||||
|
```yaml
|
||||||
|
redisCacheHost: "redis:6379"
|
||||||
|
```
|
||||||
|
Here **redis** is the hostname of a container located in the same network as Traefik and **6379** the default port of redis
|
||||||
|
|
||||||
|
To run the demo environnement run:
|
||||||
|
```bash
|
||||||
|
make run_cacheredis
|
||||||
|
```
|
||||||
|
|
||||||
### About
|
### About
|
||||||
|
|
||||||
|
|||||||
+14
-10
@@ -43,7 +43,9 @@ type Config struct {
|
|||||||
ForwardedHeadersCustomName string `json:"forwardedheaderscustomheader,omitempty"`
|
ForwardedHeadersCustomName string `json:"forwardedheaderscustomheader,omitempty"`
|
||||||
UpdateIntervalSeconds int64 `json:"updateIntervalSeconds,omitempty"`
|
UpdateIntervalSeconds int64 `json:"updateIntervalSeconds,omitempty"`
|
||||||
DefaultDecisionSeconds int64 `json:"defaultDecisionSeconds,omitempty"`
|
DefaultDecisionSeconds int64 `json:"defaultDecisionSeconds,omitempty"`
|
||||||
ForwardedHeadersTrustedIPs []string `json:"forwardedheaderstrustedips,omitempty"`
|
ForwardedHeadersTrustedIPs []string `json:"forwardedHeadersTrustedIps,omitempty"`
|
||||||
|
RedisCacheEnabled bool `json:"redisCacheEnabled,omitempty"`
|
||||||
|
RedisCacheHost string `json:"redisCacheHost,omitempty"`
|
||||||
}
|
}
|
||||||
|
|
||||||
// CreateConfig creates the default plugin configuration.
|
// CreateConfig creates the default plugin configuration.
|
||||||
@@ -59,6 +61,8 @@ func CreateConfig() *Config {
|
|||||||
DefaultDecisionSeconds: 60,
|
DefaultDecisionSeconds: 60,
|
||||||
ForwardedHeadersTrustedIPs: []string{},
|
ForwardedHeadersTrustedIPs: []string{},
|
||||||
ForwardedHeadersCustomName: "X-Forwarded-For",
|
ForwardedHeadersCustomName: "X-Forwarded-For",
|
||||||
|
RedisCacheEnabled: false,
|
||||||
|
RedisCacheHost: "redis:6379",
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -115,16 +119,16 @@ func New(ctx context.Context, next http.Handler, config *Config, name string) (h
|
|||||||
Timeout: 5 * time.Second,
|
Timeout: 5 * time.Second,
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
if config.CrowdsecMode == streamMode {
|
if config.RedisCacheEnabled {
|
||||||
go func() {
|
cache.InitRedisClient(config.RedisCacheHost)
|
||||||
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
|
return bouncer, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -2,7 +2,7 @@ version: "3.8"
|
|||||||
|
|
||||||
services:
|
services:
|
||||||
traefik:
|
traefik:
|
||||||
image: "traefik:v2.9.1"
|
image: "traefik:v2.9.4"
|
||||||
container_name: "traefik"
|
container_name: "traefik"
|
||||||
command:
|
command:
|
||||||
# - "--log.level=DEBUG"
|
# - "--log.level=DEBUG"
|
||||||
@@ -35,7 +35,6 @@ services:
|
|||||||
- "traefik.http.services.service1.loadbalancer.server.port=80"
|
- "traefik.http.services.service1.loadbalancer.server.port=80"
|
||||||
- "traefik.http.middlewares.crowdsec1.plugin.bouncer.enabled=true"
|
- "traefik.http.middlewares.crowdsec1.plugin.bouncer.enabled=true"
|
||||||
- "traefik.http.middlewares.crowdsec1.plugin.bouncer.crowdseclapikey=40796d93c2958f9e58345514e67740e5"
|
- "traefik.http.middlewares.crowdsec1.plugin.bouncer.crowdseclapikey=40796d93c2958f9e58345514e67740e5"
|
||||||
- "traefik.http.middlewares.crowdsec1.plugin.bouncer.forwardedheaderstrustedips=172.26.0.5"
|
|
||||||
|
|
||||||
whoami2:
|
whoami2:
|
||||||
image: traefik/whoami
|
image: traefik/whoami
|
||||||
@@ -48,7 +47,6 @@ services:
|
|||||||
- "traefik.http.services.service2.loadbalancer.server.port=80"
|
- "traefik.http.services.service2.loadbalancer.server.port=80"
|
||||||
- "traefik.http.middlewares.crowdsec2.plugin.bouncer.enabled=true"
|
- "traefik.http.middlewares.crowdsec2.plugin.bouncer.enabled=true"
|
||||||
- "traefik.http.middlewares.crowdsec2.plugin.bouncer.crowdseclapikey=44c36dac5c4140af9f06f397508e82c7"
|
- "traefik.http.middlewares.crowdsec2.plugin.bouncer.crowdseclapikey=44c36dac5c4140af9f06f397508e82c7"
|
||||||
- "traefik.http.middlewares.crowdsec2.plugin.bouncer.forwardedheaderstrustedips=172.26.0.5"
|
|
||||||
|
|
||||||
crowdsec:
|
crowdsec:
|
||||||
image: crowdsecurity/crowdsec:v1.4.1
|
image: crowdsecurity/crowdsec:v1.4.1
|
||||||
|
|||||||
+3
-3
@@ -2,7 +2,7 @@ version: "3.8"
|
|||||||
|
|
||||||
services:
|
services:
|
||||||
traefik:
|
traefik:
|
||||||
image: "traefik:v2.9.1"
|
image: "traefik:v2.9.4"
|
||||||
container_name: "traefik"
|
container_name: "traefik"
|
||||||
command:
|
command:
|
||||||
- "--accesslog"
|
- "--accesslog"
|
||||||
@@ -29,7 +29,7 @@ services:
|
|||||||
labels:
|
labels:
|
||||||
- "traefik.enable=true"
|
- "traefik.enable=true"
|
||||||
# Definition of the router
|
# Definition of the router
|
||||||
- "traefik.http.routers.router1.rule=Host(`localhost`) && Path(`/foo`)"
|
- "traefik.http.routers.router1.rule=Path(`/foo`)"
|
||||||
- "traefik.http.routers.router1.entrypoints=web"
|
- "traefik.http.routers.router1.entrypoints=web"
|
||||||
- "traefik.http.routers.router1.middlewares=crowdsec2@docker"
|
- "traefik.http.routers.router1.middlewares=crowdsec2@docker"
|
||||||
# Definition of the service
|
# Definition of the service
|
||||||
@@ -47,7 +47,7 @@ services:
|
|||||||
labels:
|
labels:
|
||||||
- "traefik.enable=true"
|
- "traefik.enable=true"
|
||||||
# Definition of the router
|
# Definition of the router
|
||||||
- "traefik.http.routers.router2.rule=Host(`localhost`) && Path(`/bar`)"
|
- "traefik.http.routers.router2.rule=Path(`/bar`)"
|
||||||
- "traefik.http.routers.router2.entrypoints=web"
|
- "traefik.http.routers.router2.entrypoints=web"
|
||||||
- "traefik.http.routers.router2.middlewares=crowdsec2@docker"
|
- "traefik.http.routers.router2.middlewares=crowdsec2@docker"
|
||||||
# Definition of the service
|
# Definition of the service
|
||||||
|
|||||||
@@ -0,0 +1,4 @@
|
|||||||
|
filenames:
|
||||||
|
- /var/log/traefik/access.log
|
||||||
|
labels:
|
||||||
|
type: traefik
|
||||||
@@ -0,0 +1,94 @@
|
|||||||
|
version: "3.8"
|
||||||
|
|
||||||
|
services:
|
||||||
|
traefik:
|
||||||
|
image: "traefik:v2.9.4"
|
||||||
|
container_name: "traefik"
|
||||||
|
command:
|
||||||
|
# - "--log.level=DEBUG"
|
||||||
|
- "--accesslog"
|
||||||
|
- "--accesslog.filepath=/var/log/traefik/access.log"
|
||||||
|
- "--api.insecure=true"
|
||||||
|
- "--providers.docker=true"
|
||||||
|
- "--providers.docker.exposedbydefault=false"
|
||||||
|
- "--entrypoints.web.address=:80"
|
||||||
|
|
||||||
|
#- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
|
||||||
|
#- "--experimental.plugins.bouncer.version=v1.0.9"
|
||||||
|
- "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
|
||||||
|
volumes:
|
||||||
|
- /var/run/docker.sock:/var/run/docker.sock:ro
|
||||||
|
- logs-redis:/var/log/traefik
|
||||||
|
- ./../../:/plugins-local/src/github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin
|
||||||
|
ports:
|
||||||
|
- 80:80
|
||||||
|
- 8080:8080
|
||||||
|
depends_on:
|
||||||
|
- crowdsec
|
||||||
|
- redis
|
||||||
|
|
||||||
|
whoami1:
|
||||||
|
image: traefik/whoami
|
||||||
|
container_name: "simple-service1"
|
||||||
|
labels:
|
||||||
|
- "traefik.enable=true"
|
||||||
|
# Definition of the router
|
||||||
|
- "traefik.http.routers.router1.rule=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"
|
||||||
|
# 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.rediscacheenabled=true"
|
||||||
|
|
||||||
|
whoami2:
|
||||||
|
image: traefik/whoami
|
||||||
|
container_name: "simple-service2"
|
||||||
|
labels:
|
||||||
|
- "traefik.enable=true"
|
||||||
|
# Definition of the router
|
||||||
|
- "traefik.http.routers.router2.rule=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"
|
||||||
|
# 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.rediscacheenabled=true"
|
||||||
|
|
||||||
|
|
||||||
|
crowdsec:
|
||||||
|
image: crowdsecurity/crowdsec:v1.4.1
|
||||||
|
container_name: "crowdsec"
|
||||||
|
environment:
|
||||||
|
COLLECTIONS: crowdsecurity/traefik
|
||||||
|
CUSTOM_HOSTNAME: crowdsec
|
||||||
|
BOUNCER_KEY_TRAEFIK_DEV_1: 40796d93c2958f9e58345514e67740e5
|
||||||
|
BOUNCER_KEY_TRAEFIK_DEV_2: 44c36dac5c4140af9f06f397508e82c7
|
||||||
|
volumes:
|
||||||
|
- ./acquis.yaml:/etc/crowdsec/acquis.yaml:ro
|
||||||
|
- logs-redis:/var/log/traefik:ro
|
||||||
|
- crowdsec-db-redis:/var/lib/crowdsec/data/
|
||||||
|
- crowdsec-config-redis:/etc/crowdsec/
|
||||||
|
labels:
|
||||||
|
- "traefik.enable=false"
|
||||||
|
|
||||||
|
redis:
|
||||||
|
image: "redis:7.0.5-alpine"
|
||||||
|
container_name: "redis"
|
||||||
|
command: "redis-server --save 60 1"
|
||||||
|
volumes:
|
||||||
|
- redis-data:/data
|
||||||
|
ports:
|
||||||
|
- 6379:6379
|
||||||
|
|
||||||
|
volumes:
|
||||||
|
logs-redis:
|
||||||
|
crowdsec-db-redis:
|
||||||
|
crowdsec-config-redis:
|
||||||
|
redis-data:
|
||||||
@@ -3,3 +3,5 @@ module github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin
|
|||||||
go 1.17
|
go 1.17
|
||||||
|
|
||||||
require github.com/leprosus/golang-ttl-map v1.1.7
|
require github.com/leprosus/golang-ttl-map v1.1.7
|
||||||
|
|
||||||
|
require github.com/tehnerd/goUtils v0.0.0-20150515130609-5a2d8fb2ded8 // indirect
|
||||||
|
|||||||
@@ -1,2 +1,4 @@
|
|||||||
github.com/leprosus/golang-ttl-map v1.1.7 h1:cF4AAFDDnJTFSV+/42sKLhmMluvLdRlCGS2UaifH6UM=
|
github.com/leprosus/golang-ttl-map v1.1.7 h1:cF4AAFDDnJTFSV+/42sKLhmMluvLdRlCGS2UaifH6UM=
|
||||||
github.com/leprosus/golang-ttl-map v1.1.7/go.mod h1:4QWHJPeVBbrkhOhXdhCv9IEiyj/YzkO04/iexy4vSe0=
|
github.com/leprosus/golang-ttl-map v1.1.7/go.mod h1:4QWHJPeVBbrkhOhXdhCv9IEiyj/YzkO04/iexy4vSe0=
|
||||||
|
github.com/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
+66
-12
@@ -6,17 +6,20 @@ import (
|
|||||||
ttl_map "github.com/leprosus/golang-ttl-map"
|
ttl_map "github.com/leprosus/golang-ttl-map"
|
||||||
|
|
||||||
logger "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/logger"
|
logger "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/logger"
|
||||||
|
simpleredis "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/redis"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
cacheBannedValue = "t"
|
cacheBannedValue = "t"
|
||||||
cacheNoBannedValue = "f"
|
cacheNoBannedValue = "f"
|
||||||
)
|
)
|
||||||
|
|
||||||
var cache = ttl_map.New()
|
var cache = ttl_map.New()
|
||||||
|
var redis simpleredis.SimpleRedis
|
||||||
|
|
||||||
// Get Decision check in the cache if the IP has the banned / not banned value.
|
var redisEnabled = false
|
||||||
// Otherwise return with an error to add the IP in cache if we are on.
|
|
||||||
func GetDecision(clientIP string) (bool, error) {
|
func getDecisionLocalCache(clientIP string) (bool, error) {
|
||||||
banned, isCached := cache.Get(clientIP)
|
banned, isCached := cache.Get(clientIP)
|
||||||
bannedString, isValid := banned.(string)
|
bannedString, isValid := banned.(string)
|
||||||
if isCached && isValid && len(bannedString) > 0 {
|
if isCached && isValid && len(bannedString) > 0 {
|
||||||
@@ -25,15 +28,66 @@ func GetDecision(clientIP string) (bool, error) {
|
|||||||
return false, fmt.Errorf("no cache data")
|
return false, fmt.Errorf("no cache data")
|
||||||
}
|
}
|
||||||
|
|
||||||
func SetDecision(clientIP string, isBanned bool, duration int64) {
|
func setDecisionLocalCache(clientIP string, value string, duration int64) {
|
||||||
if isBanned {
|
cache.Set(clientIP, value, duration)
|
||||||
logger.Debug(fmt.Sprintf("%v banned", clientIP))
|
}
|
||||||
cache.Set(clientIP, cacheBannedValue, duration)
|
|
||||||
} else {
|
func deleteDecisionLocalCache(clientIP string) {
|
||||||
cache.Set(clientIP, cacheNoBannedValue, duration)
|
cache.Del(clientIP)
|
||||||
|
}
|
||||||
|
|
||||||
|
func getDecisionRedisCache(clientIP string) (bool, error) {
|
||||||
|
banned, err := redis.Get(clientIP)
|
||||||
|
bannedString := string(banned)
|
||||||
|
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) {
|
||||||
|
redis.Set(clientIP, []byte(value), duration)
|
||||||
|
}
|
||||||
|
|
||||||
|
func deleteDecisionRedisCache(clientIP string) {
|
||||||
|
redis.Del(clientIP)
|
||||||
}
|
}
|
||||||
|
|
||||||
func DeleteDecision(clientIP string) {
|
func DeleteDecision(clientIP string) {
|
||||||
cache.Del(clientIP)
|
if redisEnabled {
|
||||||
|
deleteDecisionRedisCache(clientIP)
|
||||||
|
} else {
|
||||||
|
deleteDecisionLocalCache(clientIP)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// GetDecision 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.
|
||||||
|
func GetDecision(clientIP string) (bool, error) {
|
||||||
|
if redisEnabled {
|
||||||
|
return getDecisionRedisCache(clientIP)
|
||||||
|
} else {
|
||||||
|
return getDecisionLocalCache(clientIP)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func SetDecision(clientIP string, isBanned bool, duration int64) {
|
||||||
|
var value string
|
||||||
|
if isBanned {
|
||||||
|
logger.Debug(fmt.Sprintf("%v banned", clientIP))
|
||||||
|
value = cacheBannedValue
|
||||||
|
} else {
|
||||||
|
value = cacheNoBannedValue
|
||||||
|
}
|
||||||
|
if redisEnabled {
|
||||||
|
setDecisionRedisCache(clientIP, value, duration)
|
||||||
|
} else {
|
||||||
|
setDecisionLocalCache(clientIP, value, duration)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func InitRedisClient(host string) {
|
||||||
|
redisEnabled = true
|
||||||
|
redis.Init(host)
|
||||||
|
logger.Debug("connect to redis")
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,212 @@
|
|||||||
|
package simpleredis
|
||||||
|
|
||||||
|
import (
|
||||||
|
"fmt"
|
||||||
|
"net"
|
||||||
|
"strconv"
|
||||||
|
"strings"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/tehnerd/goUtils/netutils"
|
||||||
|
)
|
||||||
|
|
||||||
|
type RedisCmd struct {
|
||||||
|
Command string
|
||||||
|
Name string
|
||||||
|
Data []byte
|
||||||
|
Duration int64
|
||||||
|
Error error
|
||||||
|
}
|
||||||
|
|
||||||
|
type SimpleRedis struct {
|
||||||
|
redisChanRead chan RedisCmd
|
||||||
|
redisChanWrite chan RedisCmd
|
||||||
|
redisHost string
|
||||||
|
redisCmd RedisCmd
|
||||||
|
}
|
||||||
|
|
||||||
|
func genRedisArray(params ...[]byte) []byte {
|
||||||
|
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, "\r\n"}, "")
|
||||||
|
return []byte(MSG)
|
||||||
|
}
|
||||||
|
|
||||||
|
func parseResponse(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
|
||||||
|
}
|
||||||
|
}
|
||||||
|
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
|
||||||
|
}
|
||||||
|
}
|
||||||
|
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
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func initContext(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 {
|
||||||
|
select {
|
||||||
|
case cmd := <-redisCmdWrite:
|
||||||
|
switch cmd.Command {
|
||||||
|
case "SET":
|
||||||
|
data := genRedisArray([]byte("SET"), []byte(cmd.Name), cmd.Data, []byte("EX"), []byte(fmt.Sprintf("%v", cmd.Duration)))
|
||||||
|
writeChan <- data
|
||||||
|
case "GET":
|
||||||
|
data := genRedisArray([]byte("GET"), []byte(cmd.Name))
|
||||||
|
writeChan <- data
|
||||||
|
case "DEL":
|
||||||
|
data := genRedisArray([]byte("DEL"), []byte(cmd.Name))
|
||||||
|
writeChan <- data
|
||||||
|
}
|
||||||
|
case response := <-readChan:
|
||||||
|
data, dataBuf, err := parseResponse(response, dataBuf, &dataLen)
|
||||||
|
if dataLen != 0 {
|
||||||
|
for data == nil {
|
||||||
|
response = <-readChan
|
||||||
|
data, dataBuf, err = parseResponse(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):
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (sr *SimpleRedis) Init(redisHost string) {
|
||||||
|
sr.redisHost = redisHost
|
||||||
|
sr.redisChanWrite = make(chan RedisCmd)
|
||||||
|
sr.redisChanRead = make(chan RedisCmd)
|
||||||
|
go initContext(sr.redisHost, sr.redisChanWrite, sr.redisChanRead)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (sr *SimpleRedis) Get(name string) ([]byte, error) {
|
||||||
|
sr.redisCmd.Command = "GET"
|
||||||
|
sr.redisCmd.Name = name
|
||||||
|
sr.redisChanWrite <- sr.redisCmd
|
||||||
|
resp := <-sr.redisChanRead
|
||||||
|
if resp.Error != nil {
|
||||||
|
return nil, resp.Error
|
||||||
|
}
|
||||||
|
return resp.Data, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (sr *SimpleRedis) Set(name string, data []byte, duration int64) error {
|
||||||
|
sr.redisCmd.Command = "SET"
|
||||||
|
sr.redisCmd.Name = name
|
||||||
|
sr.redisCmd.Data = data
|
||||||
|
sr.redisCmd.Duration = duration
|
||||||
|
sr.redisChanWrite <- sr.redisCmd
|
||||||
|
resp := <-sr.redisChanRead
|
||||||
|
if resp.Error != nil {
|
||||||
|
return resp.Error
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (sr *SimpleRedis) Del(name string) error {
|
||||||
|
sr.redisCmd.Command = "DEL"
|
||||||
|
sr.redisCmd.Name = name
|
||||||
|
sr.redisChanWrite <- sr.redisCmd
|
||||||
|
resp := <-sr.redisChanRead
|
||||||
|
if resp.Error != nil {
|
||||||
|
return resp.Error
|
||||||
|
}
|
||||||
|
return 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
|
# github.com/leprosus/golang-ttl-map v1.1.7
|
||||||
## explicit; go 1.15
|
## explicit; go 1.15
|
||||||
github.com/leprosus/golang-ttl-map
|
github.com/leprosus/golang-ttl-map
|
||||||
|
# github.com/tehnerd/goUtils v0.0.0-20150515130609-5a2d8fb2ded8
|
||||||
|
## explicit
|
||||||
|
github.com/tehnerd/goUtils/netutils
|
||||||
|
|||||||
Reference in New Issue
Block a user