REPOSITORY / ScuroNeko/mtg

Compare commits

DIFF REPOSITORY

Compare commits

...
55 Commits
Author SHA1 Message Date
9seconds b1074193a8 Update dependencies 2019-06-11 13:16:01 +03:00
Sergey ArkhipovandGitHub 78dea9ae3f Merge pull request #76 from 9seconds/replay
Base prevention of replay attacks on proxy
2019-04-26 16:34:02 +03:00
9seconds da8dba1585 Add some notes to README 2019-04-26 16:24:39 +03:00
9seconds 33852ca481 Base prevention of replay attacks on proxy 2019-04-26 15:53:10 +03:00
9seconds eff53694a0 Update dependencies 2019-04-26 14:56:58 +03:00
9seconds 74a7c505eb Add notes on ansible 2019-03-01 10:30:08 +03:00
9seconds 9983f7fcba Add a note on secure mode 2019-03-01 10:27:33 +03:00
9seconds a90072260e Update docker image to use Go 1.12 2019-03-01 10:24:58 +03:00
9seconds b153240bc3 Update dependencies 2019-03-01 10:20:53 +03:00
9seconds 6dfbd26524 Update golangci-lint version 2019-03-01 10:18:23 +03:00
Sergey ArkhipovandGitHub 39e7663e0c Merge pull request #70 from im-kulikov/fix/deprecation-warning
DualStack is deprecated
2019-03-01 10:12:11 +03:00
Sergey ArkhipovandGitHub 71f7bf6cad Merge pull request #71 from im-kulikov/feature/add-go1.12-to-travis-ci
Add Golang 1.12 to TravisCI
2019-03-01 10:06:59 +03:00
Evgeniy Kulikov fba0e235c1 Add Golang 1.12 to TravisCI 2019-02-26 21:47:29 +03:00
Evgeniy Kulikov 70be70de3c DualStack is deprecated
```
	// DualStack previously enabled RFC 6555 Fast Fallback
	// support, also known as "Happy Eyeballs", in which IPv4 is
	// tried soon if IPv6 appears to be misconfigured and
	// hanging.
	//
	// Deprecated: Fast Fallback is enabled by default. To
	// disable, set FallbackDelay to a negative value.
	DualStack bool
```
2019-02-26 21:44:47 +03:00
Sergey ArkhipovandGitHub b565d83d40 Merge pull request #62 from Eternity-Yarr/patch-1
Update README.md
2019-01-24 08:43:34 +03:00
Sergey ArkhipovandGitHub ebc459440c Merge branch 'master' into patch-1 2019-01-24 08:28:16 +03:00
9seconds 6e6049a73a Remove critic from travis 2019-01-24 08:17:05 +03:00
9seconds 916714a909 Update dependencies 2019-01-24 08:16:47 +03:00
9seconds f68b195306 Update golangci-lint 2019-01-24 08:15:56 +03:00
Dmitry V. RoenkoandGitHub 5f1d1ce883 Update README.md
- little typo in docker command line 
- added more preferred 'secure only' example
2019-01-22 13:03:02 +03:00
9seconds 16558a7c55 Merge remote-tracking branch 'origin/master' into stable 2018-11-02 23:21:06 +03:00
Sergey ArkhipovandGitHub e7958aaf33 Merge pull request #57 from 9seconds/remove-go-generate
Remove generating of version.go
2018-11-02 20:20:18 +00:00
9seconds 7d6b661d97 More correct command for versioning 2018-11-02 19:10:21 +03:00
9seconds 28b3cbe91a Update dependencies 2018-11-02 19:02:39 +03:00
9seconds 6818364231 Update golangci-lint 2018-11-02 19:02:07 +03:00
9seconds ea97bf51c8 Remove autogeneration of version.go 2018-11-02 18:54:28 +03:00
9seconds 1678ddbd53 Reformat README 2018-11-02 18:37:16 +03:00
Sergey ArkhipovandGitHub f81f29cbb4 Merge pull request #55 from 9seconds/secure-mode-readme
Add explanation on secure mode
2018-11-02 15:21:27 +00:00
Sergey ArkhipovandGitHub 2d9259db48 Merge pull request #56 from hdid/patch-1
tiny fix README.md
2018-11-02 15:20:47 +00:00
hdidandGitHub c721c636d9 tiny fix README.md 2018-11-02 13:09:26 +03:30
Sergey ArkhipovandGitHub 60472072aa Add explanation on secure mode
Hopefully, this closes https://github.com/9seconds/mtg/issues/49
2018-11-02 09:02:45 +00:00
9seconds 6721e6fd9f Update dependencies 2018-11-01 18:34:38 +03:00
Sergey ArkhipovandGitHub 47fb5c23cb Merge pull request #53 from 9seconds/prometheus-mtg-only
Use mtg metrics only for prometheus endpoint
2018-11-01 18:32:55 +03:00
Sergey ArkhipovandGitHub fce6118c71 Merge pull request #52 from 9seconds/log-config
Log configuration on proxy start
2018-11-01 18:25:25 +03:00
9seconds d277a2975a Use mtg metrics only for prometheus endpoint 2018-11-01 18:23:36 +03:00
9seconds c3f21a7b0d Log configuration on proxy start
This has to simplify the debugging because usually we do not know all
options which user used to start mtg.
2018-11-01 18:07:16 +03:00
9seconds 61a1024264 Merge remote-tracking branch 'origin/master' into stable 2018-10-17 14:20:36 +03:00
Sergey ArkhipovandGitHub 4695a0c433 Merge pull request #47 from 9seconds/prometheus
Prometheus integration
2018-10-17 13:14:31 +03:00
9seconds a9d0ca1e90 Describe and document prometheus integration 2018-10-17 12:54:14 +03:00
9seconds 9346c11f37 Add support for prometheus 2018-10-17 12:06:46 +03:00
9seconds 5f00363da5 Add goreportcard 2018-10-17 10:52:30 +03:00
9seconds 3aec9657d7 Add comments for exported functions 2018-10-17 10:51:16 +03:00
9seconds 358d42ec61 Fix spelling issue 2018-10-17 10:49:27 +03:00
Sergey ArkhipovandGitHub ac2625209b Merge pull request #45 from 9seconds/stats
Rework stats server
2018-10-16 22:24:01 +03:00
9seconds c74e0ee359 Rework stats server 2018-10-16 22:17:46 +03:00
Sergey ArkhipovandGitHub 7c463851b2 Okay, 3 megabytes. Image grows :( 2018-10-13 11:18:36 +03:00
Sergey ArkhipovandGitHub c9f6b922c3 Merge pull request #44 from 9seconds/goroutines-number
Decrease the number of used goroutines
2018-10-12 10:56:46 +03:00
9seconds 7128e70e50 Minor changes in dockerfile 2018-10-12 10:50:52 +03:00
9seconds cf9bf56d8e Update dependencies 2018-10-12 10:40:54 +03:00
9seconds 76cfbb009a Decrease an amount of goroutines 2018-10-12 10:38:35 +03:00
9seconds cb78169d6a Update dependencies 2018-10-10 13:49:52 +03:00
9seconds d1703873c1 Merge remote-tracking branch 'origin/master' into stable 2018-09-24 18:29:24 +03:00
Sergey ArkhipovandGitHub ac33abbbb1 Merge pull request #40 from 9seconds/secure-only
"Secure only" mode
2018-09-24 18:15:36 +03:00
9seconds a6893c8df7 Update README 2018-09-24 18:08:06 +03:00
9seconds 7182c7bf65 Add new secure-only mode 2018-09-24 18:04:00 +03:00
27 changed files with 677 additions and 382 deletions
+1
View File
@@ -10,3 +10,4 @@ format = "colored-line-number"
[linters] [linters]
enable-all = true enable-all = true
disable = ["gochecknoglobals"]
+1 -1
View File
@@ -6,6 +6,7 @@ dist: trusty
go: go:
- "1.11.x" - "1.11.x"
- 1.12.x
- master - master
before_script: make prepare before_script: make prepare
@@ -13,7 +14,6 @@ before_script: make prepare
script: script:
- make all - make all
- make lint - make lint
- make critic
- make test - make test
matrix: matrix:
+4 -5
View File
@@ -1,7 +1,7 @@
############################################################################### ###############################################################################
# BUILD STAGE # BUILD STAGE
FROM golang:1.11-alpine FROM golang:1.12-alpine
RUN set -x \ RUN set -x \
&& apk --no-cache --update add \ && apk --no-cache --update add \
@@ -10,8 +10,7 @@ RUN set -x \
curl \ curl \
git \ git \
make \ make \
upx \ upx
&& update-ca-certificates
COPY . /go/src/github.com/9seconds/mtg/ COPY . /go/src/github.com/9seconds/mtg/
@@ -26,7 +25,7 @@ RUN set -x \
FROM scratch FROM scratch
ENTRYPOINT ["/usr/local/bin/mtg"] ENTRYPOINT ["/mtg"]
ENV MTG_IP=0.0.0.0 \ ENV MTG_IP=0.0.0.0 \
MTG_PORT=3128 \ MTG_PORT=3128 \
MTG_STATS_IP=0.0.0.0 \ MTG_STATS_IP=0.0.0.0 \
@@ -34,4 +33,4 @@ ENV MTG_IP=0.0.0.0 \
EXPOSE 3128 3129 EXPOSE 3128 3129
COPY --from=0 /etc/ssl/certs/ca-certificates.crt /etc/ssl/certs/ca-certificates.crt COPY --from=0 /etc/ssl/certs/ca-certificates.crt /etc/ssl/certs/ca-certificates.crt
COPY --from=0 /go/src/github.com/9seconds/mtg/mtg /usr/local/bin/mtg COPY --from=0 /go/src/github.com/9seconds/mtg/mtg /mtg
+11 -20
View File
@@ -3,26 +3,28 @@ IMAGE_NAME := mtg
APP_NAME := $(IMAGE_NAME) APP_NAME := $(IMAGE_NAME)
CC_BINARIES := $(shell bash -c "echo -n $(APP_NAME)-{linux,freebsd,openbsd}-{386,amd64} $(APP_NAME)-linux-{arm,arm64}") CC_BINARIES := $(shell bash -c "echo -n $(APP_NAME)-{linux,freebsd,openbsd}-{386,amd64} $(APP_NAME)-linux-{arm,arm64}")
APP_DEPS := version.go
GOLANGCI_LINT_VERSION := v1.10.2 GOLANGCI_LINT_VERSION := v1.15.0
COMMON_BUILD_FLAGS := -ldflags="-s -w" VERSION_GO := $(shell go version)
VERSION_DATE := $(shell date -Ru)
VERSION_TAG := $(shell git describe --tags --always)
COMMON_BUILD_FLAGS := -ldflags="-s -w -X 'main.version=$(VERSION_TAG) ($(VERSION_GO)) [$(VERSION_DATE)]'"
MOD_ON := env GO111MODULE=on MOD_ON := env GO111MODULE=on
MOD_OFF := env GO111MODULE=auto MOD_OFF := env GO111MODULE=auto
# ----------------------------------------------------------------------------- # -----------------------------------------------------------------------------
$(APP_NAME): $(APP_DEPS) $(APP_NAME):
@$(MOD_ON) go build $(COMMON_BUILD_FLAGS) -o "$(APP_NAME)" @$(MOD_ON) go build $(COMMON_BUILD_FLAGS) -o "$(APP_NAME)"
static-$(APP_NAME): $(APP_DEPS) static-$(APP_NAME):
@$(MOD_ON) env CGO_ENABLED=0 GOOS=linux go build -a -installsuffix cgo $(COMMON_BUILD_FLAGS) -o "$(APP_NAME)" @$(MOD_ON) env CGO_ENABLED=0 GOOS=linux go build -a -installsuffix cgo $(COMMON_BUILD_FLAGS) -o "$(APP_NAME)"
$(APP_NAME)-%: GOOS=$(shell echo -n "$@" | sed 's?$(APP_NAME)-??' | cut -f1 -d-) $(APP_NAME)-%: GOOS=$(shell echo -n "$@" | sed 's?$(APP_NAME)-??' | cut -f1 -d-)
$(APP_NAME)-%: GOARCH=$(shell echo -n "$@" | sed 's?$(APP_NAME)-??' | cut -f2 -d-) $(APP_NAME)-%: GOARCH=$(shell echo -n "$@" | sed 's?$(APP_NAME)-??' | cut -f2 -d-)
$(APP_NAME)-%: $(APP_DEPS) ccbuilds $(APP_NAME)-%: ccbuilds
@$(MOD_ON) env "GOOS=$(GOOS)" "GOARCH=$(GOARCH)" \ @$(MOD_ON) env "GOOS=$(GOOS)" "GOARCH=$(GOARCH)" \
go build \ go build \
$(COMMON_BUILD_FLAGS) \ $(COMMON_BUILD_FLAGS) \
@@ -31,9 +33,6 @@ $(APP_NAME)-%: $(APP_DEPS) ccbuilds
ccbuilds: ccbuilds:
@rm -rf ./ccbuilds && mkdir -p ./ccbuilds @rm -rf ./ccbuilds && mkdir -p ./ccbuilds
version.go:
@$(MOD_ON) go generate main.go
vendor: go.mod go.sum vendor: go.mod go.sum
@$(MOD_ON) go mod vendor @$(MOD_ON) go mod vendor
@@ -53,17 +52,13 @@ crosscompile-dir:
@rm -rf "$(CC_DIR)" && mkdir -p "$(CC_DIR)" @rm -rf "$(CC_DIR)" && mkdir -p "$(CC_DIR)"
.PHONY: test .PHONY: test
test: vendor $(APP_DEPS) test: vendor
@$(MOD_ON) go test -v ./... @$(MOD_ON) go test -v ./...
.PHONY: lint .PHONY: lint
lint: vendor $(APP_DEPS) lint: vendor
@$(MOD_OFF) golangci-lint run @$(MOD_OFF) golangci-lint run
.PHONY: critic
critic: vendor $(APP_DEPS)
@$(MOD_OFF) gocritic check-project "$(ROOT_DIR)"
.PHONY: clean .PHONY: clean
clean: clean:
@git clean -xfd && \ @git clean -xfd && \
@@ -75,13 +70,9 @@ docker:
@docker build --pull -t "$(IMAGE_NAME)" "$(ROOT_DIR)" @docker build --pull -t "$(IMAGE_NAME)" "$(ROOT_DIR)"
.PHONY: prepare .PHONY: prepare
prepare: install-lint install-critic prepare: install-lint
.PHONY: install-lint .PHONY: install-lint
install-lint: install-lint:
@curl -sfL https://install.goreleaser.com/github.com/golangci/golangci-lint.sh \ @curl -sfL https://install.goreleaser.com/github.com/golangci/golangci-lint.sh \
| $(MOD_OFF) bash -s -- -b $(GOPATH)/bin $(GOLANGCI_LINT_VERSION) | $(MOD_OFF) bash -s -- -b $(GOPATH)/bin $(GOLANGCI_LINT_VERSION)
.PHONY: install-critic
install-critic:
@$(MOD_OFF) go get -u github.com/go-critic/go-critic/...
+104 -26
View File
@@ -3,6 +3,7 @@
Bullshit-free MTPROTO proxy for Telegram Bullshit-free MTPROTO proxy for Telegram
[![Build Status](https://travis-ci.org/9seconds/mtg.svg?branch=master)](https://travis-ci.org/9seconds/mtg) [![Build Status](https://travis-ci.org/9seconds/mtg.svg?branch=master)](https://travis-ci.org/9seconds/mtg)
[![Go Report Card](https://goreportcard.com/badge/github.com/9seconds/mtg)](https://goreportcard.com/report/github.com/9seconds/mtg)
[![Docker Build Status](https://img.shields.io/docker/build/nineseconds/mtg.svg)](https://hub.docker.com/r/nineseconds/mtg/) [![Docker Build Status](https://img.shields.io/docker/build/nineseconds/mtg.svg)](https://hub.docker.com/r/nineseconds/mtg/)
# Rationale # Rationale
@@ -32,7 +33,7 @@ mtg is an implementation in golang which is intended to be:
software. I also believe that in case of throwout proxies, this feature software. I also believe that in case of throwout proxies, this feature
is useless luxury. is useless luxury.
* **Minimum docker image size** * **Minimum docker image size**
Official image is less than 2.5 megabytes. Literally. Official image is less than 3 megabytes. Literally.
* **No management WebUI** * **No management WebUI**
This is an implementation of simple lightweight proxy. I won't do that. This is an implementation of simple lightweight proxy. I won't do that.
@@ -94,6 +95,11 @@ docker pull nineseconds/mtg:stable
docker pull nineseconds/mtg:0.10 docker pull nineseconds/mtg:0.10
``` ```
# Ansible role
You can find unofficial Ansible role for mtg here: https://github.com/rlex/ansible-role-mtg
Also, there is another project on Ansible Galaxy: https://galaxy.ansible.com/ivansible/lin_mtproxy
# Configuration # Configuration
Basically, to run this tool you need to configure as less as possible. Basically, to run this tool you need to configure as less as possible.
@@ -112,11 +118,58 @@ head -c 512 /dev/urandom | md5sum | cut -f 1 -d ' '
## Secure mode ## Secure mode
If you want to support new secure mode, please prepend `dd` to the _tl;dr - use secret mode for all new installation of proxy; only clients
secret. For example, secret `cf18fa8ea0267057e2c61a5f7322a8e7` should with dd-secrets will be able to connect. This mode abuses attempts to
be `ddcf18fa8ea0267057e2c61a5f7322a8e7`. But pay attention that some DPI MTPROTO traffic._
old clients won't support this mode. If this is not your case, I would
suggest to go with this mode. Secure mode is not the best name and of course, it creates a lot of
confusion. To explain what it means, we need to tell you some bits on
dd-secrets.
MTPROTO proxy protocol requires 16-byte secret. You usually
propagate it as a 32 characters hexadecimal string like
`282831900f371ca182feb0e4e1e1aeef` (if you decode this string
to bytes, you will get a real secret which is used in the
protocol). Everything went quite good until the moment when
developers found an evidence that [protocol is quite weak to
DPI](https://github.com/TelegramMessenger/MTProxy/issues/35) and some
enthusiasts even created simple proofs of concepts on [detecting MTPROTO
traffic](https://github.com/darkk/poormansmtproto).
Telegram team has introduced a patch called dd-secrets. If you have
a secret `282831900f371ca182feb0e4e1e1aeef` then your dd-secret is
`dd282831900f371ca182feb0e4e1e1aeef`. That is, you just add dd prefix
to the secret, prepend it with dd. In that case, original secret
`282831900f371ca182feb0e4e1e1aeef` is used but client and server start
to act a little bit different: they start to add random noise to the
packets so they can't be detected by their length. In order to keep
backward compatibility, all proxies a quite liberal to the secrets to
use: if the client uses plain secret, without dd prefix, they fall back
to the normal behavior. If dd-secret is used (proxy can extract this
information on the handshake), then more secured, the hardened behavior
is used.
Yes, it can look like a hack but it is as it is.
Now going back to the secure mode: if you do not pass `-s` flag to the
mtg, then it checks what mode is requested by the client. If the client
uses plain secret, without dd prefix, then proxy falls back to the
original behavior and do not play with paddings. If dd-secret is used
and client demands this mode, then proxy start to add that random noise
to the packets. But if you pass `-s`, then only clients with dd-secrets
can connect. How to migrate existing clients then? If a client is new
enough, you can just prepend the secret with dd string in the settings.
If it is an old guy, then nothing to do, sorry.
Why this mode matters? We do not have evidence but there is quite a big
suspicion that some ISPs start to filter MTPROTO traffic. If they detect
the IP address which acts as a proxy, they block it and no clients can
use this proxy. This is an attempt to prevent such a situation.
General rule of thumb: with all new installation of proxies I would
advise to go with secure mode by default. But please do remember that it
means that clients, which do not pass dd-prefix to their secrets, will
not be able to connect. *Secure mode works only with dd-prefixes!*
Oneliners to generate such secrets: Oneliners to generate such secrets:
@@ -130,32 +183,46 @@ or
echo dd$(head -c 512 /dev/urandom | md5sum | cut -f 1 -d ' ') echo dd$(head -c 512 /dev/urandom | md5sum | cut -f 1 -d ' ')
``` ```
## Antireplay cache
In order to prevent replay attacks, we have internal storage of first
frames messages for connected clients. These frames are generated
randomly by design and we have negligible possibility of duplication
(probability is 1/(2^64)) but it could be quite effective in order to
prevent replays.
## Environment variables ## Environment variables
It is possible to configure this tool using environment variables. You It is possible to configure this tool using environment variables. You
can configure any flag but not secret or adtag. Here is the list of can configure any flag but not secret or adtag. Here is the list of
supported environment variables: supported environment variables:
| Environment variable | Corresponding flags | Default value | Description | | Environment variable | Corresponding flags | Default value | Description |
|--------------------------|------------------------|-----------------------------------|----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------| |-------------------------------|-----------------------------|-----------------------------------|----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|
| `MTG_DEBUG` | `-d`, `--debug` | `false` | Run in debug mode. Usually, you need to run in this mode only if you develop this tool or its maintainer is asking you to provide logs with such verbosity. | | `MTG_DEBUG` | `-d`, `--debug` | `false` | Run in debug mode. Usually, you need to run in this mode only if you develop this tool or its maintainer is asking you to provide logs with such verbosity. |
| `MTG_VERBOSE` | `-v`, `--verbose` | `false` | Run in verbose mode. This is way less chatty than debug mode. | | `MTG_VERBOSE` | `-v`, `--verbose` | `false` | Run in verbose mode. This is way less chatty than debug mode. |
| `MTG_IP` | `-b`, `--bind-ip` | `127.0.0.1` | Which IP should we bind to. As usual, `0.0.0.0` means that we want to listen on all interfaces. Also, 4 zeroes will bind to both IPv4 and IPv6. | | `MTG_IP` | `-b`, `--bind-ip` | `127.0.0.1` | Which IP should we bind to. As usual, `0.0.0.0` means that we want to listen on all interfaces. Also, 4 zeroes will bind to both IPv4 and IPv6. |
| `MTG_PORT` | `-p`, `--bind-port` | `3128` | Which port should we bind to (listen on). | | `MTG_PORT` | `-p`, `--bind-port` | `3128` | Which port should we bind to (listen on). |
| `MTG_IPV4` | `-4`, `--public-ipv4` | [Autodetect](https://ifconfig.co) | IPv4 address of this proxy. This is required if you NAT your proxy or run it in a docker container. In that case, you absolutely need to specify public IPv4 address of the proxy, otherwise either URLs will be broken or proxy could not access Telegram middle proxies. | | `MTG_IPV4` | `-4`, `--public-ipv4` | [Autodetect](https://ifconfig.co) | IPv4 address of this proxy. This is required if you NAT your proxy or run it in a docker container. In that case, you absolutely need to specify public IPv4 address of the proxy, otherwise either URLs will be broken or proxy could not access Telegram middle proxies. |
| `MTG_IPV4_PORT` | `--public-ipv4-port` | Value of `--bind-port` | Which port should be public of IPv4 interface. This affects only generated links and should be changed only if you NAT your proxy or run it in a docker container. | | `MTG_IPV4_PORT` | `--public-ipv4-port` | Value of `--bind-port` | Which port should be public of IPv4 interface. This affects only generated links and should be changed only if you NAT your proxy or run it in a docker container. |
| `MTG_IPV6` | `-6`, `--public-ipv6` | [Autodetect](https://ifconfig.co) | IPv6 address of this proxy. This is required if you NAT your proxy or run it in a docker container. In that case, you absolutely need to specify public IPv6 address of the proxy, otherwise either URLs will be broken or proxy could not access Telegram middle proxies. | | `MTG_IPV6` | `-6`, `--public-ipv6` | [Autodetect](https://ifconfig.co) | IPv6 address of this proxy. This is required if you NAT your proxy or run it in a docker container. In that case, you absolutely need to specify public IPv6 address of the proxy, otherwise either URLs will be broken or proxy could not access Telegram middle proxies. |
| `MTG_IPV6_PORT` | `--public-ipv6-port` | Value of `--bind-port` | Which port should be public of IPv6 interface. This affects only generated links and should be changed only if you NAT your proxy or run it in a docker container. | | `MTG_IPV6_PORT` | `--public-ipv6-port` | Value of `--bind-port` | Which port should be public of IPv6 interface. This affects only generated links and should be changed only if you NAT your proxy or run it in a docker container. |
| `MTG_STATS_IP` | `-t`, `--stats-ip` | `127.0.0.1` | Which IP should we bind the internal statistics HTTP server. | | `MTG_STATS_IP` | `-t`, `--stats-ip` | `127.0.0.1` | Which IP should we bind the internal statistics HTTP server. |
| `MTG_STATS_PORT` | `-q`, `--stats-port` | `3129` | Which port should we bind the internal statistics HTTP server. | | `MTG_STATS_PORT` | `-q`, `--stats-port` | `3129` | Which port should we bind the internal statistics HTTP server. |
| `MTG_STATSD_IP` | `--statsd-ip` | | IP/host addresses of statsd service. No defaults, by defaults we do not send anything there. | | `MTG_STATSD_IP` | `--statsd-ip` | | IP/host addresses of statsd service. No defaults, by defaults we do not send anything there. |
| `MTG_STATSD_PORT` | `--statsd-port` | `8125` | Which port should we use to work with statsd. | | `MTG_STATSD_PORT` | `--statsd-port` | `8125` | Which port should we use to work with statsd. |
| `MTG_STATSD_NETWORK` | `--statsd-network` | `udp` | Which protocol should we use to work with statsd. Possible options are `udp` and `tcp`. | | `MTG_STATSD_NETWORK` | `--statsd-network` | `udp` | Which protocol should we use to work with statsd. Possible options are `udp` and `tcp`. |
| `MTG_STATSD_PREFIX` | `--statsd-prefix` | `mtg` | Which bucket prefix we should use. For example, if you set `mtg`, then metric `traffic.ingress` would be send as `mtg.traffic.ingress`. | | `MTG_STATSD_PREFIX` | `--statsd-prefix` | `mtg` | Which bucket prefix we should use. For example, if you set `mtg`, then metric `traffic.ingress` would be send as `mtg.traffic.ingress`. |
| `MTG_STATSD_TAGS_FORMAT` | `--statsd-tags-format` | | Which tags format we should use. By default, we are using default vanilla statsd tags format but if you want to send directly to InfluxDB or Datadog, please specify it there. Possible options are `influxdb` and `datadog`. | | `MTG_STATSD_TAGS_FORMAT` | `--statsd-tags-format` | | Which tags format we should use. By default, we are using default vanilla statsd tags format but if you want to send directly to InfluxDB or Datadog, please specify it there. Possible options are `influxdb` and `datadog`. |
| `MTG_STATSD_TAGS` | `--statsd-tags` | | Which tags should we send to statsd with our metrics. Please specify them as `key=value` pairs. | | `MTG_STATSD_TAGS` | `--statsd-tags` | | Which tags should we send to statsd with our metrics. Please specify them as `key=value` pairs. |
| `MTG_BUFFER_WRITE` | `-w`, `--write-buffer` | `65536` | The size of TCP write buffer in bytes. Write buffer is the buffer for messages which are going from client to Telegram. | | `MTG_PROMETHEUS_PREFIX` | `--prometheus-prefix` | `mtg` | Which namespace should be used for prometheus metrics. |
| `MTG_BUFFER_READ` | `-r`, `--read-buffer` | `131072` | The size of TCP read buffer in bytes. Read buffer is the buffer for messages from Telegram to client. | | `MTG_BUFFER_WRITE` | `-w`, `--write-buffer` | `65536` | The size of TCP write buffer in bytes. Write buffer is the buffer for messages which are going from client to Telegram. |
| `MTG_BUFFER_READ` | `-r`, `--read-buffer` | `131072` | The size of TCP read buffer in bytes. Read buffer is the buffer for messages from Telegram to client. |
| `MTG_SECURE_ONLY` | `-s`, `--secure-only` | `false` | Support only clients with secure mode (i.e only clients with dd-secrets). |
| `MTG_ANTIREPLAY_MAXSIZE` | `anti-replay-max-size` | `128` | Max size of antireplay cache in megabytes. |
| `MTG_ANTIREPLAY_EVICTIONTIME` | `anti-replay-eviction-time` | `168h` | Eviction time for antireplay cache entries. |
Usually you want to modify only read/write buffer sizes. If you feel Usually you want to modify only read/write buffer sizes. If you feel
that proxy is slow, try to increase both sizes giving more priority to that proxy is slow, try to increase both sizes giving more priority to
@@ -228,3 +295,14 @@ All metrics are gauges. Here is the list of metrics and their meaning:
All metrics are prefixed with given prefix. Default prefix is `mtg`. All metrics are prefixed with given prefix. Default prefix is `mtg`.
With such prefix metric name `traffic.ingress`, for example, would be With such prefix metric name `traffic.ingress`, for example, would be
`mtg.traffic.ingress`. `mtg.traffic.ingress`.
# Prometheus integration
[Prometheus](https://prometheus.io) integration comes out of
the box, you do not need to setup anything special. Prometheus
scrape endpoint lives on the same IP/port where generic stats
service (`http://${MTG_STATS_IP}:${MTG_STATS_PORT}`) but on
`/prometheus` path. So, if you access http stats service as `curl
http://localhost:3129/`, then your prometheus endpoint is `curl
http://localhost:3129/prometheus/`.
+37
View File
@@ -0,0 +1,37 @@
package antireplay
import (
"github.com/allegro/bigcache"
"github.com/juju/errors"
"github.com/9seconds/mtg/config"
)
// Cache defines storage for obfuscated2 handshake frames.
type Cache struct {
cache *bigcache.BigCache
}
func (a Cache) Add(frame []byte) {
a.cache.Set(string(frame), nil) // nolint: errcheck
}
func (a Cache) Has(frame []byte) bool {
_, err := a.cache.Get(string(frame))
return err == nil
}
func NewCache(config *config.Config) (Cache, error) {
cache, err := bigcache.NewBigCache(bigcache.Config{
Shards: 1024,
LifeWindow: config.AntiReplayEvictionTime,
Hasher: hasher{},
HardMaxCacheSize: config.AntiReplayMaxSize,
})
if err != nil {
return Cache{}, errors.Annotate(err, "Cannot make cache")
}
return Cache{cache}, nil
}
+9
View File
@@ -0,0 +1,9 @@
package antireplay
import "github.com/cespare/xxhash"
type hasher struct{}
func (h hasher) Sum64(value string) uint64 {
return xxhash.Sum64String(value)
}
+2 -1
View File
@@ -4,6 +4,7 @@ import (
"context" "context"
"net" "net"
"github.com/9seconds/mtg/antireplay"
"github.com/9seconds/mtg/config" "github.com/9seconds/mtg/config"
"github.com/9seconds/mtg/mtproto" "github.com/9seconds/mtg/mtproto"
"github.com/9seconds/mtg/wrappers" "github.com/9seconds/mtg/wrappers"
@@ -11,4 +12,4 @@ import (
// Init defines common method for initializing client connections. // Init defines common method for initializing client connections.
type Init func(context.Context, context.CancelFunc, net.Conn, string, type Init func(context.Context, context.CancelFunc, net.Conn, string,
*config.Config) (wrappers.Wrap, *mtproto.ConnectionOpts, error) antireplay.Cache, *config.Config) (wrappers.Wrap, *mtproto.ConnectionOpts, error)
+9 -1
View File
@@ -7,6 +7,7 @@ import (
"github.com/juju/errors" "github.com/juju/errors"
"github.com/9seconds/mtg/antireplay"
"github.com/9seconds/mtg/config" "github.com/9seconds/mtg/config"
"github.com/9seconds/mtg/mtproto" "github.com/9seconds/mtg/mtproto"
"github.com/9seconds/mtg/obfuscated2" "github.com/9seconds/mtg/obfuscated2"
@@ -18,7 +19,8 @@ const handshakeTimeout = 10 * time.Second
// DirectInit initializes client connection for proxy which connects to // DirectInit initializes client connection for proxy which connects to
// Telegram directly. // Telegram directly.
func DirectInit(ctx context.Context, cancel context.CancelFunc, socket net.Conn, func DirectInit(ctx context.Context, cancel context.CancelFunc, socket net.Conn,
connID string, conf *config.Config) (wrappers.Wrap, *mtproto.ConnectionOpts, error) { connID string, antiReplayCache antireplay.Cache,
conf *config.Config) (wrappers.Wrap, *mtproto.ConnectionOpts, error) {
tcpSocket := socket.(*net.TCPConn) tcpSocket := socket.(*net.TCPConn)
if err := tcpSocket.SetNoDelay(false); err != nil { if err := tcpSocket.SetNoDelay(false); err != nil {
return nil, nil, errors.Annotate(err, "Cannot disable NO_DELAY to client socket") return nil, nil, errors.Annotate(err, "Cannot disable NO_DELAY to client socket")
@@ -42,6 +44,12 @@ func DirectInit(ctx context.Context, cancel context.CancelFunc, socket net.Conn,
if err != nil { if err != nil {
return nil, nil, errors.Annotate(err, "Cannot parse obfuscated frame") return nil, nil, errors.Annotate(err, "Cannot parse obfuscated frame")
} }
if antiReplayCache.Has([]byte(frame)) {
return nil, nil, errors.New("Replay attack is detected")
}
antiReplayCache.Add([]byte(frame))
connOpts.ConnectionProto = mtproto.ConnectionProtocolAny connOpts.ConnectionProto = mtproto.ConnectionProtocolAny
connOpts.ClientAddr = conn.RemoteAddr() connOpts.ClientAddr = conn.RemoteAddr()
+4 -2
View File
@@ -4,6 +4,7 @@ import (
"context" "context"
"net" "net"
"github.com/9seconds/mtg/antireplay"
"github.com/9seconds/mtg/config" "github.com/9seconds/mtg/config"
"github.com/9seconds/mtg/mtproto" "github.com/9seconds/mtg/mtproto"
"github.com/9seconds/mtg/wrappers" "github.com/9seconds/mtg/wrappers"
@@ -12,8 +13,9 @@ import (
// MiddleInit initializes client connection for proxy which has to // MiddleInit initializes client connection for proxy which has to
// support promoted channels, connect to Telegram middle proxies etc. // support promoted channels, connect to Telegram middle proxies etc.
func MiddleInit(ctx context.Context, cancel context.CancelFunc, socket net.Conn, func MiddleInit(ctx context.Context, cancel context.CancelFunc, socket net.Conn,
connID string, conf *config.Config) (wrappers.Wrap, *mtproto.ConnectionOpts, error) { connID string, antiReplayCache antireplay.Cache,
conn, opts, err := DirectInit(ctx, cancel, socket, connID, conf) conf *config.Config) (wrappers.Wrap, *mtproto.ConnectionOpts, error) {
conn, opts, err := DirectInit(ctx, cancel, socket, connID, antiReplayCache, conf)
if err != nil { if err != nil {
return nil, nil, err return nil, nil, err
} }
+31 -17
View File
@@ -6,6 +6,7 @@ import (
"fmt" "fmt"
"net" "net"
"strconv" "strconv"
"time"
"github.com/juju/errors" "github.com/juju/errors"
statsd "gopkg.in/alexcesaro/statsd.v2" statsd "gopkg.in/alexcesaro/statsd.v2"
@@ -16,6 +17,7 @@ type Config struct {
Debug bool Debug bool
Verbose bool Verbose bool
SecureMode bool SecureMode bool
SecureOnly bool
ReadBufferSize int ReadBufferSize int
WriteBufferSize int WriteBufferSize int
@@ -30,6 +32,9 @@ type Config struct {
PublicIPv6 net.IP PublicIPv6 net.IP
StatsIP net.IP StatsIP net.IP
AntiReplayMaxSize int
AntiReplayEvictionTime time.Duration
StatsD struct { StatsD struct {
Addr net.Addr Addr net.Addr
Prefix string Prefix string
@@ -37,6 +42,9 @@ type Config struct {
TagsFormat statsd.TagFormat TagsFormat statsd.TagFormat
Enabled bool Enabled bool
} }
Prometheus struct {
Prefix string
}
Secret []byte Secret []byte
AdTag []byte AdTag []byte
@@ -115,9 +123,11 @@ func NewConfig(debug, verbose bool, // nolint: gocyclo
bindIP, publicIPv4, publicIPv6, statsIP net.IP, bindIP, publicIPv4, publicIPv6, statsIP net.IP,
bindPort, publicIPv4Port, publicIPv6Port, statsPort, statsdPort uint16, bindPort, publicIPv4Port, publicIPv6Port, statsPort, statsdPort uint16,
statsdIP, statsdNetwork, statsdPrefix, statsdTagsFormat string, statsdIP, statsdNetwork, statsdPrefix, statsdTagsFormat string,
statsdTags map[string]string, statsdTags map[string]string, prometheusPrefix string,
secureOnly bool,
antiReplayMaxSize int, antiReplayEvictionTime time.Duration,
secret, adtag []byte) (*Config, error) { secret, adtag []byte) (*Config, error) {
secureMode := false secureMode := secureOnly
if bytes.HasPrefix(secret, []byte{0xdd}) && len(secret) == 17 { if bytes.HasPrefix(secret, []byte{0xdd}) && len(secret) == 17 {
secureMode = true secureMode = true
secret = bytes.TrimPrefix(secret, []byte{0xdd}) secret = bytes.TrimPrefix(secret, []byte{0xdd})
@@ -155,22 +165,26 @@ func NewConfig(debug, verbose bool, // nolint: gocyclo
} }
conf := &Config{ conf := &Config{
Debug: debug, Debug: debug,
Verbose: verbose, Verbose: verbose,
BindIP: bindIP, SecureOnly: secureOnly,
BindPort: bindPort, BindIP: bindIP,
PublicIPv4: publicIPv4, BindPort: bindPort,
PublicIPv4Port: publicIPv4Port, PublicIPv4: publicIPv4,
PublicIPv6: publicIPv6, PublicIPv4Port: publicIPv4Port,
PublicIPv6Port: publicIPv6Port, PublicIPv6: publicIPv6,
StatsIP: statsIP, PublicIPv6Port: publicIPv6Port,
StatsPort: statsPort, StatsIP: statsIP,
Secret: secret, StatsPort: statsPort,
AdTag: adtag, Secret: secret,
SecureMode: secureMode, AdTag: adtag,
ReadBufferSize: int(readBufferSize), SecureMode: secureMode,
WriteBufferSize: int(writeBufferSize), ReadBufferSize: int(readBufferSize),
WriteBufferSize: int(writeBufferSize),
AntiReplayMaxSize: antiReplayMaxSize,
AntiReplayEvictionTime: antiReplayEvictionTime,
} }
conf.Prometheus.Prefix = prometheusPrefix
if statsdIP != "" { if statsdIP != "" {
conf.StatsD.Enabled = true conf.StatsD.Enabled = true
+1 -1
View File
@@ -21,7 +21,7 @@ func getGlobalIPv6() (net.IP, error) {
} }
func fetchIP(network string) (net.IP, error) { func fetchIP(network string) (net.IP, error) {
dialer := &net.Dialer{DualStack: false} dialer := &net.Dialer{FallbackDelay: -1}
client := &http.Client{ client := &http.Client{
Jar: nil, Jar: nil,
Transport: &http.Transport{ Transport: &http.Transport{
+19 -15
View File
@@ -1,26 +1,30 @@
module github.com/9seconds/mtg module github.com/9seconds/mtg
replace github.com/golang/lint => github.com/golang/lint v0.0.0-20190227174305-8f45f776aaf1
require ( require (
github.com/alecthomas/template v0.0.0-20160405071501-a0175ee3bccc // indirect github.com/OneOfOne/xxhash v1.2.5 // indirect
github.com/alecthomas/units v0.0.0-20151022065526-2efee857e7cf // indirect github.com/allegro/bigcache v1.2.0
github.com/beevik/ntp v0.2.0 github.com/beevik/ntp v0.2.0
github.com/davecgh/go-spew v1.1.1 // indirect github.com/cespare/xxhash v1.1.0
github.com/dustin/go-humanize v0.0.0-20180713052910-9f541cc9db5d github.com/dustin/go-humanize v1.0.0
github.com/gofrs/uuid v3.1.0+incompatible github.com/gofrs/uuid v3.2.0+incompatible
github.com/juju/errors v0.0.0-20180806074554-22422dad46e1 github.com/juju/errors v0.0.0-20190207033735-e65537c515d7
github.com/juju/loggo v0.0.0-20180524022052-584905176618 // indirect github.com/juju/loggo v0.0.0-20190526231331-6e530bcce5d8 // indirect
github.com/juju/testing v0.0.0-20180920084828-472a3e8b2073 // indirect github.com/juju/testing v0.0.0-20190429233213-dfc56b8c09fc // indirect
github.com/kr/pretty v0.1.0 // indirect github.com/kr/pretty v0.1.0 // indirect
github.com/pkg/errors v0.8.0 // indirect github.com/pkg/errors v0.8.1 // indirect
github.com/pmezard/go-difflib v1.0.0 // indirect github.com/prometheus/client_golang v0.9.4
github.com/stretchr/testify v1.2.2 github.com/spaolacci/murmur3 v1.1.0 // indirect
go.uber.org/atomic v1.3.2 // indirect github.com/stretchr/testify v1.3.0
go.uber.org/atomic v1.4.0 // indirect
go.uber.org/multierr v1.1.0 // indirect go.uber.org/multierr v1.1.0 // indirect
go.uber.org/zap v1.9.1 go.uber.org/zap v1.10.0
golang.org/x/net v0.0.0-20180921000356-2f5d2388922f // indirect golang.org/x/net v0.0.0-20190607181551-461777fb6f67 // indirect
golang.org/x/sys v0.0.0-20190610200419-93c9922d18ae // indirect
gopkg.in/alecthomas/kingpin.v2 v2.2.6 gopkg.in/alecthomas/kingpin.v2 v2.2.6
gopkg.in/alexcesaro/statsd.v2 v2.0.0 gopkg.in/alexcesaro/statsd.v2 v2.0.0
gopkg.in/check.v1 v1.0.0-20180628173108-788fd7840127 // indirect gopkg.in/check.v1 v1.0.0-20180628173108-788fd7840127 // indirect
gopkg.in/mgo.v2 v2.0.0-20180705113604-9856a29383ce // indirect gopkg.in/mgo.v2 v2.0.0-20180705113604-9856a29383ce // indirect
gopkg.in/yaml.v2 v2.2.1 // indirect gopkg.in/yaml.v2 v2.2.2 // indirect
) )
+83 -16
View File
@@ -1,40 +1,105 @@
github.com/OneOfOne/xxhash v1.2.2/go.mod h1:HSdplMjZKSmBqAxg5vPj2TmRDmfkzw+cTzAElWljhcU=
github.com/OneOfOne/xxhash v1.2.5 h1:zl/OfRA6nftbBK9qTohYBJ5xvw6C/oNKizR7cZGl3cI=
github.com/OneOfOne/xxhash v1.2.5/go.mod h1:eZbhyaAYD41SGSSsnmcpxVoRiQ/MPUTjUdIIOT9Um7Q=
github.com/alecthomas/template v0.0.0-20160405071501-a0175ee3bccc h1:cAKDfWh5VpdgMhJosfJnn5/FoN2SRZ4p7fJNX58YPaU= github.com/alecthomas/template v0.0.0-20160405071501-a0175ee3bccc h1:cAKDfWh5VpdgMhJosfJnn5/FoN2SRZ4p7fJNX58YPaU=
github.com/alecthomas/template v0.0.0-20160405071501-a0175ee3bccc/go.mod h1:LOuyumcjzFXgccqObfd/Ljyb9UuFJ6TxHnclSeseNhc= github.com/alecthomas/template v0.0.0-20160405071501-a0175ee3bccc/go.mod h1:LOuyumcjzFXgccqObfd/Ljyb9UuFJ6TxHnclSeseNhc=
github.com/alecthomas/units v0.0.0-20151022065526-2efee857e7cf h1:qet1QNfXsQxTZqLG4oE62mJzwPIB8+Tee4RNCL9ulrY= github.com/alecthomas/units v0.0.0-20151022065526-2efee857e7cf h1:qet1QNfXsQxTZqLG4oE62mJzwPIB8+Tee4RNCL9ulrY=
github.com/alecthomas/units v0.0.0-20151022065526-2efee857e7cf/go.mod h1:ybxpYRFXyAe+OPACYpWeL0wqObRcbAqCMya13uyzqw0= github.com/alecthomas/units v0.0.0-20151022065526-2efee857e7cf/go.mod h1:ybxpYRFXyAe+OPACYpWeL0wqObRcbAqCMya13uyzqw0=
github.com/allegro/bigcache v1.2.0 h1:qDaE0QoF29wKBb3+pXFrJFy1ihe5OT9OiXhg1t85SxM=
github.com/allegro/bigcache v1.2.0/go.mod h1:Cb/ax3seSYIx7SuZdm2G2xzfwmv3TPSk2ucNfQESPXM=
github.com/beevik/ntp v0.2.0 h1:sGsd+kAXzT0bfVfzJfce04g+dSRfrs+tbQW8lweuYgw= github.com/beevik/ntp v0.2.0 h1:sGsd+kAXzT0bfVfzJfce04g+dSRfrs+tbQW8lweuYgw=
github.com/beevik/ntp v0.2.0/go.mod h1:hIHWr+l3+/clUnF44zdK+CWW7fO8dR5cIylAQ76NRpg= github.com/beevik/ntp v0.2.0/go.mod h1:hIHWr+l3+/clUnF44zdK+CWW7fO8dR5cIylAQ76NRpg=
github.com/beorn7/perks v0.0.0-20180321164747-3a771d992973 h1:xJ4a3vCFaGF/jqvzLMYoU8P317H5OQ+Via4RmuPwCS0=
github.com/beorn7/perks v0.0.0-20180321164747-3a771d992973/go.mod h1:Dwedo/Wpr24TaqPxmxbtue+5NUziq4I4S80YR8gNf3Q=
github.com/beorn7/perks v1.0.0 h1:HWo1m869IqiPhD389kmkxeTalrjNbbJTC8LXupb+sl0=
github.com/beorn7/perks v1.0.0/go.mod h1:KWe93zE9D1o94FZ5RNwFwVgaQK1VOXiVxmqh+CedLV8=
github.com/cespare/xxhash v1.1.0 h1:a6HrQnmkObjyL+Gs60czilIUGqrzKutQD6XZog3p+ko=
github.com/cespare/xxhash v1.1.0/go.mod h1:XrSqR1VqqWfGrhpAt58auRo0WTKS1nRRg3ghfAqPWnc=
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/dustin/go-humanize v0.0.0-20180713052910-9f541cc9db5d h1:lDrio3iIdNb0Gw9CgH7cQF+iuB5mOOjdJ9ERNJCBgb4= github.com/dustin/go-humanize v1.0.0 h1:VSnTsYCnlFHaM2/igO1h6X3HA71jcobQuxemgkq4zYo=
github.com/dustin/go-humanize v0.0.0-20180713052910-9f541cc9db5d/go.mod h1:HtrtbFcZ19U5GC7JDqmcUSB87Iq5E25KnS6fMYU6eOk= github.com/dustin/go-humanize v1.0.0/go.mod h1:HtrtbFcZ19U5GC7JDqmcUSB87Iq5E25KnS6fMYU6eOk=
github.com/gofrs/uuid v3.1.0+incompatible h1:q2rtkjaKT4YEr6E1kamy0Ha4RtepWlQBedyHx0uzKwA= github.com/go-kit/kit v0.8.0/go.mod h1:xBxKIO96dXMWWy0MnWVtmwkA9/13aqxPnvrjFYMA2as=
github.com/gofrs/uuid v3.1.0+incompatible/go.mod h1:b2aQJv3Z4Fp6yNu3cdSllBxTCLRxnplIgP/c0N/04lM= github.com/go-logfmt/logfmt v0.3.0/go.mod h1:Qt1PoO58o5twSAckw1HlFXLmHsOX5/0LbT9GBnD5lWE=
github.com/juju/errors v0.0.0-20180806074554-22422dad46e1 h1:wnhMXidtb70kDZCeLt/EfsVtkXS5c8zLnE9y/6DIRAU= github.com/go-stack/stack v1.8.0/go.mod h1:v0f6uXyyMGvRgIKkXu+yp6POWl0qKG85gN/melR3HDY=
github.com/juju/errors v0.0.0-20180806074554-22422dad46e1/go.mod h1:W54LbzXuIE0boCoNJfwqpmkKJ1O4TCTZMetAt6jGk7Q= github.com/gofrs/uuid v3.2.0+incompatible h1:y12jRkkFxsd7GpqdSZ+/KCs/fJbqpEXSGd4+jfEaewE=
github.com/juju/loggo v0.0.0-20180524022052-584905176618 h1:MK144iBQF9hTSwBW/9eJm034bVoG30IshVm688T2hi8= github.com/gofrs/uuid v3.2.0+incompatible/go.mod h1:b2aQJv3Z4Fp6yNu3cdSllBxTCLRxnplIgP/c0N/04lM=
github.com/juju/loggo v0.0.0-20180524022052-584905176618/go.mod h1:vgyd7OREkbtVEN/8IXZe5Ooef3LQePvuBm9UWj6ZL8U= github.com/gogo/protobuf v1.1.1 h1:72R+M5VuhED/KujmZVcIquuo8mBgX4oVda//DQb3PXo=
github.com/juju/testing v0.0.0-20180920084828-472a3e8b2073 h1:WQM1NildKThwdP7qWrNAFGzp4ijNLw8RlgENkaI4MJs= github.com/gogo/protobuf v1.1.1/go.mod h1:r8qH/GZQm5c6nD/R0oafs1akxWv10x8SbQlK7atdtwQ=
github.com/juju/testing v0.0.0-20180920084828-472a3e8b2073/go.mod h1:63prj8cnj0tU0S9OHjGJn+b1h0ZghCndfnbQolrYTwA= github.com/golang/protobuf v1.2.0 h1:P3YflyNX/ehuJFLhxviNdFxQPkGK5cDcApsge1SqnvM=
github.com/golang/protobuf v1.2.0/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U=
github.com/golang/protobuf v1.3.1 h1:YF8+flBXS5eO826T4nzqPrxfhQThhXl0YzfuUPu4SBg=
github.com/golang/protobuf v1.3.1/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U=
github.com/json-iterator/go v1.1.6/go.mod h1:+SdeFBvtyEkXs7REEP0seUULqWtbJapLOCVDaaPEHmU=
github.com/juju/errors v0.0.0-20190207033735-e65537c515d7 h1:dMIPRDg6gi7CUp0Kj2+HxqJ5kTr1iAdzsXYIrLCNSmU=
github.com/juju/errors v0.0.0-20190207033735-e65537c515d7/go.mod h1:W54LbzXuIE0boCoNJfwqpmkKJ1O4TCTZMetAt6jGk7Q=
github.com/juju/loggo v0.0.0-20190526231331-6e530bcce5d8 h1:UUHMLvzt/31azWTN/ifGWef4WUqvXk0iRqdhdy/2uzI=
github.com/juju/loggo v0.0.0-20190526231331-6e530bcce5d8/go.mod h1:vgyd7OREkbtVEN/8IXZe5Ooef3LQePvuBm9UWj6ZL8U=
github.com/juju/testing v0.0.0-20190429233213-dfc56b8c09fc h1:5xUWujf6ES9tEpFHFzI34vcHm8U07lGjxAuJML3qwqM=
github.com/juju/testing v0.0.0-20190429233213-dfc56b8c09fc/go.mod h1:63prj8cnj0tU0S9OHjGJn+b1h0ZghCndfnbQolrYTwA=
github.com/julienschmidt/httprouter v1.2.0/go.mod h1:SYymIcj16QtmaHHD7aYtjjsJG7VTCxuUUipMqKk8s4w=
github.com/konsorten/go-windows-terminal-sequences v1.0.1/go.mod h1:T0+1ngSBFLxvqU3pZ+m/2kptfBszLMUkC4ZK/EgS/cQ=
github.com/kr/logfmt v0.0.0-20140226030751-b84e30acd515/go.mod h1:+0opPa2QZZtGFBFZlji/RkVcI2GknAs/DXo4wKdlNEc=
github.com/kr/pretty v0.1.0 h1:L/CwN0zerZDmRFUapSPitk6f+Q3+0za1rQkzVuMiMFI= github.com/kr/pretty v0.1.0 h1:L/CwN0zerZDmRFUapSPitk6f+Q3+0za1rQkzVuMiMFI=
github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo= github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo=
github.com/kr/pty v1.1.1/go.mod h1:pFQYn66WHrOpPYNljwOMqo10TkYh1fy3cYio2l3bCsQ= github.com/kr/pty v1.1.1/go.mod h1:pFQYn66WHrOpPYNljwOMqo10TkYh1fy3cYio2l3bCsQ=
github.com/kr/text v0.1.0 h1:45sCR5RtlFHMR4UwH9sdQ5TC8v0qDQCHnXt+kaKSTVE= github.com/kr/text v0.1.0 h1:45sCR5RtlFHMR4UwH9sdQ5TC8v0qDQCHnXt+kaKSTVE=
github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI= github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI=
github.com/matttproud/golang_protobuf_extensions v1.0.1 h1:4hp9jkHxhMHkqkrB3Ix0jegS5sx/RkqARlsWZ6pIwiU=
github.com/matttproud/golang_protobuf_extensions v1.0.1/go.mod h1:D8He9yQNgCq6Z5Ld7szi9bcBfOoFv/3dc6xSMkL2PC0=
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q=
github.com/modern-go/reflect2 v1.0.1/go.mod h1:bx2lNnkwVCuqBIxFjflWJWanXIb3RllmbCylyMrvgv0=
github.com/mwitkow/go-conntrack v0.0.0-20161129095857-cc309e4a2223/go.mod h1:qRWi+5nqEBWmkhHvq77mSJWrCKwh8bxhgT7d/eI7P4U=
github.com/pkg/errors v0.8.0 h1:WdK/asTD0HN+q6hsWO3/vpuAkAr+tw6aNJNDFFf0+qw= github.com/pkg/errors v0.8.0 h1:WdK/asTD0HN+q6hsWO3/vpuAkAr+tw6aNJNDFFf0+qw=
github.com/pkg/errors v0.8.0/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pkg/errors v0.8.0/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
github.com/pkg/errors v0.8.1 h1:iURUrRGxPUNPdy5/HRSm+Yj6okJ6UtLINN0Q9M4+h3I=
github.com/pkg/errors v0.8.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/prometheus/client_golang v0.9.1/go.mod h1:7SWBe2y4D6OKWSNQJUaRYU/AaXPKyh/dDVn+NZz0KFw=
github.com/prometheus/client_golang v0.9.4 h1:Y8E/JaaPbmFSW2V81Ab/d8yZFYQQGbni1b1jPcG9Y6A=
github.com/prometheus/client_golang v0.9.4/go.mod h1:oCXIBxdI62A4cR6aTRJCgetEjecSIYzOEaeAn4iYEpM=
github.com/prometheus/client_model v0.0.0-20180712105110-5c3871d89910 h1:idejC8f05m9MGOsuEi1ATq9shN03HrxNkD/luQvxCv8=
github.com/prometheus/client_model v0.0.0-20180712105110-5c3871d89910/go.mod h1:MbSGuTsp3dbXC40dX6PRTWyKYBIrTGTE9sqQNg2J8bo=
github.com/prometheus/client_model v0.0.0-20190129233127-fd36f4220a90 h1:S/YWwWx/RA8rT8tKFRuGUZhuA90OyIBpPCXkcbwU8DE=
github.com/prometheus/client_model v0.0.0-20190129233127-fd36f4220a90/go.mod h1:xMI15A0UPsDsEKsMN9yxemIoYk6Tm2C1GtYGdfGttqA=
github.com/prometheus/common v0.4.1 h1:K0MGApIoQvMw27RTdJkPbr3JZ7DNbtxQNyi5STVM6Kw=
github.com/prometheus/common v0.4.1/go.mod h1:TNfzLD0ON7rHzMJeJkieUDPYmFC7Snx/y86RQel1bk4=
github.com/prometheus/procfs v0.0.0-20181005140218-185b4288413d h1:GoAlyOgbOEIFdaDqxJVlbOQ1DtGmZWs/Qau0hIlk+WQ=
github.com/prometheus/procfs v0.0.0-20181005140218-185b4288413d/go.mod h1:c3At6R/oaqEKCNdg8wHV1ftS6bRYblBhIjjI8uT2IGk=
github.com/prometheus/procfs v0.0.2 h1:6LJUbpNm42llc4HRCuvApCSWB/WfhuNo9K98Q9sNGfs=
github.com/prometheus/procfs v0.0.2/go.mod h1:TjEm7ze935MbeOT/UhFTIMYKhuLP4wbCsTZCD3I8kEA=
github.com/sirupsen/logrus v1.2.0/go.mod h1:LxeOpSwHxABJmUn/MG1IvRgCAasNZTLOkJPxbbu5VWo=
github.com/spaolacci/murmur3 v0.0.0-20180118202830-f09979ecbc72/go.mod h1:JwIasOWyU6f++ZhiEuf87xNszmSA2myDM2Kzu9HwQUA=
github.com/spaolacci/murmur3 v1.1.0 h1:7c1g84S4BPRrfL5Xrdp6fOJ206sU9y293DDHaoy0bLI=
github.com/spaolacci/murmur3 v1.1.0/go.mod h1:JwIasOWyU6f++ZhiEuf87xNszmSA2myDM2Kzu9HwQUA=
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
github.com/stretchr/objx v0.1.1/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
github.com/stretchr/testify v1.2.2 h1:bSDNvY7ZPG5RlJ8otE/7V6gMiyenm9RtJ7IUVIAoJ1w= github.com/stretchr/testify v1.2.2 h1:bSDNvY7ZPG5RlJ8otE/7V6gMiyenm9RtJ7IUVIAoJ1w=
github.com/stretchr/testify v1.2.2/go.mod h1:a8OnRcib4nhh0OaRAV+Yts87kKdq0PP7pXfy6kDkUVs= github.com/stretchr/testify v1.2.2/go.mod h1:a8OnRcib4nhh0OaRAV+Yts87kKdq0PP7pXfy6kDkUVs=
go.uber.org/atomic v1.3.2 h1:2Oa65PReHzfn29GpvgsYwloV9AVFHPDk8tYxt2c2tr4= github.com/stretchr/testify v1.3.0 h1:TivCn/peBQ7UY8ooIcPgZFpTNSz0Q2U6UrFlUfqbe0Q=
go.uber.org/atomic v1.3.2/go.mod h1:gD2HeocX3+yG+ygLZcrzQJaqmWj9AIm7n08wl/qW/PE= github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI=
go.uber.org/atomic v1.4.0 h1:cxzIVoETapQEqDhQu3QfnvXAV4AlzcvUCxkVUFw3+EU=
go.uber.org/atomic v1.4.0/go.mod h1:gD2HeocX3+yG+ygLZcrzQJaqmWj9AIm7n08wl/qW/PE=
go.uber.org/multierr v1.1.0 h1:HoEmRHQPVSqub6w2z2d2EOVs2fjyFRGyofhKuyDq0QI= go.uber.org/multierr v1.1.0 h1:HoEmRHQPVSqub6w2z2d2EOVs2fjyFRGyofhKuyDq0QI=
go.uber.org/multierr v1.1.0/go.mod h1:wR5kodmAFQ0UK8QlbwjlSNy0Z68gJhDJUG5sjR94q/0= go.uber.org/multierr v1.1.0/go.mod h1:wR5kodmAFQ0UK8QlbwjlSNy0Z68gJhDJUG5sjR94q/0=
go.uber.org/zap v1.9.1 h1:XCJQEf3W6eZaVwhRBof6ImoYGJSITeKWsyeh3HFu/5o= go.uber.org/zap v1.10.0 h1:ORx85nbTijNz8ljznvCMR1ZBIPKFn3jQrag10X2AsuM=
go.uber.org/zap v1.9.1/go.mod h1:vwi/ZaCAaUcBkycHslxD9B2zi4UTXhF60s6SWpuDF0Q= go.uber.org/zap v1.10.0/go.mod h1:vwi/ZaCAaUcBkycHslxD9B2zi4UTXhF60s6SWpuDF0Q=
golang.org/x/net v0.0.0-20180921000356-2f5d2388922f h1:QM2QVxvDoW9PFSPp/zy9FgxJLfaWTZlS61KEPtBwacM= golang.org/x/crypto v0.0.0-20180904163835-0709b304e793/go.mod h1:6SG95UA2DQfeDnfUPMdvaQW0Q7yPrPDi9nlGo2tz2b4=
golang.org/x/net v0.0.0-20180921000356-2f5d2388922f/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w=
golang.org/x/net v0.0.0-20181114220301-adae6a3d119a/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4=
golang.org/x/net v0.0.0-20190607181551-461777fb6f67 h1:rJJxsykSlULwd2P2+pg/rtnwN2FrWp4IuCxOSyS0V00=
golang.org/x/net v0.0.0-20190607181551-461777fb6f67/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s=
golang.org/x/sync v0.0.0-20181108010431-42b317875d0f/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sync v0.0.0-20181221193216-37e7f081c4d4 h1:YUO/7uOKsKeq9UokNS62b8FYywz3ker1l1vDZRCRefw=
golang.org/x/sync v0.0.0-20181221193216-37e7f081c4d4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sys v0.0.0-20180905080454-ebe1bf3edb33/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
golang.org/x/sys v0.0.0-20181116152217-5ac8a444bdc5/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
golang.org/x/sys v0.0.0-20190610200419-93c9922d18ae h1:xiXzMMEQdQcric9hXtr1QU98MHunKK7OTtsoU6bYWs4=
golang.org/x/sys v0.0.0-20190610200419-93c9922d18ae/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
gopkg.in/alecthomas/kingpin.v2 v2.2.6 h1:jMFz6MfLP0/4fUyZle81rXUoxOBFi19VUFKVDOQfozc= gopkg.in/alecthomas/kingpin.v2 v2.2.6 h1:jMFz6MfLP0/4fUyZle81rXUoxOBFi19VUFKVDOQfozc=
gopkg.in/alecthomas/kingpin.v2 v2.2.6/go.mod h1:FMv+mEhP44yOT+4EoQTLFTRgOQ1FBLkstjWtayDeSgw= gopkg.in/alecthomas/kingpin.v2 v2.2.6/go.mod h1:FMv+mEhP44yOT+4EoQTLFTRgOQ1FBLkstjWtayDeSgw=
gopkg.in/alexcesaro/statsd.v2 v2.0.0 h1:FXkZSCZIH17vLCO5sO2UucTHsH9pc+17F6pl3JVCwMc= gopkg.in/alexcesaro/statsd.v2 v2.0.0 h1:FXkZSCZIH17vLCO5sO2UucTHsH9pc+17F6pl3JVCwMc=
@@ -46,3 +111,5 @@ gopkg.in/mgo.v2 v2.0.0-20180705113604-9856a29383ce h1:xcEWjVhvbDy+nHP67nPDDpbYrY
gopkg.in/mgo.v2 v2.0.0-20180705113604-9856a29383ce/go.mod h1:yeKp02qBN3iKW1OzL3MGk2IdtZzaj7SFntXj72NppTA= gopkg.in/mgo.v2 v2.0.0-20180705113604-9856a29383ce/go.mod h1:yeKp02qBN3iKW1OzL3MGk2IdtZzaj7SFntXj72NppTA=
gopkg.in/yaml.v2 v2.2.1 h1:mUhvW9EsL+naU5Q3cakzfE91YhliOondGd6ZrsDBHQE= gopkg.in/yaml.v2 v2.2.1 h1:mUhvW9EsL+naU5Q3cakzfE91YhliOondGd6ZrsDBHQE=
gopkg.in/yaml.v2 v2.2.1/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= gopkg.in/yaml.v2 v2.2.1/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
gopkg.in/yaml.v2 v2.2.2 h1:ZCJp+EgiOT7lHqUV2J862kp8Qj64Jo6az82+3Td9dZw=
gopkg.in/yaml.v2 v2.2.2/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
+33 -8
View File
@@ -1,7 +1,5 @@
package main package main
//go:generate scripts/generate_version.sh
import ( import (
"encoding/json" "encoding/json"
"fmt" "fmt"
@@ -22,6 +20,8 @@ import (
"github.com/9seconds/mtg/stats" "github.com/9seconds/mtg/stats"
) )
var version = "dev" // this has to be set by build ld flags
var ( var (
app = kingpin.New("mtg", "Simple MTPROTO proxy.") app = kingpin.New("mtg", "Simple MTPROTO proxy.")
@@ -110,6 +110,12 @@ var (
Envar("MTG_STATSD_TAGS"). Envar("MTG_STATSD_TAGS").
StringMap() StringMap()
prometheusPrefix = app.Flag("prometheus-prefix",
"Which namespace to use to send stats to Prometheus.").
Envar("MTG_PROMETHEUS_PREFIX").
Default("mtg").
String()
writeBufferSize = app.Flag("write-buffer", writeBufferSize = app.Flag("write-buffer",
"Write buffer size in bytes. You can think about it as a buffer from client to Telegram."). "Write buffer size in bytes. You can think about it as a buffer from client to Telegram.").
Short('w'). Short('w').
@@ -122,18 +128,32 @@ var (
Envar("MTG_BUFFER_READ"). Envar("MTG_BUFFER_READ").
Default("131072"). Default("131072").
Uint32() Uint32()
secureOnly = app.Flag("secure-only",
"Support clients with dd-secrets only.").
Short('s').
Envar("MTG_SECURE_ONLY").
Bool()
antiReplayMaxSize = app.Flag("anti-replay-max-size",
"Max size of antireplay cache in megabytes.").
Envar("MTG_ANTIREPLAY_MAXSIZE").
Default("128").
Int()
antiReplayEvictionTime = app.Flag("anti-replay-eviction-time",
"Eviction time period for obfuscated2 handshakes").
Envar("MTG_ANTIREPLAY_EVICTIONTIME").
Default("168h").
Duration()
secret = app.Arg("secret", "Secret of this proxy.").Required().HexBytes() secret = app.Arg("secret", "Secret of this proxy.").Required().HexBytes()
adtag = app.Arg("adtag", "ADTag of the proxy.").HexBytes() adtag = app.Arg("adtag", "ADTag of the proxy.").HexBytes()
) )
func init() { func main() { // nolint: gocyclo
rand.Seed(time.Now().UTC().UnixNano()) rand.Seed(time.Now().UTC().UnixNano())
app.Version(version) app.Version(version)
app.HelpFlag.Short('h') app.HelpFlag.Short('h')
}
func main() { // nolint: gocyclo
kingpin.MustParse(app.Parse(os.Args[1:])) kingpin.MustParse(app.Parse(os.Args[1:]))
err := setRLimit() err := setRLimit()
@@ -146,7 +166,8 @@ func main() { // nolint: gocyclo
*bindIP, *publicIPv4, *publicIPv6, *statsIP, *bindIP, *publicIPv4, *publicIPv6, *statsIP,
*bindPort, *publicIPv4Port, *publicIPv6Port, *statsPort, *statsdPort, *bindPort, *publicIPv4Port, *publicIPv6Port, *statsPort, *statsdPort,
*statsdIP, *statsdNetwork, *statsdPrefix, *statsdTagsFormat, *statsdIP, *statsdNetwork, *statsdPrefix, *statsdTagsFormat,
*statsdTags, *statsdTags, *prometheusPrefix, *secureOnly,
*antiReplayMaxSize, *antiReplayEvictionTime,
*secret, *adtag, *secret, *adtag,
) )
if err != nil { if err != nil {
@@ -172,6 +193,7 @@ func main() { // nolint: gocyclo
defer logger.Sync() // nolint: errcheck defer logger.Sync() // nolint: errcheck
printURLs(conf.GetURLs()) printURLs(conf.GetURLs())
zap.S().Debugw("Configuration", "config", conf)
if conf.UseMiddleProxy() { if conf.UseMiddleProxy() {
zap.S().Infow("Use middle proxy connection to Telegram") zap.S().Infow("Use middle proxy connection to Telegram")
@@ -188,11 +210,14 @@ func main() { // nolint: gocyclo
zap.S().Infow("Use direct connection to Telegram") zap.S().Infow("Use direct connection to Telegram")
} }
if err := stats.Start(conf); err != nil { if err := stats.Init(conf); err != nil {
panic(err) panic(err)
} }
server := proxy.NewProxy(conf) server, err := proxy.NewProxy(conf)
if err != nil {
panic(err)
}
if err := server.Serve(); err != nil { if err := server.Serve(); err != nil {
zap.S().Fatalw("Server stopped", "error", err) zap.S().Fatalw("Server stopped", "error", err)
} }
+2 -7
View File
@@ -23,11 +23,6 @@ var (
ProxyRequestExtraSize = []byte{0x18, 0x00, 0x00, 0x00} ProxyRequestExtraSize = []byte{0x18, 0x00, 0x00, 0x00}
ProxyRequestProxyTag = []byte{0xae, 0x26, 0x1e, 0xdb} ProxyRequestProxyTag = []byte{0xae, 0x26, 0x1e, 0xdb}
HandshakeSenderPID []byte
HandshakePeerPID []byte
)
func init() {
HandshakeSenderPID = []byte("IPIPPRPDTIME") HandshakeSenderPID = []byte("IPIPPRPDTIME")
HandshakePeerPID = []byte("IPIPPRPDTIME") HandshakePeerPID = []byte("IPIPPRPDTIME")
} )
+3 -3
View File
@@ -78,10 +78,10 @@ func TestFrameGenerateValid(t *testing.T) {
} }
for _, test := range validTests { for _, test := range validTests {
t.Run(strconv.Itoa(int(test)), func(tt *testing.T) { t.Run(strconv.Itoa(int(test)), func(tt *testing.T) {
frame := generateFrame(test) frame := generateFrame(test) // nolint: scopelint
conType, err := frame.ConnectionType() conType, err := frame.ConnectionType()
assert.Nil(t, err) assert.Nil(tt, err)
assert.Equal(t, conType, test) assert.Equal(tt, conType, test) // nolint: scopelint
}) })
} }
} }
+28 -25
View File
@@ -10,6 +10,7 @@ import (
"github.com/juju/errors" "github.com/juju/errors"
"go.uber.org/zap" "go.uber.org/zap"
"github.com/9seconds/mtg/antireplay"
"github.com/9seconds/mtg/client" "github.com/9seconds/mtg/client"
"github.com/9seconds/mtg/config" "github.com/9seconds/mtg/config"
"github.com/9seconds/mtg/mtproto" "github.com/9seconds/mtg/mtproto"
@@ -20,9 +21,10 @@ import (
// Proxy is a core of this program. // Proxy is a core of this program.
type Proxy struct { type Proxy struct {
clientInit client.Init antiReplayCache antireplay.Cache
tg telegram.Telegram clientInit client.Init
conf *config.Config tg telegram.Telegram
conf *config.Config
} }
// Serve runs TCP proxy server. // Serve runs TCP proxy server.
@@ -58,13 +60,18 @@ func (p *Proxy) accept(conn net.Conn) {
log.Infow("Client connected", "addr", conn.RemoteAddr()) log.Infow("Client connected", "addr", conn.RemoteAddr())
clientConn, opts, err := p.clientInit(ctx, cancel, conn, connID, p.conf) clientConn, opts, err := p.clientInit(ctx, cancel, conn, connID, p.antiReplayCache, p.conf)
if err != nil { if err != nil {
log.Errorw("Cannot initialize client connection", "error", err) log.Errorw("Cannot initialize client connection", "error", err)
return return
} }
defer clientConn.(io.Closer).Close() // nolint: errcheck defer clientConn.(io.Closer).Close() // nolint: errcheck
if p.conf.SecureOnly && opts.ConnectionType != mtproto.ConnectionTypeSecure {
log.Errorw("Proxy supports only secure connections", "connection_type", opts.ConnectionType)
return
}
stats.ClientConnected(opts.ConnectionType, clientConn.RemoteAddr()) stats.ClientConnected(opts.ConnectionType, clientConn.RemoteAddr())
defer stats.ClientDisconnected(opts.ConnectionType, clientConn.RemoteAddr()) defer stats.ClientDisconnected(opts.ConnectionType, clientConn.RemoteAddr())
@@ -88,12 +95,12 @@ func (p *Proxy) accept(conn net.Conn) {
clientPacket := clientConn.(wrappers.PacketReadWriteCloser) clientPacket := clientConn.(wrappers.PacketReadWriteCloser)
serverPacket := serverConn.(wrappers.PacketReadWriteCloser) serverPacket := serverConn.(wrappers.PacketReadWriteCloser)
go p.middlePipe(clientPacket, serverPacket, wait, &opts.ReadHacks) go p.middlePipe(clientPacket, serverPacket, wait, &opts.ReadHacks)
go p.middlePipe(serverPacket, clientPacket, wait, &opts.WriteHacks) p.middlePipe(serverPacket, clientPacket, wait, &opts.WriteHacks)
} else { } else {
clientStream := clientConn.(wrappers.StreamReadWriteCloser) clientStream := clientConn.(wrappers.StreamReadWriteCloser)
serverStream := serverConn.(wrappers.StreamReadWriteCloser) serverStream := serverConn.(wrappers.StreamReadWriteCloser)
go p.directPipe(clientStream, serverStream, wait, p.conf.ReadBufferSize) go p.directPipe(clientStream, serverStream, wait, p.conf.ReadBufferSize)
go p.directPipe(serverStream, clientStream, wait, p.conf.WriteBufferSize) p.directPipe(serverStream, clientStream, wait, p.conf.WriteBufferSize)
} }
wait.Wait() wait.Wait()
@@ -116,13 +123,8 @@ func (p *Proxy) getTelegramConn(ctx context.Context, cancel context.CancelFunc,
return packetConn, nil return packetConn, nil
} }
func (p *Proxy) middlePipe(src wrappers.PacketReadCloser, dst io.WriteCloser, func (p *Proxy) middlePipe(src wrappers.PacketReadCloser, dst io.Writer, wait *sync.WaitGroup, hacks *mtproto.Hacks) {
wait *sync.WaitGroup, hacks *mtproto.Hacks) { defer wait.Done()
defer func() {
src.Close() // nolint: errcheck, gosec
dst.Close() // nolint: errcheck, gosec
wait.Done()
}()
for { for {
hacks.SimpleAck = false hacks.SimpleAck = false
@@ -140,13 +142,8 @@ func (p *Proxy) middlePipe(src wrappers.PacketReadCloser, dst io.WriteCloser,
} }
} }
func (p *Proxy) directPipe(src wrappers.StreamReadCloser, dst io.WriteCloser, func (p *Proxy) directPipe(src wrappers.StreamReadCloser, dst io.Writer, wait *sync.WaitGroup, bufferSize int) {
wait *sync.WaitGroup, bufferSize int) { defer wait.Done()
defer func() {
src.Close() // nolint: errcheck, gosec
dst.Close() // nolint: errcheck, gosec
wait.Done()
}()
buffer := make([]byte, bufferSize) buffer := make([]byte, bufferSize)
if _, err := io.CopyBuffer(dst, src, buffer); err != nil { if _, err := io.CopyBuffer(dst, src, buffer); err != nil {
@@ -155,10 +152,15 @@ func (p *Proxy) directPipe(src wrappers.StreamReadCloser, dst io.WriteCloser,
} }
// NewProxy returns new proxy instance. // NewProxy returns new proxy instance.
func NewProxy(conf *config.Config) *Proxy { func NewProxy(conf *config.Config) (*Proxy, error) {
var clientInit client.Init var clientInit client.Init
var tg telegram.Telegram var tg telegram.Telegram
cache, err := antireplay.NewCache(conf)
if err != nil {
return nil, errors.Annotate(err, "Cannot make proxy")
}
if conf.UseMiddleProxy() { if conf.UseMiddleProxy() {
clientInit = client.MiddleInit clientInit = client.MiddleInit
tg = telegram.NewMiddleTelegram(conf) tg = telegram.NewMiddleTelegram(conf)
@@ -168,8 +170,9 @@ func NewProxy(conf *config.Config) *Proxy {
} }
return &Proxy{ return &Proxy{
conf: conf, antiReplayCache: cache,
clientInit: clientInit, conf: conf,
tg: tg, clientInit: clientInit,
} tg: tg,
}, nil
} }
-12
View File
@@ -1,12 +0,0 @@
#!/bin/sh
set -eu
PROJECT_DIR="$(git rev-parse --show-toplevel)"
OUTPUT_FILE="${PROJECT_DIR}/version.go"
cat > "$OUTPUT_FILE" <<EOF
package main
// autogenerated by $(basename "$0") on $(date -Ru)
const version = "$(git describe --long --always) ($(go version)) [$(date -Ru)]"
EOF
+17 -86
View File
@@ -2,21 +2,20 @@ package stats
import ( import (
"net" "net"
"time"
"github.com/9seconds/mtg/mtproto" "github.com/9seconds/mtg/mtproto"
) )
const ( const (
crashesChanLength = 1 connectionsChanLength = 10
connectionsChanLength = 20 trafficChanLength = 10
trafficChanLength = 5000
) )
var ( var (
crashesChan = make(chan struct{}, crashesChanLength) crashesChan = make(chan struct{})
connectionsChan = make(chan *connectionData, connectionsChanLength) statsChan = make(chan chan<- Stats)
trafficChan = make(chan *trafficData, trafficChanLength) connectionsChan = make(chan connectionData, connectionsChanLength)
trafficChan = make(chan trafficData, trafficChanLength)
) )
type connectionData struct { type connectionData struct {
@@ -30,81 +29,6 @@ type trafficData struct {
ingress bool ingress bool
} }
func crashManager() {
for range crashesChan {
instance.mutex.RLock()
instance.Crashes++
instance.mutex.RUnlock()
}
}
func connectionManager() {
for event := range connectionsChan {
instance.mutex.RLock()
isIPv4 := event.addr.IP.To4() != nil
var inc uint32 = 1
if !event.connected {
inc = ^uint32(0)
}
switch event.connectionType {
case mtproto.ConnectionTypeAbridged:
if isIPv4 {
instance.Connections.Abridged.IPv4 += inc
} else {
instance.Connections.Abridged.IPv6 += inc
}
case mtproto.ConnectionTypeSecure:
if isIPv4 {
instance.Connections.Secure.IPv4 += inc
} else {
instance.Connections.Secure.IPv6 += inc
}
default:
if isIPv4 {
instance.Connections.Intermediate.IPv4 += inc
} else {
instance.Connections.Intermediate.IPv6 += inc
}
}
instance.mutex.RUnlock()
}
}
func trafficManager() {
speedChan := time.Tick(time.Second)
for {
select {
case event := <-trafficChan:
instance.mutex.RLock()
if event.ingress {
instance.Traffic.Ingress += trafficValue(event.traffic)
instance.speedCurrent.Ingress += trafficSpeedValue(event.traffic)
} else {
instance.Traffic.Egress += trafficValue(event.traffic)
instance.speedCurrent.Egress += trafficSpeedValue(event.traffic)
}
instance.mutex.RUnlock()
case <-speedChan:
instance.mutex.RLock()
instance.Speed.Ingress = instance.speedCurrent.Ingress
instance.Speed.Egress = instance.speedCurrent.Egress
instance.speedCurrent.Ingress = trafficSpeedValue(0)
instance.speedCurrent.Egress = trafficSpeedValue(0)
instance.mutex.RUnlock()
}
}
}
// NewCrash indicates new crash. // NewCrash indicates new crash.
func NewCrash() { func NewCrash() {
crashesChan <- struct{}{} crashesChan <- struct{}{}
@@ -112,7 +36,7 @@ func NewCrash() {
// ClientConnected indicates that new client was connected. // ClientConnected indicates that new client was connected.
func ClientConnected(connectionType mtproto.ConnectionType, addr *net.TCPAddr) { func ClientConnected(connectionType mtproto.ConnectionType, addr *net.TCPAddr) {
connectionsChan <- &connectionData{ connectionsChan <- connectionData{
connectionType: connectionType, connectionType: connectionType,
addr: addr, addr: addr,
connected: true, connected: true,
@@ -121,7 +45,7 @@ func ClientConnected(connectionType mtproto.ConnectionType, addr *net.TCPAddr) {
// ClientDisconnected indicates that client was disconnected. // ClientDisconnected indicates that client was disconnected.
func ClientDisconnected(connectionType mtproto.ConnectionType, addr *net.TCPAddr) { func ClientDisconnected(connectionType mtproto.ConnectionType, addr *net.TCPAddr) {
connectionsChan <- &connectionData{ connectionsChan <- connectionData{
connectionType: connectionType, connectionType: connectionType,
addr: addr, addr: addr,
connected: false, connected: false,
@@ -130,7 +54,7 @@ func ClientDisconnected(connectionType mtproto.ConnectionType, addr *net.TCPAddr
// IngressTraffic accounts new ingress traffic. // IngressTraffic accounts new ingress traffic.
func IngressTraffic(traffic int) { func IngressTraffic(traffic int) {
trafficChan <- &trafficData{ trafficChan <- trafficData{
traffic: traffic, traffic: traffic,
ingress: true, ingress: true,
} }
@@ -138,8 +62,15 @@ func IngressTraffic(traffic int) {
// EgressTraffic accounts new ingress traffic. // EgressTraffic accounts new ingress traffic.
func EgressTraffic(traffic int) { func EgressTraffic(traffic int) {
trafficChan <- &trafficData{ trafficChan <- trafficData{
traffic: traffic, traffic: traffic,
ingress: false, ingress: false,
} }
} }
// GetStats returns a snapshot of Stats instance.
func GetStats() Stats {
rpcChan := make(chan Stats)
statsChan <- rpcChan
return <-rpcChan
}
+28
View File
@@ -0,0 +1,28 @@
package stats
import (
"github.com/juju/errors"
"github.com/9seconds/mtg/config"
)
// Init initializes stats subsystem.
func Init(conf *config.Config) error {
if conf.StatsD.Enabled {
client, err := newStatsd(conf)
if err != nil {
return errors.Annotate(err, "Cannot initialize statsd client")
}
go client.run()
}
prometheus, err := newPrometheus(conf)
if err != nil {
return errors.Annotate(err, "Cannot initialize prometheus client")
}
go prometheus.run()
go NewStats(conf).start()
go startServer(conf, prometheus.getHTTPHandler())
return nil
}
+91
View File
@@ -0,0 +1,91 @@
package stats
import (
"net/http"
"time"
"github.com/juju/errors"
"github.com/prometheus/client_golang/prometheus"
"github.com/prometheus/client_golang/prometheus/promhttp"
"github.com/9seconds/mtg/config"
)
const prometheusPollTime = time.Second
type prometheusExporter struct {
registry prometheus.Gatherer
connections *prometheus.GaugeVec
traffic *prometheus.GaugeVec
speed *prometheus.GaugeVec
crashes prometheus.Gauge
}
func (p *prometheusExporter) run() {
for range time.Tick(prometheusPollTime) {
instance := GetStats()
p.connections.WithLabelValues("abridged", "v4").Set(float64(instance.Connections.Abridged.IPv4))
p.connections.WithLabelValues("abridged", "v6").Set(float64(instance.Connections.Abridged.IPv6))
p.connections.WithLabelValues("intermediate", "v4").Set(float64(instance.Connections.Intermediate.IPv4))
p.connections.WithLabelValues("intermediate", "v6").Set(float64(instance.Connections.Intermediate.IPv6))
p.connections.WithLabelValues("secure", "v4").Set(float64(instance.Connections.Secure.IPv4))
p.connections.WithLabelValues("secure", "v6").Set(float64(instance.Connections.Secure.IPv6))
p.traffic.WithLabelValues("ingress").Set(float64(instance.Traffic.ingress))
p.traffic.WithLabelValues("egress").Set(float64(instance.Traffic.egress))
p.speed.WithLabelValues("ingress").Set(float64(instance.Speed.ingress))
p.speed.WithLabelValues("egress").Set(float64(instance.Speed.egress))
p.crashes.Set(float64(instance.Crashes))
}
}
func (p *prometheusExporter) getHTTPHandler() http.Handler {
return promhttp.HandlerFor(p.registry, promhttp.HandlerOpts{})
}
func newPrometheus(conf *config.Config) (*prometheusExporter, error) {
registry := prometheus.NewRegistry()
connections := prometheus.NewGaugeVec(prometheus.GaugeOpts{
Namespace: conf.Prometheus.Prefix,
Name: "connections",
Help: "Current number of connections to the proxy.",
}, []string{"type", "protocol"})
traffic := prometheus.NewGaugeVec(prometheus.GaugeOpts{
Namespace: conf.Prometheus.Prefix,
Name: "traffic",
Help: "Traffic passed through the proxy in bytes.",
}, []string{"direction"})
speed := prometheus.NewGaugeVec(prometheus.GaugeOpts{
Namespace: conf.Prometheus.Prefix,
Name: "speed",
Help: "Current throughput in bytes per second.",
}, []string{"direction"})
crashes := prometheus.NewGauge(prometheus.GaugeOpts{
Namespace: conf.Prometheus.Prefix,
Name: "crashes",
Help: "How many crashes happened.",
})
if err := registry.Register(connections); err != nil {
return nil, errors.Annotate(err, "Cannot register connections collector")
}
if err := registry.Register(traffic); err != nil {
return nil, errors.Annotate(err, "cannot register traffic collector")
}
if err := registry.Register(speed); err != nil {
return nil, errors.Annotate(err, "cannot register speed collector")
}
if err := registry.Register(crashes); err != nil {
return nil, errors.Annotate(err, "cannot register crashes collector")
}
return &prometheusExporter{
registry: registry,
connections: connections,
traffic: traffic,
speed: speed,
crashes: crashes,
}, nil
}
+10 -39
View File
@@ -3,67 +3,38 @@ package stats
import ( import (
"encoding/json" "encoding/json"
"net/http" "net/http"
"sync"
"time"
"go.uber.org/zap" "go.uber.org/zap"
"github.com/9seconds/mtg/config" "github.com/9seconds/mtg/config"
) )
var instance *stats func startServer(conf *config.Config, prometheusHandler http.Handler) {
// Start starts new statistics server.
func Start(conf *config.Config) error {
log := zap.S().Named("stats") log := zap.S().Named("stats")
instance = &stats{ http.HandleFunc("/", func(w http.ResponseWriter, _ *http.Request) {
URLs: conf.GetURLs(),
Uptime: uptime(time.Now()),
mutex: &sync.RWMutex{},
}
if conf.StatsD.Enabled {
client, err := newStatsd(conf)
if err != nil {
return err
}
go client.run()
}
go crashManager()
go connectionManager()
go trafficManager()
http.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json") w.Header().Set("Content-Type", "application/json")
instance.mutex.Lock() first, err := json.Marshal(GetStats())
first, err := json.Marshal(instance)
instance.mutex.Unlock()
if err != nil { if err != nil {
log.Errorw("Cannot encode json", "error", err) log.Errorw("Cannot encode json", "error", err)
http.Error(w, "Internal server error", 500) http.Error(w, "Internal server error", 500)
return return
} }
interm := map[string]interface{}{} interim := map[string]interface{}{}
json.Unmarshal(first, &interm) // nolint: errcheck, gosec json.Unmarshal(first, &interim) // nolint: errcheck, gosec
encoder := json.NewEncoder(w) encoder := json.NewEncoder(w)
encoder.SetEscapeHTML(false) encoder.SetEscapeHTML(false)
encoder.SetIndent("", " ") encoder.SetIndent("", " ")
if err = encoder.Encode(interm); err != nil { if err = encoder.Encode(interim); err != nil {
log.Errorw("Cannot encode json", "error", err) log.Errorw("Cannot encode json", "error", err)
} }
}) })
http.Handle("/prometheus/", prometheusHandler)
go func() { if err := http.ListenAndServe(conf.StatAddr(), nil); err != nil {
if err := http.ListenAndServe(conf.StatAddr(), nil); err != nil { log.Fatalw("Stats server has been stopped", "error", err)
log.Fatalw("Stats server has been stopped", "error", err) }
}
}()
return nil
} }
+122 -49
View File
@@ -4,12 +4,12 @@ import (
"encoding/json" "encoding/json"
"fmt" "fmt"
"strconv" "strconv"
"sync"
"time" "time"
humanize "github.com/dustin/go-humanize" humanize "github.com/dustin/go-humanize"
"github.com/9seconds/mtg/config" "github.com/9seconds/mtg/config"
"github.com/9seconds/mtg/mtproto"
) )
type uptime time.Time type uptime time.Time
@@ -24,72 +24,73 @@ func (u uptime) MarshalJSON() ([]byte, error) {
return json.Marshal(value) return json.Marshal(value)
} }
type trafficValue uint64 type connectionType struct {
IPv6 uint32 `json:"ipv6"`
func (t trafficValue) MarshalJSON() ([]byte, error) { IPv4 uint32 `json:"ipv4"`
tv := uint64(t)
value := map[string]interface{}{
"bytes": tv,
"human": humanize.Bytes(tv),
}
return json.Marshal(value)
} }
type trafficSpeedValue uint64 type baseConnections struct {
func (t trafficSpeedValue) MarshalJSON() ([]byte, error) {
speed := uint64(t)
value := map[string]interface{}{
"bytes/s": speed,
"human": fmt.Sprintf("%s/S", humanize.Bytes(speed)),
}
return json.Marshal(value)
}
type connections struct {
All connectionType `json:"all"` All connectionType `json:"all"`
Abridged connectionType `json:"abridged"` Abridged connectionType `json:"abridged"`
Intermediate connectionType `json:"intermediate"` Intermediate connectionType `json:"intermediate"`
Secure connectionType `json:"secure"` Secure connectionType `json:"secure"`
} }
type connections struct {
baseConnections
}
func (c connections) MarshalJSON() ([]byte, error) { func (c connections) MarshalJSON() ([]byte, error) {
c.All.IPv4 = c.Abridged.IPv4 + c.Intermediate.IPv4 + c.Secure.IPv4 c.All.IPv4 = c.Abridged.IPv4 + c.Intermediate.IPv4 + c.Secure.IPv4
c.All.IPv6 = c.Abridged.IPv6 + c.Intermediate.IPv6 + c.Secure.IPv6 c.All.IPv6 = c.Abridged.IPv6 + c.Intermediate.IPv6 + c.Secure.IPv6
value := struct { return json.Marshal(c.baseConnections)
All connectionType `json:"all"` }
Abridged connectionType `json:"abridged"`
Intermediate connectionType `json:"intermediate"` type traffic struct {
Secure connectionType `json:"secure"` ingress uint64
}{ egress uint64
All: c.All, }
Abridged: c.Abridged,
Intermediate: c.Intermediate, func (t *traffic) dumpValue(value uint64) map[string]interface{} {
Secure: c.Secure, return map[string]interface{}{
"bytes": value,
"human": humanize.Bytes(value),
}
}
func (t traffic) MarshalJSON() ([]byte, error) {
value := map[string]map[string]interface{}{
"ingress": t.dumpValue(t.ingress),
"egress": t.dumpValue(t.egress),
} }
return json.Marshal(value) return json.Marshal(value)
} }
type connectionType struct {
IPv6 uint32 `json:"ipv6"`
IPv4 uint32 `json:"ipv4"`
}
type traffic struct {
Ingress trafficValue `json:"ingress"`
Egress trafficValue `json:"egress"`
}
type speed struct { type speed struct {
Ingress trafficSpeedValue `json:"ingress"` ingress uint64
Egress trafficSpeedValue `json:"egress"` egress uint64
} }
type stats struct { func (s *speed) dumpValue(value uint64) map[string]interface{} {
return map[string]interface{}{
"bytes/s": value,
"human": fmt.Sprintf("%s/s", humanize.Bytes(value)),
}
}
func (s speed) MarshalJSON() ([]byte, error) {
value := map[string]map[string]interface{}{
"ingress": s.dumpValue(s.ingress),
"egress": s.dumpValue(s.egress),
}
return json.Marshal(value)
}
// Stats represents a statistics of the proxy.
type Stats struct {
URLs config.IPURLs `json:"urls"` URLs config.IPURLs `json:"urls"`
Connections connections `json:"connections"` Connections connections `json:"connections"`
Traffic traffic `json:"traffic"` Traffic traffic `json:"traffic"`
@@ -97,6 +98,78 @@ type stats struct {
Uptime uptime `json:"uptime"` Uptime uptime `json:"uptime"`
Crashes uint32 `json:"crashes"` Crashes uint32 `json:"crashes"`
speedCurrent speed previousTraffic traffic
mutex *sync.RWMutex }
func (s *Stats) start() {
speedChan := time.Tick(time.Second)
for {
select {
case <-speedChan:
s.handleSpeed()
case event := <-trafficChan:
s.handleTraffic(event)
case event := <-connectionsChan:
s.handleConnection(event)
case getStatsChan := <-statsChan:
s.handleGetStats(getStatsChan)
case <-crashesChan:
s.handleCrash()
}
}
}
func (s *Stats) handleTraffic(evt trafficData) {
if evt.ingress {
s.Traffic.ingress += uint64(evt.traffic)
} else {
s.Traffic.egress += uint64(evt.traffic)
}
}
func (s *Stats) handleSpeed() {
s.Speed.ingress = s.Traffic.ingress - s.previousTraffic.ingress
s.Speed.egress = s.Traffic.egress - s.previousTraffic.egress
s.previousTraffic.ingress = s.Traffic.ingress
s.previousTraffic.egress = s.Traffic.egress
}
func (s *Stats) handleConnection(evt connectionData) {
var inc uint32 = 1
if !evt.connected {
inc = ^uint32(0)
}
var conn *connectionType
switch evt.connectionType {
case mtproto.ConnectionTypeAbridged:
conn = &s.Connections.Abridged
case mtproto.ConnectionTypeSecure:
conn = &s.Connections.Secure
default:
conn = &s.Connections.Intermediate
}
if evt.addr.IP.To4() != nil {
conn.IPv4 += inc
} else {
conn.IPv6 += inc
}
}
func (s *Stats) handleGetStats(getStatsChan chan<- Stats) {
getStatsChan <- *s
}
func (s *Stats) handleCrash() {
s.Crashes++
}
// NewStats creates a new instance of Stats structure.
func NewStats(conf *config.Config) *Stats {
return &Stats{
URLs: conf.GetURLs(),
Uptime: uptime(time.Now()),
}
} }
+5 -7
View File
@@ -36,7 +36,7 @@ type statsdExporter struct {
func (s *statsdExporter) run() { func (s *statsdExporter) run() {
for range time.Tick(statsdPollTime) { for range time.Tick(statsdPollTime) {
instance.mutex.Lock() instance := GetStats()
s.client.Gauge(statsdConnectionsAbridgedV4, instance.Connections.Abridged.IPv4) s.client.Gauge(statsdConnectionsAbridgedV4, instance.Connections.Abridged.IPv4)
s.client.Gauge(statsdConnectionsAbridgedV6, instance.Connections.Abridged.IPv6) s.client.Gauge(statsdConnectionsAbridgedV6, instance.Connections.Abridged.IPv6)
@@ -44,13 +44,11 @@ func (s *statsdExporter) run() {
s.client.Gauge(statsdConnectionsIntermediateV6, instance.Connections.Intermediate.IPv6) s.client.Gauge(statsdConnectionsIntermediateV6, instance.Connections.Intermediate.IPv6)
s.client.Gauge(statsdConnectionsSecureV4, instance.Connections.Secure.IPv4) s.client.Gauge(statsdConnectionsSecureV4, instance.Connections.Secure.IPv4)
s.client.Gauge(statsdConnectionsSecureV6, instance.Connections.Secure.IPv6) s.client.Gauge(statsdConnectionsSecureV6, instance.Connections.Secure.IPv6)
s.client.Gauge(statsdTrafficIngress, uint64(instance.Traffic.Ingress)) s.client.Gauge(statsdTrafficIngress, instance.Traffic.ingress)
s.client.Gauge(statsdTrafficEgress, uint64(instance.Traffic.Egress)) s.client.Gauge(statsdTrafficEgress, instance.Traffic.egress)
s.client.Gauge(statsdSpeedIngress, uint64(instance.Speed.Ingress)) s.client.Gauge(statsdSpeedIngress, instance.Speed.ingress)
s.client.Gauge(statsdSpeedEgress, uint64(instance.Speed.Egress)) s.client.Gauge(statsdSpeedEgress, instance.Speed.egress)
s.client.Gauge(statsdCrashes, instance.Crashes) s.client.Gauge(statsdCrashes, instance.Crashes)
instance.mutex.Unlock()
} }
} }
+20 -37
View File
@@ -38,13 +38,6 @@ const (
connTimeoutWrite = 2 * time.Minute connTimeoutWrite = 2 * time.Minute
) )
type ioResult struct {
n int
err error
}
type ioFunc func([]byte) (int, error)
// Conn is a basic wrapper for net.Conn providing the most low-level // Conn is a basic wrapper for net.Conn providing the most low-level
// logic and management as possible. // logic and management as possible.
type Conn struct { type Conn struct {
@@ -61,12 +54,20 @@ type Conn struct {
func (c *Conn) Write(p []byte) (int, error) { func (c *Conn) Write(p []byte) (int, error) {
select { select {
case <-c.ctx.Done(): case <-c.ctx.Done():
c.Close() // nolint: gosec
return 0, errors.Annotate(c.ctx.Err(), "Cannot write because context was closed") return 0, errors.Annotate(c.ctx.Err(), "Cannot write because context was closed")
default: default:
n, err := c.doIO(c.conn.Write, p, connTimeoutWrite) if err := c.conn.SetWriteDeadline(time.Now().Add(connTimeoutWrite)); err != nil {
c.Close() // nolint: gosec
return 0, errors.Annotate(err, "Cannot set write deadline to the socket")
}
n, err := c.conn.Write(p)
c.logger.Debugw("Write to stream", "bytes", n, "error", err) c.logger.Debugw("Write to stream", "bytes", n, "error", err)
stats.EgressTraffic(n) stats.EgressTraffic(n)
if err != nil {
c.Close() // nolint: gosec
}
return n, err return n, err
} }
@@ -75,48 +76,30 @@ func (c *Conn) Write(p []byte) (int, error) {
func (c *Conn) Read(p []byte) (int, error) { func (c *Conn) Read(p []byte) (int, error) {
select { select {
case <-c.ctx.Done(): case <-c.ctx.Done():
c.Close() // nolint: gosec
return 0, errors.Annotate(c.ctx.Err(), "Cannot read because context was closed") return 0, errors.Annotate(c.ctx.Err(), "Cannot read because context was closed")
default: default:
n, err := c.doIO(c.conn.Read, p, connTimeoutRead) if err := c.conn.SetReadDeadline(time.Now().Add(connTimeoutRead)); err != nil {
c.Close() // nolint: gosec
return 0, errors.Annotate(err, "Cannot set read deadline to the socket")
}
n, err := c.conn.Read(p)
c.logger.Debugw("Read from stream", "bytes", n, "error", err) c.logger.Debugw("Read from stream", "bytes", n, "error", err)
stats.IngressTraffic(n) stats.IngressTraffic(n)
if err != nil {
c.Close() // nolint: gosec
}
return n, err return n, err
} }
} }
func (c *Conn) doIO(callback ioFunc, p []byte, timeout time.Duration) (int, error) {
resChan := make(chan ioResult, 1)
timer := time.NewTimer(timeout)
go func() {
n, err := callback(p)
resChan <- ioResult{n: n, err: err}
}()
select {
case res := <-resChan:
timer.Stop()
if res.err != nil {
c.Close() // nolint: gosec
}
return res.n, res.err
case <-c.ctx.Done():
timer.Stop()
c.Close() // nolint: gosec
return 0, errors.Annotate(c.ctx.Err(), "Cannot do IO because context is closed")
case <-timer.C:
c.Close() // nolint: gosec
return 0, errors.Annotate(c.ctx.Err(), "Timeout on IO operation")
}
}
// Close closes underlying net.Conn instance. // Close closes underlying net.Conn instance.
func (c *Conn) Close() error { func (c *Conn) Close() error {
defer c.logger.Debugw("Close connection") c.logger.Debugw("Close connection")
c.cancel() c.cancel()
return c.conn.Close() return c.conn.Close()
} }
+2 -4
View File
@@ -5,12 +5,10 @@ import (
"sync" "sync"
) )
var streamCipherBufferPool sync.Pool var (
func init() {
streamCipherBufferPool = sync.Pool{ streamCipherBufferPool = sync.Pool{
New: func() interface{} { New: func() interface{} {
return &bytes.Buffer{} return &bytes.Buffer{}
}, },
} }
} )