tuic_job.go 3.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111
  1. package job
  2. import (
  3. "fmt"
  4. "time"
  5. "github.com/mhsanaei/3x-ui/v3/internal/logger"
  6. "github.com/mhsanaei/3x-ui/v3/internal/tuic"
  7. "github.com/mhsanaei/3x-ui/v3/internal/web/service"
  8. "github.com/mhsanaei/3x-ui/v3/internal/xray"
  9. )
  10. type TuicJob struct {
  11. inboundService service.InboundService
  12. }
  13. func NewTuicJob() *TuicJob {
  14. return new(TuicJob)
  15. }
  16. func (j *TuicJob) Run() {
  17. tuicJournalMu.Lock()
  18. journalErr := j.replayTuicJournal()
  19. tuicJournalMu.Unlock()
  20. if journalErr != nil {
  21. logger.Warning("tuic job: recover traffic journal failed:", journalErr)
  22. }
  23. desired, err := j.inboundService.DesiredTuicInstances()
  24. if err != nil {
  25. logger.Warning("tuic job: get desired instances failed:", err)
  26. return
  27. }
  28. activeTags := make([]string, 0, len(desired))
  29. for _, inst := range desired {
  30. activeTags = append(activeTags, inst.Tag)
  31. }
  32. mgr := tuic.GetManager()
  33. mgr.Reconcile(desired)
  34. _, clientDeltas := mgr.CollectAllTraffic()
  35. onlineEmails, _ := mgr.GetActiveClients(30 * time.Second)
  36. clientTraffics := aggregateTuicClientTraffic(clientDeltas, onlineEmails)
  37. // Inbound total traffic is already metered through the loopback SOCKS relay
  38. // by xray_traffic_job (matching mtproto); only per-client deltas are submitted here.
  39. if len(clientTraffics) > 0 {
  40. needRestart, _, err := j.inboundService.AddTraffic(nil, clientTraffics)
  41. if err != nil {
  42. logger.Warning("tuic job: add traffic failed:", err)
  43. mgr.RequeueClientTraffic(clientDeltas)
  44. } else if needRestart {
  45. if desired, err := j.inboundService.DesiredTuicInstances(); err == nil {
  46. mgr.Reconcile(desired)
  47. }
  48. }
  49. }
  50. if len(onlineEmails) > 0 {
  51. if err := j.inboundService.BumpClientsLastOnline(onlineEmails); err != nil {
  52. logger.Warning("tuic job: bump last online for tuic clients failed:", err)
  53. }
  54. }
  55. j.inboundService.RefreshLocalOnlineClients(onlineEmails, activeTags)
  56. }
  57. // FlushStoppedTraffic persists counters drained when the TUIC manager stops its
  58. // listeners. Call it after scheduled jobs have stopped and before the traffic
  59. // writer shuts down.
  60. func (j *TuicJob) FlushStoppedTraffic() error {
  61. return j.flushTuicJournal()
  62. }
  63. func aggregateTuicClientTraffic(clientDeltas []tuic.ClientTrafficDelta, onlineEmails []string) []*xray.ClientTraffic {
  64. clientTrafficMap := make(map[string]*xray.ClientTraffic, len(clientDeltas)+len(onlineEmails))
  65. for _, cd := range clientDeltas {
  66. key := cd.Email
  67. if cd.TrafficID > 0 {
  68. key = fmt.Sprintf("traffic:%d", cd.TrafficID)
  69. }
  70. if cd.TrafficID == 0 && cd.InboundID > 0 && cd.UUID != "" {
  71. key = fmt.Sprintf("tuic:%d:%s", cd.InboundID, cd.UUID)
  72. }
  73. traffic := clientTrafficMap[key]
  74. if traffic == nil {
  75. traffic = &xray.ClientTraffic{Email: cd.Email, TuicTrafficID: cd.TrafficID, TuicUUID: cd.UUID, TuicInboundId: cd.InboundID}
  76. clientTrafficMap[key] = traffic
  77. }
  78. traffic.Up += cd.Up
  79. traffic.Down += cd.Down
  80. }
  81. for _, email := range onlineEmails {
  82. if _, exists := clientTrafficMap[email]; !exists {
  83. clientTrafficMap[email] = &xray.ClientTraffic{
  84. Email: email,
  85. Up: 0,
  86. Down: 0,
  87. }
  88. }
  89. }
  90. clientTraffics := make([]*xray.ClientTraffic, 0, len(clientTrafficMap))
  91. for _, ct := range clientTrafficMap {
  92. clientTraffics = append(clientTraffics, ct)
  93. }
  94. return clientTraffics
  95. }