mirror of
https://github.com/ScuroNeko/mtg.git
synced 2026-09-01 01:24:02 +03:00
REPOSITORY / ScuroNeko/mtg
Compare commits
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1a81efcb6e | ||
|
|
2544c521ed | ||
|
|
3a68ea5f2d | ||
|
|
dbced77566 | ||
|
|
f4f969e702 | ||
|
|
e8368f7645 | ||
|
|
38abee7d7f | ||
|
|
a3663fe8b5 | ||
|
|
2aa3321bd4 | ||
|
|
3793558c4c | ||
|
|
b6427ee321 | ||
|
|
0c9fa5e710 | ||
|
|
1fcec38aea | ||
|
|
89930631cf | ||
|
|
eedee63143 | ||
|
|
73c6a3aa37 | ||
|
|
e54d9d60d3 | ||
|
|
db2e6031a3 | ||
|
|
4b8da719ae | ||
|
|
a2de52f071 | ||
|
|
b926a0590c | ||
|
|
58e8c8f982 | ||
|
|
4642546b35 | ||
|
|
4627910238 | ||
|
|
8b3c622ea6 | ||
|
|
9d43c2d759 | ||
|
|
0840c7e3e5 | ||
|
|
de48e177b1 | ||
|
|
3ed09146b9 | ||
|
|
018bd2fdc1 | ||
|
|
0ad3a06863 | ||
|
|
d3a090d6b4 | ||
|
|
1725a0d721 | ||
|
|
87988326ab | ||
|
|
ec271baab0 | ||
|
|
139db15e83 | ||
|
|
0c030646f9 | ||
|
|
822560bede | ||
|
|
9917f61bc8 | ||
|
|
735466b90d | ||
|
|
6b51de8305 | ||
|
|
7b62e06e36 | ||
|
|
2b07c0037e | ||
|
|
450381ee16 | ||
|
|
46c33f1532 | ||
|
|
f355512aa6 | ||
|
|
289bb283b1 | ||
|
|
836090ebdf | ||
|
|
01402bdba2 | ||
|
|
026ec74dfd | ||
|
|
cc4b6ce2f4 | ||
|
|
9dfd992c1d | ||
|
|
80213ad35d | ||
|
|
d32e8e8b97 | ||
|
|
60c57c2306 | ||
|
|
0edd5e6f92 | ||
|
|
a8e4acb6f8 | ||
|
|
27d10e6820 | ||
|
|
b47e13556e | ||
|
|
de81ed565d |
@@ -119,6 +119,39 @@ jobs:
|
||||
- name: Run linter
|
||||
run: mise tasks run lint
|
||||
|
||||
artifacts:
|
||||
name: Build release artifacts
|
||||
runs-on: ubuntu-latest
|
||||
timeout-minutes: 20
|
||||
steps:
|
||||
- name: Checkout
|
||||
uses: actions/checkout@v6
|
||||
with:
|
||||
submodules: recursive
|
||||
|
||||
- uses: jdx/mise-action@v3
|
||||
name: Install mise
|
||||
|
||||
- name: Cache Go modules
|
||||
uses: actions/cache@v5
|
||||
with:
|
||||
path: ~/go/pkg/mod
|
||||
key: ${{ runner.os }}-gomod-${{ hashFiles('go.sum') }}
|
||||
restore-keys: |
|
||||
${{ runner.os }}-gomod-
|
||||
|
||||
- name: Cache cross-compilation build
|
||||
uses: actions/cache@v5
|
||||
with:
|
||||
path: ~/.cache/go-build
|
||||
key: ${{ runner.os }}-goreleaser-${{ hashFiles('go.sum') }}-${{ hashFiles('**/*.go') }}
|
||||
restore-keys: |
|
||||
${{ runner.os }}-goreleaser-${{ hashFiles('go.sum') }}-
|
||||
${{ runner.os }}-goreleaser-
|
||||
|
||||
- name: Run release
|
||||
run: mise tasks run release
|
||||
|
||||
docker:
|
||||
name: Docker
|
||||
runs-on: ubuntu-latest
|
||||
|
||||
@@ -45,11 +45,13 @@ jobs:
|
||||
|
||||
steps:
|
||||
- name: Checkout repository
|
||||
uses: actions/checkout@v2
|
||||
uses: actions/checkout@v6
|
||||
with:
|
||||
submodules: recursive
|
||||
|
||||
# Initializes the CodeQL tools for scanning.
|
||||
- name: Initialize CodeQL
|
||||
uses: github/codeql-action/init@v1
|
||||
uses: github/codeql-action/init@v4
|
||||
with:
|
||||
languages: ${{ matrix.language }}
|
||||
# If you wish to specify custom queries, you can do so here or in a config file.
|
||||
@@ -60,7 +62,7 @@ jobs:
|
||||
# Autobuild attempts to build any compiled languages (C/C++, C#, or Java).
|
||||
# If this step fails, then you should remove it and run the build manually (see below)
|
||||
- name: Autobuild
|
||||
uses: github/codeql-action/autobuild@v1
|
||||
uses: github/codeql-action/autobuild@v4
|
||||
|
||||
# ℹ️ Command-line programs to run using the OS shell.
|
||||
# 📚 https://git.io/JvXDl
|
||||
@@ -74,4 +76,4 @@ jobs:
|
||||
# make release
|
||||
|
||||
- name: Perform CodeQL Analysis
|
||||
uses: github/codeql-action/analyze@v1
|
||||
uses: github/codeql-action/analyze@v4
|
||||
|
||||
@@ -16,6 +16,12 @@ sources = ["**/*.go", "go.mod", "go.sum"]
|
||||
outputs = ["mtg"]
|
||||
run = "go build"
|
||||
|
||||
[tasks."build:prof"]
|
||||
description = "Build binary with profiling enabled"
|
||||
sources = ["**/*.go", "go.mod", "go.sum"]
|
||||
outputs = ["mtg"]
|
||||
run = "go build -tags prof"
|
||||
|
||||
[tasks.update]
|
||||
description = "Update dependencies"
|
||||
run = [
|
||||
|
||||
+1
-1
@@ -33,7 +33,7 @@ RUN go mod download
|
||||
COPY . /app
|
||||
|
||||
RUN set -x \
|
||||
&& version="$(git describe --exact-match HEAD 2>/dev/null || git describe --tags --always)" \
|
||||
&& version="$(git describe --exact-match HEAD 2>/dev/null || git describe --tags --always 2>/dev/null || echo dev)" \
|
||||
&& go build \
|
||||
-trimpath \
|
||||
-mod=readonly \
|
||||
|
||||
@@ -29,6 +29,8 @@ are the most notable:
|
||||
* [Official](https://github.com/TelegramMessenger/MTProxy)
|
||||
* [Python](https://github.com/alexbers/mtprotoproxy)
|
||||
* [Erlang](https://github.com/seriyps/mtproto_proxy)
|
||||
* [Teleproxy (C)](https://github.com/teleproxy/teleproxy)
|
||||
* [mtproto.zig (Zig)](https://github.com/sleep3r/mtproto.zig)
|
||||
* [Telemt (Rust)](https://github.com/telemt/telemt)
|
||||
|
||||
You can use any of these. They work great and all implementations have
|
||||
@@ -91,6 +93,9 @@ that probably matter.
|
||||
software. I also believe that in the case of throwout proxies, this
|
||||
the feature is a useless luxury.
|
||||
|
||||
This is very controversial topic. Please read [rationale (in russian)](https://github.com/9seconds/mtg/issues/376#issuecomment-4118726699)
|
||||
and use [mtg-multi](https://github.com/dolonet/mtg-multi) fork if you are disagree with.
|
||||
|
||||
* **No adtag support**
|
||||
|
||||
Please read [Version 2](#version-2) chapter.
|
||||
|
||||
@@ -12,7 +12,7 @@ type StableBloomFilterTestSuite struct {
|
||||
}
|
||||
|
||||
func (suite *StableBloomFilterTestSuite) TestOp() {
|
||||
filter := antireplay.NewStableBloomFilter(500, 0.001)
|
||||
filter := antireplay.NewStableBloomFilter(100000, 0.001)
|
||||
|
||||
suite.False(filter.SeenBefore([]byte{1, 2, 3}))
|
||||
suite.False(filter.SeenBefore([]byte{4, 5, 6}))
|
||||
|
||||
BIN
Binary file not shown.
@@ -38,7 +38,7 @@ func (e EventStream) Send(ctx context.Context, evt mtglib.Event) {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
case <-e.ctx.Done():
|
||||
case e.chans[int(chanNo)%len(e.chans)] <- evt:
|
||||
case e.chans[chanNo%uint32(len(e.chans))] <- evt:
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -48,6 +48,13 @@ concurrency = 8192
|
||||
# Only ipv4 connectivity is used
|
||||
prefer-ip = "prefer-ipv6"
|
||||
|
||||
# Public IP addresses of this server. Used by 'mtg access' to generate
|
||||
# proxy links and by 'mtg doctor' to validate SNI-DNS match.
|
||||
# If not set, mtg tries to detect them automatically via ifconfig.co.
|
||||
# Set these if ifconfig.co is unreachable from your server.
|
||||
# public-ipv4 = "1.2.3.4"
|
||||
# public-ipv6 = "2001:db8::1"
|
||||
|
||||
# If this setting is set, then mtg will try to get proxy updates from Telegram
|
||||
# Usually this is completely fine to have it disabled, because mtg has a list
|
||||
# of some core proxies hardcoded.
|
||||
|
||||
@@ -15,7 +15,7 @@ require (
|
||||
github.com/prometheus/client_golang v1.23.2
|
||||
github.com/prometheus/common v0.67.5 // indirect
|
||||
github.com/prometheus/procfs v0.20.1 // indirect
|
||||
github.com/rs/zerolog v1.34.0
|
||||
github.com/rs/zerolog v1.35.0
|
||||
github.com/smira/go-statsd v1.3.4
|
||||
github.com/stretchr/objx v0.5.2 // indirect
|
||||
github.com/stretchr/testify v1.11.1
|
||||
@@ -29,7 +29,7 @@ require (
|
||||
require (
|
||||
github.com/beevik/ntp v1.5.0
|
||||
github.com/ncruces/go-dns v1.3.2
|
||||
github.com/pelletier/go-toml/v2 v2.2.4
|
||||
github.com/pelletier/go-toml/v2 v2.3.0
|
||||
github.com/pires/go-proxyproto v0.11.0
|
||||
github.com/things-go/go-socks5 v0.1.0
|
||||
github.com/txthinking/socks5 v0.0.0-20251011041537-5c31f201a10e
|
||||
|
||||
@@ -18,14 +18,12 @@ github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM=
|
||||
github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw=
|
||||
github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs=
|
||||
github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs=
|
||||
github.com/coreos/go-systemd/v22 v22.5.0/go.mod h1:Y58oyj3AT4RCenI/lSvhwexgC+NSVTIJ3seZv2GcEnc=
|
||||
github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E=
|
||||
github.com/d4l3k/messagediff v1.2.1 h1:ZcAIMYsUg0EAp9X+tt8/enBE/Q8Yd5kzPynLyKptt9U=
|
||||
github.com/d4l3k/messagediff v1.2.1/go.mod h1:Oozbb1TVXFac9FtSIxHBMnBCq2qeH/2KkEQxENCrlLo=
|
||||
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/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
|
||||
github.com/godbus/dbus/v5 v5.0.4/go.mod h1:xhWf0FNVPg57R7Z0UbKHbJfkEywrmjJnf7w5xrFpKfA=
|
||||
github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8=
|
||||
github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU=
|
||||
github.com/hexops/gotextdiff v1.0.3 h1:gitA9+qJrrTCsiCl7+kh75nPqQt1cx4ZkudSTLoUqJM=
|
||||
@@ -40,11 +38,8 @@ github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY=
|
||||
github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE=
|
||||
github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0SNc=
|
||||
github.com/kylelemons/godebug v1.1.0/go.mod h1:9/0rRGxNHcop5bhtWyNeEfOS8JIWk580+fNqagV/RAw=
|
||||
github.com/mattn/go-colorable v0.1.13/go.mod h1:7S9/ev0klgBDR4GtXTXX8a3vIGJpMovkB8vQcUbaXHg=
|
||||
github.com/mattn/go-colorable v0.1.14 h1:9A9LHSqF/7dyVVX6g0U9cwm9pG3kP9gSzcuIPHPsaIE=
|
||||
github.com/mattn/go-colorable v0.1.14/go.mod h1:6LmQG8QLFO4G5z1gPvYEzlUgJ2wF+stgPZH1UqBm1s8=
|
||||
github.com/mattn/go-isatty v0.0.16/go.mod h1:kYGgaQfpe5nmfYZH+SKPsOc2e4SrIfOl2e/yFXSvRLM=
|
||||
github.com/mattn/go-isatty v0.0.19/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y=
|
||||
github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY=
|
||||
github.com/mattn/go-isatty v0.0.20/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y=
|
||||
github.com/mccutchen/go-httpbin v1.1.1 h1:aEws49HEJEyXHLDnshQVswfUlCVoS8g6h9YaDyaW7RE=
|
||||
@@ -59,11 +54,10 @@ github.com/panjf2000/ants/v2 v2.12.0 h1:u9JhESo83i/GkZnhfTNuFMMWcNt7mnV1bGJ6FT4w
|
||||
github.com/panjf2000/ants/v2 v2.12.0/go.mod h1:tSQuaNQ6r6NRhPt+IZVUevvDyFMTs+eS4ztZc52uJTY=
|
||||
github.com/patrickmn/go-cache v2.1.0+incompatible h1:HRMgzkcYKYpi3C8ajMPV8OFXaaRUnok+kx1WdO15EQc=
|
||||
github.com/patrickmn/go-cache v2.1.0+incompatible/go.mod h1:3Qf8kWWT7OJRJbdiICTKqZju1ZixQ/KpMGzzAfe6+WQ=
|
||||
github.com/pelletier/go-toml/v2 v2.2.4 h1:mye9XuhQ6gvn5h28+VilKrrPoQVanw5PMw/TB0t5Ec4=
|
||||
github.com/pelletier/go-toml/v2 v2.2.4/go.mod h1:2gIqNv+qfxSVS7cM2xJQKtLSTLUE9V8t9Stt+h56mCY=
|
||||
github.com/pelletier/go-toml/v2 v2.3.0 h1:k59bC/lIZREW0/iVaQR8nDHxVq8OVlIzYCOJf421CaM=
|
||||
github.com/pelletier/go-toml/v2 v2.3.0/go.mod h1:2gIqNv+qfxSVS7cM2xJQKtLSTLUE9V8t9Stt+h56mCY=
|
||||
github.com/pires/go-proxyproto v0.11.0 h1:gUQpS85X/VJMdUsYyEgyn59uLJvGqPhJV5YvG68wXH4=
|
||||
github.com/pires/go-proxyproto v0.11.0/go.mod h1:ZKAAyp3cgy5Y5Mo4n9AlScrkCZwUy0g3Jf+slqQVcuU=
|
||||
github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
|
||||
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
|
||||
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
|
||||
github.com/prometheus/client_golang v1.23.2 h1:Je96obch5RDVy3FDMndoUsjAhG5Edi49h0RJWRi/o0o=
|
||||
@@ -76,9 +70,8 @@ github.com/prometheus/procfs v0.20.1 h1:XwbrGOIplXW/AU3YhIhLODXMJYyC1isLFfYCsTEy
|
||||
github.com/prometheus/procfs v0.20.1/go.mod h1:o9EMBZGRyvDrSPH1RqdxhojkuXstoe4UlK79eF5TGGo=
|
||||
github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0tI/otEQ=
|
||||
github.com/rogpeppe/go-internal v1.14.1/go.mod h1:MaRKkUm5W0goXpeCfT7UZI6fk/L7L7so1lCWt35ZSgc=
|
||||
github.com/rs/xid v1.6.0/go.mod h1:7XoLgs4eV+QndskICGsho+ADou8ySMSjJKDIan90Nz0=
|
||||
github.com/rs/zerolog v1.34.0 h1:k43nTLIwcTVQAncfCw4KZ2VY6ukYoZaBPNOE8txlOeY=
|
||||
github.com/rs/zerolog v1.34.0/go.mod h1:bJsvje4Z08ROH4Nhs5iH600c3IkWhwp44iRc54W6wYQ=
|
||||
github.com/rs/zerolog v1.35.0 h1:VD0ykx7HMiMJytqINBsKcbLS+BJ4WYjz+05us+LRTdI=
|
||||
github.com/rs/zerolog v1.35.0/go.mod h1:EjML9kdfa/RMA7h/6z6pYmq1ykOuA8/mjWaEvGI+jcw=
|
||||
github.com/smira/go-statsd v1.3.4 h1:kBYWcLSGT+qC6JVbvfz48kX7mQys32fjDOPrfmsSx2c=
|
||||
github.com/smira/go-statsd v1.3.4/go.mod h1:RjdsESPgDODtg1VpVVf9MJrEW2Hw0wtRNbmB1CAhu6A=
|
||||
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
|
||||
@@ -133,10 +126,8 @@ golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7w
|
||||
golang.org/x/sys v0.0.0-20210615035016-665e8c7367d1/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
golang.org/x/sys v0.0.0-20220520151302-bc2c85ada10a/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
golang.org/x/sys v0.0.0-20220722155257-8c9f86f7a55f/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
golang.org/x/sys v0.0.0-20220811171246-fbc7d0a398ab/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
golang.org/x/sys v0.2.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
golang.org/x/sys v0.12.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
golang.org/x/sys v0.42.0 h1:omrd2nAlyT5ESRdCLYdm3+fMfNFE/+Rf4bDIQImRJeo=
|
||||
golang.org/x/sys v0.42.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
|
||||
golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo=
|
||||
|
||||
@@ -58,6 +58,9 @@ func (a *Access) Run(cli *CLI, version string) error {
|
||||
|
||||
wg.Go(func() {
|
||||
ip := a.PublicIPv4
|
||||
if ip == nil {
|
||||
ip = conf.PublicIPv4.Get(nil)
|
||||
}
|
||||
if ip == nil {
|
||||
ip = getIP(ntw, "tcp4")
|
||||
}
|
||||
@@ -70,6 +73,9 @@ func (a *Access) Run(cli *CLI, version string) error {
|
||||
})
|
||||
wg.Go(func() {
|
||||
ip := a.PublicIPv6
|
||||
if ip == nil {
|
||||
ip = conf.PublicIPv6.Get(nil)
|
||||
}
|
||||
if ip == nil {
|
||||
ip = getIP(ntw, "tcp6")
|
||||
}
|
||||
|
||||
+10
-3
@@ -332,13 +332,20 @@ func (d *Doctor) checkSecretHost(resolver *net.Resolver, ntw mtglib.Network) boo
|
||||
return false
|
||||
}
|
||||
|
||||
ourIP4 := getIP(ntw, "tcp4")
|
||||
ourIP6 := getIP(ntw, "tcp6")
|
||||
ourIP4 := d.conf.PublicIPv4.Get(nil)
|
||||
if ourIP4 == nil {
|
||||
ourIP4 = getIP(ntw, "tcp4")
|
||||
}
|
||||
|
||||
ourIP6 := d.conf.PublicIPv6.Get(nil)
|
||||
if ourIP6 == nil {
|
||||
ourIP6 = getIP(ntw, "tcp6")
|
||||
}
|
||||
|
||||
if ourIP4 == nil && ourIP6 == nil {
|
||||
tplError.Execute(os.Stdout, map[string]any{ //nolint: errcheck
|
||||
"description": "cannot detect public IP address",
|
||||
"error": errors.New("ifconfig.co is unreachable for both IPv4 and IPv6"),
|
||||
"error": errors.New("cannot detect automatically and public-ipv4/public-ipv6 are not set in config"),
|
||||
})
|
||||
return false
|
||||
}
|
||||
|
||||
@@ -5,6 +5,7 @@ import (
|
||||
"fmt"
|
||||
"net"
|
||||
"os"
|
||||
"time"
|
||||
|
||||
"github.com/9seconds/mtg/v2/antireplay"
|
||||
"github.com/9seconds/mtg/v2/events"
|
||||
@@ -262,6 +263,7 @@ func runProxy(conf *config.Config, version string) error { //nolint: funlen
|
||||
|
||||
AllowFallbackOnUnknownDC: conf.AllowFallbackOnUnknownDC.Get(false),
|
||||
TolerateTimeSkewness: conf.TolerateTimeSkewness.Value,
|
||||
IdleTimeout: conf.Network.Timeout.Idle.Get(time.Minute),
|
||||
|
||||
DoppelGangerURLs: doppelGangerURLs,
|
||||
DoppelGangerPerRaid: conf.Defense.Doppelganger.Repeats.Get(mtglib.DoppelGangerPerRaid),
|
||||
|
||||
@@ -35,6 +35,8 @@ type Config struct {
|
||||
DomainFrontingProxyProtocol TypeBool `json:"domainFrontingProxyProtocol"`
|
||||
TolerateTimeSkewness TypeDuration `json:"tolerateTimeSkewness"`
|
||||
Concurrency TypeConcurrency `json:"concurrency"`
|
||||
PublicIPv4 TypeIP `json:"publicIpv4"`
|
||||
PublicIPv6 TypeIP `json:"publicIpv6"`
|
||||
DomainFronting struct {
|
||||
IP TypeIP `json:"ip"`
|
||||
Port TypePort `json:"port"`
|
||||
|
||||
@@ -42,6 +42,32 @@ func (suite *ConfigTestSuite) TestParseMinimalConfig() {
|
||||
suite.Equal("0.0.0.0:3128", conf.BindTo.String())
|
||||
}
|
||||
|
||||
func (suite *ConfigTestSuite) TestParsePublicIP() {
|
||||
conf, err := config.Parse(suite.ReadConfig("public_ip.toml"))
|
||||
suite.NoError(err)
|
||||
suite.Equal("203.0.113.1", conf.PublicIPv4.Get(nil).String())
|
||||
suite.Equal("2001:db8::1", conf.PublicIPv6.Get(nil).String())
|
||||
}
|
||||
|
||||
func (suite *ConfigTestSuite) TestParsePublicIPv4Only() {
|
||||
conf, err := config.Parse(suite.ReadConfig("public_ip_v4_only.toml"))
|
||||
suite.NoError(err)
|
||||
suite.Equal("203.0.113.1", conf.PublicIPv4.Get(nil).String())
|
||||
suite.Nil(conf.PublicIPv6.Get(nil))
|
||||
}
|
||||
|
||||
func (suite *ConfigTestSuite) TestParsePublicIPInvalid() {
|
||||
_, err := config.Parse(suite.ReadConfig("public_ip_invalid.toml"))
|
||||
suite.Error(err)
|
||||
}
|
||||
|
||||
func (suite *ConfigTestSuite) TestParsePublicIPNotSet() {
|
||||
conf, err := config.Parse(suite.ReadConfig("minimal.toml"))
|
||||
suite.NoError(err)
|
||||
suite.Nil(conf.PublicIPv4.Get(nil))
|
||||
suite.Nil(conf.PublicIPv6.Get(nil))
|
||||
}
|
||||
|
||||
func (suite *ConfigTestSuite) TestString() {
|
||||
conf, err := config.Parse(suite.ReadConfig("minimal.toml"))
|
||||
suite.NoError(err)
|
||||
|
||||
@@ -21,6 +21,8 @@ type tomlConfig struct {
|
||||
DomainFrontingProxyProtocol bool `toml:"domain-fronting-proxy-protocol" json:"domainFrontingProxyProtocol,omitempty"`
|
||||
TolerateTimeSkewness string `toml:"tolerate-time-skewness" json:"tolerateTimeSkewness,omitempty"`
|
||||
Concurrency uint `toml:"concurrency" json:"concurrency,omitempty"`
|
||||
PublicIPv4 string `toml:"public-ipv4" json:"publicIpv4,omitempty"`
|
||||
PublicIPv6 string `toml:"public-ipv6" json:"publicIpv6,omitempty"`
|
||||
DomainFronting struct {
|
||||
IP string `toml:"ip" json:"ip,omitempty"`
|
||||
Port uint `toml:"port" json:"port,omitempty"`
|
||||
|
||||
+4
@@ -0,0 +1,4 @@
|
||||
secret = "7oe1GqLy6TBc38CV3jx7q09nb29nbGUuY29t"
|
||||
bind-to = "0.0.0.0:3128"
|
||||
public-ipv4 = "203.0.113.1"
|
||||
public-ipv6 = "2001:db8::1"
|
||||
@@ -0,0 +1,3 @@
|
||||
secret = "7oe1GqLy6TBc38CV3jx7q09nb29nbGUuY29t"
|
||||
bind-to = "0.0.0.0:3128"
|
||||
public-ipv4 = "not-an-ip"
|
||||
@@ -0,0 +1,3 @@
|
||||
secret = "7oe1GqLy6TBc38CV3jx7q09nb29nbGUuY29t"
|
||||
bind-to = "0.0.0.0:3128"
|
||||
public-ipv4 = "203.0.113.1"
|
||||
@@ -82,40 +82,40 @@ checksum = "sha256:4932cfca5e75bf60fe1c576edf459e5e809e6644664a068185d64b84af3fa
|
||||
url = "https://github.com/golangci/golangci-lint/releases/download/v2.11.4/golangci-lint-2.11.4-windows-amd64.zip"
|
||||
|
||||
[[tools.goreleaser]]
|
||||
version = "2.14.3"
|
||||
version = "2.15.2"
|
||||
backend = "aqua:goreleaser/goreleaser"
|
||||
|
||||
[tools.goreleaser."platforms.linux-arm64"]
|
||||
checksum = "sha256:581a10e53c1176b3e81ee45cf531e02dbf899db0bc7b795669347df4276ce948"
|
||||
url = "https://github.com/goreleaser/goreleaser/releases/download/v2.14.3/goreleaser_Linux_arm64.tar.gz"
|
||||
checksum = "sha256:5db66761a98f6693161e49e1a95d28d2673a892ba60cb4a5e16736cafd41c4c9"
|
||||
url = "https://github.com/goreleaser/goreleaser/releases/download/v2.15.2/goreleaser_Linux_arm64.tar.gz"
|
||||
provenance = "cosign"
|
||||
|
||||
[tools.goreleaser."platforms.linux-arm64-musl"]
|
||||
checksum = "sha256:581a10e53c1176b3e81ee45cf531e02dbf899db0bc7b795669347df4276ce948"
|
||||
url = "https://github.com/goreleaser/goreleaser/releases/download/v2.14.3/goreleaser_Linux_arm64.tar.gz"
|
||||
checksum = "sha256:5db66761a98f6693161e49e1a95d28d2673a892ba60cb4a5e16736cafd41c4c9"
|
||||
url = "https://github.com/goreleaser/goreleaser/releases/download/v2.15.2/goreleaser_Linux_arm64.tar.gz"
|
||||
provenance = "cosign"
|
||||
|
||||
[tools.goreleaser."platforms.linux-x64"]
|
||||
checksum = "sha256:dc7faeeeb6da8bdfda788626263a4ae725892a8c7504b975c3234127d4a44579"
|
||||
url = "https://github.com/goreleaser/goreleaser/releases/download/v2.14.3/goreleaser_Linux_x86_64.tar.gz"
|
||||
checksum = "sha256:0ebdbf0353aba566b969dde746cc4e4806f96c27aa2f3971b229a9df7611fedc"
|
||||
url = "https://github.com/goreleaser/goreleaser/releases/download/v2.15.2/goreleaser_Linux_x86_64.tar.gz"
|
||||
provenance = "cosign"
|
||||
|
||||
[tools.goreleaser."platforms.linux-x64-musl"]
|
||||
checksum = "sha256:dc7faeeeb6da8bdfda788626263a4ae725892a8c7504b975c3234127d4a44579"
|
||||
url = "https://github.com/goreleaser/goreleaser/releases/download/v2.14.3/goreleaser_Linux_x86_64.tar.gz"
|
||||
checksum = "sha256:0ebdbf0353aba566b969dde746cc4e4806f96c27aa2f3971b229a9df7611fedc"
|
||||
url = "https://github.com/goreleaser/goreleaser/releases/download/v2.15.2/goreleaser_Linux_x86_64.tar.gz"
|
||||
provenance = "cosign"
|
||||
|
||||
[tools.goreleaser."platforms.macos-arm64"]
|
||||
checksum = "sha256:3507798489e107a78aff36b169de48148a335ac26eb3161608d905f3f3a957bd"
|
||||
url = "https://github.com/goreleaser/goreleaser/releases/download/v2.14.3/goreleaser_Darwin_all.tar.gz"
|
||||
provenance = "cosign"
|
||||
checksum = "sha256:0e6bd67688ac949780bf1166813a91f89856898ef4c40d7d46c2c74ebaa4b9ee"
|
||||
url = "https://github.com/goreleaser/goreleaser/releases/download/v2.15.2/goreleaser_Darwin_all.tar.gz"
|
||||
provenance = "github-attestations"
|
||||
|
||||
[tools.goreleaser."platforms.macos-x64"]
|
||||
checksum = "sha256:3507798489e107a78aff36b169de48148a335ac26eb3161608d905f3f3a957bd"
|
||||
url = "https://github.com/goreleaser/goreleaser/releases/download/v2.14.3/goreleaser_Darwin_all.tar.gz"
|
||||
checksum = "sha256:0e6bd67688ac949780bf1166813a91f89856898ef4c40d7d46c2c74ebaa4b9ee"
|
||||
url = "https://github.com/goreleaser/goreleaser/releases/download/v2.15.2/goreleaser_Darwin_all.tar.gz"
|
||||
provenance = "cosign"
|
||||
|
||||
[tools.goreleaser."platforms.windows-x64"]
|
||||
checksum = "sha256:3deea8ff471aa258a2d99f3e5302971d7028647ae8ddaf103257a8113e485a31"
|
||||
url = "https://github.com/goreleaser/goreleaser/releases/download/v2.14.3/goreleaser_Windows_x86_64.zip"
|
||||
checksum = "sha256:7459832946dbe122c144f8d7f87484d8572ca005b779310aa6bb03346e8de17a"
|
||||
url = "https://github.com/goreleaser/goreleaser/releases/download/v2.15.2/goreleaser_Windows_x86_64.zip"
|
||||
provenance = "cosign"
|
||||
|
||||
@@ -3,9 +3,12 @@ package mtglib
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"github.com/9seconds/mtg/v2/essentials"
|
||||
"github.com/pires/go-proxyproto"
|
||||
@@ -95,3 +98,64 @@ func newConnProxyProtocol(source, target essentials.Conn) *connProxyProtocol {
|
||||
sourceAddr: source.RemoteAddr(),
|
||||
}
|
||||
}
|
||||
|
||||
// idleTracker is a shared idle tracker for a pair of relay connections.
|
||||
// Both directions update the same timestamp so that activity in one direction
|
||||
// prevents the other (idle) direction from timing out.
|
||||
type idleTracker struct {
|
||||
lastActive atomic.Pointer[time.Time]
|
||||
timeout time.Duration
|
||||
}
|
||||
|
||||
func newIdleTracker(timeout time.Duration) *idleTracker {
|
||||
t := &idleTracker{timeout: timeout}
|
||||
t.touch()
|
||||
|
||||
return t
|
||||
}
|
||||
|
||||
func (t *idleTracker) touch() {
|
||||
stamp := time.Now()
|
||||
t.lastActive.Store(&stamp)
|
||||
}
|
||||
|
||||
func (t *idleTracker) isIdle() bool {
|
||||
return time.Since(*t.lastActive.Load()) >= t.timeout
|
||||
}
|
||||
|
||||
type connIdleTimeout struct {
|
||||
essentials.Conn
|
||||
|
||||
tracker *idleTracker
|
||||
}
|
||||
|
||||
func (c connIdleTimeout) Read(b []byte) (int, error) {
|
||||
var netErr net.Error
|
||||
|
||||
for {
|
||||
c.SetReadDeadline(time.Now().Add(c.tracker.timeout)) //nolint: errcheck
|
||||
|
||||
n, err := c.Conn.Read(b)
|
||||
|
||||
switch {
|
||||
case err == nil:
|
||||
c.tracker.touch()
|
||||
return n, nil
|
||||
case errors.As(err, &netErr) && netErr.Timeout() && !c.tracker.isIdle():
|
||||
continue
|
||||
}
|
||||
|
||||
return n, err
|
||||
}
|
||||
}
|
||||
|
||||
func (c connIdleTimeout) Write(b []byte) (int, error) {
|
||||
c.SetWriteDeadline(time.Now().Add(c.tracker.timeout)) //nolint: errcheck
|
||||
|
||||
n, err := c.Conn.Write(b)
|
||||
if n > 0 {
|
||||
c.tracker.touch()
|
||||
}
|
||||
|
||||
return n, err //nolint: wrapcheck
|
||||
}
|
||||
|
||||
@@ -16,6 +16,12 @@ import (
|
||||
"github.com/stretchr/testify/suite"
|
||||
)
|
||||
|
||||
type netTimeoutError struct{}
|
||||
|
||||
func (e netTimeoutError) Error() string { return "i/o timeout" }
|
||||
func (e netTimeoutError) Timeout() bool { return true }
|
||||
func (e netTimeoutError) Temporary() bool { return true }
|
||||
|
||||
type ConnRewindBaseConn struct {
|
||||
testlib.EssentialsConnMock
|
||||
|
||||
@@ -291,6 +297,141 @@ func (suite *ConnProxyProtocolTestSuite) TearDownTest() {
|
||||
suite.targetConnMock.AssertExpectations(suite.T())
|
||||
}
|
||||
|
||||
type IdleTrackerTestSuite struct {
|
||||
suite.Suite
|
||||
}
|
||||
|
||||
func (suite *IdleTrackerTestSuite) TestNewNotIdle() {
|
||||
tracker := newIdleTracker(time.Second)
|
||||
suite.False(tracker.isIdle())
|
||||
}
|
||||
|
||||
func (suite *IdleTrackerTestSuite) TestIdleAfterTimeout() {
|
||||
tracker := newIdleTracker(10 * time.Millisecond)
|
||||
time.Sleep(20 * time.Millisecond)
|
||||
|
||||
suite.True(tracker.isIdle())
|
||||
}
|
||||
|
||||
func (suite *IdleTrackerTestSuite) TestTouchResetsIdle() {
|
||||
tracker := newIdleTracker(50 * time.Millisecond)
|
||||
time.Sleep(30 * time.Millisecond)
|
||||
|
||||
tracker.touch()
|
||||
|
||||
suite.False(tracker.isIdle())
|
||||
}
|
||||
|
||||
type ConnIdleTimeoutTestSuite struct {
|
||||
suite.Suite
|
||||
|
||||
connMock *testlib.EssentialsConnMock
|
||||
tracker *idleTracker
|
||||
conn connIdleTimeout
|
||||
}
|
||||
|
||||
func (suite *ConnIdleTimeoutTestSuite) SetupTest() {
|
||||
suite.connMock = &testlib.EssentialsConnMock{}
|
||||
suite.tracker = newIdleTracker(time.Second)
|
||||
suite.conn = connIdleTimeout{
|
||||
Conn: suite.connMock,
|
||||
tracker: suite.tracker,
|
||||
}
|
||||
}
|
||||
|
||||
func (suite *ConnIdleTimeoutTestSuite) TearDownTest() {
|
||||
suite.connMock.AssertExpectations(suite.T())
|
||||
}
|
||||
|
||||
func (suite *ConnIdleTimeoutTestSuite) TestReadOk() {
|
||||
suite.connMock.On("SetReadDeadline", mock.Anything).Return(nil)
|
||||
suite.connMock.On("Read", mock.Anything).Once().Return(5, nil)
|
||||
|
||||
n, err := suite.conn.Read(make([]byte, 10))
|
||||
suite.NoError(err)
|
||||
suite.Equal(5, n)
|
||||
}
|
||||
|
||||
func (suite *ConnIdleTimeoutTestSuite) TestReadNonTimeoutErr() {
|
||||
suite.connMock.On("SetReadDeadline", mock.Anything).Return(nil)
|
||||
suite.connMock.On("Read", mock.Anything).Once().Return(0, io.EOF)
|
||||
|
||||
n, err := suite.conn.Read(make([]byte, 10))
|
||||
suite.True(errors.Is(err, io.EOF))
|
||||
suite.Equal(0, n)
|
||||
}
|
||||
|
||||
func (suite *ConnIdleTimeoutTestSuite) TestReadTimeoutRetriesWhenNotIdle() {
|
||||
suite.connMock.On("SetReadDeadline", mock.Anything).Return(nil)
|
||||
suite.connMock.On("Read", mock.Anything).Once().Return(0, netTimeoutError{})
|
||||
suite.connMock.On("Read", mock.Anything).Once().Return(5, nil)
|
||||
|
||||
n, err := suite.conn.Read(make([]byte, 10))
|
||||
suite.NoError(err)
|
||||
suite.Equal(5, n)
|
||||
}
|
||||
|
||||
func (suite *ConnIdleTimeoutTestSuite) TestReadTimeoutClosesWhenIdle() {
|
||||
suite.tracker = newIdleTracker(time.Millisecond)
|
||||
suite.conn = connIdleTimeout{
|
||||
Conn: suite.connMock,
|
||||
tracker: suite.tracker,
|
||||
}
|
||||
|
||||
time.Sleep(5 * time.Millisecond)
|
||||
|
||||
suite.connMock.On("SetReadDeadline", mock.Anything).Return(nil)
|
||||
suite.connMock.On("Read", mock.Anything).Once().Return(0, netTimeoutError{})
|
||||
|
||||
n, err := suite.conn.Read(make([]byte, 10))
|
||||
suite.Equal(0, n)
|
||||
|
||||
netErr, ok := err.(net.Error) //nolint: errorlint
|
||||
suite.True(ok)
|
||||
suite.True(netErr.Timeout())
|
||||
}
|
||||
|
||||
func (suite *ConnIdleTimeoutTestSuite) TestSharedTrackerPreventsFalseTimeout() {
|
||||
connMock2 := &testlib.EssentialsConnMock{}
|
||||
conn2 := connIdleTimeout{
|
||||
Conn: connMock2,
|
||||
tracker: suite.tracker,
|
||||
}
|
||||
|
||||
connMock2.On("SetWriteDeadline", mock.Anything).Return(nil)
|
||||
connMock2.On("Write", mock.Anything).Once().Return(5, nil)
|
||||
|
||||
_, _ = conn2.Write(make([]byte, 5))
|
||||
|
||||
suite.connMock.On("SetReadDeadline", mock.Anything).Return(nil)
|
||||
suite.connMock.On("Read", mock.Anything).Once().Return(0, netTimeoutError{})
|
||||
suite.connMock.On("Read", mock.Anything).Once().Return(3, nil)
|
||||
|
||||
n, err := suite.conn.Read(make([]byte, 10))
|
||||
suite.NoError(err)
|
||||
suite.Equal(3, n)
|
||||
|
||||
connMock2.AssertExpectations(suite.T())
|
||||
}
|
||||
|
||||
func (suite *ConnIdleTimeoutTestSuite) TestWriteOk() {
|
||||
suite.connMock.On("SetWriteDeadline", mock.Anything).Return(nil)
|
||||
suite.connMock.On("Write", mock.Anything).Once().Return(5, nil)
|
||||
|
||||
n, err := suite.conn.Write(make([]byte, 5))
|
||||
suite.NoError(err)
|
||||
suite.Equal(5, n)
|
||||
}
|
||||
|
||||
func (suite *ConnIdleTimeoutTestSuite) TestWriteErr() {
|
||||
suite.connMock.On("SetWriteDeadline", mock.Anything).Return(nil)
|
||||
suite.connMock.On("Write", mock.Anything).Once().Return(0, io.EOF)
|
||||
|
||||
n, err := suite.conn.Write(make([]byte, 5))
|
||||
suite.True(errors.Is(err, io.EOF))
|
||||
suite.Equal(0, n)
|
||||
}
|
||||
|
||||
func TestConnTraffic(t *testing.T) {
|
||||
t.Parallel()
|
||||
suite.Run(t, &ConnTrafficTestSuite{})
|
||||
@@ -305,3 +446,13 @@ func TestConnProxyProtocol(t *testing.T) {
|
||||
t.Parallel()
|
||||
suite.Run(t, &ConnProxyProtocolTestSuite{})
|
||||
}
|
||||
|
||||
func TestIdleTracker(t *testing.T) {
|
||||
t.Parallel()
|
||||
suite.Run(t, &IdleTrackerTestSuite{})
|
||||
}
|
||||
|
||||
func TestConnIdleTimeout(t *testing.T) {
|
||||
t.Parallel()
|
||||
suite.Run(t, &ConnIdleTimeoutTestSuite{})
|
||||
}
|
||||
|
||||
@@ -5,15 +5,19 @@ type dcView struct {
|
||||
}
|
||||
|
||||
func (d dcView) getV4(dc int) []Addr {
|
||||
addrs := d.publicConfigs.getV4(dc)
|
||||
var addrs []Addr
|
||||
|
||||
addrs = append(addrs, defaultDCAddrSet.getV4(dc)...)
|
||||
addrs = append(addrs, d.publicConfigs.getV4(dc)...)
|
||||
|
||||
return addrs
|
||||
}
|
||||
|
||||
func (d dcView) getV6(dc int) []Addr {
|
||||
addrs := d.publicConfigs.getV6(dc)
|
||||
var addrs []Addr
|
||||
|
||||
addrs = append(addrs, defaultDCAddrSet.getV6(dc)...)
|
||||
addrs = append(addrs, d.publicConfigs.getV6(dc)...)
|
||||
|
||||
return addrs
|
||||
}
|
||||
|
||||
@@ -1,35 +0,0 @@
|
||||
package doppel
|
||||
|
||||
import (
|
||||
"context"
|
||||
"time"
|
||||
)
|
||||
|
||||
type Clock struct {
|
||||
stats *Stats
|
||||
tick chan struct{}
|
||||
}
|
||||
|
||||
func (c Clock) Start(ctx context.Context) {
|
||||
tickTock := time.NewTimer(c.stats.Delay())
|
||||
defer func() {
|
||||
tickTock.Stop()
|
||||
select {
|
||||
case <-tickTock.C:
|
||||
default:
|
||||
}
|
||||
}()
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-tickTock.C:
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
case c.tick <- struct{}{}:
|
||||
}
|
||||
tickTock.Reset(c.stats.Delay())
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,80 +0,0 @@
|
||||
package doppel
|
||||
|
||||
import (
|
||||
"context"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/stretchr/testify/suite"
|
||||
)
|
||||
|
||||
type ClockTestSuite struct {
|
||||
suite.Suite
|
||||
|
||||
clock Clock
|
||||
wg sync.WaitGroup
|
||||
ctx context.Context
|
||||
ctxCancel context.CancelFunc
|
||||
}
|
||||
|
||||
func (suite *ClockTestSuite) SetupTest() {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
|
||||
suite.ctx = ctx
|
||||
suite.ctxCancel = cancel
|
||||
suite.clock = Clock{
|
||||
stats: &Stats{
|
||||
k: StatsDefaultK,
|
||||
lambda: StatsDefaultLambda,
|
||||
},
|
||||
tick: make(chan struct{}),
|
||||
}
|
||||
|
||||
suite.wg.Go(func() {
|
||||
suite.clock.Start(suite.ctx)
|
||||
})
|
||||
}
|
||||
|
||||
func (suite *ClockTestSuite) TearDownTest() {
|
||||
suite.ctxCancel()
|
||||
suite.wg.Wait()
|
||||
}
|
||||
|
||||
func (suite *ClockTestSuite) TestTicks() {
|
||||
received := 0
|
||||
|
||||
for range 3 {
|
||||
select {
|
||||
case <-suite.clock.tick:
|
||||
received++
|
||||
case <-time.After(2 * time.Second):
|
||||
suite.Fail("timed out waiting for tick")
|
||||
}
|
||||
}
|
||||
|
||||
suite.Equal(3, received)
|
||||
}
|
||||
|
||||
func (suite *ClockTestSuite) TestStopsOnCancel() {
|
||||
select {
|
||||
case <-suite.clock.tick:
|
||||
case <-time.After(2 * time.Second):
|
||||
suite.Fail("timed out waiting for first tick")
|
||||
}
|
||||
|
||||
suite.ctxCancel()
|
||||
|
||||
time.Sleep(50 * time.Millisecond)
|
||||
|
||||
select {
|
||||
case <-suite.clock.tick:
|
||||
suite.Fail("received tick after cancel")
|
||||
default:
|
||||
}
|
||||
}
|
||||
|
||||
func TestClock(t *testing.T) {
|
||||
t.Parallel()
|
||||
suite.Run(t, &ClockTestSuite{})
|
||||
}
|
||||
@@ -4,11 +4,19 @@ import (
|
||||
"bytes"
|
||||
"context"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/9seconds/mtg/v2/essentials"
|
||||
"github.com/9seconds/mtg/v2/mtglib/internal/tls"
|
||||
)
|
||||
|
||||
var doppelBufPool = sync.Pool{
|
||||
New: func() any {
|
||||
b := make([]byte, tls.MaxRecordSize)
|
||||
return &b
|
||||
},
|
||||
}
|
||||
|
||||
type Conn struct {
|
||||
essentials.Conn
|
||||
|
||||
@@ -18,7 +26,7 @@ type Conn struct {
|
||||
type connPayload struct {
|
||||
ctx context.Context
|
||||
ctxCancel context.CancelCauseFunc
|
||||
clock Clock
|
||||
stats Stats
|
||||
wg sync.WaitGroup
|
||||
writeStream bytes.Buffer
|
||||
writtenCond sync.Cond
|
||||
@@ -39,23 +47,23 @@ func (c Conn) Write(p []byte) (int, error) {
|
||||
return len(p), context.Cause(c.p.ctx)
|
||||
}
|
||||
|
||||
func (c Conn) Start() {
|
||||
c.p.wg.Go(func() {
|
||||
c.start()
|
||||
})
|
||||
}
|
||||
|
||||
func (c Conn) start() {
|
||||
buf := [tls.MaxRecordSize]byte{}
|
||||
bp := doppelBufPool.Get().(*[]byte)
|
||||
buf := *bp
|
||||
defer doppelBufPool.Put(bp)
|
||||
|
||||
timer := time.NewTimer(c.p.stats.Delay())
|
||||
defer timer.Stop()
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-c.p.ctx.Done():
|
||||
return
|
||||
case <-c.p.clock.tick:
|
||||
case <-timer.C:
|
||||
timer.Reset(c.p.stats.Delay())
|
||||
}
|
||||
|
||||
size := c.p.clock.stats.Size()
|
||||
size := c.p.stats.Size()
|
||||
|
||||
c.p.writtenCond.L.Lock()
|
||||
for c.p.writeStream.Len() == 0 && !c.p.done {
|
||||
@@ -68,7 +76,7 @@ func (c Conn) start() {
|
||||
continue
|
||||
}
|
||||
|
||||
if err := tls.WriteRecordInPlace(c.Conn, buf[:], n); err != nil {
|
||||
if err := tls.WriteRecordInPlace(c.Conn, buf, n); err != nil {
|
||||
c.p.ctxCancel(err)
|
||||
return
|
||||
}
|
||||
@@ -86,28 +94,22 @@ func (c Conn) Stop() {
|
||||
c.p.wg.Wait()
|
||||
}
|
||||
|
||||
func NewConn(ctx context.Context, conn essentials.Conn, stats *Stats) Conn {
|
||||
func NewConn(ctx context.Context, conn essentials.Conn, stats Stats) Conn {
|
||||
ctx, cancel := context.WithCancelCause(ctx)
|
||||
rv := Conn{
|
||||
Conn: conn,
|
||||
p: &connPayload{
|
||||
ctx: ctx,
|
||||
ctxCancel: cancel,
|
||||
stats: stats,
|
||||
writtenCond: sync.Cond{
|
||||
L: &sync.Mutex{},
|
||||
},
|
||||
clock: Clock{
|
||||
stats: stats,
|
||||
tick: make(chan struct{}),
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
rv.p.writeStream.Grow(tls.DefaultBufferSize)
|
||||
|
||||
rv.p.wg.Go(func() {
|
||||
rv.p.clock.Start(ctx)
|
||||
})
|
||||
rv.p.wg.Go(func() {
|
||||
rv.start()
|
||||
})
|
||||
|
||||
@@ -63,7 +63,7 @@ func (suite *ConnTestSuite) TearDownTest() {
|
||||
}
|
||||
|
||||
func (suite *ConnTestSuite) makeConn() Conn {
|
||||
return NewConn(suite.ctx, suite.connMock, &Stats{
|
||||
return NewConn(suite.ctx, suite.connMock, Stats{
|
||||
k: 2.0,
|
||||
lambda: 0.01,
|
||||
})
|
||||
@@ -152,7 +152,7 @@ func (suite *ConnTestSuite) TestStopDoesNotDeadlockWhenStartIsWaiting() {
|
||||
ctx, cancel := context.WithCancel(suite.ctx)
|
||||
defer cancel()
|
||||
|
||||
c := NewConn(ctx, suite.connMock, &Stats{
|
||||
c := NewConn(ctx, suite.connMock, Stats{
|
||||
k: 2.0,
|
||||
lambda: 0.01,
|
||||
})
|
||||
|
||||
@@ -2,7 +2,9 @@ package doppel
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"github.com/9seconds/mtg/v2/essentials"
|
||||
@@ -12,8 +14,22 @@ const (
|
||||
DoppelGangerMaxDurations = 4096
|
||||
DoppelGangerScoutRaidEach = 6 * time.Hour
|
||||
DoppelGangerScoutRepeats = 10
|
||||
|
||||
MinCertSizesToCalculate = 3
|
||||
)
|
||||
|
||||
// NoiseParams holds the measured cert chain size for FakeTLS noise calibration.
|
||||
// If Mean is 0, the caller should use a legacy fallback.
|
||||
type NoiseParams struct {
|
||||
Mean int
|
||||
Jitter int
|
||||
}
|
||||
|
||||
type scoutRaidResult struct {
|
||||
durations []time.Duration
|
||||
certSizes []int
|
||||
}
|
||||
|
||||
type gangerConnRequest struct {
|
||||
ret chan<- Conn
|
||||
payload essentials.Conn
|
||||
@@ -31,8 +47,11 @@ type Ganger struct {
|
||||
|
||||
drs bool
|
||||
|
||||
stats *Stats
|
||||
stats Stats
|
||||
durations []time.Duration
|
||||
certSizes []int
|
||||
|
||||
noiseParams atomic.Pointer[NoiseParams]
|
||||
|
||||
connRequests chan gangerConnRequest
|
||||
}
|
||||
@@ -48,6 +67,16 @@ func (g *Ganger) Run() {
|
||||
})
|
||||
}
|
||||
|
||||
// NoiseParams returns the current cert-size-based noise parameters.
|
||||
// Returns zero-value NoiseParams if not yet measured (caller should use fallback).
|
||||
func (g *Ganger) NoiseParams() NoiseParams {
|
||||
if p := g.noiseParams.Load(); p != nil {
|
||||
return *p
|
||||
}
|
||||
|
||||
return NoiseParams{}
|
||||
}
|
||||
|
||||
func (g *Ganger) NewConn(conn essentials.Conn) (Conn, error) {
|
||||
rvChan := make(chan Conn)
|
||||
req := gangerConnRequest{
|
||||
@@ -81,10 +110,10 @@ func (g *Ganger) run() {
|
||||
}
|
||||
}()
|
||||
|
||||
scoutCollectedChan := make(chan []time.Duration)
|
||||
scoutCollectedChan := make(chan scoutRaidResult)
|
||||
currentScoutCollectedChan := scoutCollectedChan
|
||||
|
||||
updatedStatsChan := make(chan *Stats)
|
||||
updatedStatsChan := make(chan Stats)
|
||||
|
||||
g.wg.Go(func() {
|
||||
g.runScoutRaid(scoutCollectedChan)
|
||||
@@ -94,18 +123,29 @@ func (g *Ganger) run() {
|
||||
select {
|
||||
case <-g.ctx.Done():
|
||||
return
|
||||
case durations := <-currentScoutCollectedChan:
|
||||
g.durations = append(g.durations, durations...)
|
||||
case result := <-currentScoutCollectedChan:
|
||||
g.durations = append(g.durations, result.durations...)
|
||||
|
||||
if len(g.durations) > DoppelGangerMaxDurations {
|
||||
copy(g.durations, g.durations[len(g.durations)-DoppelGangerMaxDurations:])
|
||||
g.durations = g.durations[:DoppelGangerMaxDurations]
|
||||
}
|
||||
|
||||
// Update cert sizes and recompute noise params.
|
||||
g.certSizes = append(g.certSizes, result.certSizes...)
|
||||
if len(g.certSizes) > DoppelGangerMaxDurations {
|
||||
g.certSizes = g.certSizes[len(g.certSizes)-DoppelGangerMaxDurations:]
|
||||
}
|
||||
|
||||
if len(g.certSizes) >= MinCertSizesToCalculate {
|
||||
g.updateNoiseParams()
|
||||
}
|
||||
|
||||
if len(g.durations) < MinDurationsToCalculate {
|
||||
continue
|
||||
}
|
||||
|
||||
durations := g.durations
|
||||
currentScoutCollectedChan = nil
|
||||
g.wg.Go(func() {
|
||||
select {
|
||||
@@ -129,8 +169,45 @@ func (g *Ganger) run() {
|
||||
}
|
||||
}
|
||||
|
||||
func (g *Ganger) runScoutRaid(rvChan chan<- []time.Duration) {
|
||||
durations := []time.Duration{}
|
||||
func (g *Ganger) updateNoiseParams() {
|
||||
if len(g.certSizes) == 0 {
|
||||
return
|
||||
}
|
||||
|
||||
sum := 0
|
||||
for _, s := range g.certSizes {
|
||||
sum += s
|
||||
}
|
||||
|
||||
mean := sum / len(g.certSizes)
|
||||
|
||||
maxDev := 0
|
||||
for _, s := range g.certSizes {
|
||||
d := s - mean
|
||||
if d < 0 {
|
||||
d = -d
|
||||
}
|
||||
|
||||
if d > maxDev {
|
||||
maxDev = d
|
||||
}
|
||||
}
|
||||
|
||||
if maxDev < 100 {
|
||||
maxDev = 100
|
||||
}
|
||||
|
||||
np := &NoiseParams{Mean: mean, Jitter: maxDev}
|
||||
g.noiseParams.Store(np)
|
||||
|
||||
g.logger.Info(fmt.Sprintf(
|
||||
"updated noise params: mean=%d jitter=%d samples=%d",
|
||||
mean, maxDev, len(g.certSizes),
|
||||
))
|
||||
}
|
||||
|
||||
func (g *Ganger) runScoutRaid(rvChan chan<- scoutRaidResult) {
|
||||
var result scoutRaidResult
|
||||
|
||||
for range g.scoutRaidRepeats {
|
||||
learned, err := g.scout.Learn(g.ctx)
|
||||
@@ -138,13 +215,18 @@ func (g *Ganger) runScoutRaid(rvChan chan<- []time.Duration) {
|
||||
g.logger.WarningError("cannot learn", err)
|
||||
continue
|
||||
}
|
||||
durations = append(durations, learned...)
|
||||
|
||||
result.durations = append(result.durations, learned.Durations...)
|
||||
|
||||
if learned.CertSize > 0 {
|
||||
result.certSizes = append(result.certSizes, learned.CertSize)
|
||||
}
|
||||
}
|
||||
|
||||
select {
|
||||
case <-g.ctx.Done():
|
||||
return
|
||||
case rvChan <- durations:
|
||||
case rvChan <- result:
|
||||
}
|
||||
}
|
||||
|
||||
@@ -174,7 +256,7 @@ func NewGanger(
|
||||
scoutRaidEach: scoutEach,
|
||||
scoutRaidRepeats: scoutRepeats,
|
||||
drs: drs,
|
||||
stats: &Stats{
|
||||
stats: Stats{
|
||||
k: StatsDefaultK,
|
||||
lambda: StatsDefaultLambda,
|
||||
drs: drs,
|
||||
|
||||
@@ -12,36 +12,46 @@ import (
|
||||
"github.com/9seconds/mtg/v2/mtglib/internal/tls"
|
||||
)
|
||||
|
||||
// ScoutResult holds measurements from a single scout HTTP request.
|
||||
type ScoutResult struct {
|
||||
Durations []time.Duration
|
||||
CertSize int // total ApplicationData bytes during TLS handshake; 0 if unknown
|
||||
}
|
||||
|
||||
type Scout struct {
|
||||
network Network
|
||||
urls []string
|
||||
}
|
||||
|
||||
func (s Scout) Learn(ctx context.Context) ([]time.Duration, error) {
|
||||
var durations []time.Duration
|
||||
func (s Scout) Learn(ctx context.Context) (ScoutResult, error) {
|
||||
var combined ScoutResult
|
||||
|
||||
for _, url := range s.urls {
|
||||
learned, err := s.learn(ctx, url)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
return ScoutResult{}, err
|
||||
}
|
||||
|
||||
durations = append(durations, learned...)
|
||||
combined.Durations = append(combined.Durations, learned.Durations...)
|
||||
|
||||
if learned.CertSize > 0 && combined.CertSize == 0 {
|
||||
combined.CertSize = learned.CertSize
|
||||
}
|
||||
}
|
||||
|
||||
return durations, nil
|
||||
return combined, nil
|
||||
}
|
||||
|
||||
func (s Scout) learn(ctx context.Context, url string) ([]time.Duration, error) {
|
||||
func (s Scout) learn(ctx context.Context, url string) (ScoutResult, error) {
|
||||
client, results := s.makeClient()
|
||||
|
||||
if !strings.HasPrefix(url, "https://") {
|
||||
return nil, fmt.Errorf("url %s must be https", url)
|
||||
return ScoutResult{}, fmt.Errorf("url %s must be https", url)
|
||||
}
|
||||
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
return ScoutResult{}, err
|
||||
}
|
||||
|
||||
resp, err := client.Do(req)
|
||||
@@ -51,31 +61,62 @@ func (s Scout) learn(ctx context.Context, url string) ([]time.Duration, error) {
|
||||
client.CloseIdleConnections()
|
||||
}
|
||||
|
||||
if err != nil || len(results.data) == 0 {
|
||||
return nil, err
|
||||
if err != nil {
|
||||
return ScoutResult{}, err
|
||||
}
|
||||
|
||||
durations := []time.Duration{}
|
||||
data, writeIndex := results.Snapshot()
|
||||
|
||||
if len(data) == 0 {
|
||||
return ScoutResult{}, nil
|
||||
}
|
||||
|
||||
var result ScoutResult
|
||||
|
||||
// Compute inter-record durations (existing logic).
|
||||
lastTimestamp := time.Time{}
|
||||
|
||||
for i, v := range results.data {
|
||||
for i, v := range data {
|
||||
if v.recordType != tls.TypeApplicationData {
|
||||
continue
|
||||
}
|
||||
|
||||
if lastTimestamp.IsZero() {
|
||||
if i > 0 {
|
||||
lastTimestamp = results.data[i-1].timestamp
|
||||
lastTimestamp = data[i-1].timestamp
|
||||
} else {
|
||||
lastTimestamp = v.timestamp
|
||||
}
|
||||
}
|
||||
|
||||
durations = append(durations, v.timestamp.Sub(lastTimestamp))
|
||||
result.Durations = append(result.Durations, v.timestamp.Sub(lastTimestamp))
|
||||
lastTimestamp = v.timestamp
|
||||
}
|
||||
|
||||
return durations, nil
|
||||
// Compute cert size: sum of ApplicationData payload between CCS and
|
||||
// the first client Write (which marks the end of server handshake).
|
||||
seenCCS := false
|
||||
boundary := writeIndex
|
||||
if boundary < 0 {
|
||||
boundary = len(data)
|
||||
}
|
||||
|
||||
for i, v := range data {
|
||||
if i >= boundary {
|
||||
break
|
||||
}
|
||||
|
||||
if v.recordType == tls.TypeChangeCipherSpec {
|
||||
seenCCS = true
|
||||
continue
|
||||
}
|
||||
|
||||
if seenCCS && v.recordType == tls.TypeApplicationData {
|
||||
result.CertSize += v.payloadLen
|
||||
}
|
||||
}
|
||||
|
||||
return result, nil
|
||||
}
|
||||
|
||||
func (s Scout) makeClient() (*http.Client, *ScoutConnCollected) {
|
||||
|
||||
@@ -14,9 +14,10 @@ type ScoutConn struct {
|
||||
|
||||
results *ScoutConnCollected
|
||||
rawBuf *bytes.Buffer
|
||||
seenCCS bool
|
||||
}
|
||||
|
||||
func (s ScoutConn) Read(p []byte) (int, error) {
|
||||
func (s *ScoutConn) Read(p []byte) (int, error) {
|
||||
buf := &bytes.Buffer{}
|
||||
|
||||
for {
|
||||
@@ -31,7 +32,11 @@ func (s ScoutConn) Read(p []byte) (int, error) {
|
||||
return 0, err
|
||||
}
|
||||
|
||||
s.results.Add(recordType)
|
||||
if recordType == tls.TypeChangeCipherSpec {
|
||||
s.seenCCS = true
|
||||
}
|
||||
|
||||
s.results.Add(recordType, int(length))
|
||||
s.rawBuf.Write([]byte{recordType})
|
||||
s.rawBuf.Write(tls.TLSVersion[:])
|
||||
|
||||
@@ -45,11 +50,19 @@ func (s ScoutConn) Read(p []byte) (int, error) {
|
||||
}
|
||||
}
|
||||
|
||||
func NewScoutConn(conn essentials.Conn, results *ScoutConnCollected) ScoutConn {
|
||||
func (s *ScoutConn) Write(p []byte) (int, error) {
|
||||
if s.seenCCS {
|
||||
s.results.MarkWrite()
|
||||
}
|
||||
|
||||
return s.Conn.Write(p)
|
||||
}
|
||||
|
||||
func NewScoutConn(conn essentials.Conn, results *ScoutConnCollected) *ScoutConn {
|
||||
rawBuf := &bytes.Buffer{}
|
||||
rawBuf.Grow(tls.MaxRecordSize)
|
||||
|
||||
return ScoutConn{
|
||||
return &ScoutConn{
|
||||
Conn: tls.New(conn, false, false),
|
||||
results: results,
|
||||
rawBuf: rawBuf,
|
||||
|
||||
@@ -1,6 +1,10 @@
|
||||
package doppel
|
||||
|
||||
import "time"
|
||||
import (
|
||||
"slices"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
const (
|
||||
ScoutConnCollectedPreallocSize = 100
|
||||
@@ -9,21 +13,47 @@ const (
|
||||
type ScoutConnResult struct {
|
||||
timestamp time.Time
|
||||
recordType byte
|
||||
payloadLen int
|
||||
}
|
||||
|
||||
type ScoutConnCollected struct {
|
||||
data []ScoutConnResult
|
||||
mu sync.Mutex
|
||||
data []ScoutConnResult
|
||||
writeIndex int // index at which client first wrote post-handshake data; -1 if not set
|
||||
}
|
||||
|
||||
func (s *ScoutConnCollected) Add(record byte) {
|
||||
func (s *ScoutConnCollected) Add(record byte, payloadLen int) {
|
||||
s.mu.Lock()
|
||||
s.data = append(s.data, ScoutConnResult{
|
||||
timestamp: time.Now(),
|
||||
recordType: record,
|
||||
payloadLen: payloadLen,
|
||||
})
|
||||
s.mu.Unlock()
|
||||
}
|
||||
|
||||
// MarkWrite records the current data length as the handshake boundary.
|
||||
func (s *ScoutConnCollected) MarkWrite() {
|
||||
s.mu.Lock()
|
||||
if s.writeIndex < 0 {
|
||||
s.writeIndex = len(s.data)
|
||||
}
|
||||
s.mu.Unlock()
|
||||
}
|
||||
|
||||
// Snapshot returns a copy of the collected data and the write index.
|
||||
func (s *ScoutConnCollected) Snapshot() ([]ScoutConnResult, int) {
|
||||
s.mu.Lock()
|
||||
snapshot := slices.Clone(s.data)
|
||||
writeIndex := s.writeIndex
|
||||
s.mu.Unlock()
|
||||
|
||||
return snapshot, writeIndex
|
||||
}
|
||||
|
||||
func NewScoutConnCollected() *ScoutConnCollected {
|
||||
return &ScoutConnCollected{
|
||||
data: make([]ScoutConnResult, 0, ScoutConnCollectedPreallocSize),
|
||||
data: make([]ScoutConnResult, 0, ScoutConnCollectedPreallocSize),
|
||||
writeIndex: -1,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package doppel
|
||||
|
||||
import (
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
@@ -14,28 +15,71 @@ type ScoutConnCollectedTestSuite struct {
|
||||
|
||||
func (suite *ScoutConnCollectedTestSuite) TestAddSingle() {
|
||||
collected := NewScoutConnCollected()
|
||||
collected.Add(tls.TypeApplicationData)
|
||||
collected.Add(tls.TypeApplicationData, 100)
|
||||
|
||||
suite.Len(collected.data, 1)
|
||||
suite.Equal(byte(tls.TypeApplicationData), collected.data[0].recordType)
|
||||
data, _ := collected.Snapshot()
|
||||
|
||||
suite.Len(data, 1)
|
||||
suite.Equal(byte(tls.TypeApplicationData), data[0].recordType)
|
||||
}
|
||||
|
||||
func (suite *ScoutConnCollectedTestSuite) TestAddTimestampsAreMonotonic() {
|
||||
collected := NewScoutConnCollected()
|
||||
|
||||
collected.Add(tls.TypeApplicationData)
|
||||
collected.Add(tls.TypeApplicationData, 100)
|
||||
|
||||
time.Sleep(time.Microsecond)
|
||||
collected.Add(tls.TypeApplicationData)
|
||||
collected.Add(tls.TypeApplicationData, 100)
|
||||
|
||||
time.Sleep(time.Microsecond)
|
||||
collected.Add(tls.TypeApplicationData)
|
||||
collected.Add(tls.TypeApplicationData, 100)
|
||||
|
||||
for i := 1; i < len(collected.data); i++ {
|
||||
suite.True(collected.data[i].timestamp.After(collected.data[i-1].timestamp))
|
||||
data, _ := collected.Snapshot()
|
||||
|
||||
for i := 1; i < len(data); i++ {
|
||||
suite.True(data[i].timestamp.After(data[i-1].timestamp))
|
||||
}
|
||||
}
|
||||
|
||||
func (suite *ScoutConnCollectedTestSuite) TestConcurrentAddSnapshot() {
|
||||
collected := NewScoutConnCollected()
|
||||
|
||||
var wg sync.WaitGroup
|
||||
|
||||
wg.Add(3)
|
||||
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
|
||||
for i := 0; i < 1000; i++ {
|
||||
collected.Add(tls.TypeApplicationData, i)
|
||||
}
|
||||
}()
|
||||
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
|
||||
for i := 0; i < 100; i++ {
|
||||
collected.MarkWrite()
|
||||
}
|
||||
}()
|
||||
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
|
||||
for i := 0; i < 1000; i++ {
|
||||
// call Snapshot concurrently to exercise the lock under -race
|
||||
collected.Snapshot() //nolint:errcheck
|
||||
}
|
||||
}()
|
||||
|
||||
wg.Wait()
|
||||
|
||||
data, writeIndex := collected.Snapshot()
|
||||
suite.Len(data, 1000)
|
||||
suite.GreaterOrEqual(writeIndex, 0)
|
||||
}
|
||||
|
||||
func TestScoutConnCollected(t *testing.T) {
|
||||
t.Parallel()
|
||||
suite.Run(t, &ScoutConnCollectedTestSuite{})
|
||||
|
||||
@@ -22,9 +22,9 @@ func (suite *ScoutTestSuite) SetupSuite() {
|
||||
}
|
||||
|
||||
func (suite *ScoutTestSuite) TestCollectResults() {
|
||||
durations, err := suite.scout.Learn(suite.ctx)
|
||||
result, err := suite.scout.Learn(suite.ctx)
|
||||
suite.NoError(err)
|
||||
suite.Less(3, len(durations))
|
||||
suite.Less(3, len(result.Durations))
|
||||
}
|
||||
|
||||
func (suite *ScoutTestSuite) TestCollectNothing() {
|
||||
|
||||
@@ -112,7 +112,7 @@ func (d *Stats) Size() int {
|
||||
return TLSRecordSizeMax
|
||||
}
|
||||
|
||||
func NewStats(durations []time.Duration, drs bool) *Stats {
|
||||
func NewStats(durations []time.Duration, drs bool) Stats {
|
||||
n := float64(len(durations))
|
||||
|
||||
// in milliseconds
|
||||
@@ -162,7 +162,7 @@ func NewStats(durations []time.Duration, drs bool) *Stats {
|
||||
// λ = (Σxᵢᵏ / n)^(1/k)
|
||||
lambda := math.Pow(sumXK/n, 1.0/k)
|
||||
|
||||
return &Stats{
|
||||
return Stats{
|
||||
k: k,
|
||||
lambda: lambda,
|
||||
drs: drs,
|
||||
|
||||
@@ -0,0 +1,13 @@
|
||||
//go:build mips || mipsle
|
||||
|
||||
package relay
|
||||
|
||||
import "github.com/9seconds/mtg/v2/mtglib/internal/tls"
|
||||
|
||||
const (
|
||||
// MIPS is quite short in resources, and usually it means that it will run
|
||||
// on Microtiks, OpenWRT-based routers or similar hardware. I think it worth
|
||||
// to sacrifice a number of read syscalls (read, CPU load) to shrink
|
||||
// limited RAM resources.
|
||||
bufPoolSize = tls.MaxRecordPayloadSize / 2
|
||||
)
|
||||
@@ -0,0 +1,9 @@
|
||||
//go:build !mips && !mipsle
|
||||
|
||||
package relay
|
||||
|
||||
import "github.com/9seconds/mtg/v2/mtglib/internal/tls"
|
||||
|
||||
const (
|
||||
bufPoolSize = tls.MaxRecordPayloadSize
|
||||
)
|
||||
@@ -0,0 +1,18 @@
|
||||
package relay
|
||||
|
||||
import "sync"
|
||||
|
||||
var bufPool = sync.Pool{
|
||||
New: func() any {
|
||||
b := make([]byte, bufPoolSize)
|
||||
return &b
|
||||
},
|
||||
}
|
||||
|
||||
func acquireBuffer() *[]byte {
|
||||
return bufPool.Get().(*[]byte)
|
||||
}
|
||||
|
||||
func releaseBuffer(p *[]byte) {
|
||||
bufPool.Put(p)
|
||||
}
|
||||
@@ -6,7 +6,6 @@ import (
|
||||
"io"
|
||||
|
||||
"github.com/9seconds/mtg/v2/essentials"
|
||||
"github.com/9seconds/mtg/v2/mtglib/internal/tls"
|
||||
)
|
||||
|
||||
func Relay(ctx context.Context, log Logger, telegramConn, clientConn essentials.Conn) {
|
||||
@@ -16,11 +15,11 @@ func Relay(ctx context.Context, log Logger, telegramConn, clientConn essentials.
|
||||
ctx, cancel := context.WithCancel(ctx)
|
||||
defer cancel()
|
||||
|
||||
go func() {
|
||||
<-ctx.Done()
|
||||
stop := context.AfterFunc(ctx, func() {
|
||||
telegramConn.Close() //nolint: errcheck
|
||||
clientConn.Close() //nolint: errcheck
|
||||
}()
|
||||
})
|
||||
defer stop()
|
||||
|
||||
closeChan := make(chan struct{})
|
||||
|
||||
@@ -36,12 +35,13 @@ func Relay(ctx context.Context, log Logger, telegramConn, clientConn essentials.
|
||||
}
|
||||
|
||||
func pump(log Logger, src, dst essentials.Conn, direction string) {
|
||||
var buf [tls.MaxRecordPayloadSize]byte
|
||||
buf := acquireBuffer()
|
||||
defer releaseBuffer(buf)
|
||||
|
||||
defer src.CloseRead() //nolint: errcheck
|
||||
defer dst.CloseWrite() //nolint: errcheck
|
||||
|
||||
n, err := io.CopyBuffer(src, dst, buf[:])
|
||||
n, err := io.CopyBuffer(src, dst, *buf)
|
||||
|
||||
switch {
|
||||
case err == nil:
|
||||
|
||||
@@ -11,8 +11,6 @@ import (
|
||||
"net"
|
||||
"slices"
|
||||
"time"
|
||||
|
||||
"github.com/9seconds/mtg/v2/mtglib/internal/tls"
|
||||
)
|
||||
|
||||
const (
|
||||
@@ -56,25 +54,17 @@ func ReadClientHello(
|
||||
// 4. New digest should be all 0 except of last 4 bytes
|
||||
// 5. Last 4 bytes are little endian uint32 of UNIX timestamp when
|
||||
// this message was created.
|
||||
handshakeCopyBuf := &bytes.Buffer{}
|
||||
reader := io.TeeReader(conn, handshakeCopyBuf)
|
||||
|
||||
reader, err := parseTLSHeader(reader)
|
||||
clientHelloCopy, handshakeReader, err := parseClientHello(conn)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("cannot parse tls header: %w", err)
|
||||
return nil, fmt.Errorf("cannot read client hello: %w", err)
|
||||
}
|
||||
|
||||
reader, err = parseHandshakeHeader(reader)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("cannot parse handshake header: %w", err)
|
||||
}
|
||||
|
||||
hello, err := parseHandshake(reader)
|
||||
hello, err := parseHandshake(handshakeReader)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("cannot parse handshake: %w", err)
|
||||
}
|
||||
|
||||
sniHostnames, err := parseSNI(reader)
|
||||
sniHostnames, err := parseSNI(handshakeReader)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("cannot parse SNI: %w", err)
|
||||
}
|
||||
@@ -85,10 +75,10 @@ func ReadClientHello(
|
||||
|
||||
digest := hmac.New(sha256.New, secret)
|
||||
// we write a copy of the handshake with client random all nullified.
|
||||
digest.Write(handshakeCopyBuf.Next(RandomOffset))
|
||||
handshakeCopyBuf.Next(RandomLen)
|
||||
digest.Write(clientHelloCopy.Next(RandomOffset))
|
||||
clientHelloCopy.Next(RandomLen)
|
||||
digest.Write(emptyRandom[:])
|
||||
digest.Write(handshakeCopyBuf.Bytes())
|
||||
digest.Write(clientHelloCopy.Bytes())
|
||||
|
||||
computed := digest.Sum(nil)
|
||||
|
||||
@@ -110,58 +100,6 @@ func ReadClientHello(
|
||||
return hello, nil
|
||||
}
|
||||
|
||||
func parseTLSHeader(r io.Reader) (io.Reader, error) {
|
||||
// record_type(1) + version(2) + size(2)
|
||||
// 16 - type is 0x16 (handshake record)
|
||||
// 03 01 - protocol version is "3,1" (also known as TLS 1.0)
|
||||
// 00 f8 - 0xF8 (248) bytes of handshake message follows
|
||||
header := [1 + 2 + 2]byte{}
|
||||
|
||||
if _, err := io.ReadFull(r, header[:]); err != nil {
|
||||
return nil, fmt.Errorf("cannot read record header: %w", err)
|
||||
}
|
||||
|
||||
if header[0] != tls.TypeHandshake {
|
||||
return nil, fmt.Errorf("unexpected record type %#x", header[0])
|
||||
}
|
||||
|
||||
if header[1] != 3 || header[2] != 1 {
|
||||
return nil, fmt.Errorf("unexpected protocol version %#x %#x", header[1], header[2])
|
||||
}
|
||||
|
||||
length := int64(binary.BigEndian.Uint16(header[3:]))
|
||||
buf := &bytes.Buffer{}
|
||||
|
||||
_, err := io.CopyN(buf, r, length)
|
||||
|
||||
return buf, err
|
||||
}
|
||||
|
||||
func parseHandshakeHeader(r io.Reader) (io.Reader, error) {
|
||||
// type(1) + size(3 / uint24)
|
||||
// 01 - handshake message type 0x01 (client hello)
|
||||
// 00 00 f4 - 0xF4 (244) bytes of client hello data follows
|
||||
header := [1 + 3]byte{}
|
||||
|
||||
if _, err := io.ReadFull(r, header[:]); err != nil {
|
||||
return nil, fmt.Errorf("cannot read handshake header: %w", err)
|
||||
}
|
||||
|
||||
if header[0] != TypeHandshakeClient {
|
||||
return nil, fmt.Errorf("incorrect handshake type: %#x", header[0])
|
||||
}
|
||||
|
||||
// unfortunately there is not uint24 in golang, so we just reust header
|
||||
header[0] = 0
|
||||
|
||||
length := int64(binary.BigEndian.Uint32(header[:]))
|
||||
buf := &bytes.Buffer{}
|
||||
|
||||
_, err := io.CopyN(buf, r, length)
|
||||
|
||||
return buf, err
|
||||
}
|
||||
|
||||
func parseHandshake(r io.Reader) (*ClientHello, error) {
|
||||
// A protocol version of "3,3" (meaning TLS 1.2) is given.
|
||||
header := [2]byte{}
|
||||
|
||||
@@ -3,8 +3,10 @@ package fake_test
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/binary"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"io"
|
||||
"os"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
@@ -393,3 +395,273 @@ func TestParseClientHelloSNI(t *testing.T) {
|
||||
t.Parallel()
|
||||
suite.Run(t, &ParseClientHelloSNITestSuite{})
|
||||
}
|
||||
|
||||
// fragmentTLSRecord splits a single TLS record into n TLS records by
|
||||
// dividing the payload into roughly equal parts. Each part gets its own
|
||||
// TLS record header with the same record type and version.
|
||||
func fragmentTLSRecord(t testing.TB, full []byte, n int) []byte {
|
||||
t.Helper()
|
||||
|
||||
recordType := full[0]
|
||||
version := full[1:3]
|
||||
payload := full[tls.SizeHeader:]
|
||||
|
||||
chunkSize := len(payload) / n
|
||||
result := &bytes.Buffer{}
|
||||
|
||||
for i := 0; i < n; i++ {
|
||||
start := i * chunkSize
|
||||
end := start + chunkSize
|
||||
|
||||
if i == n-1 {
|
||||
end = len(payload)
|
||||
}
|
||||
|
||||
chunk := payload[start:end]
|
||||
result.WriteByte(recordType)
|
||||
result.Write(version)
|
||||
require.NoError(t, binary.Write(result, binary.BigEndian, uint16(len(chunk))))
|
||||
result.Write(chunk)
|
||||
}
|
||||
|
||||
return result.Bytes()
|
||||
}
|
||||
|
||||
// splitPayloadAt creates two TLS records from a single record by splitting
|
||||
// the payload at the given byte position.
|
||||
func splitPayloadAt(t testing.TB, full []byte, pos int) []byte {
|
||||
t.Helper()
|
||||
|
||||
payload := full[tls.SizeHeader:]
|
||||
buf := &bytes.Buffer{}
|
||||
|
||||
buf.WriteByte(tls.TypeHandshake)
|
||||
buf.Write(full[1:3])
|
||||
require.NoError(t, binary.Write(buf, binary.BigEndian, uint16(pos)))
|
||||
buf.Write(payload[:pos])
|
||||
|
||||
buf.WriteByte(tls.TypeHandshake)
|
||||
buf.Write(full[1:3])
|
||||
require.NoError(t, binary.Write(buf, binary.BigEndian, uint16(len(payload)-pos)))
|
||||
buf.Write(payload[pos:])
|
||||
|
||||
return buf.Bytes()
|
||||
}
|
||||
|
||||
type ParseClientHelloFragmentedTestSuite struct {
|
||||
suite.Suite
|
||||
|
||||
secret mtglib.Secret
|
||||
snapshot *clientHelloSnapshot
|
||||
}
|
||||
|
||||
func (s *ParseClientHelloFragmentedTestSuite) SetupSuite() {
|
||||
parsed, err := mtglib.ParseSecret(
|
||||
"ee367a189aee18fa31c190054efd4a8e9573746f726167652e676f6f676c65617069732e636f6d",
|
||||
)
|
||||
require.NoError(s.T(), err)
|
||||
|
||||
s.secret = parsed
|
||||
|
||||
fileData, err := os.ReadFile("testdata/client-hello-ok-19dfe38384b9884b.json")
|
||||
require.NoError(s.T(), err)
|
||||
|
||||
s.snapshot = &clientHelloSnapshot{}
|
||||
require.NoError(s.T(), json.Unmarshal(fileData, s.snapshot))
|
||||
}
|
||||
|
||||
func (s *ParseClientHelloFragmentedTestSuite) makeConn(data []byte) *parseClientHelloConnMock {
|
||||
readBuf := &bytes.Buffer{}
|
||||
readBuf.Write(data)
|
||||
|
||||
connMock := &parseClientHelloConnMock{
|
||||
readBuf: readBuf,
|
||||
}
|
||||
|
||||
connMock.
|
||||
On("SetReadDeadline", mock.AnythingOfType("time.Time")).
|
||||
Twice().
|
||||
Return(nil)
|
||||
|
||||
return connMock
|
||||
}
|
||||
|
||||
func (s *ParseClientHelloFragmentedTestSuite) TestReassemblySuccess() {
|
||||
full := s.snapshot.GetFull()
|
||||
|
||||
tests := []struct {
|
||||
name string
|
||||
data []byte
|
||||
}{
|
||||
{"two equal fragments", fragmentTLSRecord(s.T(), full, 2)},
|
||||
{"three equal fragments", fragmentTLSRecord(s.T(), full, 3)},
|
||||
{"single byte first fragment", splitPayloadAt(s.T(), full, 1)},
|
||||
{"three byte first fragment", splitPayloadAt(s.T(), full, 3)},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
s.Run(tt.name, func() {
|
||||
connMock := s.makeConn(tt.data)
|
||||
defer connMock.AssertExpectations(s.T())
|
||||
|
||||
hello, err := fake.ReadClientHello(
|
||||
connMock,
|
||||
s.secret.Key[:],
|
||||
s.secret.Host,
|
||||
TolerateTime,
|
||||
)
|
||||
s.Require().NoError(err)
|
||||
|
||||
s.Equal(s.snapshot.GetRandom(), hello.Random[:])
|
||||
s.Equal(s.snapshot.GetSessionID(), hello.SessionID)
|
||||
s.Equal(uint16(s.snapshot.CipherSuite), hello.CipherSuite)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func (s *ParseClientHelloFragmentedTestSuite) TestReassemblyErrors() {
|
||||
full := s.snapshot.GetFull()
|
||||
payload := full[tls.SizeHeader:]
|
||||
|
||||
tests := []struct {
|
||||
name string
|
||||
buildData func() []byte
|
||||
errMsg string
|
||||
}{
|
||||
{
|
||||
name: "wrong continuation record type",
|
||||
buildData: func() []byte {
|
||||
buf := &bytes.Buffer{}
|
||||
buf.WriteByte(tls.TypeHandshake)
|
||||
buf.Write(full[1:3])
|
||||
require.NoError(s.T(), binary.Write(buf, binary.BigEndian, uint16(10)))
|
||||
buf.Write(payload[:10])
|
||||
// Wrong type: application data instead of handshake
|
||||
buf.WriteByte(tls.TypeApplicationData)
|
||||
buf.Write(full[1:3])
|
||||
require.NoError(s.T(), binary.Write(buf, binary.BigEndian, uint16(len(payload)-10)))
|
||||
buf.Write(payload[10:])
|
||||
return buf.Bytes()
|
||||
},
|
||||
errMsg: "unexpected record type",
|
||||
},
|
||||
{
|
||||
name: "too many continuation records",
|
||||
buildData: func() []byte {
|
||||
// Handshake header claiming 256 bytes, but we only send 1 byte per continuation
|
||||
handshakePayload := []byte{0x01, 0x00, 0x01, 0x00}
|
||||
buf := &bytes.Buffer{}
|
||||
buf.WriteByte(tls.TypeHandshake)
|
||||
buf.Write([]byte{3, 1})
|
||||
require.NoError(s.T(), binary.Write(buf, binary.BigEndian, uint16(len(handshakePayload))))
|
||||
buf.Write(handshakePayload)
|
||||
for range 11 {
|
||||
buf.WriteByte(tls.TypeHandshake)
|
||||
buf.Write([]byte{3, 1})
|
||||
require.NoError(s.T(), binary.Write(buf, binary.BigEndian, uint16(1)))
|
||||
buf.WriteByte(0xAB)
|
||||
}
|
||||
return buf.Bytes()
|
||||
},
|
||||
errMsg: "too many fragments",
|
||||
},
|
||||
{
|
||||
name: "zero-length continuation record",
|
||||
buildData: func() []byte {
|
||||
buf := &bytes.Buffer{}
|
||||
buf.WriteByte(tls.TypeHandshake)
|
||||
buf.Write(full[1:3])
|
||||
require.NoError(s.T(), binary.Write(buf, binary.BigEndian, uint16(10)))
|
||||
buf.Write(payload[:10])
|
||||
// Valid header but zero-length payload
|
||||
buf.WriteByte(tls.TypeHandshake)
|
||||
buf.Write(full[1:3])
|
||||
require.NoError(s.T(), binary.Write(buf, binary.BigEndian, uint16(0)))
|
||||
return buf.Bytes()
|
||||
},
|
||||
errMsg: "cannot read record header",
|
||||
},
|
||||
{
|
||||
name: "wrong continuation record version",
|
||||
buildData: func() []byte {
|
||||
buf := &bytes.Buffer{}
|
||||
buf.WriteByte(tls.TypeHandshake)
|
||||
buf.Write(full[1:3])
|
||||
require.NoError(s.T(), binary.Write(buf, binary.BigEndian, uint16(10)))
|
||||
buf.Write(payload[:10])
|
||||
// Wrong version: 3.3 instead of 3.1
|
||||
buf.WriteByte(tls.TypeHandshake)
|
||||
buf.Write([]byte{3, 3})
|
||||
require.NoError(s.T(), binary.Write(buf, binary.BigEndian, uint16(len(payload)-10)))
|
||||
buf.Write(payload[10:])
|
||||
return buf.Bytes()
|
||||
},
|
||||
errMsg: "unexpected protocol version",
|
||||
},
|
||||
{
|
||||
name: "handshake message too large",
|
||||
buildData: func() []byte {
|
||||
// Handshake header claiming 0x010000 (65536) bytes — exceeds 0xFFFF limit
|
||||
handshakePayload := []byte{0x01, 0x01, 0x00, 0x00}
|
||||
buf := &bytes.Buffer{}
|
||||
buf.WriteByte(tls.TypeHandshake)
|
||||
buf.Write([]byte{3, 1})
|
||||
require.NoError(s.T(), binary.Write(buf, binary.BigEndian, uint16(len(handshakePayload))))
|
||||
buf.Write(handshakePayload)
|
||||
return buf.Bytes()
|
||||
},
|
||||
errMsg: "cannot read record header",
|
||||
},
|
||||
{
|
||||
name: "truncated continuation record header",
|
||||
buildData: func() []byte {
|
||||
buf := &bytes.Buffer{}
|
||||
buf.WriteByte(tls.TypeHandshake)
|
||||
buf.Write(full[1:3])
|
||||
require.NoError(s.T(), binary.Write(buf, binary.BigEndian, uint16(10)))
|
||||
buf.Write(payload[:10])
|
||||
// Connection ends mid-header (only 2 bytes)
|
||||
buf.WriteByte(tls.TypeHandshake)
|
||||
buf.WriteByte(3)
|
||||
return buf.Bytes()
|
||||
},
|
||||
errMsg: "cannot read record header",
|
||||
},
|
||||
{
|
||||
name: "truncated continuation record payload",
|
||||
buildData: func() []byte {
|
||||
buf := &bytes.Buffer{}
|
||||
buf.WriteByte(tls.TypeHandshake)
|
||||
buf.Write(full[1:3])
|
||||
require.NoError(s.T(), binary.Write(buf, binary.BigEndian, uint16(10)))
|
||||
buf.Write(payload[:10])
|
||||
// Claims 100 bytes but no payload follows
|
||||
buf.WriteByte(tls.TypeHandshake)
|
||||
buf.Write(full[1:3])
|
||||
require.NoError(s.T(), binary.Write(buf, binary.BigEndian, uint16(100)))
|
||||
return buf.Bytes()
|
||||
},
|
||||
errMsg: "EOF",
|
||||
},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
s.Run(tt.name, func() {
|
||||
connMock := s.makeConn(tt.buildData())
|
||||
defer connMock.AssertExpectations(s.T())
|
||||
|
||||
_, err := fake.ReadClientHello(
|
||||
connMock,
|
||||
s.secret.Key[:],
|
||||
s.secret.Host,
|
||||
TolerateTime,
|
||||
)
|
||||
s.ErrorContains(err, tt.errMsg)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseClientHelloFragmented(t *testing.T) {
|
||||
t.Parallel()
|
||||
suite.Run(t, &ParseClientHelloFragmentedTestSuite{})
|
||||
}
|
||||
|
||||
@@ -9,11 +9,18 @@ import (
|
||||
"io"
|
||||
rnd "math/rand/v2"
|
||||
|
||||
"github.com/9seconds/mtg/v2/mtglib/internal/doppel"
|
||||
"github.com/9seconds/mtg/v2/mtglib/internal/tls"
|
||||
"golang.org/x/crypto/curve25519"
|
||||
)
|
||||
|
||||
// NoiseParams controls the size of the fake ApplicationData record
|
||||
// in ServerHello. If Mean is 0, the legacy random range (2500-4700)
|
||||
// is used.
|
||||
type NoiseParams struct {
|
||||
Mean int
|
||||
Jitter int
|
||||
}
|
||||
|
||||
const (
|
||||
TypeHandshakeServer = 0x02
|
||||
ChangeCipherValue = 0x01
|
||||
@@ -33,13 +40,13 @@ var serverHelloSuffix = []byte{
|
||||
0x00, 0x20, // 32 bytes of key
|
||||
}
|
||||
|
||||
func SendServerHello(w io.Writer, secret []byte, clientHello *ClientHello) error {
|
||||
func SendServerHello(w io.Writer, secret []byte, clientHello *ClientHello, noise NoiseParams) error {
|
||||
buf := &bytes.Buffer{}
|
||||
buf.Grow(tls.MaxRecordSize)
|
||||
|
||||
generateServerHello(buf, clientHello)
|
||||
generateChangeCipherValue(buf)
|
||||
generateNoise(buf)
|
||||
generateNoise(buf, noise)
|
||||
|
||||
packet := buf.Bytes()
|
||||
digest := hmac.New(sha256.New, secret)
|
||||
@@ -125,19 +132,31 @@ func generateChangeCipherValue(buf *bytes.Buffer) {
|
||||
buf.WriteByte(ChangeCipherValue)
|
||||
}
|
||||
|
||||
func generateNoise(buf *bytes.Buffer) {
|
||||
data := make(
|
||||
[]byte,
|
||||
int64(
|
||||
doppel.TLSRecordSizeStart+rnd.IntN(
|
||||
doppel.TLSRecordSizeAccel-doppel.TLSRecordSizeStart,
|
||||
),
|
||||
),
|
||||
)
|
||||
// generateNoise writes a single ApplicationData record mimicking the combined
|
||||
// size of a real TLS 1.3 encrypted server handshake (EncryptedExtensions +
|
||||
// Certificate chain + CertificateVerify + Finished).
|
||||
//
|
||||
// NOTE: Must be exactly ONE ApplicationData record — the Telegram client reads
|
||||
// ServerHello + CCS + 1 ApplicationData and computes HMAC over all three.
|
||||
// Multiple records would cause HMAC mismatch and connection failure.
|
||||
func generateNoise(buf *bytes.Buffer, noise NoiseParams) {
|
||||
var size int
|
||||
|
||||
if _, err := rand.Read(data[:]); err != nil {
|
||||
if noise.Mean > 0 && noise.Jitter > 0 {
|
||||
// Calibrated: use measured cert chain size ± jitter.
|
||||
size = noise.Mean - noise.Jitter + rnd.IntN(2*noise.Jitter)
|
||||
if size < 1000 {
|
||||
size = 1000
|
||||
}
|
||||
} else {
|
||||
// Legacy fallback: random in 2500-4700 range.
|
||||
size = 2500 + rnd.IntN(2200)
|
||||
}
|
||||
|
||||
data := make([]byte, size)
|
||||
if _, err := rand.Read(data); err != nil {
|
||||
panic(err)
|
||||
}
|
||||
|
||||
tls.WriteRecord(buf, data[:]) //nolint: errcheck
|
||||
tls.WriteRecord(buf, data) //nolint: errcheck
|
||||
}
|
||||
|
||||
@@ -8,7 +8,6 @@ import (
|
||||
"testing"
|
||||
|
||||
"github.com/9seconds/mtg/v2/mtglib"
|
||||
"github.com/9seconds/mtg/v2/mtglib/internal/doppel"
|
||||
"github.com/9seconds/mtg/v2/mtglib/internal/tls"
|
||||
"github.com/9seconds/mtg/v2/mtglib/internal/tls/fake"
|
||||
"github.com/stretchr/testify/suite"
|
||||
@@ -39,7 +38,7 @@ func (suite *SendServerHelloTestSuite) SetupTest() {
|
||||
}
|
||||
|
||||
func (suite *SendServerHelloTestSuite) TestRecordStructure() {
|
||||
err := fake.SendServerHello(suite.buf, suite.secret.Key[:], suite.hello)
|
||||
err := fake.SendServerHello(suite.buf, suite.secret.Key[:], suite.hello, fake.NoiseParams{})
|
||||
suite.NoError(err)
|
||||
|
||||
var rec bytes.Buffer
|
||||
@@ -59,13 +58,13 @@ func (suite *SendServerHelloTestSuite) TestRecordStructure() {
|
||||
recordType, length, err := tls.ReadRecord(suite.buf, &rec)
|
||||
suite.NoError(err)
|
||||
suite.Equal(byte(tls.TypeApplicationData), recordType)
|
||||
suite.Greater(length, int64(doppel.TLSRecordSizeStart))
|
||||
suite.Greater(length, int64(2500))
|
||||
|
||||
suite.Empty(suite.buf.Bytes())
|
||||
}
|
||||
|
||||
func (suite *SendServerHelloTestSuite) TestHMAC() {
|
||||
err := fake.SendServerHello(suite.buf, suite.secret.Key[:], suite.hello)
|
||||
err := fake.SendServerHello(suite.buf, suite.secret.Key[:], suite.hello, fake.NoiseParams{})
|
||||
suite.NoError(err)
|
||||
|
||||
packet := make([]byte, suite.buf.Len())
|
||||
@@ -83,7 +82,7 @@ func (suite *SendServerHelloTestSuite) TestHMAC() {
|
||||
}
|
||||
|
||||
func (suite *SendServerHelloTestSuite) TestHandshakePayload() {
|
||||
err := fake.SendServerHello(suite.buf, suite.secret.Key[:], suite.hello)
|
||||
err := fake.SendServerHello(suite.buf, suite.secret.Key[:], suite.hello, fake.NoiseParams{})
|
||||
suite.NoError(err)
|
||||
|
||||
packet := suite.buf.Bytes()
|
||||
@@ -105,7 +104,7 @@ func (suite *SendServerHelloTestSuite) TestHandshakePayload() {
|
||||
}
|
||||
|
||||
func (suite *SendServerHelloTestSuite) TestChangeCipherSpec() {
|
||||
err := fake.SendServerHello(suite.buf, suite.secret.Key[:], suite.hello)
|
||||
err := fake.SendServerHello(suite.buf, suite.secret.Key[:], suite.hello, fake.NoiseParams{})
|
||||
suite.NoError(err)
|
||||
|
||||
// Skip first record
|
||||
@@ -124,6 +123,33 @@ func (suite *SendServerHelloTestSuite) TestChangeCipherSpec() {
|
||||
suite.Equal([]byte{fake.ChangeCipherValue}, rec.Bytes())
|
||||
}
|
||||
|
||||
func (suite *SendServerHelloTestSuite) TestCalibratedNoiseSize() {
|
||||
noise := fake.NoiseParams{Mean: 6480, Jitter: 100}
|
||||
err := fake.SendServerHello(suite.buf, suite.secret.Key[:], suite.hello, noise)
|
||||
suite.NoError(err)
|
||||
|
||||
var rec bytes.Buffer
|
||||
|
||||
// Skip ServerHello
|
||||
_, _, err = tls.ReadRecord(suite.buf, &rec)
|
||||
suite.NoError(err)
|
||||
|
||||
// Skip ChangeCipherSpec
|
||||
rec.Reset()
|
||||
_, _, err = tls.ReadRecord(suite.buf, &rec)
|
||||
suite.NoError(err)
|
||||
|
||||
// Read noise ApplicationData
|
||||
rec.Reset()
|
||||
recordType, length, err := tls.ReadRecord(suite.buf, &rec)
|
||||
suite.NoError(err)
|
||||
suite.Equal(byte(tls.TypeApplicationData), recordType)
|
||||
|
||||
// Should be within mean ± jitter range.
|
||||
suite.GreaterOrEqual(length, int64(noise.Mean-noise.Jitter))
|
||||
suite.LessOrEqual(length, int64(noise.Mean+noise.Jitter))
|
||||
}
|
||||
|
||||
func TestSendServerHello(t *testing.T) {
|
||||
t.Parallel()
|
||||
suite.Run(t, &SendServerHelloTestSuite{})
|
||||
|
||||
@@ -0,0 +1,158 @@
|
||||
package fake
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/binary"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
|
||||
"github.com/9seconds/mtg/v2/mtglib/internal/tls"
|
||||
)
|
||||
|
||||
const (
|
||||
maxFragmentsCount = 10
|
||||
)
|
||||
|
||||
var ErrTooManyFragments = errors.New("too many fragments")
|
||||
|
||||
// https://datatracker.ietf.org/doc/html/rfc5246#section-6.2.1
|
||||
// client hello can be fragmented in a series of packets:
|
||||
//
|
||||
// Bytes on the wire:
|
||||
//
|
||||
// 16 03 01 00 F8 01 00 00 F4 03 03 [32 bytes random] [session_id] [ciphers] [SNI...]
|
||||
// ├─────────────┤├──────────────────────────────────────────────────────────────────┤
|
||||
//
|
||||
// TLS record Payload (248 bytes)
|
||||
// header (5B)
|
||||
//
|
||||
// 16 = Handshake
|
||||
// 03 01 = TLS 1.0 (record layer version)
|
||||
// 00 F8 = 248 bytes follow
|
||||
//
|
||||
// 01 = ClientHello (handshake type)
|
||||
// 00 00 F4 = 244 bytes of handshake body
|
||||
// 03 03 = TLS 1.2 (actual protocol version)
|
||||
// ...rest of ClientHello...
|
||||
//
|
||||
// Fragmented record look like:
|
||||
//
|
||||
// Record 1:
|
||||
//
|
||||
// 16 03 01 00 03 01 00 00
|
||||
// ├─────────────┤├──────┤
|
||||
//
|
||||
// TLS header 3 bytes of payload
|
||||
//
|
||||
// 16 = Handshake
|
||||
// 03 01 = TLS 1.0
|
||||
// 00 03 = only 3 bytes follow
|
||||
//
|
||||
// 01 = ClientHello type
|
||||
// 00 00 = first 2 bytes of the uint24 length (INCOMPLETE!)
|
||||
//
|
||||
// Record 2:
|
||||
// 16 03 01 00 F5 F4 03 03 [32 bytes random] [session_id] [ciphers] [SNI...]
|
||||
// ├─────────────┤├────────────────────────────────────────────────────────────┤
|
||||
//
|
||||
// TLS header remaining 245 bytes of payload
|
||||
//
|
||||
// 16 = Handshake
|
||||
// 03 01 = TLS 1.0
|
||||
// 00 F5 = 245 bytes follow
|
||||
//
|
||||
// F4 = last byte of uint24 length (now complete: 00 00 F4 = 244)
|
||||
// 03 03 = TLS 1.2
|
||||
// ...rest of ClientHello continues...
|
||||
//
|
||||
// So it means that there could be a series of handshake packets of different
|
||||
// lengths. The goal of this function is to concatenate these fragments.
|
||||
type fragmentedHandshakeReader struct {
|
||||
r io.Reader
|
||||
buf bytes.Buffer
|
||||
readFragments int
|
||||
}
|
||||
|
||||
func (f *fragmentedHandshakeReader) Read(p []byte) (int, error) {
|
||||
if n, err := f.buf.Read(p); err == nil {
|
||||
return n, nil
|
||||
}
|
||||
|
||||
f.buf.Reset()
|
||||
|
||||
for f.buf.Len() == 0 {
|
||||
if f.readFragments > maxFragmentsCount {
|
||||
return 0, ErrTooManyFragments
|
||||
}
|
||||
|
||||
if err := f.parseNextFragment(); err != nil {
|
||||
return 0, err
|
||||
}
|
||||
|
||||
f.readFragments++
|
||||
}
|
||||
|
||||
return f.buf.Read(p)
|
||||
}
|
||||
|
||||
func (f *fragmentedHandshakeReader) parseNextFragment() error {
|
||||
// record_type(1) + version(2) + size(2)
|
||||
// 16 - type is 0x16 (handshake record)
|
||||
// 03 01 - protocol version is "3,1" (also known as TLS 1.0)
|
||||
// 00 f8 - 0xF8 (248) bytes of handshake message follows
|
||||
header := [1 + 2 + 2]byte{}
|
||||
|
||||
if _, err := io.ReadFull(f.r, header[:]); err != nil {
|
||||
return fmt.Errorf("cannot read record header: %w", err)
|
||||
}
|
||||
|
||||
if header[0] != tls.TypeHandshake {
|
||||
return fmt.Errorf("unexpected record type %#x", header[0])
|
||||
}
|
||||
|
||||
if header[1] != 3 || header[2] != 1 {
|
||||
return fmt.Errorf("unexpected protocol version %#x %#x", header[1], header[2])
|
||||
}
|
||||
|
||||
length := int64(binary.BigEndian.Uint16(header[3:]))
|
||||
_, err := io.CopyN(&f.buf, f.r, length)
|
||||
|
||||
return err
|
||||
}
|
||||
|
||||
func parseClientHello(r io.Reader) (*bytes.Buffer, *bytes.Buffer, error) {
|
||||
r = &fragmentedHandshakeReader{r: r}
|
||||
header := [1 + 3]byte{}
|
||||
|
||||
if _, err := io.ReadFull(r, header[:]); err != nil {
|
||||
return nil, nil, fmt.Errorf("cannot read handshake header: %w", err)
|
||||
}
|
||||
|
||||
if header[0] != TypeHandshakeClient {
|
||||
return nil, nil, fmt.Errorf("incorrect handshake type: %#x", header[0])
|
||||
}
|
||||
|
||||
// unfortunately there is not uint24 in golang, so we just reuse header
|
||||
header[0] = 0
|
||||
length := int64(binary.BigEndian.Uint32(header[:]))
|
||||
|
||||
clientHelloCopy := &bytes.Buffer{}
|
||||
clientHelloCopy.Write([]byte{tls.TypeHandshake, 3, 1})
|
||||
binary.Write( //nolint: errcheck
|
||||
clientHelloCopy,
|
||||
binary.BigEndian,
|
||||
// 1 for handshake type
|
||||
// 3 for handshake length
|
||||
uint16(1+3+length),
|
||||
)
|
||||
clientHelloCopy.WriteByte(TypeHandshakeClient)
|
||||
clientHelloCopy.Write(header[1:])
|
||||
|
||||
handshakeCopy := &bytes.Buffer{}
|
||||
writer := io.MultiWriter(clientHelloCopy, handshakeCopy)
|
||||
|
||||
_, err := io.CopyN(writer, r, length)
|
||||
|
||||
return clientHelloCopy, handshakeCopy, err
|
||||
}
|
||||
+18
-8
@@ -27,6 +27,7 @@ type Proxy struct {
|
||||
|
||||
allowFallbackOnUnknownDC bool
|
||||
tolerateTimeSkewness time.Duration
|
||||
idleTimeout time.Duration
|
||||
domainFrontingPort int
|
||||
domainFrontingIP string
|
||||
domainFrontingProxyProtocol bool
|
||||
@@ -65,10 +66,10 @@ func (p *Proxy) ServeConn(conn essentials.Conn) {
|
||||
ctx := newStreamContext(p.ctx, p.logger, conn)
|
||||
defer ctx.Close()
|
||||
|
||||
go func() {
|
||||
<-ctx.Done()
|
||||
stop := context.AfterFunc(ctx, func() {
|
||||
ctx.Close()
|
||||
}()
|
||||
})
|
||||
defer stop()
|
||||
|
||||
p.eventStream.Send(ctx, NewEventStart(ctx.streamID, ctx.ClientIP()))
|
||||
ctx.logger.Info("Stream has been started")
|
||||
@@ -101,11 +102,13 @@ func (p *Proxy) ServeConn(conn essentials.Conn) {
|
||||
return
|
||||
}
|
||||
|
||||
tracker := newIdleTracker(p.idleTimeout)
|
||||
|
||||
relay.Relay(
|
||||
ctx,
|
||||
ctx.logger.Named("relay"),
|
||||
ctx.telegramConn,
|
||||
ctx.clientConn,
|
||||
connIdleTimeout{Conn: ctx.telegramConn, tracker: tracker},
|
||||
connIdleTimeout{Conn: ctx.clientConn, tracker: tracker},
|
||||
)
|
||||
}
|
||||
|
||||
@@ -151,6 +154,7 @@ func (p *Proxy) Serve(listener net.Listener) error {
|
||||
case errors.Is(err, ants.ErrPoolClosed):
|
||||
return nil
|
||||
case errors.Is(err, ants.ErrPoolOverload):
|
||||
conn.Close() //nolint: errcheck
|
||||
logger.Info("connection was concurrency limited")
|
||||
p.eventStream.Send(p.ctx, NewEventConcurrencyLimited())
|
||||
}
|
||||
@@ -192,7 +196,10 @@ func (p *Proxy) doFakeTLSHandshake(ctx *streamContext) bool {
|
||||
return false
|
||||
}
|
||||
|
||||
if err := fake.SendServerHello(ctx.clientConn, p.secret.Key[:], clientHello); err != nil {
|
||||
gangerNoise := p.doppelGanger.NoiseParams()
|
||||
noiseParams := fake.NoiseParams{Mean: gangerNoise.Mean, Jitter: gangerNoise.Jitter}
|
||||
|
||||
if err := fake.SendServerHello(ctx.clientConn, p.secret.Key[:], clientHello, noiseParams); err != nil {
|
||||
p.logger.InfoError("cannot send welcome packet", err)
|
||||
return false
|
||||
}
|
||||
@@ -300,11 +307,13 @@ func (p *Proxy) doDomainFronting(ctx *streamContext, conn *connRewind) {
|
||||
stream: p.eventStream,
|
||||
}
|
||||
|
||||
tracker := newIdleTracker(p.idleTimeout)
|
||||
|
||||
relay.Relay(
|
||||
ctx,
|
||||
ctx.logger.Named("domain-fronting"),
|
||||
frontConn,
|
||||
conn,
|
||||
connIdleTimeout{Conn: frontConn, tracker: tracker},
|
||||
connIdleTimeout{Conn: conn, tracker: tracker},
|
||||
)
|
||||
}
|
||||
|
||||
@@ -336,6 +345,7 @@ func NewProxy(opts ProxyOpts) (*Proxy, error) {
|
||||
domainFrontingPort: opts.getDomainFrontingPort(),
|
||||
domainFrontingIP: opts.DomainFrontingIP,
|
||||
tolerateTimeSkewness: opts.getTolerateTimeSkewness(),
|
||||
idleTimeout: opts.getIdleTimeout(),
|
||||
allowFallbackOnUnknownDC: opts.AllowFallbackOnUnknownDC,
|
||||
telegram: tg,
|
||||
doppelGanger: doppel.NewGanger(
|
||||
|
||||
@@ -215,6 +215,14 @@ func (p ProxyOpts) getPreferIP() string {
|
||||
return p.PreferIP
|
||||
}
|
||||
|
||||
func (p ProxyOpts) getIdleTimeout() time.Duration {
|
||||
if p.IdleTimeout == 0 {
|
||||
return time.Minute
|
||||
}
|
||||
|
||||
return p.IdleTimeout
|
||||
}
|
||||
|
||||
func (p ProxyOpts) getLogger(name string) Logger {
|
||||
return p.Logger.Named(name)
|
||||
}
|
||||
|
||||
@@ -175,7 +175,7 @@ func (suite *ProxyTestSuite) TestHTTPSRequest() {
|
||||
addr := fmt.Sprintf("https://%s/headers", suite.ProxyAddress())
|
||||
|
||||
resp, err := client.Get(addr) //nolint: noctx
|
||||
suite.NoError(err)
|
||||
suite.Require().NoError(err)
|
||||
|
||||
defer resp.Body.Close() //nolint: errcheck
|
||||
|
||||
|
||||
@@ -3,7 +3,6 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"net"
|
||||
"net/http"
|
||||
_ "net/http/pprof" //nolint: gosec
|
||||
|
||||
Reference in New Issue
Block a user