inbound_traffic.go 41 KB

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