client_traffic.go 5.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219
  1. package service
  2. import (
  3. "time"
  4. "github.com/mhsanaei/3x-ui/v3/internal/database"
  5. "github.com/mhsanaei/3x-ui/v3/internal/database/model"
  6. "github.com/mhsanaei/3x-ui/v3/internal/logger"
  7. "github.com/mhsanaei/3x-ui/v3/internal/util/common"
  8. "github.com/mhsanaei/3x-ui/v3/internal/xray"
  9. "gorm.io/gorm"
  10. )
  11. func (s *ClientService) ResetTrafficByEmail(inboundSvc *InboundService, email string) (bool, error) {
  12. if email == "" {
  13. return false, common.NewError("client email is required")
  14. }
  15. rec, err := s.GetRecordByEmail(nil, email)
  16. if err != nil {
  17. return false, err
  18. }
  19. inboundIds, err := s.GetInboundIdsForRecord(rec.Id)
  20. if err != nil {
  21. return false, err
  22. }
  23. needRestart := false
  24. if len(inboundIds) == 0 {
  25. if rErr := inboundSvc.ResetClientTrafficByEmail(email); rErr != nil {
  26. return false, rErr
  27. }
  28. } else {
  29. applies := make([]inboundApply, 0, len(inboundIds))
  30. for _, ibId := range inboundIds {
  31. applies = append(applies, inboundApply{id: ibId, run: func() (bool, error) {
  32. return inboundSvc.ResetClientTraffic(ibId, email)
  33. }})
  34. }
  35. nr, applyErr := fanoutInboundApplies(applies)
  36. if applyErr != nil {
  37. return nr, applyErr
  38. }
  39. needRestart = nr
  40. }
  41. // Enable only once the counters are zero: a still-depleted client enabled
  42. // first is switched off again by the next traffic tick.
  43. if !rec.Enable {
  44. updated := rec.ToClient()
  45. updated.Enable = true
  46. nr, uErr := s.Update(inboundSvc, rec.Id, *updated, rec.LimitHwid)
  47. if uErr != nil {
  48. logger.Warning("Failed to auto-enable client during traffic reset:", uErr)
  49. }
  50. if nr {
  51. needRestart = true
  52. }
  53. }
  54. return needRestart, nil
  55. }
  56. func (s *ClientService) BulkResetTraffic(inboundSvc *InboundService, emails []string) (int, error) {
  57. if len(emails) == 0 {
  58. return 0, nil
  59. }
  60. cleanEmails := trimmedUniqueEmails(emails)
  61. if len(cleanEmails) == 0 {
  62. return 0, nil
  63. }
  64. recordsByEmail, err := clientRecordsByEmail(nil, cleanEmails)
  65. if err != nil {
  66. return 0, err
  67. }
  68. affected := 0
  69. err = submitTrafficWrite(func() error {
  70. db := database.GetDB()
  71. return db.Transaction(func(tx *gorm.DB) error {
  72. if err := adjustGroupBaselinesForRemovedTraffic(tx, cleanEmails); err != nil {
  73. return err
  74. }
  75. for _, batch := range chunkStrings(cleanEmails, sqlInChunk) {
  76. res := tx.Model(xray.ClientTraffic{}).
  77. Where("email IN ?", batch).
  78. Updates(map[string]any{"enable": true, "up": 0, "down": 0})
  79. if res.Error != nil {
  80. return res.Error
  81. }
  82. affected += int(res.RowsAffected)
  83. }
  84. if err := clearGlobalTraffic(tx, cleanEmails...); err != nil {
  85. return err
  86. }
  87. for _, batch := range chunkStrings(cleanEmails, sqlInChunk) {
  88. if err := tx.Where("email IN ?", batch).Delete(&model.NodeClientTraffic{}).Error; err != nil {
  89. return err
  90. }
  91. }
  92. return nil
  93. })
  94. })
  95. if err != nil {
  96. return 0, err
  97. }
  98. // After the zeroing, as in ResetTrafficByEmail: enabling a still-depleted
  99. // client first lets the next traffic tick switch it off again.
  100. for _, e := range cleanEmails {
  101. rec := recordsByEmail[e]
  102. if rec == nil || rec.Enable {
  103. continue
  104. }
  105. updated := rec.ToClient()
  106. updated.Enable = true
  107. if _, uErr := s.Update(inboundSvc, rec.Id, *updated, rec.LimitHwid); uErr != nil {
  108. logger.Warning("Failed to auto-enable client during bulk traffic reset:", uErr)
  109. }
  110. }
  111. return affected, nil
  112. }
  113. func (s *ClientService) ResetAllClientTraffics(inboundSvc *InboundService, id int) error {
  114. err := submitTrafficWrite(func() error {
  115. return s.resetAllClientTrafficsLocked(id)
  116. })
  117. if err == nil {
  118. inboundSvc.resetAllMtprotoQuotas()
  119. }
  120. return err
  121. }
  122. func (s *ClientService) resetAllClientTrafficsLocked(id int) error {
  123. db := database.GetDB()
  124. now := time.Now().Unix() * 1000
  125. if err := db.Transaction(func(tx *gorm.DB) error {
  126. // client_traffics.inbound_id is stale: it reflects the inbound the row was
  127. // first inserted under and is never refreshed. Use the client_inbounds join
  128. // as the authoritative source for which emails belong to a given inbound.
  129. var resetEmails []string
  130. if id == -1 {
  131. if err := tx.Model(xray.ClientTraffic{}).Pluck("email", &resetEmails).Error; err != nil {
  132. return err
  133. }
  134. } else {
  135. if err := tx.Table("client_inbounds ci").
  136. Select("c.email").
  137. Joins("JOIN clients c ON c.id = ci.client_id").
  138. Where("ci.inbound_id = ?", id).
  139. Pluck("c.email", &resetEmails).Error; err != nil {
  140. return err
  141. }
  142. }
  143. if len(resetEmails) == 0 {
  144. return nil
  145. }
  146. if err := adjustGroupBaselinesForRemovedTraffic(tx, resetEmails); err != nil {
  147. return err
  148. }
  149. result := tx.Model(xray.ClientTraffic{}).
  150. Where("email IN ?", resetEmails).
  151. Updates(map[string]any{"enable": true, "up": 0, "down": 0})
  152. if result.Error != nil {
  153. return result.Error
  154. }
  155. if err := clearGlobalTraffic(tx, resetEmails...); err != nil {
  156. return err
  157. }
  158. for _, batch := range chunkStrings(resetEmails, sqlInChunk) {
  159. if err := tx.Where("email IN ?", batch).Delete(&model.NodeClientTraffic{}).Error; err != nil {
  160. return err
  161. }
  162. }
  163. inboundWhereText := "id "
  164. if id == -1 {
  165. inboundWhereText += " > ?"
  166. } else {
  167. inboundWhereText += " = ?"
  168. }
  169. result = tx.Model(model.Inbound{}).
  170. Where(inboundWhereText, id).
  171. Update("last_traffic_reset_time", now)
  172. return result.Error
  173. }); err != nil {
  174. return err
  175. }
  176. return nil
  177. }
  178. func (s *ClientService) ResetAllTraffics() (bool, error) {
  179. var affected int64
  180. err := submitTrafficWrite(func() error {
  181. return database.GetDB().Transaction(func(tx *gorm.DB) error {
  182. res := tx.Model(&xray.ClientTraffic{}).
  183. Where("1 = 1").
  184. Updates(map[string]any{"enable": true, "up": 0, "down": 0})
  185. if res.Error != nil {
  186. return res.Error
  187. }
  188. affected = res.RowsAffected
  189. if err := tx.Where("1 = 1").Delete(&model.ClientGlobalTraffic{}).Error; err != nil {
  190. return err
  191. }
  192. return tx.Where("1 = 1").Delete(&model.NodeClientTraffic{}).Error
  193. })
  194. })
  195. if err != nil {
  196. return false, err
  197. }
  198. return affected > 0, nil
  199. }