diff --git a/internal/cli/simple_run.go b/internal/cli/simple_run.go index b0eebde..192cf3c 100644 --- a/internal/cli/simple_run.go +++ b/internal/cli/simple_run.go @@ -18,7 +18,7 @@ type SimpleRun struct { TCPBuffer string `kong:"name='tcp-buffer',short='b',default='4KB',help='Deprecated and ignored'"` //nolint: lll PreferIP string `kong:"name='prefer-ip',short='i',default='prefer-ipv6',help='IP preference. By default we prefer IPv6 with fallback to IPv4.'"` //nolint: lll DomainFrontingPort uint64 `kong:"name='domain-fronting-port',short='p',default='443',help='A port to access for domain fronting.'"` //nolint: lll - DomainFrontingIP string `kong:"name='domain-fronting-ip',help='An IP address to use for domain fronting instead of resolving the hostname via DNS.'"` //nolint: lll + DomainFrontingIP string `kong:"name='domain-fronting-ip',help='An IP address to use for domain fronting instead of resolving the hostname via DNS.'"` //nolint: lll DOHIP net.IP `kong:"name='doh-ip',short='n',default='1.1.1.1',help='IP address of DNS-over-HTTP to use.'"` //nolint: lll Timeout time.Duration `kong:"name='timeout',short='t',default='10s',help='Network timeout to use'"` //nolint: lll Socks5Proxies []string `kong:"name='socks5-proxy',short='s',help='Socks5 proxies to use for network access.'"` //nolint: lll diff --git a/mtglib/internal/dc/init.go b/mtglib/internal/dc/init.go index 2c0f42c..fd8befa 100644 --- a/mtglib/internal/dc/init.go +++ b/mtglib/internal/dc/init.go @@ -38,43 +38,41 @@ type Updater interface { Run(ctx context.Context) } -var ( - // https://github.com/telegramdesktop/tdesktop/blob/master/Telegram/SourceFiles/mtproto/mtproto_dc_options.cpp#L30 - defaultDCAddrSet = dcAddrSet{ - v4: map[int][]Addr{ - 1: { - {Network: "tcp4", Address: "149.154.175.50:443"}, - }, - 2: { - {Network: "tcp4", Address: "149.154.167.51:443"}, - {Network: "tcp4", Address: "95.161.76.100:443"}, - }, - 3: { - {Network: "tcp4", Address: "149.154.175.100:443"}, - }, - 4: { - {Network: "tcp4", Address: "149.154.167.91:443"}, - }, - 5: { - {Network: "tcp4", Address: "149.154.171.5:443"}, - }, +// https://github.com/telegramdesktop/tdesktop/blob/master/Telegram/SourceFiles/mtproto/mtproto_dc_options.cpp#L30 +var defaultDCAddrSet = dcAddrSet{ + v4: map[int][]Addr{ + 1: { + {Network: "tcp4", Address: "149.154.175.50:443"}, }, - v6: map[int][]Addr{ - 1: { - {Network: "tcp6", Address: "[2001:b28:f23d:f001::a]:443"}, - }, - 2: { - {Network: "tcp6", Address: "[2001:67c:04e8:f002::a]:443"}, - }, - 3: { - {Network: "tcp6", Address: "[2001:b28:f23d:f003::a]:443"}, - }, - 4: { - {Network: "tcp6", Address: "[2001:67c:04e8:f004::a]:443"}, - }, - 5: { - {Network: "tcp6", Address: "[2001:b28:f23f:f005::a]:443"}, - }, + 2: { + {Network: "tcp4", Address: "149.154.167.51:443"}, + {Network: "tcp4", Address: "95.161.76.100:443"}, }, - } -) + 3: { + {Network: "tcp4", Address: "149.154.175.100:443"}, + }, + 4: { + {Network: "tcp4", Address: "149.154.167.91:443"}, + }, + 5: { + {Network: "tcp4", Address: "149.154.171.5:443"}, + }, + }, + v6: map[int][]Addr{ + 1: { + {Network: "tcp6", Address: "[2001:b28:f23d:f001::a]:443"}, + }, + 2: { + {Network: "tcp6", Address: "[2001:67c:04e8:f002::a]:443"}, + }, + 3: { + {Network: "tcp6", Address: "[2001:b28:f23d:f003::a]:443"}, + }, + 4: { + {Network: "tcp6", Address: "[2001:67c:04e8:f004::a]:443"}, + }, + 5: { + {Network: "tcp6", Address: "[2001:b28:f23f:f005::a]:443"}, + }, + }, +} diff --git a/mtglib/internal/dc/public_config_updater.go b/mtglib/internal/dc/public_config_updater.go index 8f22832..a02b87c 100644 --- a/mtglib/internal/dc/public_config_updater.go +++ b/mtglib/internal/dc/public_config_updater.go @@ -19,7 +19,7 @@ type PublicConfigUpdater struct { tg *Telegram } -func (p PublicConfigUpdater) Run(ctx context.Context, url, network string) { +func (p *PublicConfigUpdater) Run(ctx context.Context, url, network string) { p.run(ctx, func() error { req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil) if err != nil { @@ -81,8 +81,8 @@ func (p PublicConfigUpdater) Run(ctx context.Context, url, network string) { }) } -func NewPublicConfigUpdater(tg *Telegram, logger Logger, client *http.Client) PublicConfigUpdater { - return PublicConfigUpdater{ +func NewPublicConfigUpdater(tg *Telegram, logger Logger, client *http.Client) *PublicConfigUpdater { + return &PublicConfigUpdater{ updater: updater{ logger: logger, period: PublicConfigUpdateEach, diff --git a/mtglib/internal/dc/public_config_updater_test.go b/mtglib/internal/dc/public_config_updater_test.go index f51f3a2..bc1d148 100644 --- a/mtglib/internal/dc/public_config_updater_test.go +++ b/mtglib/internal/dc/public_config_updater_test.go @@ -14,7 +14,7 @@ import ( type PublicConfigUpdaterTestSuite struct { UpdaterTestSuiteBase - u PublicConfigUpdater + u *PublicConfigUpdater lock sync.Mutex srv *httptest.Server responseHandler func(w http.ResponseWriter) @@ -42,39 +42,27 @@ func (s *PublicConfigUpdaterTestSuite) SetupTest() { } func (s *PublicConfigUpdaterTestSuite) Test502StatusCode() { - done := false - s.responseHandler = func(w http.ResponseWriter) { w.WriteHeader(http.StatusBadGateway) - done = true } - go s.u.Run(s.ctx, s.srv.URL, "tcp4") + s.u.Run(s.ctx, s.srv.URL, "tcp4") - s.Eventually(func() bool { - s.lock.Lock() - defer s.lock.Unlock() - - return done - }, time.Second, 10*time.Millisecond) + time.Sleep(100 * time.Millisecond) + s.ctxCancel() + s.u.Wait() s.Len(s.u.tg.view.publicConfigs.v4, 0) } func (s *PublicConfigUpdaterTestSuite) TestEmptyFile() { - done := false - s.responseHandler = func(w http.ResponseWriter) { - done = true w.WriteHeader(http.StatusOK) } - go s.u.Run(s.ctx, s.srv.URL, "tcp4") + s.u.Run(s.ctx, s.srv.URL, "tcp4") - s.Eventually(func() bool { - s.lock.Lock() - defer s.lock.Unlock() - - return done - }, time.Second, 10*time.Millisecond) + time.Sleep(100 * time.Millisecond) + s.ctxCancel() + s.u.Wait() s.Len(s.u.tg.view.publicConfigs.v4, 0) } @@ -85,21 +73,16 @@ proxy_for -1 -1; proxy_for 100 100.10.0.0:3333; lala 0 0 ` - done := false s.responseHandler = func(w http.ResponseWriter) { - done = true w.WriteHeader(http.StatusOK) w.Write([]byte(result)) //nolint: errcheck } - go s.u.Run(s.ctx, s.srv.URL, "tcp4") + s.u.Run(s.ctx, s.srv.URL, "tcp4") - s.Eventually(func() bool { - s.lock.Lock() - defer s.lock.Unlock() - - return done - }, time.Second, 10*time.Millisecond) + time.Sleep(100 * time.Millisecond) + s.ctxCancel() + s.u.Wait() s.Len(s.u.tg.view.publicConfigs.v4, 0) } @@ -109,21 +92,16 @@ func (s *PublicConfigUpdaterTestSuite) TestOk() { proxy_for 203 100.10.0.0:3333; proxy_for -100 101.10.0.0:3333; ` - done := false s.responseHandler = func(w http.ResponseWriter) { - done = true w.WriteHeader(http.StatusOK) w.Write([]byte(result)) //nolint: errcheck } - go s.u.Run(s.ctx, s.srv.URL, "tcp4") + s.u.Run(s.ctx, s.srv.URL, "tcp4") - s.Eventually(func() bool { - s.lock.Lock() - defer s.lock.Unlock() - - return done - }, time.Second, 10*time.Millisecond) + time.Sleep(100 * time.Millisecond) + s.ctxCancel() + s.u.Wait() s.Len(s.u.tg.view.publicConfigs.v4, 1) s.Len(s.u.tg.view.publicConfigs.v4[203], 1) diff --git a/mtglib/internal/dc/updater.go b/mtglib/internal/dc/updater.go index 5d57af4..72adf38 100644 --- a/mtglib/internal/dc/updater.go +++ b/mtglib/internal/dc/updater.go @@ -2,38 +2,46 @@ package dc import ( "context" + "sync" "time" ) type updater struct { + wg sync.WaitGroup logger Logger period time.Duration } -func (u updater) run(ctx context.Context, callback func() error) { - ticker := time.NewTicker(u.period) +func (u *updater) Wait() { + u.wg.Wait() +} - defer func() { - ticker.Stop() +func (u *updater) run(ctx context.Context, callback func() error) { + u.wg.Go(func() { + ticker := time.NewTicker(u.period) - select { - case <-ticker.C: - default: - } - }() + defer func() { + ticker.Stop() + + select { + case <-ticker.C: + default: + } + }() - for { - u.logger.Info("start update") - if err := callback(); err != nil { - u.logger.WarningError("cannot update: %w", err) - } - u.logger.Info("updated") + for { + u.logger.Info("start update") + if err := callback(); err != nil { + u.logger.WarningError("cannot update: %w", err) + } + u.logger.Info("updated") - select { - case <-ctx.Done(): - u.logger.Info("stop updating") - return - case <-ticker.C: + select { + case <-ctx.Done(): + u.logger.Info("stop updating") + return + case <-ticker.C: + } } - } + }) } diff --git a/mtglib/proxy.go b/mtglib/proxy.go index 567b88e..468fa10 100644 --- a/mtglib/proxy.go +++ b/mtglib/proxy.go @@ -30,6 +30,7 @@ type Proxy struct { domainFrontingIP string workerPool *ants.PoolWithFunc telegram *dc.Telegram + configUpdater *dc.PublicConfigUpdater clientObfuscatror obfuscation.Obfuscator secret Secret @@ -152,6 +153,7 @@ func (p *Proxy) Shutdown() { p.ctxCancel() p.streamWaitGroup.Wait() p.workerPool.Release() + p.configUpdater.Wait() p.allowlist.Shutdown() p.blocklist.Shutdown() @@ -328,18 +330,18 @@ func NewProxy(opts ProxyOpts) (*Proxy, error) { tolerateTimeSkewness: opts.getTolerateTimeSkewness(), allowFallbackOnUnknownDC: opts.AllowFallbackOnUnknownDC, telegram: tg, + configUpdater: dc.NewPublicConfigUpdater( + tg, + updatersLogger.Named("public-config"), + opts.Network.MakeHTTPClient(nil), + ), clientObfuscatror: obfuscation.Obfuscator{ Secret: opts.Secret.Key[:], }, } - publicConfigUpdater := dc.NewPublicConfigUpdater( - tg, - updatersLogger.Named("public-config"), - opts.Network.MakeHTTPClient(nil), - ) - go publicConfigUpdater.Run(ctx, dc.PublicConfigUpdateURLv4, "tcp4") - go publicConfigUpdater.Run(ctx, dc.PublicConfigUpdateURLv6, "tcp6") + proxy.configUpdater.Run(ctx, dc.PublicConfigUpdateURLv4, "tcp4") + proxy.configUpdater.Run(ctx, dc.PublicConfigUpdateURLv6, "tcp6") pool, err := ants.NewPoolWithFunc(opts.getConcurrency(), func(arg any) {