1
0

tuic_job.go 4.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147
  1. package job
  2. import (
  3. "fmt"
  4. "sync"
  5. "time"
  6. "github.com/mhsanaei/3x-ui/v3/internal/logger"
  7. "github.com/mhsanaei/3x-ui/v3/internal/tuic"
  8. "github.com/mhsanaei/3x-ui/v3/internal/web/service"
  9. "github.com/mhsanaei/3x-ui/v3/internal/web/websocket"
  10. "github.com/mhsanaei/3x-ui/v3/internal/xray"
  11. )
  12. const defaultTuicSpeedSampleInterval = 10 * time.Second
  13. type TuicJob struct {
  14. inboundService service.InboundService
  15. runMu sync.Mutex
  16. lastSpeedSample time.Time
  17. }
  18. func NewTuicJob() *TuicJob {
  19. return new(TuicJob)
  20. }
  21. func (j *TuicJob) Run() {
  22. j.runMu.Lock()
  23. defer j.runMu.Unlock()
  24. tuicJournalMu.Lock()
  25. journalErr := j.replayTuicJournal()
  26. tuicJournalMu.Unlock()
  27. if journalErr != nil {
  28. logger.Warning("tuic job: recover traffic journal failed:", journalErr)
  29. }
  30. desired, err := j.inboundService.DesiredTuicInstances()
  31. if err != nil {
  32. logger.Warning("tuic job: get desired instances failed:", err)
  33. return
  34. }
  35. activeTags := make([]string, 0, len(desired))
  36. for _, inst := range desired {
  37. activeTags = append(activeTags, inst.Tag)
  38. }
  39. mgr := tuic.GetManager()
  40. mgr.Reconcile(desired)
  41. _, clientDeltas := mgr.CollectAllTraffic()
  42. onlineEmails, _ := mgr.GetActiveClients(30 * time.Second)
  43. clientTraffics := aggregateTuicClientTraffic(clientDeltas, onlineEmails)
  44. sampledAt := time.Now()
  45. sampleInterval := tuicSpeedSampleInterval(j.lastSpeedSample, sampledAt)
  46. // Inbound total traffic is already metered through the loopback SOCKS relay
  47. // by xray_traffic_job (matching mtproto); only per-client deltas are submitted here.
  48. persisted := true
  49. if len(clientTraffics) > 0 {
  50. needRestart, _, err := j.inboundService.AddTraffic(nil, clientTraffics)
  51. if err != nil {
  52. logger.Warning("tuic job: add traffic failed:", err)
  53. mgr.RequeueClientTraffic(clientDeltas)
  54. persisted = false
  55. } else if needRestart {
  56. if desired, err := j.inboundService.DesiredTuicInstances(); err == nil {
  57. mgr.Reconcile(desired)
  58. }
  59. }
  60. }
  61. if persisted {
  62. websocket.BroadcastTraffic(tuicSpeedPayload(clientTraffics, sampleInterval))
  63. j.lastSpeedSample = sampledAt
  64. }
  65. if len(onlineEmails) > 0 {
  66. if err := j.inboundService.BumpClientsLastOnline(onlineEmails); err != nil {
  67. logger.Warning("tuic job: bump last online for tuic clients failed:", err)
  68. }
  69. }
  70. j.inboundService.RefreshLocalOnlineClients(onlineEmails, activeTags)
  71. }
  72. func tuicSpeedSampleInterval(previous, current time.Time) time.Duration {
  73. if previous.IsZero() || !current.After(previous) {
  74. return defaultTuicSpeedSampleInterval
  75. }
  76. return current.Sub(previous)
  77. }
  78. func tuicSpeedPayload(clientTraffics []*xray.ClientTraffic, sampleInterval time.Duration) map[string]any {
  79. intervalMs := sampleInterval.Milliseconds()
  80. if intervalMs < 1 {
  81. intervalMs = 1
  82. }
  83. return map[string]any{
  84. "clientTraffics": clientTraffics,
  85. "clientTrafficSource": "tuic",
  86. "clientTrafficIntervalMs": intervalMs,
  87. }
  88. }
  89. // FlushStoppedTraffic persists counters drained when the TUIC manager stops its
  90. // listeners. Call it after scheduled jobs have stopped and before the traffic
  91. // writer shuts down.
  92. func (j *TuicJob) FlushStoppedTraffic() error {
  93. return j.flushTuicJournal()
  94. }
  95. func aggregateTuicClientTraffic(clientDeltas []tuic.ClientTrafficDelta, onlineEmails []string) []*xray.ClientTraffic {
  96. clientTrafficMap := make(map[string]*xray.ClientTraffic, len(clientDeltas)+len(onlineEmails))
  97. for _, cd := range clientDeltas {
  98. key := cd.Email
  99. if cd.TrafficID > 0 {
  100. key = fmt.Sprintf("traffic:%d", cd.TrafficID)
  101. }
  102. if cd.TrafficID == 0 && cd.InboundID > 0 && cd.UUID != "" {
  103. key = fmt.Sprintf("tuic:%d:%s", cd.InboundID, cd.UUID)
  104. }
  105. traffic := clientTrafficMap[key]
  106. if traffic == nil {
  107. traffic = &xray.ClientTraffic{Email: cd.Email, TuicTrafficID: cd.TrafficID, TuicUUID: cd.UUID, TuicInboundId: cd.InboundID}
  108. clientTrafficMap[key] = traffic
  109. }
  110. traffic.Up += cd.Up
  111. traffic.Down += cd.Down
  112. }
  113. for _, email := range onlineEmails {
  114. if _, exists := clientTrafficMap[email]; !exists {
  115. clientTrafficMap[email] = &xray.ClientTraffic{
  116. Email: email,
  117. Up: 0,
  118. Down: 0,
  119. }
  120. }
  121. }
  122. clientTraffics := make([]*xray.ClientTraffic, 0, len(clientTrafficMap))
  123. for _, ct := range clientTrafficMap {
  124. clientTraffics = append(clientTraffics, ct)
  125. }
  126. return clientTraffics
  127. }