inbound_traffic.go 42 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311
  1. package service
  2. import (
  3. "context"
  4. "encoding/json"
  5. "errors"
  6. "fmt"
  7. "maps"
  8. "slices"
  9. "strconv"
  10. "strings"
  11. "time"
  12. "github.com/mhsanaei/3x-ui/v3/internal/database"
  13. "github.com/mhsanaei/3x-ui/v3/internal/database/model"
  14. "github.com/mhsanaei/3x-ui/v3/internal/logger"
  15. "github.com/mhsanaei/3x-ui/v3/internal/web/runtime"
  16. "github.com/mhsanaei/3x-ui/v3/internal/xray"
  17. "gorm.io/gorm"
  18. "gorm.io/gorm/clause"
  19. )
  20. // A client with a renewal day set auto-renews too, so it must not read as
  21. // depleted — otherwise the operator's purge deletes it between cycles (#6239).
  22. const depletedClientsClause = "reset = 0 and reset_day = 0 and reset_weekday = 0 and ((total > 0 and up + down >= total) or (expiry_time > 0 and expiry_time <= ?))"
  23. func (s *InboundService) AddTraffic(inboundTraffics []*xray.Traffic, clientTraffics []*xray.ClientTraffic) (needRestart bool, clientsDisabled bool, err error) {
  24. var disabledNodeIDs []int
  25. var remotePlans []trafficInboundUpdatePlan
  26. var renewed []string
  27. err = submitTrafficWrite(func() error {
  28. var inner error
  29. needRestart, clientsDisabled, disabledNodeIDs, remotePlans, renewed, inner = s.addTrafficLocked(inboundTraffics, clientTraffics)
  30. return inner
  31. })
  32. if err != nil {
  33. return
  34. }
  35. s.resetMtprotoClientQuotas(renewed)
  36. // Off the serial writer: a hanging node must not stall traffic accounting.
  37. needRestart = s.applyTrafficRemotePlans(remotePlans) || needRestart
  38. if len(disabledNodeIDs) > 0 {
  39. s.restartRemoteNodesOnDisable(disabledNodeIDs)
  40. }
  41. return
  42. }
  43. func (s *InboundService) addTrafficLocked(inboundTraffics []*xray.Traffic, clientTraffics []*xray.ClientTraffic) (bool, bool, []int, []trafficInboundUpdatePlan, []string, error) {
  44. db := database.GetDB()
  45. // Commit durable traffic before best-effort lifecycle maintenance so helper
  46. // failures cannot discard usage already reported by Xray.
  47. if err := db.Transaction(func(tx *gorm.DB) error {
  48. if err := s.addInboundTraffic(tx, inboundTraffics); err != nil {
  49. return err
  50. }
  51. return s.addClientTraffic(tx, clientTraffics)
  52. }); err != nil {
  53. return false, false, nil, nil, nil, err
  54. }
  55. var (
  56. needRestart bool
  57. clientsDisabled bool
  58. disabledNodeIDs []int
  59. disabledClientsCount int64
  60. )
  61. batch := newTrafficMutationBatch()
  62. err := db.Transaction(func(tx *gorm.DB) error {
  63. needRestart0, count, err := s.autoRenewClients(tx, batch)
  64. if err != nil {
  65. return fmt.Errorf("renew clients: %w", err)
  66. }
  67. if count > 0 {
  68. logger.Debugf("%v clients renewed", count)
  69. }
  70. needRestart1, count, nodeIDs, err := s.disableInvalidClients(tx, batch)
  71. if err != nil {
  72. return fmt.Errorf("disable invalid clients: %w", err)
  73. }
  74. if count > 0 {
  75. logger.Debugf("%v clients disabled", count)
  76. disabledClientsCount = count
  77. }
  78. needRestart2, count, err := s.disableInvalidInbounds(tx, batch)
  79. if err != nil {
  80. return fmt.Errorf("disable invalid inbounds: %w", err)
  81. }
  82. if count > 0 {
  83. logger.Debugf("%v inbounds disabled", count)
  84. }
  85. if err := batch.markNodesTx(tx); err != nil {
  86. return err
  87. }
  88. needRestart = needRestart0 || needRestart1 || needRestart2
  89. clientsDisabled = disabledClientsCount > 0
  90. disabledNodeIDs = nodeIDs
  91. return nil
  92. })
  93. if err != nil {
  94. logger.Warning("traffic lifecycle maintenance failed after traffic commit:", err)
  95. return false, false, nil, nil, nil, nil
  96. }
  97. needRestart = needRestart || s.applyTrafficMutationBatch(batch)
  98. return needRestart, clientsDisabled, disabledNodeIDs, batch.remotePlans, batch.renewedEmails, nil
  99. }
  100. func (s *InboundService) addInboundTraffic(tx *gorm.DB, traffics []*xray.Traffic) error {
  101. if len(traffics) == 0 {
  102. return nil
  103. }
  104. var err error
  105. for _, traffic := range traffics {
  106. if traffic.IsInbound {
  107. err = tx.Model(&model.Inbound{}).Where("tag = ? AND node_id IS NULL", traffic.Tag).
  108. Updates(map[string]any{
  109. "up": gorm.Expr(database.ClampedAddExpr("up"), traffic.Up),
  110. "down": gorm.Expr(database.ClampedAddExpr("down"), traffic.Down),
  111. }).Error
  112. if err != nil {
  113. return err
  114. }
  115. }
  116. }
  117. return nil
  118. }
  119. func (s *InboundService) addClientTraffic(tx *gorm.DB, traffics []*xray.ClientTraffic) (err error) {
  120. if len(traffics) == 0 {
  121. return nil
  122. }
  123. traffics, err = canonicalizeClientTraffic(tx, traffics)
  124. if err != nil {
  125. return fmt.Errorf("resolve client traffic identities: %w", err)
  126. }
  127. emails := make([]string, 0, len(traffics))
  128. for _, traffic := range traffics {
  129. emails = append(emails, traffic.Email)
  130. }
  131. dbClientTraffics := make([]*xray.ClientTraffic, 0, len(traffics))
  132. // Match purely by email. client_traffics is email-keyed (one shared row per
  133. // email regardless of how many inbounds the client is attached to), and these
  134. // emails come from the local xray's report, so they always belong to a client
  135. // attached to a local inbound. The old `inbound_id NOT IN (node inbounds)`
  136. // filter dropped the local traffic of a client attached to both a node and the
  137. // mother inbound whenever the node inbound happened to be attached first — its
  138. // shared row then carried the node inbound's id (AddClientStat used to use
  139. // OnConflict DoNothing and never refreshed it; it now refreshes inbound_id on
  140. // conflict, but this filter was removed rather than relying on that ordering).
  141. err = tx.Model(xray.ClientTraffic{}).
  142. Where("email IN (?)", emails).
  143. Order("id").
  144. Find(&dbClientTraffics).Error
  145. if err != nil {
  146. return err
  147. }
  148. // Avoid empty slice error
  149. if len(dbClientTraffics) == 0 {
  150. return nil
  151. }
  152. dbClientTraffics, convertedExpiryByEmail, err := s.adjustTraffics(tx, dbClientTraffics)
  153. if err != nil {
  154. return err
  155. }
  156. // Index by email for O(N) merge.
  157. trafficByEmail := make(map[string]*xray.ClientTraffic, len(traffics))
  158. for i := range traffics {
  159. if traffics[i] != nil {
  160. trafficByEmail[traffics[i].Email] = traffics[i]
  161. }
  162. }
  163. now := time.Now().UnixMilli()
  164. // Use atomic per-row UPDATE instead of read-modify-write Save. tx.Save
  165. // issues UPDATEs in slice order, which varies between concurrent callers;
  166. // on PostgreSQL two transactions locking the same rows in opposite order
  167. // deadlock. An atomic "SET up = up + ?" never holds a row lock across a
  168. // subsequent lock acquisition, so concurrent writers cannot deadlock.
  169. for _, ct := range dbClientTraffics {
  170. t, ok := trafficByEmail[ct.Email]
  171. if !ok || (t.Up == 0 && t.Down == 0) {
  172. continue
  173. }
  174. if err = tx.Exec(
  175. fmt.Sprintf(
  176. `UPDATE client_traffics SET up = %s, down = %s, last_online = %s WHERE email = ?`,
  177. database.ClampedAddExpr("up"),
  178. database.ClampedAddExpr("down"),
  179. database.GreatestExpr("last_online", "?"),
  180. ),
  181. t.Up, t.Down, now, ct.Email,
  182. ).Error; err != nil {
  183. return fmt.Errorf("update traffic for %s: %w", ct.Email, err)
  184. }
  185. }
  186. // adjustTraffics converts delayed-start rows (negative ExpiryTime → absolute
  187. // deadline) in-memory. Persist that conversion now since the traffic UPDATE
  188. // above only touches up/down/last_online. Only converted emails are written:
  189. // updating every polled row issued one no-op UPDATE per active client per
  190. // poll. Sorted order keeps concurrent writers lock-compatible on Postgres.
  191. for _, email := range slices.Sorted(maps.Keys(convertedExpiryByEmail)) {
  192. if err = tx.Exec(
  193. `UPDATE client_traffics SET expiry_time = ? WHERE email = ? AND expiry_time < 0`,
  194. convertedExpiryByEmail[email], email,
  195. ).Error; err != nil {
  196. return fmt.Errorf("update expiry time for %s: %w", email, err)
  197. }
  198. }
  199. return nil
  200. }
  201. func canonicalizeClientTraffic(tx *gorm.DB, traffics []*xray.ClientTraffic) ([]*xray.ClientTraffic, error) {
  202. byEmail := make(map[string]*xray.ClientTraffic, len(traffics))
  203. for _, traffic := range traffics {
  204. if traffic == nil {
  205. continue
  206. }
  207. email := traffic.Email
  208. if traffic.TuicTrafficID > 0 {
  209. var owner xray.ClientTraffic
  210. err := tx.Select("email").Where("id = ?", traffic.TuicTrafficID).Take(&owner).Error
  211. if errors.Is(err, gorm.ErrRecordNotFound) {
  212. continue
  213. }
  214. if err != nil {
  215. return nil, err
  216. }
  217. email = owner.Email
  218. }
  219. current := byEmail[email]
  220. if current == nil {
  221. copy := *traffic
  222. copy.Email = email
  223. copy.TuicUUID = ""
  224. copy.TuicInboundId = 0
  225. byEmail[email] = &copy
  226. continue
  227. }
  228. current.Up += traffic.Up
  229. current.Down += traffic.Down
  230. current.Enable = current.Enable || traffic.Enable
  231. }
  232. result := make([]*xray.ClientTraffic, 0, len(byEmail))
  233. for _, traffic := range byEmail {
  234. result = append(result, traffic)
  235. }
  236. return result, nil
  237. }
  238. func (s *InboundService) adjustTraffics(tx *gorm.DB, dbClientTraffics []*xray.ClientTraffic) ([]*xray.ClientTraffic, map[string]int64, error) {
  239. now := time.Now().UnixMilli()
  240. // "Start After First Use" stores a negative expiry (the duration). On the
  241. // first traffic tick it becomes an absolute deadline of now+duration. Compute
  242. // it once per email so every inbound the client is attached to lands on the
  243. // same value (recomputing per inbound would skip all but the first one).
  244. newExpiryByEmail := make(map[string]int64, len(dbClientTraffics))
  245. for traffic_index := range dbClientTraffics {
  246. if dbClientTraffics[traffic_index].ExpiryTime < 0 {
  247. newExpiryByEmail[dbClientTraffics[traffic_index].Email] = now - dbClientTraffics[traffic_index].ExpiryTime
  248. }
  249. }
  250. if len(newExpiryByEmail) == 0 {
  251. return dbClientTraffics, nil, nil
  252. }
  253. delayedEmails := make([]string, 0, len(newExpiryByEmail))
  254. for email := range newExpiryByEmail {
  255. delayedEmails = append(delayedEmails, email)
  256. }
  257. // Resolve the owning inbounds through the client_inbounds link, which is
  258. // authoritative. client_traffics.inbound_id goes stale when an inbound is
  259. // deleted and recreated, which would leave the negative expiry unconverted.
  260. var inboundIds []int
  261. err := tx.Table("client_inbounds").
  262. Joins("JOIN clients ON clients.id = client_inbounds.client_id").
  263. Where("clients.email IN (?)", delayedEmails).
  264. Distinct().
  265. Pluck("client_inbounds.inbound_id", &inboundIds).Error
  266. if err != nil {
  267. return nil, nil, err
  268. }
  269. if len(inboundIds) == 0 {
  270. return dbClientTraffics, nil, nil
  271. }
  272. var inbounds []*model.Inbound
  273. err = tx.Model(model.Inbound{}).Where("id IN (?)", inboundIds).Find(&inbounds).Error
  274. if err != nil {
  275. return nil, nil, err
  276. }
  277. for inbound_index := range inbounds {
  278. settings := map[string]any{}
  279. _ = json.Unmarshal([]byte(inbounds[inbound_index].Settings), &settings)
  280. clients, ok := settings["clients"].([]any)
  281. if ok {
  282. var newClients []any
  283. for client_index := range clients {
  284. c := clients[client_index].(map[string]any)
  285. email, _ := c["email"].(string)
  286. if newExpiry, ok := newExpiryByEmail[email]; ok {
  287. c["expiryTime"] = newExpiry
  288. c["updated_at"] = now
  289. }
  290. if _, ok := c["created_at"]; !ok {
  291. c["created_at"] = now
  292. }
  293. if _, ok := c["updated_at"]; !ok {
  294. c["updated_at"] = now
  295. }
  296. newClients = append(newClients, any(c))
  297. }
  298. settings["clients"] = newClients
  299. modifiedSettings, err := json.MarshalIndent(settings, "", " ")
  300. if err != nil {
  301. return nil, nil, err
  302. }
  303. inbounds[inbound_index].Settings = string(modifiedSettings)
  304. }
  305. }
  306. for traffic_index := range dbClientTraffics {
  307. if newExpiry, ok := newExpiryByEmail[dbClientTraffics[traffic_index].Email]; ok {
  308. dbClientTraffics[traffic_index].ExpiryTime = newExpiry
  309. }
  310. }
  311. err = tx.Save(inbounds).Error
  312. if err != nil {
  313. logger.Warning("AddClientTraffic update inbounds ", err)
  314. logger.Error(inbounds)
  315. } else {
  316. for _, ib := range inbounds {
  317. if ib == nil {
  318. continue
  319. }
  320. cs, gcErr := s.GetClients(ib)
  321. if gcErr != nil {
  322. logger.Warning("AddClientTraffic sync clients: GetClients failed", gcErr)
  323. continue
  324. }
  325. if syncErr := s.clientService.SyncInbound(tx, ib.Id, cs); syncErr != nil {
  326. logger.Warning("AddClientTraffic sync clients: SyncInbound failed", syncErr)
  327. }
  328. }
  329. }
  330. return dbClientTraffics, newExpiryByEmail, nil
  331. }
  332. // apiUserFromClient prepares a stored client object for the runtime AddUser
  333. // call. The copy matters twice over: the stored object keeps being mutated and
  334. // marshalled back into the inbound's settings, which must not gain an API-only
  335. // key, and shadowsocks clients carry no cipher of their own — it lives on the
  336. // inbound, and without it the API cannot tell which of xray's two shadowsocks
  337. // account types the running inbound expects.
  338. func apiUserFromClient(client map[string]any, cipher string) map[string]any {
  339. user := maps.Clone(client)
  340. if user == nil {
  341. user = map[string]any{}
  342. }
  343. if cipher != "" {
  344. user["cipher"] = cipher
  345. }
  346. return user
  347. }
  348. // Candidates and renewals are not the same set: a skipped candidate keeps its
  349. // counters, so only the clients actually reset may lose their cross-panel rows.
  350. func (s *InboundService) autoRenewClients(tx *gorm.DB, mutationBatch *trafficMutationBatch) (bool, int64, error) {
  351. // check for time expired
  352. var traffics []*xray.ClientTraffic
  353. now := time.Now().Unix() * 1000
  354. var err error
  355. // Filter to clients that have at least one local inbound. Using
  356. // client_traffics.inbound_id is wrong: it goes stale after an inbound is
  357. // deleted/recreated and always points to the first inbound the client was
  358. // attached to, so it could be a node inbound even when the client also has
  359. // local inbounds. The email-based join through client_inbounds is authoritative.
  360. err = tx.Model(xray.ClientTraffic{}).
  361. Where("(reset > 0 or reset_day > 0 or reset_weekday > 0) and expiry_time > 0 and expiry_time <= ?", now).
  362. // A prepaid plan stops itself: once as many renewals have fired as the
  363. // operator allowed, the client is left to expire like any other.
  364. Where("reset_max <= 0 or reset_count < reset_max").
  365. Where("email IN (?)", tx.Table("client_inbounds ci").
  366. Select("c.email").
  367. Joins("JOIN clients c ON c.id = ci.client_id").
  368. Joins("JOIN inbounds i ON i.id = ci.inbound_id").
  369. Where("i.node_id IS NULL")).
  370. Find(&traffics).Error
  371. if err != nil {
  372. return false, 0, err
  373. }
  374. // return if there is no client to renew
  375. if len(traffics) == 0 {
  376. return false, 0, nil
  377. }
  378. renewLocation, locErr := (&SettingService{}).GetTimeLocation()
  379. if locErr != nil || renewLocation == nil {
  380. // Falling back to UTC keeps renewals happening; the alternative is
  381. // skipping them entirely because a setting could not be read.
  382. logger.Warning("autoRenewClients: could not read the panel time zone, using UTC:", locErr)
  383. renewLocation = time.UTC
  384. }
  385. var inbound_ids []int
  386. var inbounds []*model.Inbound
  387. needRestart := false
  388. type inboundClientKey struct {
  389. inboundID int
  390. email string
  391. }
  392. var clientsToAdd []struct {
  393. inbound model.Inbound
  394. client map[string]any
  395. }
  396. clientsToAddSet := make(map[inboundClientKey]struct{})
  397. // Resolve the inbounds to renew through the client_inbounds link rather than
  398. // client_traffics.inbound_id, which goes stale after an inbound is deleted and
  399. // recreated and would otherwise skip the renew entirely.
  400. renewEmails := make([]string, 0, len(traffics))
  401. for _, traffic := range traffics {
  402. renewEmails = append(renewEmails, traffic.Email)
  403. }
  404. for _, batch := range chunkStrings(renewEmails, sqliteMaxVars) {
  405. var ids []int
  406. if err = tx.Table("client_inbounds").
  407. Joins("JOIN clients ON clients.id = client_inbounds.client_id").
  408. Where("clients.email IN ?", batch).
  409. Distinct().
  410. Pluck("client_inbounds.inbound_id", &ids).Error; err != nil {
  411. return false, 0, err
  412. }
  413. inbound_ids = append(inbound_ids, ids...)
  414. }
  415. // Dedupe so an inbound hosting N expired clients is fetched and saved once
  416. // per tick instead of N times across chunk boundaries.
  417. inbound_ids = uniqueInts(inbound_ids)
  418. // Chunked to stay under SQLite's bind-variable limit when many inbounds
  419. // are touched in a single tick.
  420. for _, batch := range chunkInts(inbound_ids, sqliteMaxVars) {
  421. var page []*model.Inbound
  422. if err = tx.Model(model.Inbound{}).Where("id IN ?", batch).Find(&page).Error; err != nil {
  423. return false, 0, err
  424. }
  425. inbounds = append(inbounds, page...)
  426. }
  427. // Index the expired traffics by email so each client is an O(1) lookup
  428. // instead of a linear scan of every expired row (O(clients × expired) per
  429. // inbound, quadratic at scale). Pointers keep the in-place mutation below.
  430. trafficByEmail := make(map[string]*xray.ClientTraffic, len(traffics))
  431. // Keep the pre-renewal quota state: the shared pointer becomes enabled while
  432. // processing the first inbound, while an already-enabled row paired with
  433. // disabled settings represents an operator-disabled client we must preserve.
  434. trafficWasEnabled := make(map[string]bool, len(traffics))
  435. for i := range traffics {
  436. trafficByEmail[traffics[i].Email] = traffics[i]
  437. trafficWasEnabled[traffics[i].Email] = traffics[i].Enable
  438. }
  439. renewedEmails := make([]string, 0, len(traffics))
  440. for inbound_index := range inbounds {
  441. settings := map[string]any{}
  442. _ = json.Unmarshal([]byte(inbounds[inbound_index].Settings), &settings)
  443. clients, _ := settings["clients"].([]any)
  444. if len(clients) == 0 {
  445. continue
  446. }
  447. cipher := ""
  448. if inbounds[inbound_index].Protocol == model.Shadowsocks {
  449. cipher, _ = settings["method"].(string)
  450. }
  451. for client_index := range clients {
  452. c := clients[client_index].(map[string]any)
  453. email, _ := c["email"].(string)
  454. traffic, ok := trafficByEmail[email]
  455. if !ok {
  456. continue
  457. }
  458. // One allowance per period, not per tick: a client away for three
  459. // cycles must not catch up three of them against a prepaid cap.
  460. newExpiryTime, renewals := catchUpClientRenewal(traffic, now, renewLocation)
  461. if renewals > 0 {
  462. traffic.ExpiryTime = newExpiryTime
  463. traffic.ResetCount += renewals
  464. }
  465. c["expiryTime"] = traffic.ExpiryTime
  466. if traffic.ExpiryTime <= now {
  467. // Cap ran out mid-catch-up and the client is still expired: enabling it
  468. // for disableInvalidClients to undo adds and removes an xray user for nothing.
  469. clients[client_index] = any(c)
  470. continue
  471. }
  472. if renewals > 0 {
  473. traffic.Down = 0
  474. traffic.Up = 0
  475. renewedEmails = append(renewedEmails, email)
  476. }
  477. if !trafficWasEnabled[email] {
  478. traffic.Enable = true
  479. c["enable"] = true
  480. key := inboundClientKey{inboundID: inbounds[inbound_index].Id, email: email}
  481. if _, planned := clientsToAddSet[key]; !planned {
  482. clientsToAddSet[key] = struct{}{}
  483. clientsToAdd = append(clientsToAdd,
  484. struct {
  485. inbound model.Inbound
  486. client map[string]any
  487. }{
  488. inbound: *inbounds[inbound_index],
  489. client: apiUserFromClient(c, cipher),
  490. })
  491. }
  492. }
  493. clients[client_index] = any(c)
  494. }
  495. settings["clients"] = clients
  496. newSettings, err := json.MarshalIndent(settings, "", " ")
  497. if err != nil {
  498. return false, 0, err
  499. }
  500. inbounds[inbound_index].Settings = string(newSettings)
  501. }
  502. err = tx.Save(inbounds).Error
  503. if err != nil {
  504. return false, 0, err
  505. }
  506. for _, ib := range inbounds {
  507. if ib == nil {
  508. continue
  509. }
  510. cs, gcErr := s.GetClients(ib)
  511. if gcErr != nil {
  512. logger.Warning("autoRenewClients sync clients: GetClients failed", gcErr)
  513. continue
  514. }
  515. if syncErr := s.clientService.SyncInbound(tx, ib.Id, cs); syncErr != nil {
  516. logger.Warning("autoRenewClients sync clients: SyncInbound failed", syncErr)
  517. }
  518. }
  519. err = tx.Save(traffics).Error
  520. if err != nil {
  521. return false, 0, err
  522. }
  523. // A renewed client starts a fresh quota window: drop the cross-panel rows
  524. // too, or the stale pushed totals would re-deplete it immediately.
  525. if err = clearGlobalTraffic(tx, renewedEmails...); err != nil {
  526. return false, 0, err
  527. }
  528. mutationBatch.renewedEmails = append(mutationBatch.renewedEmails, renewedEmails...)
  529. for _, clientToAdd := range clientsToAdd {
  530. if clientToAdd.inbound.NodeID != nil {
  531. mutationBatch.addNode(*clientToAdd.inbound.NodeID)
  532. continue
  533. }
  534. mutationBatch.localPlans = append(mutationBatch.localPlans, trafficLocalApplyPlan{
  535. action: trafficAddUser, inbound: clientToAdd.inbound, client: clientToAdd.client,
  536. })
  537. }
  538. return needRestart, int64(len(renewedEmails)), nil
  539. }
  540. // AddClientStat inserts a per-client accounting row, or refreshes the
  541. // config-derived columns on an email conflict. Xray reports traffic per
  542. // email, so the surviving row also acts as the shared accumulator for
  543. // inbounds that re-use the same identity — every call for that identity
  544. // (one per attached inbound) carries the same enable/expiry/reset/total,
  545. // so re-asserting them here is idempotent for that legitimate case.
  546. //
  547. // The conflict path matters on its own for a second reason: an inbound
  548. // delete detaches its clients (InboundService.DelInbound) without deleting
  549. // their client_traffics row, by design — mirroring ClientService.Detach,
  550. // which intentionally leaves a fully-detached client's row in place so a
  551. // later Attach can resume it with its accumulated traffic intact. If that
  552. // same email is instead reused for a freshly (re)created client, the new
  553. // config's enable/expiry/reset/total must win over whatever the orphaned
  554. // row still holds; DoNothing left them stale indefinitely (#5958).
  555. //
  556. // up/down are deliberately excluded from the refresh: they are the
  557. // accumulated traffic totals, and zeroing them here would erase real usage
  558. // every time an existing, actively-used client is attached to one more
  559. // inbound. One tradeoff this does not resolve: a genuinely new client that
  560. // happens to reuse an orphaned email still inherits that row's leftover
  561. // up/down, since nothing at this call site can tell the two cases apart.
  562. func (s *InboundService) AddClientStat(tx *gorm.DB, inboundId int, client *model.Client) error {
  563. if err := validateClientRenewal(*client); err != nil {
  564. return err
  565. }
  566. clientTraffic := xray.ClientTraffic{
  567. InboundId: inboundId,
  568. Email: client.Email,
  569. Total: client.TotalGB,
  570. ExpiryTime: client.ExpiryTime,
  571. Enable: client.Enable,
  572. Reset: client.Reset,
  573. ResetDay: client.ResetDay,
  574. ResetWeekday: client.ResetWeekday,
  575. ResetMax: client.ResetMax,
  576. }
  577. return tx.Clauses(clause.OnConflict{
  578. Columns: []clause.Column{{Name: "email"}},
  579. DoUpdates: clause.AssignmentColumns([]string{"inbound_id", "total", "expiry_time", "enable", "reset", "reset_day", "reset_weekday", "reset_max"}),
  580. }).Create(&clientTraffic).Error
  581. }
  582. func (s *InboundService) UpdateClientStat(tx *gorm.DB, email string, client *model.Client) error {
  583. if err := validateClientRenewal(*client); err != nil {
  584. return err
  585. }
  586. result := tx.Model(xray.ClientTraffic{}).
  587. Where("email = ?", email).
  588. Updates(map[string]any{
  589. "enable": client.Enable,
  590. "email": client.Email,
  591. "total": client.TotalGB,
  592. "expiry_time": client.ExpiryTime,
  593. "reset": client.Reset,
  594. "reset_day": client.ResetDay,
  595. "reset_weekday": client.ResetWeekday,
  596. "reset_max": client.ResetMax,
  597. })
  598. err := result.Error
  599. return err
  600. }
  601. func (s *InboundService) DelClientStat(tx *gorm.DB, email string) error {
  602. if err := adjustGroupBaselinesForRemovedTraffic(tx, []string{email}); err != nil {
  603. return err
  604. }
  605. if err := tx.Where("email = ?", email).Delete(xray.ClientTraffic{}).Error; err != nil {
  606. return err
  607. }
  608. if err := clearGlobalTraffic(tx, email); err != nil {
  609. return err
  610. }
  611. return tx.Where("email = ?", email).Delete(&model.NodeClientTraffic{}).Error
  612. }
  613. func (s *InboundService) delClientStatsByEmails(tx *gorm.DB, emails []string) error {
  614. if err := adjustGroupBaselinesForRemovedTraffic(tx, emails); err != nil {
  615. return err
  616. }
  617. const chunk = 400
  618. for start := 0; start < len(emails); start += chunk {
  619. end := min(start+chunk, len(emails))
  620. batch := emails[start:end]
  621. if err := tx.Where("email IN ?", batch).Delete(xray.ClientTraffic{}).Error; err != nil {
  622. return err
  623. }
  624. if err := tx.Where("email IN ?", batch).Delete(&model.ClientGlobalTraffic{}).Error; err != nil {
  625. return err
  626. }
  627. if err := tx.Where("email IN ?", batch).Delete(&model.NodeClientTraffic{}).Error; err != nil {
  628. return err
  629. }
  630. }
  631. return nil
  632. }
  633. func (s *InboundService) ResetClientTrafficByEmail(clientEmail string) error {
  634. err := submitTrafficWrite(func() error {
  635. return database.GetDB().Transaction(func(tx *gorm.DB) error {
  636. if err := adjustGroupBaselinesForRemovedTraffic(tx, []string{clientEmail}); err != nil {
  637. return err
  638. }
  639. if err := clearGlobalTraffic(tx, clientEmail); err != nil {
  640. return err
  641. }
  642. if err := tx.Model(xray.ClientTraffic{}).
  643. Where("email = ?", clientEmail).
  644. Updates(map[string]any{"enable": true, "up": 0, "down": 0}).Error; err != nil {
  645. return err
  646. }
  647. return tx.Where("email = ?", clientEmail).Delete(&model.NodeClientTraffic{}).Error
  648. })
  649. })
  650. if err == nil {
  651. s.resetMtprotoClientQuota(clientEmail)
  652. }
  653. return err
  654. }
  655. func (s *InboundService) ResetClientTraffic(id int, clientEmail string) (needRestart bool, err error) {
  656. var ownNode *int
  657. err = submitTrafficWrite(func() error {
  658. var inner error
  659. needRestart, ownNode, inner = s.resetClientTrafficLocked(id, clientEmail)
  660. return inner
  661. })
  662. if err == nil {
  663. s.resetMtprotoClientQuota(clientEmail)
  664. // Siblings on other nodes are delivered by their own inbound's reset.
  665. if ownNode != nil {
  666. s.deliverNodeResetsNow([]int{*ownNode})
  667. }
  668. }
  669. return
  670. }
  671. func (s *InboundService) resetClientTrafficLocked(id int, clientEmail string) (bool, *int, error) {
  672. needRestart := false
  673. var reenablePlan *trafficLocalApplyPlan
  674. var reenableNodeID *int
  675. traffic, err := s.GetClientTrafficByEmail(clientEmail)
  676. if err != nil {
  677. return false, nil, err
  678. }
  679. if !traffic.Enable {
  680. inbound, err := s.GetInbound(id)
  681. if err != nil {
  682. return false, nil, err
  683. }
  684. clients, err := s.GetClients(inbound)
  685. if err != nil {
  686. return false, nil, err
  687. }
  688. for _, client := range clients {
  689. if client.Email == clientEmail && client.Enable {
  690. cipher := ""
  691. if string(inbound.Protocol) == "shadowsocks" {
  692. var oldSettings map[string]any
  693. err = json.Unmarshal([]byte(inbound.Settings), &oldSettings)
  694. if err != nil {
  695. return false, nil, err
  696. }
  697. cipher, _ = oldSettings["method"].(string)
  698. }
  699. clientMap := map[string]any{
  700. "email": client.Email,
  701. "id": client.ID,
  702. "auth": client.Auth,
  703. "security": client.Security,
  704. "flow": client.Flow,
  705. "password": client.Password,
  706. "cipher": cipher,
  707. "reverse": client.Reverse,
  708. }
  709. if inbound.NodeID != nil {
  710. reenableNodeID = inbound.NodeID
  711. } else {
  712. reenablePlan = &trafficLocalApplyPlan{action: trafficAddUser, inbound: *inbound, client: clientMap}
  713. }
  714. break
  715. }
  716. }
  717. }
  718. traffic.Up = 0
  719. traffic.Down = 0
  720. traffic.Enable = true
  721. db := database.GetDB()
  722. now := time.Now().UnixMilli()
  723. inbound, err := s.GetInbound(id)
  724. if err != nil {
  725. return false, nil, err
  726. }
  727. if err := db.Transaction(func(tx *gorm.DB) error {
  728. if err := adjustGroupBaselinesForRemovedTraffic(tx, []string{clientEmail}); err != nil {
  729. return err
  730. }
  731. if err := tx.Save(traffic).Error; err != nil {
  732. return err
  733. }
  734. if err := clearGlobalTraffic(tx, clientEmail); err != nil {
  735. return err
  736. }
  737. if err := tx.Where("email = ?", clientEmail).Delete(&model.NodeClientTraffic{}).Error; err != nil {
  738. return err
  739. }
  740. if _, err := queueNodeResets(tx, []string{clientEmail}); err != nil {
  741. return err
  742. }
  743. if err := tx.Model(model.Inbound{}).
  744. Where("id = ?", id).
  745. Update("last_traffic_reset_time", now).Error; err != nil {
  746. return err
  747. }
  748. if reenableNodeID != nil {
  749. return (&NodeService{}).MarkNodeDirtyTx(tx, *reenableNodeID)
  750. }
  751. if inbound != nil && inbound.NodeID != nil {
  752. return (&NodeService{}).MarkNodeDirtyTx(tx, *inbound.NodeID)
  753. }
  754. return nil
  755. }); err != nil {
  756. return false, nil, err
  757. }
  758. if reenablePlan != nil {
  759. rt, err := s.runtimeFor(&reenablePlan.inbound)
  760. if err != nil {
  761. needRestart = true
  762. } else if err := rt.AddUser(context.Background(), &reenablePlan.inbound, reenablePlan.client); err != nil {
  763. logger.Debug("Error in enabling client on", rt.Name(), ":", err)
  764. needRestart = true
  765. } else {
  766. logger.Debug("Client enabled on", rt.Name(), "due to reset traffic:", clientEmail)
  767. }
  768. }
  769. if inbound != nil {
  770. return needRestart, inbound.NodeID, nil
  771. }
  772. return needRestart, nil, nil
  773. }
  774. func (s *InboundService) ResetAllTraffics() error {
  775. err := submitTrafficWrite(func() error {
  776. return s.resetAllTrafficsLocked()
  777. })
  778. if err == nil {
  779. s.propagateResetAllTrafficsToNodes()
  780. }
  781. return err
  782. }
  783. func (s *InboundService) resetAllTrafficsLocked() error {
  784. db := database.GetDB()
  785. now := time.Now().UnixMilli()
  786. return db.Model(model.Inbound{}).
  787. Where("user_id > ?", 0).
  788. Updates(map[string]any{
  789. "up": 0,
  790. "down": 0,
  791. "last_traffic_reset_time": now,
  792. }).Error
  793. }
  794. // propagateResetAllTrafficsToNodes tells every node to zero its own counters.
  795. // Kept OUT of the traffic-writer transaction: each remote call can block up to
  796. // remoteHTTPTimeout, and holding the single serial writer across N such calls
  797. // stalls traffic accounting and drops the deltas of every concurrent poll.
  798. func (s *InboundService) propagateResetAllTrafficsToNodes() {
  799. nodes, err := (&NodeService{}).GetAll()
  800. if err != nil {
  801. return
  802. }
  803. ids := make([]int, len(nodes))
  804. for i, node := range nodes {
  805. ids[i] = node.Id
  806. }
  807. fanoutInboundResults(ids, nodeFanoutConcurrency, func(i int) struct{} {
  808. if rt, err := runtime.GetManager().RuntimeFor(&ids[i]); err == nil {
  809. if e := rt.ResetAllTraffics(context.Background()); e != nil {
  810. logger.Warning("ResetAllTraffics: remote propagation to", rt.Name(), "failed:", e)
  811. }
  812. }
  813. return struct{}{}
  814. })
  815. }
  816. func (s *InboundService) ResetInboundTraffic(id int) error {
  817. var inbound *model.Inbound
  818. if err := submitTrafficWrite(func() error {
  819. db := database.GetDB()
  820. if err := db.Model(model.Inbound{}).
  821. Where("id = ?", id).
  822. Updates(map[string]any{"up": 0, "down": 0}).Error; err != nil {
  823. return err
  824. }
  825. var err error
  826. inbound, err = s.GetInbound(id)
  827. if err != nil {
  828. return err
  829. }
  830. return nil
  831. }); err != nil {
  832. return err
  833. }
  834. if inbound != nil && inbound.NodeID != nil {
  835. if rt, rterr := s.runtimeFor(inbound); rterr == nil {
  836. if e := rt.ResetInboundTraffic(context.Background(), inbound); e != nil {
  837. logger.Warning("ResetInboundTraffic: remote propagation to", rt.Name(), "failed:", e)
  838. }
  839. } else {
  840. logger.Warning("ResetInboundTraffic: runtime lookup failed:", rterr)
  841. }
  842. }
  843. return nil
  844. }
  845. func (s *InboundService) DelDepletedClients(id int) (err error) {
  846. db := database.GetDB()
  847. var deletedInbounds []model.Inbound
  848. err = db.Transaction(func(tx *gorm.DB) error {
  849. // Collect depleted emails globally — a shared-email row owned by one
  850. // inbound depletes every sibling that lists the email.
  851. now := time.Now().Unix() * 1000
  852. depletedClause := depletedClientsClause
  853. var depletedRows []xray.ClientTraffic
  854. if err := tx.Model(xray.ClientTraffic{}).
  855. Where(depletedClause, now).
  856. Find(&depletedRows).Error; err != nil {
  857. return err
  858. }
  859. if len(depletedRows) == 0 {
  860. return nil
  861. }
  862. depletedEmails := make(map[string]struct{}, len(depletedRows))
  863. for _, r := range depletedRows {
  864. if r.Email == "" {
  865. continue
  866. }
  867. depletedEmails[strings.ToLower(r.Email)] = struct{}{}
  868. }
  869. if len(depletedEmails) == 0 {
  870. return nil
  871. }
  872. var inbounds []*model.Inbound
  873. inboundQuery := tx.Model(model.Inbound{})
  874. if id >= 0 {
  875. inboundQuery = inboundQuery.Where("id = ?", id)
  876. }
  877. if err := inboundQuery.Find(&inbounds).Error; err != nil {
  878. return err
  879. }
  880. for _, inbound := range inbounds {
  881. var settings map[string]any
  882. if err := json.Unmarshal([]byte(inbound.Settings), &settings); err != nil {
  883. return err
  884. }
  885. rawClients, ok := settings["clients"].([]any)
  886. if !ok {
  887. continue
  888. }
  889. newClients := make([]any, 0, len(rawClients))
  890. removed := 0
  891. for _, client := range rawClients {
  892. c, ok := client.(map[string]any)
  893. if !ok {
  894. newClients = append(newClients, client)
  895. continue
  896. }
  897. email, _ := c["email"].(string)
  898. if _, isDepleted := depletedEmails[strings.ToLower(email)]; isDepleted {
  899. removed++
  900. continue
  901. }
  902. newClients = append(newClients, client)
  903. }
  904. if removed == 0 {
  905. continue
  906. }
  907. if len(newClients) == 0 {
  908. deletedInbounds = append(deletedInbounds, *inbound)
  909. if err := s.clientService.DetachInbound(tx, inbound.Id); err != nil {
  910. return err
  911. }
  912. if err := tx.Where("inbound_id = ?", inbound.Id).Delete(&model.Host{}).Error; err != nil {
  913. return err
  914. }
  915. if err := tx.Delete(model.Inbound{}, inbound.Id).Error; err != nil {
  916. return err
  917. }
  918. if inbound.NodeID != nil {
  919. if err := (&NodeService{}).MarkNodeDirtyTx(tx, *inbound.NodeID); err != nil {
  920. return err
  921. }
  922. }
  923. continue
  924. }
  925. settings["clients"] = newClients
  926. ns, mErr := json.MarshalIndent(settings, "", " ")
  927. if mErr != nil {
  928. return mErr
  929. }
  930. inbound.Settings = string(ns)
  931. if err := tx.Save(inbound).Error; err != nil {
  932. return err
  933. }
  934. survivingClients, gcErr := s.GetClients(inbound)
  935. if gcErr != nil {
  936. return gcErr
  937. }
  938. if err := s.clientService.SyncInbound(tx, inbound.Id, survivingClients); err != nil {
  939. return err
  940. }
  941. if inbound.NodeID != nil {
  942. if err := (&NodeService{}).MarkNodeDirtyTx(tx, *inbound.NodeID); err != nil {
  943. return err
  944. }
  945. }
  946. }
  947. // Drop now-orphaned rows. With id >= 0, a row is safe to drop only when
  948. // no out-of-scope inbound still references the email.
  949. if id < 0 {
  950. return tx.Where(depletedClause, now).Delete(xray.ClientTraffic{}).Error
  951. }
  952. emails := make([]string, 0, len(depletedEmails))
  953. for e := range depletedEmails {
  954. emails = append(emails, e)
  955. }
  956. var stillReferenced []string
  957. emailExpr := database.JSONFieldText("client.value", "email")
  958. stillQuery := fmt.Sprintf(
  959. "SELECT DISTINCT LOWER(%s) %s WHERE LOWER(%s) IN ?",
  960. emailExpr,
  961. database.JSONClientsFromInbound(),
  962. emailExpr,
  963. )
  964. if err := tx.Raw(stillQuery, emails).Scan(&stillReferenced).Error; err != nil {
  965. return err
  966. }
  967. stillSet := make(map[string]struct{}, len(stillReferenced))
  968. for _, e := range stillReferenced {
  969. stillSet[e] = struct{}{}
  970. }
  971. toDelete := make([]string, 0, len(emails))
  972. for _, e := range emails {
  973. if _, kept := stillSet[e]; !kept {
  974. toDelete = append(toDelete, e)
  975. }
  976. }
  977. if len(toDelete) > 0 {
  978. if err := tx.Where("LOWER(email) IN ?", toDelete).Delete(xray.ClientTraffic{}).Error; err != nil {
  979. return err
  980. }
  981. }
  982. return nil
  983. })
  984. if err != nil {
  985. return err
  986. }
  987. for i := range deletedInbounds {
  988. inbound := &deletedInbounds[i]
  989. if rt, rtErr := s.runtimeFor(inbound); rtErr != nil {
  990. logger.Warning("DelDepletedClients: runtime lookup failed after commit:", rtErr)
  991. } else if rtErr = rt.DelInbound(context.Background(), inbound); rtErr != nil && !xray.IsMissingHandlerErr(rtErr) {
  992. logger.Warning("DelDepletedClients: runtime cleanup failed after commit:", rtErr)
  993. }
  994. if inbound.Tag != "" {
  995. if _, syncErr := (&XraySettingService{}).RemoveInboundTagReferences(inbound.Tag); syncErr != nil {
  996. logger.Warning("DelDepletedClients: routing cleanup failed after commit:", syncErr)
  997. }
  998. }
  999. }
  1000. return nil
  1001. }
  1002. func (s *InboundService) GetClientTrafficTgBot(tgId int64) ([]*xray.ClientTraffic, error) {
  1003. db := database.GetDB()
  1004. idQuery := fmt.Sprintf(
  1005. "SELECT DISTINCT inbounds.id %s WHERE %s = ?",
  1006. database.JSONClientsFromInbound(),
  1007. database.JSONFieldText("client.value", "tgId"),
  1008. )
  1009. var inboundIds []int
  1010. if err := db.Raw(idQuery, strconv.FormatInt(tgId, 10)).Scan(&inboundIds).Error; err != nil {
  1011. logger.Errorf("Error retrieving inbounds with tgId %d: %v", tgId, err)
  1012. return nil, err
  1013. }
  1014. var inbounds []*model.Inbound
  1015. if len(inboundIds) > 0 {
  1016. err := db.Model(model.Inbound{}).Where("id IN ?", inboundIds).Find(&inbounds).Error
  1017. if err != nil && !errors.Is(err, gorm.ErrRecordNotFound) {
  1018. logger.Errorf("Error retrieving inbounds with tgId %d: %v", tgId, err)
  1019. return nil, err
  1020. }
  1021. }
  1022. var emails []string
  1023. for _, inbound := range inbounds {
  1024. clients, err := s.GetClients(inbound)
  1025. if err != nil {
  1026. logger.Errorf("Error retrieving clients for inbound %d: %v", inbound.Id, err)
  1027. continue
  1028. }
  1029. for _, client := range clients {
  1030. if client.TgID == tgId {
  1031. emails = append(emails, client.Email)
  1032. }
  1033. }
  1034. }
  1035. // Chunked to stay under SQLite's bind-variable limit when a single Telegram
  1036. // account owns thousands of clients across inbounds.
  1037. uniqEmails := uniqueNonEmptyStrings(emails)
  1038. traffics := make([]*xray.ClientTraffic, 0, len(uniqEmails))
  1039. for _, batch := range chunkStrings(uniqEmails, sqliteMaxVars) {
  1040. var page []*xray.ClientTraffic
  1041. if err := db.Model(xray.ClientTraffic{}).Where("email IN ?", batch).Find(&page).Error; err != nil {
  1042. if errors.Is(err, gorm.ErrRecordNotFound) {
  1043. continue
  1044. }
  1045. logger.Errorf("Error retrieving ClientTraffic for emails %v: %v", batch, err)
  1046. return nil, err
  1047. }
  1048. traffics = append(traffics, page...)
  1049. }
  1050. if len(traffics) == 0 {
  1051. logger.Warning("No ClientTraffic records found for emails:", emails)
  1052. return nil, nil
  1053. }
  1054. // Populate UUID and other client data for each traffic record
  1055. for i := range traffics {
  1056. if ct, client, e := s.GetClientByEmail(traffics[i].Email); e == nil && ct != nil && client != nil {
  1057. traffics[i].Enable = client.Enable
  1058. traffics[i].UUID = client.ID
  1059. traffics[i].SubId = client.SubID
  1060. }
  1061. }
  1062. return traffics, nil
  1063. }
  1064. // BumpClientsLastOnline sets client_traffics.last_online to now for the given
  1065. // emails. Used in online-API mode for clients that hold a live connection but
  1066. // moved no bytes this poll — the traffic path (addClientTraffic) only bumps
  1067. // last_online on a non-zero delta, so idle-but-connected clients would
  1068. // otherwise show a stale "last online" while being reported online.
  1069. func (s *InboundService) BumpClientsLastOnline(emails []string) error {
  1070. uniq := uniqueNonEmptyStrings(emails)
  1071. if len(uniq) == 0 {
  1072. return nil
  1073. }
  1074. now := time.Now().UnixMilli()
  1075. return submitTrafficWrite(func() error {
  1076. db := database.GetDB()
  1077. for _, batch := range chunkStrings(uniq, sqliteMaxVars) {
  1078. if err := db.Model(xray.ClientTraffic{}).Where("email IN ?", batch).Update("last_online", now).Error; err != nil {
  1079. return err
  1080. }
  1081. }
  1082. return nil
  1083. })
  1084. }
  1085. func (s *InboundService) GetActiveClientTraffics(emails []string) ([]*xray.ClientTraffic, error) {
  1086. uniq := uniqueNonEmptyStrings(emails)
  1087. if len(uniq) == 0 {
  1088. return nil, nil
  1089. }
  1090. db := database.GetDB()
  1091. traffics := make([]*xray.ClientTraffic, 0, len(uniq))
  1092. for _, batch := range chunkStrings(uniq, sqliteMaxVars) {
  1093. var page []*xray.ClientTraffic
  1094. if err := db.Model(xray.ClientTraffic{}).Where("email IN ?", batch).Find(&page).Error; err != nil {
  1095. return nil, err
  1096. }
  1097. traffics = append(traffics, page...)
  1098. }
  1099. overlayGlobalTraffic(db, traffics)
  1100. return traffics, nil
  1101. }
  1102. // GetAllClientTraffics returns the full set of client_traffics rows so the
  1103. // websocket broadcasters can ship a complete snapshot every cycle. A pure
  1104. // delta path silently dropped the per-client section whenever no client moved
  1105. // bytes in the cycle or a node sync failed, leaving client rows in the UI
  1106. // stuck at stale numbers — so small installs broadcast this snapshot, and only
  1107. // above the traffic job's snapshot threshold (where the marshaled snapshot
  1108. // would exceed the hub's payload cap and be dropped wholesale) does the job
  1109. // fall back to active-row deltas.
  1110. func (s *InboundService) GetAllClientTraffics() ([]*xray.ClientTraffic, error) {
  1111. db := database.GetDB()
  1112. var traffics []*xray.ClientTraffic
  1113. if err := db.Model(xray.ClientTraffic{}).Find(&traffics).Error; err != nil {
  1114. return nil, err
  1115. }
  1116. overlayGlobalTraffic(db, traffics)
  1117. return traffics, nil
  1118. }
  1119. func (s *InboundService) CountClientTraffics() (int64, error) {
  1120. db := database.GetDB()
  1121. var count int64
  1122. err := db.Model(xray.ClientTraffic{}).Count(&count).Error
  1123. return count, err
  1124. }
  1125. type InboundTrafficSummary struct {
  1126. Id int `json:"id" example:"1"`
  1127. Up int64 `json:"up" example:"1048576"`
  1128. Down int64 `json:"down" example:"2097152"`
  1129. Total int64 `json:"total" example:"10737418240"`
  1130. Enable bool `json:"enable" example:"true"`
  1131. }
  1132. func (s *InboundService) GetInboundsTrafficSummary() ([]InboundTrafficSummary, error) {
  1133. db := database.GetDB()
  1134. var summaries []InboundTrafficSummary
  1135. if err := db.Model(&model.Inbound{}).
  1136. Select("id, up, down, total, enable").
  1137. Find(&summaries).Error; err != nil {
  1138. return nil, err
  1139. }
  1140. return summaries, nil
  1141. }
  1142. func (s *InboundService) GetClientTrafficByEmail(email string) (traffic *xray.ClientTraffic, err error) {
  1143. db := database.GetDB()
  1144. var traffics []*xray.ClientTraffic
  1145. if err := db.Model(xray.ClientTraffic{}).Where("email = ?", email).Find(&traffics).Error; err != nil {
  1146. logger.Warningf("Error retrieving ClientTraffic with email %s: %v", email, err)
  1147. return nil, err
  1148. }
  1149. if len(traffics) == 0 {
  1150. return nil, nil
  1151. }
  1152. overlayGlobalTraffic(db, traffics)
  1153. t := traffics[0]
  1154. if rec, rErr := s.clientService.GetRecordByEmail(db, email); rErr == nil && rec != nil {
  1155. c := rec.ToClient()
  1156. t.UUID = c.ID
  1157. t.SubId = c.SubID
  1158. return t, nil
  1159. }
  1160. t2, client, err := s.GetClientByEmail(email)
  1161. if err != nil {
  1162. logger.Warningf("Error retrieving ClientTraffic with email %s: %v", email, err)
  1163. return nil, err
  1164. }
  1165. if t2 != nil && client != nil {
  1166. t2.UUID = client.ID
  1167. t2.SubId = client.SubID
  1168. return t2, nil
  1169. }
  1170. return nil, nil
  1171. }
  1172. func (s *InboundService) UpdateClientTrafficByEmail(email string, upload int64, download int64) error {
  1173. return submitTrafficWrite(func() error {
  1174. db := database.GetDB()
  1175. err := db.Model(xray.ClientTraffic{}).
  1176. Where("email = ?", email).
  1177. Updates(map[string]any{
  1178. "up": upload,
  1179. "down": download,
  1180. }).Error
  1181. if err != nil {
  1182. logger.Warningf("Error updating ClientTraffic with email %s: %v", email, err)
  1183. }
  1184. return err
  1185. })
  1186. }
  1187. func (s *InboundService) SearchClientTraffic(query string) (traffic *xray.ClientTraffic, err error) {
  1188. db := database.GetDB()
  1189. inbound := &model.Inbound{}
  1190. traffic = &xray.ClientTraffic{}
  1191. // Search for inbound settings that contain the query
  1192. err = db.Model(model.Inbound{}).Where("settings LIKE ?", "%\""+query+"\"%").First(inbound).Error
  1193. if err != nil {
  1194. if errors.Is(err, gorm.ErrRecordNotFound) {
  1195. logger.Warningf("Inbound settings containing query %s not found: %v", query, err)
  1196. return nil, err
  1197. }
  1198. logger.Errorf("Error searching for inbound settings with query %s: %v", query, err)
  1199. return nil, err
  1200. }
  1201. traffic.InboundId = inbound.Id
  1202. clients, err := ParseInboundSettingsClients(inbound.Settings)
  1203. if err != nil {
  1204. logger.Errorf("Error unmarshalling inbound settings for inbound ID %d: %v", inbound.Id, err)
  1205. return nil, err
  1206. }
  1207. for _, client := range clients {
  1208. if (client.ID == query || client.Password == query) && client.Email != "" {
  1209. traffic.Email = client.Email
  1210. break
  1211. }
  1212. }
  1213. if traffic.Email == "" {
  1214. logger.Warningf("No client found with query %s in inbound ID %d", query, inbound.Id)
  1215. return nil, gorm.ErrRecordNotFound
  1216. }
  1217. // Retrieve ClientTraffic based on the found email
  1218. err = db.Model(xray.ClientTraffic{}).Where("email = ?", traffic.Email).First(traffic).Error
  1219. if err != nil {
  1220. if errors.Is(err, gorm.ErrRecordNotFound) {
  1221. logger.Warningf("ClientTraffic for email %s not found: %v", traffic.Email, err)
  1222. return nil, err
  1223. }
  1224. logger.Errorf("Error retrieving ClientTraffic for email %s: %v", traffic.Email, err)
  1225. return nil, err
  1226. }
  1227. return traffic, nil
  1228. }