Highly-opinionated (ex-bullshit-free) MTPROTO proxy for Telegram. If you use v1.0 or upgrade broke you proxy, please read the chapter Version 2
You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

ganger.go 3.2KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173
  1. package doppel
  2. import (
  3. "context"
  4. "sync"
  5. "time"
  6. "github.com/9seconds/mtg/v2/essentials"
  7. )
  8. const (
  9. DoppelGangerMaxDurations = 4096
  10. DoppelGangerScoutMissionEach = 30 * time.Minute
  11. DoppelGangerScoutRepeats = 10
  12. )
  13. type gangerConnRequest struct {
  14. ret chan Conn
  15. payload essentials.Conn
  16. }
  17. type Ganger struct {
  18. ctx context.Context
  19. ctxCancel context.CancelFunc
  20. logger Logger
  21. wg sync.WaitGroup
  22. scout Scout
  23. scoutMissionEach time.Duration
  24. scoutMissionRepeats int
  25. stats *Stats
  26. durations []time.Duration
  27. connRequests chan gangerConnRequest
  28. }
  29. func (g *Ganger) Shutdown() {
  30. g.ctxCancel()
  31. g.wg.Wait()
  32. }
  33. func (g *Ganger) Run() {
  34. g.wg.Go(func() {
  35. g.run()
  36. })
  37. }
  38. func (g *Ganger) NewConn(conn essentials.Conn) (Conn, error) {
  39. req := gangerConnRequest{
  40. ret: make(chan Conn),
  41. payload: conn,
  42. }
  43. defer close(req.ret)
  44. select {
  45. case <-g.ctx.Done():
  46. return Conn{}, context.Cause(g.ctx)
  47. case g.connRequests <- req:
  48. }
  49. select {
  50. case <-g.ctx.Done():
  51. return Conn{}, context.Cause(g.ctx)
  52. case conn := <-req.ret:
  53. return conn, nil
  54. }
  55. }
  56. func (g *Ganger) run() {
  57. scoutTicker := time.NewTicker(g.scoutMissionEach)
  58. defer func() {
  59. scoutTicker.Stop()
  60. select {
  61. case <-scoutTicker.C:
  62. default:
  63. }
  64. }()
  65. scoutCollectedChan := make(chan []time.Duration)
  66. currentScoutCollectedChan := scoutCollectedChan
  67. updatedStatsChan := make(chan *Stats)
  68. g.wg.Go(func() {
  69. g.runScoutMission(scoutCollectedChan)
  70. })
  71. for {
  72. select {
  73. case <-g.ctx.Done():
  74. return
  75. case durations := <-currentScoutCollectedChan:
  76. g.durations = append(g.durations, durations...)
  77. if len(g.durations) > DoppelGangerMaxDurations {
  78. g.durations = g.durations[len(g.durations)-DoppelGangerMaxDurations:]
  79. }
  80. currentScoutCollectedChan = nil
  81. g.wg.Go(func() {
  82. select {
  83. case <-g.ctx.Done():
  84. case updatedStatsChan <- NewStats(durations):
  85. }
  86. })
  87. case stats := <-updatedStatsChan:
  88. g.stats = stats
  89. currentScoutCollectedChan = scoutCollectedChan
  90. case <-scoutTicker.C:
  91. g.wg.Go(func() {
  92. g.runScoutMission(scoutCollectedChan)
  93. })
  94. case req := <-g.connRequests:
  95. select {
  96. case <-g.ctx.Done():
  97. case req.ret <- NewConn(g.ctx, req.payload, g.stats):
  98. }
  99. }
  100. }
  101. }
  102. func (g *Ganger) runScoutMission(rvChan chan<- []time.Duration) {
  103. durations := []time.Duration{}
  104. for range g.scoutMissionRepeats {
  105. learned, err := g.scout.Learn(g.ctx)
  106. if err != nil {
  107. g.logger.WarningError("cannot learn", err)
  108. continue
  109. }
  110. durations = append(durations, learned...)
  111. }
  112. select {
  113. case <-g.ctx.Done():
  114. return
  115. case rvChan <- durations:
  116. }
  117. }
  118. func NewGanger(
  119. ctx context.Context,
  120. network Network,
  121. logger Logger,
  122. scoutEach time.Duration,
  123. scoutRepeats int,
  124. urls []string,
  125. ) *Ganger {
  126. ctx, cancel := context.WithCancel(ctx)
  127. if scoutEach == 0 {
  128. scoutEach = DoppelGangerScoutMissionEach
  129. }
  130. if scoutRepeats == 0 {
  131. scoutRepeats = DoppelGangerScoutRepeats
  132. }
  133. return &Ganger{
  134. ctx: ctx,
  135. ctxCancel: cancel,
  136. logger: logger,
  137. scoutMissionEach: scoutEach,
  138. scoutMissionRepeats: scoutRepeats,
  139. stats: &Stats{
  140. k: StatsDefaultK,
  141. lambda: StatsDefaultLambda,
  142. },
  143. scout: NewScout(network, urls),
  144. connRequests: make(chan gangerConnRequest),
  145. }
  146. }