Compare commits

..
16 Commits
Author SHA1 Message Date
maxlerebourg 2d0bea8eec Merge pull request #32 from maxlerebourg/28-bypassing-crowdsec-for-local-ips
28 bypassing crowdsec for local ips
2022-11-19 15:47:30 +01:00
MathieuHa bb714b5dfd Change debug to info log 2022-11-19 15:17:56 +01:00
MathieuHa bba5620187 Update all documentation and exemple 2022-11-19 13:46:15 +01:00
MathieuHa 73374ccefe Rework variables name, and logic to check IPs with strategy 2022-11-19 13:27:51 +01:00
MathieuHa 6b7f8655ac Remove Fixme In Doc 2022-11-19 12:28:50 +01:00
MathieuHa c9a3f1f8a8 Remove duration for remove csli crowdsec 2022-11-19 12:26:49 +01:00
MathieuHa 490e5e934f Merge remote-tracking branch 'origin' into 28-bypassing-crowdsec-for-local-ips 2022-11-19 12:22:56 +01:00
MathieuHa 5b0d4b533f Add logic return to bypass when IP is unfiltered, update documentation 2022-11-19 12:16:20 +01:00
Max Lerebourg 400de0d552 add redis timeout 2022-11-18 20:36:42 +01:00
MathieuHa 3742b6b540 Fix nil trusted checker 2022-11-17 09:34:16 +01:00
MathieuHa d2f0a9416d Add check to load param 2022-11-17 09:21:57 +01:00
MathieuHa c1f2131bc2 Add doc for trusted IP 2022-11-16 23:19:57 +01:00
MathieuHa 3f70e2b256 First implementation of the bypass feature for trusted IP 2022-11-16 23:17:41 +01:00
Max Lerebourg bf76ca9ef5 🍱 change logging 2022-11-15 19:50:22 +01:00
maxlerebourg 1baa5d7667 Merge pull request #31 from maxlerebourg/30-significant-slow-down-after-updating-to-v111-and-configuring-to-use-redis
 rewrite redis communication
2022-11-11 00:14:00 +01:00
Max Lerebourg 18f68de196 rewrite redis communication 2022-11-11 00:09:06 +01:00
18 changed files with 403 additions and 836 deletions
+14 -1
View File
@@ -31,6 +31,9 @@ run_behindproxy:
run_cacheredis: run_cacheredis:
docker-compose -f exemples/redis-cache/docker-compose.redis.yml up -d --remove-orphans docker-compose -f exemples/redis-cache/docker-compose.redis.yml up -d --remove-orphans
run_trustedips:
docker-compose -f exemples/trusted-ips/docker-compose.trusted.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
@@ -43,6 +46,15 @@ restart_local:
restart: restart:
docker-compose -f docker-compose.yml restart docker-compose -f docker-compose.yml restart
restart_behindproxy:
docker-compose -f exemples/behind-proxy/docker-compose.cloudflare.yml restart
restart_cacheredis:
docker-compose -f exemples/redis-cache/docker-compose.redis.yml restart
restart_trustedips:
docker-compose -f exemples/trusted-ips/docker-compose.trusted.yml restart
show_logs: show_logs:
docker-compose -f docker-compose.yml restart docker-compose -f docker-compose.yml restart
@@ -54,7 +66,8 @@ 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 exemples/redis-cache/docker-compose.redis.yml down --remove-orphans
docker-compose -f exemples/trusted-ips/docker-compose.trusted.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
+49 -4
View File
@@ -89,7 +89,10 @@ make run
- string - string
- default: "redis:6379" - default: "redis:6379"
- hostname and port for the redis service - hostname and port for the redis service
- ClientTrustedIPs
- string
- default: []
- List of client IPs we trust, they will bypass any more check to the bouncer or cache (usefull for LAN or VPN IP)
### Configuration ### Configuration
@@ -139,6 +142,8 @@ http:
forwardedHeadersTrustedIPs: forwardedHeadersTrustedIPs:
- 10.0.10.23/32 - 10.0.10.23/32
- 10.0.20.0/24 - 10.0.20.0/24
clientTrustedIPs:
- 192.168.1.0/24
forwardedHeadersCustomName: X-Custom-Header forwardedHeadersCustomName: X-Custom-Header
redisCacheEnabled: false redisCacheEnabled: false
redisCacheHost: "redis:6379" redisCacheHost: "redis:6379"
@@ -178,7 +183,7 @@ docker-compose up -d
```bash ```bash
docker-compose up -d crowdsec docker-compose up -d crowdsec
docker exec crowdsec cscli decisions add --ip 10.0.0.10 # this will be effective 4h docker exec crowdsec cscli decisions add --ip 10.0.0.10 -d 10m # this will be effective 10min
docker exec crowdsec cscli decisions remove --ip 10.0.0.10 docker exec crowdsec cscli decisions remove --ip 10.0.0.10
``` ```
@@ -224,6 +229,7 @@ You need to configure your Traefik to trust Forwarded headers by your front prox
In the exemple we use another instance of traefik with the container named cloudflare to simulate a front proxy In the exemple we use another instance of traefik with the container named cloudflare to simulate a front proxy
The "internal" Traefik instance is configured to trust the cloudflare forward headers The "internal" Traefik instance is configured to trust the cloudflare forward headers
This helps Traefik choose the right IP of the client: see https://doc.traefik.io/traefik/routing/entrypoints/#forwarded-headers
```yaml ```yaml
- "--entrypoints.web.forwardedheaders.trustedips=172.21.0.5" - "--entrypoints.web.forwardedheaders.trustedips=172.21.0.5"
``` ```
@@ -233,7 +239,7 @@ We configure the middleware to trust as well the IP:
- "traefik.http.middlewares.crowdsec1.plugin.bouncer.forwardedheaderstrustedips=172.21.0.5" - "traefik.http.middlewares.crowdsec1.plugin.bouncer.forwardedheaderstrustedips=172.21.0.5"
``` ```
To run the environnement run: To play the demo environnement run:
```bash ```bash
make run_behindproxy make run_behindproxy
``` ```
@@ -246,11 +252,50 @@ The plugin must be configured to connect to a redis instance
``` ```
Here **redis** is the hostname of a container located in the same network as Traefik and **6379** the default port of redis 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: To play the demo environnement run:
```bash ```bash
make run_cacheredis make run_cacheredis
``` ```
3. Using Trusted IP (ex: LAN OR VPN) that won't get filtered by crowdsec
You need to configure your Traefik to trust Forwarded headers by your front proxy
In the exemple we use a whoami container protected by crowdsec, and we ban or IP before allowing using TrustedIPs
If you are using another proxy in front, you need to add it's IP in the trusted IP for the forwarded headers.
This helps Traefik choose the right IP of the client: see https://doc.traefik.io/traefik/routing/entrypoints/#forwarded-headers
The "internal" Traefik instance is configured to trust the forward headers
```yaml
- "--entrypoints.web.forwardedheaders.trustedips=172.21.0.5"
```
We configure the middleware to trust as well the IP of the intermediate proxy if needed:
```yaml
- "traefik.http.middlewares.crowdsec.plugin.bouncer.forwardedheaderstrustedips=172.21.0.5"
```
Add your IP to the ban list
```bash
docker exec crowdsec cscli decisions add --ip 10.0.10.30 -d 10m
```
You should get a 403 on http://localhost/foo
> Replace *10.0.10.30* by your IP
Add the IPs that will not be filtered by the plugin
```yaml
- "traefik.http.middlewares.crowdsec.plugin.bouncer.clientTrustedips=10.0.10.30/32"
```
> Replace *10.0.10.30/32* by your IP or IP range, so it's not getting checked against ban cache of crowdsec
You should get a 200 on http://localhost/foo even if you are on the ban cache
To play the demo environnement run:
```bash
make run_trustedips
```
### About ### About
Me and [mathieuHa](https://github.com/mathieuHa) have been using traefik since 2020 at [Primadviz](https://primadviz.com). Me and [mathieuHa](https://github.com/mathieuHa) have been using traefik since 2020 at [Primadviz](https://primadviz.com).
+48 -15
View File
@@ -15,6 +15,7 @@ import (
cache "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/cache" cache "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/cache"
ip "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/ip" ip "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/ip"
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 (
@@ -40,10 +41,11 @@ type Config struct {
CrowdsecLapiScheme string `json:"crowdsecLapiScheme,omitempty"` CrowdsecLapiScheme string `json:"crowdsecLapiScheme,omitempty"`
CrowdsecLapiHost string `json:"crowdsecLapiHost,omitempty"` CrowdsecLapiHost string `json:"crowdsecLapiHost,omitempty"`
CrowdsecLapiKey string `json:"crowdsecLapiKey,omitempty"` CrowdsecLapiKey string `json:"crowdsecLapiKey,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"`
ForwardedHeadersCustomName string `json:"forwardedheaderscustomheader,omitempty"`
ForwardedHeadersTrustedIPs []string `json:"forwardedHeadersTrustedIps,omitempty"` ForwardedHeadersTrustedIPs []string `json:"forwardedHeadersTrustedIps,omitempty"`
ClientTrustedIPs []string `json:"clientTrustedIps,omitempty"`
RedisCacheEnabled bool `json:"redisCacheEnabled,omitempty"` RedisCacheEnabled bool `json:"redisCacheEnabled,omitempty"`
RedisCacheHost string `json:"redisCacheHost,omitempty"` RedisCacheHost string `json:"redisCacheHost,omitempty"`
} }
@@ -60,6 +62,7 @@ func CreateConfig() *Config {
UpdateIntervalSeconds: 60, UpdateIntervalSeconds: 60,
DefaultDecisionSeconds: 60, DefaultDecisionSeconds: 60,
ForwardedHeadersTrustedIPs: []string{}, ForwardedHeadersTrustedIPs: []string{},
ClientTrustedIPs: []string{},
ForwardedHeadersCustomName: "X-Forwarded-For", ForwardedHeadersCustomName: "X-Forwarded-For",
RedisCacheEnabled: false, RedisCacheEnabled: false,
RedisCacheHost: "redis:6379", RedisCacheHost: "redis:6379",
@@ -80,7 +83,8 @@ type Bouncer struct {
updateInterval int64 updateInterval int64
defaultDecisionTimeout int64 defaultDecisionTimeout int64
customHeader string customHeader string
poolStrategy *ip.PoolStrategy serverPoolStrategy *ip.PoolStrategy
clientPoolStrategy *ip.PoolStrategy
client *http.Client client *http.Client
} }
@@ -89,11 +93,12 @@ func New(ctx context.Context, next http.Handler, config *Config, name string) (h
logger.Init(config.LogLevel) logger.Init(config.LogLevel)
err := validateParams(config) err := validateParams(config)
if err != nil { if err != nil {
logger.Info(fmt.Sprintf("%w", err)) logger.Info(err.Error())
return nil, err return nil, err
} }
checker, _ := ip.NewChecker(config.ForwardedHeadersTrustedIPs) serverChecker, _ := ip.NewChecker(config.ForwardedHeadersTrustedIPs)
clientChecker, _ := ip.NewChecker(config.ClientTrustedIPs)
bouncer := &Bouncer{ bouncer := &Bouncer{
next: next, next: next,
@@ -108,8 +113,11 @@ func New(ctx context.Context, next http.Handler, config *Config, name string) (h
updateInterval: config.UpdateIntervalSeconds, updateInterval: config.UpdateIntervalSeconds,
customHeader: config.ForwardedHeadersCustomName, customHeader: config.ForwardedHeadersCustomName,
defaultDecisionTimeout: config.DefaultDecisionSeconds, defaultDecisionTimeout: config.DefaultDecisionSeconds,
poolStrategy: &ip.PoolStrategy{ serverPoolStrategy: &ip.PoolStrategy{
Checker: checker, Checker: serverChecker,
},
clientPoolStrategy: &ip.PoolStrategy{
Checker: clientChecker,
}, },
client: &http.Client{ client: &http.Client{
Transport: &http.Transport{ Transport: &http.Transport{
@@ -138,19 +146,36 @@ func (bouncer *Bouncer) ServeHTTP(rw http.ResponseWriter, req *http.Request) {
bouncer.next.ServeHTTP(rw, req) bouncer.next.ServeHTTP(rw, req)
return return
} }
// Here we check for the trusted IPs in the customHeader
remoteHost, err := ip.GetRemoteIP(req, bouncer.poolStrategy, bouncer.customHeader) remoteHost, err := ip.GetRemoteIP(req, bouncer.serverPoolStrategy, bouncer.customHeader)
if err != nil { if err != nil {
logger.Info(fmt.Sprintf("%w", err)) logger.Info(err.Error())
bouncer.next.ServeHTTP(rw, req) bouncer.next.ServeHTTP(rw, req)
return return
} }
logger.Debug(fmt.Sprintf("ServeHTTP ip:%v", remoteHost))
trusted, err := bouncer.clientPoolStrategy.Checker.Contains(remoteHost)
if err != nil {
logger.Info(err.Error())
return
}
// if our IP is in the trusted list we bypass the next checks
logger.Debug(fmt.Sprintf("ServeHTTP ip:%v isTrusted:%v", remoteHost, trusted))
if trusted {
bouncer.next.ServeHTTP(rw, req)
return
}
healthy := crowdsecStreamHealthy
if bouncer.crowdsecMode != noneMode { if bouncer.crowdsecMode != noneMode {
isBanned, err := cache.GetDecision(remoteHost) isBanned, err := cache.GetDecision(remoteHost)
if err == nil { if err != nil {
logger.Debug(fmt.Sprintf("ServeHTTP cacheHit isBanned:%v", isBanned)) logger.Debug(err.Error())
if err.Error() == simpleredis.RedisUnreachable {
healthy = false
}
} else {
logger.Debug(fmt.Sprintf("ServeHTTP cache:hit isBanned:%v", isBanned))
if isBanned { if isBanned {
rw.WriteHeader(http.StatusForbidden) rw.WriteHeader(http.StatusForbidden)
} else { } else {
@@ -162,7 +187,7 @@ func (bouncer *Bouncer) ServeHTTP(rw http.ResponseWriter, req *http.Request) {
// Right here if we cannot join the stream we forbid the request to go on. // Right here if we cannot join the stream we forbid the request to go on.
if bouncer.crowdsecMode == streamMode { if bouncer.crowdsecMode == streamMode {
if crowdsecStreamHealthy { if healthy {
bouncer.next.ServeHTTP(rw, req) bouncer.next.ServeHTTP(rw, req)
} else { } else {
rw.WriteHeader(http.StatusForbidden) rw.WriteHeader(http.StatusForbidden)
@@ -229,7 +254,7 @@ func handleNoStreamCache(bouncer *Bouncer, rw http.ResponseWriter, req *http.Req
} }
body, err := crowdsecQuery(bouncer, routeURL.String()) body, err := crowdsecQuery(bouncer, routeURL.String())
if err != nil { if err != nil {
logger.Info(fmt.Sprintf("%w", err)) logger.Info(err.Error())
rw.WriteHeader(http.StatusForbidden) rw.WriteHeader(http.StatusForbidden)
return return
} }
@@ -286,7 +311,7 @@ func handleStreamCache(bouncer *Bouncer) {
} }
body, err := crowdsecQuery(bouncer, streamRouteURL.String()) body, err := crowdsecQuery(bouncer, streamRouteURL.String())
if err != nil { if err != nil {
logger.Info(fmt.Sprintf("%w", err)) logger.Info(err.Error())
crowdsecStreamHealthy = false crowdsecStreamHealthy = false
return return
} }
@@ -376,6 +401,14 @@ func validateParams(config *Config) error {
} else { } else {
logger.Debug("No IP provided for ForwardedHeadersTrustedIPs") logger.Debug("No IP provided for ForwardedHeadersTrustedIPs")
} }
if len(config.ClientTrustedIPs) > 0 {
_, err = ip.NewChecker(config.ClientTrustedIPs)
if err != nil {
return fmt.Errorf("TrustedIPs must be a list of IP/CIDR :%w", err)
}
} else {
logger.Debug("No IP provided for TrustedIPs")
}
return nil return nil
} }
+20 -16
View File
@@ -4,8 +4,9 @@ services:
traefik: traefik:
image: "traefik:v2.9.4" image: "traefik:v2.9.4"
container_name: "traefik" container_name: "traefik"
restart: unless-stopped
command: command:
# - "--log.level=DEBUG" - "--log.level=DEBUG"
- "--accesslog" - "--accesslog"
- "--accesslog.filepath=/var/log/traefik/access.log" - "--accesslog.filepath=/var/log/traefik/access.log"
- "--api.insecure=true" - "--api.insecure=true"
@@ -24,33 +25,36 @@ services:
depends_on: depends_on:
- crowdsec - crowdsec
whoami1: whoami-foo:
image: traefik/whoami image: traefik/whoami
container_name: "simple-service1" container_name: "simple-service-foo"
restart: unless-stopped
labels: labels:
- "traefik.enable=true" - "traefik.enable=true"
- "traefik.http.routers.router1.rule=Host(`localhost`) && Path(`/foo`)" - "traefik.http.routers.router-foo.rule=Path(`/foo`)"
- "traefik.http.routers.router1.entrypoints=web" - "traefik.http.routers.router-foo.entrypoints=web"
- "traefik.http.routers.router1.middlewares=crowdsec1@docker" - "traefik.http.routers.router-foo.middlewares=crowdsec-foo@docker"
- "traefik.http.services.service1.loadbalancer.server.port=80" - "traefik.http.services.service-foo.loadbalancer.server.port=80"
- "traefik.http.middlewares.crowdsec1.plugin.bouncer.enabled=true" - "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.enabled=true"
- "traefik.http.middlewares.crowdsec1.plugin.bouncer.crowdseclapikey=40796d93c2958f9e58345514e67740e5" - "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.crowdseclapikey=40796d93c2958f9e58345514e67740e5"
whoami2: whoami2:
image: traefik/whoami image: traefik/whoami
container_name: "simple-service2" container_name: "simple-service-bar"
restart: unless-stopped
labels: labels:
- "traefik.enable=true" - "traefik.enable=true"
- "traefik.http.routers.router2.rule=Host(`localhost`) && Path(`/bar`)" - "traefik.http.routers.router-bar.rule=Path(`/bar`)"
- "traefik.http.routers.router2.entrypoints=web" - "traefik.http.routers.router-bar.entrypoints=web"
- "traefik.http.routers.router2.middlewares=crowdsec2@docker" - "traefik.http.routers.router-bar.middlewares=crowdsec-bar@docker"
- "traefik.http.services.service2.loadbalancer.server.port=80" - "traefik.http.services.service-bar.loadbalancer.server.port=80"
- "traefik.http.middlewares.crowdsec2.plugin.bouncer.enabled=true" - "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.enabled=true"
- "traefik.http.middlewares.crowdsec2.plugin.bouncer.crowdseclapikey=44c36dac5c4140af9f06f397508e82c7" - "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.crowdseclapikey=44c36dac5c4140af9f06f397508e82c7"
crowdsec: crowdsec:
image: crowdsecurity/crowdsec:v1.4.1 image: crowdsecurity/crowdsec:v1.4.1
container_name: "crowdsec" container_name: "crowdsec"
restart: unless-stopped
environment: environment:
COLLECTIONS: crowdsecurity/traefik COLLECTIONS: crowdsecurity/traefik
CUSTOM_HOSTNAME: crowdsec CUSTOM_HOSTNAME: crowdsec
+21 -17
View File
@@ -4,6 +4,7 @@ services:
traefik: traefik:
image: "traefik:v2.9.4" image: "traefik:v2.9.4"
container_name: "traefik" container_name: "traefik"
restart: unless-stopped
command: command:
- "--accesslog" - "--accesslog"
- "--accesslog.filepath=/var/log/traefik/access.log" - "--accesslog.filepath=/var/log/traefik/access.log"
@@ -13,7 +14,7 @@ services:
- "--entrypoints.web.address=:80" - "--entrypoints.web.address=:80"
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin" - "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
- "--experimental.plugins.bouncer.version=v1.1.0" - "--experimental.plugins.bouncer.version=v1.1.2"
volumes: volumes:
- "/var/run/docker.sock:/var/run/docker.sock:ro" - "/var/run/docker.sock:/var/run/docker.sock:ro"
- "logs:/var/log/traefik" - "logs:/var/log/traefik"
@@ -25,43 +26,46 @@ services:
whoami1: whoami1:
image: traefik/whoami image: traefik/whoami
container_name: "simple-service-1" container_name: "simple-service-foo"
restart: unless-stopped
labels: labels:
- "traefik.enable=true" - "traefik.enable=true"
# Definition of the router # Definition of the router
- "traefik.http.routers.router1.rule=Path(`/foo`)" - "traefik.http.routers.router-foo.rule=Path(`/foo`)"
- "traefik.http.routers.router1.entrypoints=web" - "traefik.http.routers.router-foo.entrypoints=web"
- "traefik.http.routers.router1.middlewares=crowdsec2@docker" - "traefik.http.routers.router-foo.middlewares=crowdsec-foo@docker"
# Definition of the service # Definition of the service
- "traefik.http.services.service1.loadbalancer.server.port=80" - "traefik.http.services.service-foo.loadbalancer.server.port=80"
# Definition of the middleware # Definition of the middleware
- "traefik.http.middlewares.crowdsec1.plugin.bouncer.enabled=true" - "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.enabled=true"
# crowdseclapikey must be unique to the middleware attached to the service # crowdseclapikey must be unique to the middleware attached to the service
- "traefik.http.middlewares.crowdsec1.plugin.bouncer.crowdseclapikey=FIXME-LAPI-KEY-1" - "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.crowdseclapikey=FIXME-LAPI-KEY-1"
# forwardedheaderstrustedips should be the IP of the proxy that is in front of traefik (if any) # forwardedheaderstrustedips should be the IP of the proxy that is in front of traefik (if any)
- "traefik.http.middlewares.crowdsec1.plugin.bouncer.forwardedheaderstrustedips=172.21.0.5" - "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.forwardedheaderstrustedips=172.21.0.5"
whoami2: whoami2:
image: traefik/whoami image: traefik/whoami
container_name: "simple-service-2" container_name: "simple-service-bar"
restart: unless-stopped
labels: labels:
- "traefik.enable=true" - "traefik.enable=true"
# Definition of the router # Definition of the router
- "traefik.http.routers.router2.rule=Path(`/bar`)" - "traefik.http.routers.router-bar.rule=Path(`/bar`)"
- "traefik.http.routers.router2.entrypoints=web" - "traefik.http.routers.router-bar.entrypoints=web"
- "traefik.http.routers.router2.middlewares=crowdsec2@docker" - "traefik.http.routers.router-bar.middlewares=crowdsec-bar@docker"
# Definition of the service # Definition of the service
- "traefik.http.services.service2.loadbalancer.server.port=80" - "traefik.http.services.service-bar.loadbalancer.server.port=80"
# Definitin of the middleware # Definitin of the middleware
- "traefik.http.middlewares.crowdsec2.plugin.bouncer.enabled=true" - "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.enabled=true"
# crowdseclapikey must be unique to the middleware attached to the service # crowdseclapikey must be unique to the middleware attached to the service
- "traefik.http.middlewares.crowdsec2.plugin.bouncer.crowdseclapikey=FIXME-LAPI-KEY-2" - "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.crowdseclapikey=FIXME-LAPI-KEY-2"
# forwardedheaderstrustedips should be the IP of the proxy that is in front of traefik (if any) # forwardedheaderstrustedips should be the IP of the proxy that is in front of traefik (if any)
- "traefik.http.middlewares.crowdsec2.plugin.bouncer.forwardedheaderstrustedips=172.21.0.5" - "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.forwardedheaderstrustedips=172.21.0.5"
crowdsec: crowdsec:
image: crowdsecurity/crowdsec:v1.4.1 image: crowdsecurity/crowdsec:v1.4.1
container_name: "crowdsec" container_name: "crowdsec"
restart: unless-stopped
environment: environment:
COLLECTIONS: crowdsecurity/traefik COLLECTIONS: crowdsecurity/traefik
CUSTOM_HOSTNAME: crowdsec CUSTOM_HOSTNAME: crowdsec
@@ -2,8 +2,9 @@ version: "3.8"
services: services:
cloudflare: cloudflare:
image: "traefik:v2.9.1" image: "traefik:v2.9.4"
container_name: "cloudflare" container_name: "cloudflare"
restart: unless-stopped
command: command:
# - "--log.level=DEBUG" # - "--log.level=DEBUG"
- "--accesslog" - "--accesslog"
@@ -21,8 +22,9 @@ services:
- 8080:8080 - 8080:8080
traefik: traefik:
image: "traefik:v2.9.1" image: "traefik:v2.9.4"
container_name: "traefik" container_name: "traefik"
restart: unless-stopped
command: command:
# - "--log.level=DEBUG" # - "--log.level=DEBUG"
- "--accesslog" - "--accesslog"
@@ -34,10 +36,10 @@ services:
- "--entrypoints.web.forwardedheaders.trustedips=172.21.0.5" - "--entrypoints.web.forwardedheaders.trustedips=172.21.0.5"
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin" - "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
- "--experimental.plugins.bouncer.version=v1.1.0" - "--experimental.plugins.bouncer.version=v1.1.3"
volumes: volumes:
- /var/run/docker.sock:/var/run/docker.sock:ro - /var/run/docker.sock:/var/run/docker.sock:ro
- logs-dev:/var/log/traefik - logs-traefik:/var/log/traefik
ports: ports:
- 90:80 - 90:80
- 9080:8080 - 9080:8080
@@ -46,26 +48,49 @@ services:
whoami1: whoami1:
image: traefik/whoami image: traefik/whoami
container_name: "simple-service1" container_name: "simple-service-foo"
restart: unless-stopped
labels: labels:
- "traefik.enable=true" - "traefik.enable=true"
# Definition of the router # Definition of the router
- "traefik.http.routers.router1.rule=Path(`/foo`)" - "traefik.http.routers.router-foo.rule=Path(`/foo`)"
- "traefik.http.routers.router1.entrypoints=web" - "traefik.http.routers.router-foo.entrypoints=web"
- "traefik.http.routers.router1.middlewares=crowdsec1@docker" - "traefik.http.routers.router-foo.middlewares=crowdsec-foo@docker"
# Definition of the service # Definition of the service
- "traefik.http.services.service1.loadbalancer.server.port=80" - "traefik.http.services.service-foo.loadbalancer.server.port=80"
# Definitin of the middleware # Definitin of the middleware
- "traefik.http.middlewares.crowdsec1.plugin.bouncer.enabled=true" - "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.enabled=true"
# crowdseclapikey must be uniq to the middleware attached to the service # crowdseclapikey must be uniq to the middleware attached to the service
- "traefik.http.middlewares.crowdsec1.plugin.bouncer.crowdseclapikey=40796d93c2958f9e58345514e67740e5" - "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.crowdseclapikey=40796d93c2958f9e58345514e67740e5"
- "traefik.http.middlewares.crowdsec1.plugin.bouncer.crowdsecmode=live" - "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.crowdsecmode=live"
- "traefik.http.middlewares.crowdsec1.plugin.bouncer.forwardedheaderstrustedips=172.21.0.5" - "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.forwardedheaderstrustedips=172.21.0.5"
- "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.loglevel=DEBUG"
whoami2:
image: traefik/whoami
container_name: "simple-service-bar"
restart: unless-stopped
labels:
- "traefik.enable=true"
# Definition of the router
- "traefik.http.routers.router-bar.rule=Path(`/bar`)"
- "traefik.http.routers.router-bar.entrypoints=web"
- "traefik.http.routers.router-bar.middlewares=crowdsec-bar@docker"
# Definition of the service
- "traefik.http.services.service-bar.loadbalancer.server.port=80"
# Definitin of the middleware
- "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.enabled=true"
# crowdseclapikey must be uniq to the middleware attached to the service
- "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.crowdseclapikey=44c36dac5c4140af9f06f397508e82c7"
- "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.crowdsecmode=live"
- "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.forwardedheaderstrustedips=172.21.0.5"
- "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.loglevel=DEBUG"
crowdsec: crowdsec:
image: crowdsecurity/crowdsec:v1.4.1 image: crowdsecurity/crowdsec:v1.4.1
container_name: "crowdsec" container_name: "crowdsec"
restart: unless-stopped
environment: environment:
COLLECTIONS: crowdsecurity/traefik COLLECTIONS: crowdsecurity/traefik
CUSTOM_HOSTNAME: crowdsec CUSTOM_HOSTNAME: crowdsec
@@ -73,14 +98,14 @@ services:
BOUNCER_KEY_TRAEFIK_DEV_2: 44c36dac5c4140af9f06f397508e82c7 BOUNCER_KEY_TRAEFIK_DEV_2: 44c36dac5c4140af9f06f397508e82c7
volumes: volumes:
- ./acquis.yaml:/etc/crowdsec/acquis.yaml:ro - ./acquis.yaml:/etc/crowdsec/acquis.yaml:ro
- logs-dev:/var/log/traefik:ro - logs-cloudflare:/var/log/traefik:ro
- crowdsec-db-dev:/var/lib/crowdsec/data/ - crowdsec-db-cloudflare:/var/lib/crowdsec/data/
- crowdsec-config-dev:/etc/crowdsec/ - crowdsec-config-cloudflare:/etc/crowdsec/
labels: labels:
- "traefik.enable=false" - "traefik.enable=false"
volumes: volumes:
logs-dev: logs-traefik:
logs-cloudflare: logs-cloudflare:
crowdsec-db-dev: crowdsec-db-cloudflare:
crowdsec-config-dev: crowdsec-config-cloudflare:
+29 -22
View File
@@ -4,6 +4,7 @@ services:
traefik: traefik:
image: "traefik:v2.9.4" image: "traefik:v2.9.4"
container_name: "traefik" container_name: "traefik"
restart: unless-stopped
command: command:
# - "--log.level=DEBUG" # - "--log.level=DEBUG"
- "--accesslog" - "--accesslog"
@@ -13,9 +14,9 @@ services:
- "--providers.docker.exposedbydefault=false" - "--providers.docker.exposedbydefault=false"
- "--entrypoints.web.address=:80" - "--entrypoints.web.address=:80"
#- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin" - "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
#- "--experimental.plugins.bouncer.version=v1.0.9" - "--experimental.plugins.bouncer.version=v1.1.3"
- "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin" # - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
volumes: volumes:
- /var/run/docker.sock:/var/run/docker.sock:ro - /var/run/docker.sock:/var/run/docker.sock:ro
- logs-redis:/var/log/traefik - logs-redis:/var/log/traefik
@@ -27,44 +28,49 @@ services:
- crowdsec - crowdsec
- redis - redis
whoami1: whoami-foo:
image: traefik/whoami image: traefik/whoami
container_name: "simple-service1" container_name: "simple-service-foo"
restart: unless-stopped
labels: labels:
- "traefik.enable=true" - "traefik.enable=true"
# Definition of the router # Definition of the router
- "traefik.http.routers.router1.rule=Path(`/foo`)" - "traefik.http.routers.router-foo.rule=Path(`/foo`)"
- "traefik.http.routers.router1.entrypoints=web" - "traefik.http.routers.router-foo.entrypoints=web"
- "traefik.http.routers.router1.middlewares=crowdsec1@docker" - "traefik.http.routers.router-foo.middlewares=crowdsec-foo@docker"
# Definition of the service # Definition of the service
- "traefik.http.services.service1.loadbalancer.server.port=80" - "traefik.http.services.service-foo.loadbalancer.server.port=80"
# Definition of the middleware # Definition of the middleware
- "traefik.http.middlewares.crowdsec1.plugin.bouncer.enabled=true" - "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.enabled=true"
# crowdseclapikey must be uniq to the middleware attached to the service # crowdseclapikey must be uniq to the middleware attached to the service
- "traefik.http.middlewares.crowdsec1.plugin.bouncer.crowdseclapikey=40796d93c2958f9e58345514e67740e5" - "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.crowdseclapikey=40796d93c2958f9e58345514e67740e5"
- "traefik.http.middlewares.crowdsec1.plugin.bouncer.rediscacheenabled=true" - "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.rediscacheenabled=true"
- "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.loglevel=DEBUG"
whoami2: whoami-bar:
image: traefik/whoami image: traefik/whoami
container_name: "simple-service2" container_name: "simple-service-bar"
restart: unless-stopped
labels: labels:
- "traefik.enable=true" - "traefik.enable=true"
# Definition of the router # Definition of the router
- "traefik.http.routers.router2.rule=Path(`/bar`)" - "traefik.http.routers.router-bar.rule=Path(`/bar`)"
- "traefik.http.routers.router2.entrypoints=web" - "traefik.http.routers.router-bar.entrypoints=web"
- "traefik.http.routers.router2.middlewares=crowdsec1@docker" - "traefik.http.routers.router-bar.middlewares=crowdsec-bar@docker"
# Definition of the service # Definition of the service
- "traefik.http.services.service2.loadbalancer.server.port=80" - "traefik.http.services.service-bar.loadbalancer.server.port=80"
# Definition of the middleware # Definition of the middleware
- "traefik.http.middlewares.crowdsec2.plugin.bouncer.enabled=true" - "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.enabled=true"
# crowdseclapikey must be uniq to the middleware attached to the service # crowdseclapikey must be uniq to the middleware attached to the service
- "traefik.http.middlewares.crowdsec2.plugin.bouncer.crowdseclapikey=44c36dac5c4140af9f06f397508e82c7" - "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.crowdseclapikey=44c36dac5c4140af9f06f397508e82c7"
- "traefik.http.middlewares.crowdsec2.plugin.bouncer.rediscacheenabled=true" - "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.rediscacheenabled=true"
- "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.loglevel=DEBUG"
crowdsec: crowdsec:
image: crowdsecurity/crowdsec:v1.4.1 image: crowdsecurity/crowdsec:v1.4.1
container_name: "crowdsec" container_name: "crowdsec"
restart: unless-stopped
environment: environment:
COLLECTIONS: crowdsecurity/traefik COLLECTIONS: crowdsecurity/traefik
CUSTOM_HOSTNAME: crowdsec CUSTOM_HOSTNAME: crowdsec
@@ -77,10 +83,11 @@ services:
- crowdsec-config-redis:/etc/crowdsec/ - crowdsec-config-redis:/etc/crowdsec/
labels: labels:
- "traefik.enable=false" - "traefik.enable=false"
redis: redis:
image: "redis:7.0.5-alpine" image: "redis:7.0.5-alpine"
container_name: "redis" container_name: "redis"
restart: unless-stopped
command: "redis-server --save 60 1" command: "redis-server --save 60 1"
volumes: volumes:
- redis-data:/data - redis-data:/data
+4
View File
@@ -0,0 +1,4 @@
filenames:
- /var/log/traefik/access.log
labels:
type: traefik
@@ -0,0 +1,91 @@
version: "3.8"
services:
traefik:
image: "traefik:v2.9.4"
container_name: "traefik"
restart: unless-stopped
command:
# - "--log.level=DEBUG"
- "--accesslog"
- "--accesslog.filepath=/var/log/traefik/access.log"
- "--api.insecure=true"
- "--providers.docker=true"
- "--providers.docker.exposedbydefault=false"
- "--entrypoints.web.address=:80"
- "--experimental.plugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
- "--experimental.plugins.bouncer.version=v1.1.3"
# - "--experimental.localplugins.bouncer.modulename=github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin"
volumes:
- /var/run/docker.sock:/var/run/docker.sock:ro
- logs-trustedips:/var/log/traefik
- ./../../:/plugins-local/src/github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin
ports:
- 80:80
- 8080:8080
depends_on:
- crowdsec
whoami1:
image: traefik/whoami
container_name: "simple-service-foo"
restart: unless-stopped
labels:
- "traefik.enable=true"
# Definition of the router
- "traefik.http.routers.router-foo.rule=Path(`/foo`)"
- "traefik.http.routers.router-foo.entrypoints=web"
- "traefik.http.routers.router-foo.middlewares=crowdsec-foo@docker"
# Definition of the service
- "traefik.http.services.service-foo.loadbalancer.server.port=80"
# Definition of the middleware
- "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.enabled=true"
# crowdseclapikey must be uniq to the middleware attached to the service
- "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.crowdseclapikey=40796d93c2958f9e58345514e67740e5"
# Replace 10.0.10.30/32 by your IP range which is "trusted"
- "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.clienttrustedips=10.0.10.30/32"
- "traefik.http.middlewares.crowdsec-foo.plugin.bouncer.loglevel=DEBUG"
whoami2:
image: traefik/whoami
container_name: "simple-service-bar"
restart: unless-stopped
labels:
- "traefik.enable=true"
# Definition of the router
- "traefik.http.routers.router-bar.rule=Path(`/bar`)"
- "traefik.http.routers.router-bar.entrypoints=web"
- "traefik.http.routers.router-bar.middlewares=crowdsec-bar@docker"
# Definition of the service
- "traefik.http.services.service-bar.loadbalancer.server.port=80"
# Definition of the middleware
- "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.enabled=true"
# crowdseclapikey must be uniq to the middleware attached to the service
- "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.crowdseclapikey=44c36dac5c4140af9f06f397508e82c7"
# Replace 10.0.10.30/32 by your IP range which is "trusted"
- "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.clienttrustedips=10.0.10.30/32"
- "traefik.http.middlewares.crowdsec-bar.plugin.bouncer.loglevel=DEBUG"
crowdsec:
image: crowdsecurity/crowdsec:v1.4.1
container_name: "crowdsec"
restart: unless-stopped
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-trustedips:/var/log/traefik:ro
- crowdsec-db-trustedips:/var/lib/crowdsec/data/
- crowdsec-config-trustedips:/etc/crowdsec/
labels:
- "traefik.enable=false"
volumes:
logs-trustedips:
crowdsec-db-trustedips:
crowdsec-config-trustedips:
+1 -3
View File
@@ -1,7 +1,5 @@
module github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin module github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin
go 1.17 go 1.18
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
-2
View File
@@ -1,4 +1,2 @@
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=
+8 -3
View File
@@ -19,13 +19,15 @@ var redis simpleredis.SimpleRedis
var redisEnabled = false var redisEnabled = false
// CLASSIC
func getDecisionLocalCache(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 {
return bannedString == cacheBannedValue, nil return bannedString == cacheBannedValue, nil
} }
return false, fmt.Errorf("no cache data") return false, fmt.Errorf("cache:miss")
} }
func setDecisionLocalCache(clientIP string, value string, duration int64) { func setDecisionLocalCache(clientIP string, value string, duration int64) {
@@ -36,13 +38,15 @@ func deleteDecisionLocalCache(clientIP string) {
cache.Del(clientIP) cache.Del(clientIP)
} }
// REDIS
func getDecisionRedisCache(clientIP string) (bool, error) { func getDecisionRedisCache(clientIP string) (bool, error) {
banned, err := redis.Get(clientIP) banned, err := redis.Get(clientIP)
bannedString := string(banned) bannedString := string(banned)
if err == nil && len(bannedString) > 0 { if err == nil && len(bannedString) > 0 {
return bannedString == cacheBannedValue, nil return bannedString == cacheBannedValue, nil
} }
return false, fmt.Errorf("no cache data") return false, err
} }
func setDecisionRedisCache(clientIP string, value string, duration int64) { func setDecisionRedisCache(clientIP string, value string, duration int64) {
@@ -53,6 +57,7 @@ func deleteDecisionRedisCache(clientIP string) {
redis.Del(clientIP) redis.Del(clientIP)
} }
// DeleteDecision delete decision in cache
func DeleteDecision(clientIP string) { func DeleteDecision(clientIP string) {
if redisEnabled { if redisEnabled {
deleteDecisionRedisCache(clientIP) deleteDecisionRedisCache(clientIP)
@@ -89,5 +94,5 @@ func SetDecision(clientIP string, isBanned bool, duration int64) {
func InitRedisClient(host string) { func InitRedisClient(host string) {
redisEnabled = true redisEnabled = true
redis.Init(host) redis.Init(host)
logger.Debug("connect to redis") logger.Debug("Redis initialized")
} }
+4 -4
View File
@@ -6,6 +6,8 @@ import (
"net" "net"
"net/http" "net/http"
"strings" "strings"
"github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/logger"
) )
// CHECKER // CHECKER
@@ -18,15 +20,12 @@ type Checker struct {
// NewChecker builds a new Checker given a list of CIDR-Strings to trusted IPs. // NewChecker builds a new Checker given a list of CIDR-Strings to trusted IPs.
func NewChecker(trustedIPs []string) (*Checker, error) { func NewChecker(trustedIPs []string) (*Checker, error) {
if len(trustedIPs) == 0 {
return nil, errors.New("no trusted IPs provided")
}
checker := &Checker{} checker := &Checker{}
for _, ipMask := range trustedIPs { for _, ipMask := range trustedIPs {
if ipAddr := net.ParseIP(ipMask); ipAddr != nil { if ipAddr := net.ParseIP(ipMask); ipAddr != nil {
checker.authorizedIPs = append(checker.authorizedIPs, &ipAddr) checker.authorizedIPs = append(checker.authorizedIPs, &ipAddr)
logger.Debug(fmt.Sprintf("IP %v is trusted", ipAddr))
continue continue
} }
@@ -35,6 +34,7 @@ func NewChecker(trustedIPs []string) (*Checker, error) {
return nil, fmt.Errorf("parsing CIDR trusted IPs %s: %w", ipAddr, err) return nil, fmt.Errorf("parsing CIDR trusted IPs %s: %w", ipAddr, err)
} }
checker.authorizedIPsNet = append(checker.authorizedIPsNet, ipAddr) checker.authorizedIPsNet = append(checker.authorizedIPsNet, ipAddr)
logger.Debug(fmt.Sprintf("IP network %v is trusted", ipAddr))
} }
return checker, nil return checker, nil
+7 -7
View File
@@ -14,13 +14,13 @@ var (
// Init Set Default log level to info in case log level to defined // Init Set Default log level to info in case log level to defined
func Init(logLevel string) { func Init(logLevel string) {
switch logLevel { switch logLevel {
case "INFO": case "INFO":
loggerInfo.SetOutput(os.Stdout) loggerInfo.SetOutput(os.Stdout)
case "DEBUG": case "DEBUG":
loggerInfo.SetOutput(os.Stdout) loggerInfo.SetOutput(os.Stdout)
loggerDebug.SetOutput(os.Stdout) loggerDebug.SetOutput(os.Stdout)
default: default:
loggerInfo.SetOutput(os.Stdout) loggerInfo.SetOutput(os.Stdout)
} }
} }
+63 -153
View File
@@ -1,13 +1,20 @@
package simpleredis package simpleredis
import ( import (
"bufio"
"fmt" "fmt"
"net" "net"
"strconv" "net/textproto"
"strings" "strings"
"time" "time"
"github.com/tehnerd/goUtils/netutils" logger "github.com/maxlerebourg/crowdsec-bouncer-traefik-plugin/pkg/logger"
)
const (
RedisUnreachable = "redis:unreachable"
RedisMiss = "redis:miss"
RedisTimeout = "redis:timeout"
) )
type RedisCmd struct { type RedisCmd struct {
@@ -19,10 +26,7 @@ type RedisCmd struct {
} }
type SimpleRedis struct { type SimpleRedis struct {
redisChanRead chan RedisCmd redisHost string
redisChanWrite chan RedisCmd
redisHost string
redisCmd RedisCmd
} }
func genRedisArray(params ...[]byte) []byte { func genRedisArray(params ...[]byte) []byte {
@@ -35,135 +39,45 @@ func genRedisArray(params ...[]byte) []byte {
return []byte(MSG) return []byte(MSG)
} }
func parseResponse(response []byte, dataBuf []byte, Len *int) ([]byte, []byte, error) { func askRedis(hostnamePort string, cmd RedisCmd, channel chan RedisCmd) {
dataBuf = append(dataBuf, response...) dialer := net.Dialer{Timeout: 2 * time.Second}
lenCRLF := 2 conn, err := dialer.Dial("tcp", hostnamePort)
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 { if err != nil {
panic("cant resolve remote redis address") channel <- RedisCmd{Error: fmt.Errorf(RedisUnreachable)}
return
} }
var ladr *net.TCPAddr defer conn.Close()
msgBuf := make([]byte, 65000)
initMsg := []byte("*1\r\n$4\r\nPING\r\n") writer := textproto.NewWriter(bufio.NewWriter(conn))
writeChan := make(chan []byte) reader := textproto.NewReader(bufio.NewReader(conn))
readChan := make(chan []byte)
flushChan := make(chan int) switch cmd.Command {
go netutils.AutoRecoonectedTCP(ladr, tcpRemoteAddress, msgBuf, initMsg, writeChan, readChan, flushChan) case "SET":
<-readChan data := genRedisArray([]byte("SET"), []byte(cmd.Name), []byte(cmd.Data), []byte("EX"), []byte(fmt.Sprintf("%v", cmd.Duration)))
dataBuf := make([]byte, 0) writer.PrintfLine(string(data))
dataLen := 0 logger.Debug("redis:set")
for { case "DEL":
select { data := genRedisArray([]byte("DEL"), []byte(cmd.Name))
case cmd := <-redisCmdWrite: writer.PrintfLine(string(data))
switch cmd.Command { logger.Debug("redis:del")
case "SET": case "GET":
data := genRedisArray([]byte("SET"), []byte(cmd.Name), cmd.Data, []byte("EX"), []byte(fmt.Sprintf("%v", cmd.Duration))) data := genRedisArray([]byte("GET"), []byte(cmd.Name))
writeChan <- data writer.PrintfLine(string(data))
case "GET": logger.Debug("redis:get")
data := genRedisArray([]byte("GET"), []byte(cmd.Name)) for {
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 { select {
case redisCmdRead <- RedisCmd{}: case <-time.After(time.Second * 1):
case <-time.After(time.Second * 5): channel <- RedisCmd{Error: fmt.Errorf(RedisTimeout)}
return
default:
read, _ := reader.ReadLineBytes()
if string(read) != "$1" {
channel <- RedisCmd{Error: fmt.Errorf(RedisMiss)}
return
}
read, _ = reader.ReadLineBytes()
channel <- RedisCmd{Data: read}
return
} }
} }
} }
@@ -171,16 +85,16 @@ func initContext(hostnamePort string, redisCmdWrite, redisCmdRead chan RedisCmd)
func (sr *SimpleRedis) Init(redisHost string) { func (sr *SimpleRedis) Init(redisHost string) {
sr.redisHost = redisHost 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) { func (sr *SimpleRedis) Get(name string) ([]byte, error) {
sr.redisCmd.Command = "GET" redisCmd := RedisCmd{
sr.redisCmd.Name = name Command: "GET",
sr.redisChanWrite <- sr.redisCmd Name: name,
resp := <-sr.redisChanRead }
channel := make(chan RedisCmd)
go askRedis(sr.redisHost, redisCmd, channel)
resp := <-channel
if resp.Error != nil { if resp.Error != nil {
return nil, resp.Error return nil, resp.Error
} }
@@ -188,25 +102,21 @@ func (sr *SimpleRedis) Get(name string) ([]byte, error) {
} }
func (sr *SimpleRedis) Set(name string, data []byte, duration int64) error { func (sr *SimpleRedis) Set(name string, data []byte, duration int64) error {
sr.redisCmd.Command = "SET" redisCmd := RedisCmd{
sr.redisCmd.Name = name Command: "SET",
sr.redisCmd.Data = data Name: name,
sr.redisCmd.Duration = duration Data: data,
sr.redisChanWrite <- sr.redisCmd Duration: duration,
resp := <-sr.redisChanRead
if resp.Error != nil {
return resp.Error
} }
go askRedis(sr.redisHost, redisCmd, nil)
return nil return nil
} }
func (sr *SimpleRedis) Del(name string) error { func (sr *SimpleRedis) Del(name string) error {
sr.redisCmd.Command = "DEL" redisCmd := RedisCmd{
sr.redisCmd.Name = name Command: "DEL",
sr.redisChanWrite <- sr.redisCmd Name: name,
resp := <-sr.redisChanRead
if resp.Error != nil {
return resp.Error
} }
go askRedis(sr.redisHost, redisCmd, nil)
return nil return nil
} }
-525
View File
@@ -1,525 +0,0 @@
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
View File
@@ -1,42 +0,0 @@
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)
}
-3
View File
@@ -1,6 +1,3 @@
# 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