1
0

manager.go 6.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262
  1. package tuic
  2. import (
  3. "fmt"
  4. "sync"
  5. "time"
  6. "github.com/mhsanaei/3x-ui/v3/internal/logger"
  7. )
  8. type managed struct {
  9. server *Server
  10. tag string
  11. structuralFP string
  12. usersFP string
  13. }
  14. type Manager struct {
  15. mu sync.Mutex
  16. servers map[int]*managed
  17. lastStartErr map[int]string
  18. pendingTraffic map[string]ClientTrafficDelta
  19. }
  20. var (
  21. managerInstance *Manager
  22. managerOnce sync.Once
  23. )
  24. func GetManager() *Manager {
  25. managerOnce.Do(func() {
  26. managerInstance = &Manager{
  27. servers: make(map[int]*managed),
  28. lastStartErr: make(map[int]string),
  29. pendingTraffic: make(map[string]ClientTrafficDelta),
  30. }
  31. })
  32. return managerInstance
  33. }
  34. func (m *Manager) HasRunning() bool {
  35. m.mu.Lock()
  36. defer m.mu.Unlock()
  37. for _, mg := range m.servers {
  38. if mg.server != nil && mg.server.IsRunning() {
  39. return true
  40. }
  41. }
  42. return false
  43. }
  44. func (m *Manager) Ensure(inst Instance) error {
  45. m.mu.Lock()
  46. defer m.mu.Unlock()
  47. return m.ensureLocked(inst)
  48. }
  49. func (m *Manager) ensureLocked(inst Instance) error {
  50. if err := ValidateClients(inst.Clients); err != nil {
  51. return err
  52. }
  53. if len(inst.Clients) == 0 {
  54. m.removeLocked(inst.Id)
  55. return nil
  56. }
  57. structuralFP := inst.StructuralFingerprint()
  58. usersFP := inst.UsersFingerprint()
  59. if existing, ok := m.servers[inst.Id]; ok && existing != nil {
  60. if existing.server != nil && existing.server.IsRunning() && existing.structuralFP == structuralFP {
  61. existing.tag = inst.Tag
  62. existing.server.UpdateRuntimeSettings(inst.Tag, inst.CongestionControl, inst.LogLevel)
  63. if existing.usersFP != usersFP {
  64. existing.usersFP = usersFP
  65. existing.server.UpdateUsers(inst.Clients)
  66. }
  67. return nil
  68. }
  69. m.stopAndDrainLocked(existing)
  70. delete(m.servers, inst.Id)
  71. }
  72. server, err := m.startLocked(inst)
  73. if err != nil {
  74. if m.lastStartErr[inst.Id] != err.Error() {
  75. m.lastStartErr[inst.Id] = err.Error()
  76. if tuicLogWarn >= parseLogLevel(inst.LogLevel) {
  77. logger.Warningf("tuic: inbound %d (%s): failed to start server: %v", inst.Id, inst.Tag, err)
  78. }
  79. }
  80. return err
  81. }
  82. delete(m.lastStartErr, inst.Id)
  83. m.servers[inst.Id] = &managed{
  84. server: server,
  85. tag: inst.Tag,
  86. structuralFP: structuralFP,
  87. usersFP: usersFP,
  88. }
  89. return nil
  90. }
  91. func (m *Manager) startLocked(inst Instance) (*Server, error) {
  92. relay := &SocksRelay{
  93. Addr: fmt.Sprintf("127.0.0.1:%d", SOCKSPortForInbound(inst.Id)),
  94. Password: SocksPassword(),
  95. }
  96. server, err := NewServer(inst, relay)
  97. if err != nil {
  98. return nil, fmt.Errorf("tuic: init server for %d: %w", inst.Id, err)
  99. }
  100. if err := server.Start(); err != nil {
  101. return nil, fmt.Errorf("tuic: start server on %s for %d: %w", inst.BindTo(), inst.Id, err)
  102. }
  103. return server, nil
  104. }
  105. func (m *Manager) stopAndDrainLocked(mg *managed) {
  106. if mg == nil || mg.server == nil {
  107. return
  108. }
  109. _ = mg.server.Close()
  110. m.appendPendingTrafficLocked(mg.server.CollectClientTraffic())
  111. }
  112. func (m *Manager) appendPendingTrafficLocked(deltas []ClientTrafficDelta) {
  113. if m.pendingTraffic == nil {
  114. m.pendingTraffic = make(map[string]ClientTrafficDelta)
  115. }
  116. for _, delta := range deltas {
  117. key := delta.Email
  118. if delta.TrafficID > 0 {
  119. key = fmt.Sprintf("traffic:%d", delta.TrafficID)
  120. }
  121. if delta.TrafficID == 0 && delta.InboundID > 0 && delta.UUID != "" {
  122. key = fmt.Sprintf("%d:%s", delta.InboundID, delta.UUID)
  123. }
  124. current := m.pendingTraffic[key]
  125. current.Email = delta.Email
  126. current.UUID = delta.UUID
  127. current.InboundID = delta.InboundID
  128. current.TrafficID = delta.TrafficID
  129. current.Up += delta.Up
  130. current.Down += delta.Down
  131. m.pendingTraffic[key] = current
  132. }
  133. }
  134. func (m *Manager) RequeueClientTraffic(deltas []ClientTrafficDelta) {
  135. m.mu.Lock()
  136. defer m.mu.Unlock()
  137. m.appendPendingTrafficLocked(deltas)
  138. }
  139. func (m *Manager) GetActiveClients(window time.Duration) ([]string, []string) {
  140. m.mu.Lock()
  141. defer m.mu.Unlock()
  142. var emails []string
  143. var tags []string
  144. for _, mg := range m.servers {
  145. if mg.server != nil && mg.server.IsRunning() {
  146. active := mg.server.GetActiveEmails(window)
  147. if len(active) > 0 {
  148. emails = append(emails, active...)
  149. tags = append(tags, mg.tag)
  150. }
  151. }
  152. }
  153. return emails, tags
  154. }
  155. type InboundTrafficDelta struct {
  156. Tag string
  157. Up int64
  158. Down int64
  159. }
  160. func (m *Manager) CollectClientTraffic() []ClientTrafficDelta {
  161. _, clients := m.CollectAllTraffic()
  162. return clients
  163. }
  164. func (m *Manager) CollectAllTraffic() ([]InboundTrafficDelta, []ClientTrafficDelta) {
  165. m.mu.Lock()
  166. defer m.mu.Unlock()
  167. var inbounds []InboundTrafficDelta
  168. clients := make([]ClientTrafficDelta, 0, len(m.pendingTraffic))
  169. for email, delta := range m.pendingTraffic {
  170. clients = append(clients, delta)
  171. delete(m.pendingTraffic, email)
  172. }
  173. for _, mg := range m.servers {
  174. if mg.server != nil && mg.server.IsRunning() {
  175. up, down, cDeltas := mg.server.CollectAllTraffic()
  176. if up > 0 || down > 0 {
  177. inbounds = append(inbounds, InboundTrafficDelta{
  178. Tag: mg.tag,
  179. Up: up,
  180. Down: down,
  181. })
  182. }
  183. clients = append(clients, cDeltas...)
  184. }
  185. }
  186. return inbounds, clients
  187. }
  188. func (m *Manager) AddTestTraffic(id int, email string, up, down int64) bool {
  189. m.mu.Lock()
  190. defer m.mu.Unlock()
  191. if mg, ok := m.servers[id]; ok && mg.server != nil {
  192. return mg.server.AddTestTraffic(email, up, down)
  193. }
  194. return false
  195. }
  196. func (m *Manager) Remove(id int) {
  197. m.mu.Lock()
  198. defer m.mu.Unlock()
  199. m.removeLocked(id)
  200. }
  201. func (m *Manager) removeLocked(id int) {
  202. if existing, ok := m.servers[id]; ok && existing != nil {
  203. m.stopAndDrainLocked(existing)
  204. delete(m.servers, id)
  205. delete(m.lastStartErr, id)
  206. }
  207. }
  208. func (m *Manager) Reconcile(desired []Instance) {
  209. m.mu.Lock()
  210. defer m.mu.Unlock()
  211. desiredMap := make(map[int]Instance, len(desired))
  212. for _, inst := range desired {
  213. desiredMap[inst.Id] = inst
  214. }
  215. for id := range m.servers {
  216. if _, ok := desiredMap[id]; !ok {
  217. m.removeLocked(id)
  218. }
  219. }
  220. for _, inst := range desired {
  221. _ = m.ensureLocked(inst)
  222. }
  223. }
  224. func (m *Manager) StopAll() {
  225. m.mu.Lock()
  226. defer m.mu.Unlock()
  227. for _, mg := range m.servers {
  228. m.stopAndDrainLocked(mg)
  229. }
  230. m.servers = make(map[int]*managed)
  231. }