inbound_traffic_apply.go 3.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130
  1. package service
  2. import (
  3. "context"
  4. "strings"
  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/xray"
  8. "gorm.io/gorm"
  9. )
  10. type trafficLocalApplyAction uint8
  11. const (
  12. trafficAddUser trafficLocalApplyAction = iota + 1
  13. trafficRemoveUser
  14. trafficDisableInbound
  15. )
  16. type trafficLocalApplyPlan struct {
  17. action trafficLocalApplyAction
  18. inbound model.Inbound
  19. client map[string]any
  20. email string
  21. }
  22. type trafficMutationBatch struct {
  23. localPlans []trafficLocalApplyPlan
  24. remotePlans []trafficInboundUpdatePlan
  25. nodeIDs map[int]struct{}
  26. // renewedEmails get their MTProto sidecar quota zeroed once the tick commits.
  27. renewedEmails []string
  28. }
  29. type trafficInboundUpdatePlan struct{ oldInbound, newInbound model.Inbound }
  30. func newTrafficMutationBatch() *trafficMutationBatch {
  31. return &trafficMutationBatch{nodeIDs: make(map[int]struct{})}
  32. }
  33. func (b *trafficMutationBatch) addNode(nodeID int) {
  34. if nodeID > 0 {
  35. b.nodeIDs[nodeID] = struct{}{}
  36. }
  37. }
  38. func (b *trafficMutationBatch) markNodesTx(tx *gorm.DB) error {
  39. if b == nil {
  40. return nil
  41. }
  42. nodeSvc := NodeService{}
  43. for nodeID := range b.nodeIDs {
  44. if err := nodeSvc.MarkNodeDirtyTx(tx, nodeID); err != nil {
  45. return err
  46. }
  47. }
  48. return nil
  49. }
  50. // applyTrafficRemotePlans is bounded like every per-client node push: the nodes are
  51. // already dirty, so an offline or slow one defers to the reconcile.
  52. func (s *InboundService) applyTrafficRemotePlans(plans []trafficInboundUpdatePlan) bool {
  53. ids := make([]int, len(plans))
  54. for i := range plans {
  55. ids[i] = plans[i].newInbound.Id
  56. }
  57. failed, panics := fanoutInboundResults(ids, inboundFanoutConcurrency, func(i int) bool {
  58. rt, push, _, err := s.nodePushPlan(&plans[i].newInbound)
  59. if err == nil && push {
  60. ctx, cancel := nodePushContext()
  61. err = rt.UpdateInbound(ctx, &plans[i].oldInbound, &plans[i].newInbound)
  62. cancel()
  63. }
  64. if err != nil {
  65. logger.Debug("traffic post-commit remote apply failed:", err)
  66. }
  67. return err != nil
  68. })
  69. needRestart := false
  70. for i := range failed {
  71. needRestart = needRestart || failed[i] || panics[i] != nil
  72. }
  73. return needRestart
  74. }
  75. func (s *InboundService) applyTrafficMutationBatch(b *trafficMutationBatch) bool {
  76. if b == nil {
  77. return false
  78. }
  79. needRestart := false
  80. for i := range b.localPlans {
  81. plan := &b.localPlans[i]
  82. if plan.inbound.Protocol == model.MTProto {
  83. s.applyLocalMtproto(plan.inbound.Id)
  84. continue
  85. }
  86. if plan.inbound.Protocol == model.AmneziaWG {
  87. s.applyLocalAmneziaWG(plan.inbound.Id)
  88. continue
  89. }
  90. if plan.inbound.Protocol == model.TUIC {
  91. s.applyLocalTuic(plan.inbound.Id)
  92. continue
  93. }
  94. rt, err := s.runtimeFor(&plan.inbound)
  95. if err == nil {
  96. switch plan.action {
  97. case trafficAddUser:
  98. err = rt.AddUser(context.Background(), &plan.inbound, plan.client)
  99. case trafficRemoveUser:
  100. err = rt.RemoveUser(context.Background(), &plan.inbound, plan.email)
  101. if err != nil && strings.Contains(err.Error(), "not found") {
  102. err = nil
  103. }
  104. case trafficDisableInbound:
  105. err = rt.DelInbound(context.Background(), &plan.inbound)
  106. if xray.IsMissingHandlerErr(err) {
  107. err = nil
  108. }
  109. }
  110. }
  111. if err != nil {
  112. logger.Debug("traffic post-commit runtime apply failed:", err)
  113. needRestart = true
  114. }
  115. }
  116. return needRestart
  117. }