client_bulk.go 55 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746174717481749175017511752175317541755175617571758175917601761176217631764176517661767176817691770177117721773177417751776177717781779178017811782178317841785178617871788178917901791179217931794179517961797179817991800180118021803180418051806180718081809181018111812181318141815181618171818181918201821182218231824182518261827182818291830183118321833183418351836183718381839184018411842184318441845184618471848184918501851185218531854185518561857185818591860
  1. package service
  2. import (
  3. "context"
  4. "encoding/json"
  5. "fmt"
  6. "strings"
  7. "time"
  8. "github.com/google/uuid"
  9. "github.com/mhsanaei/3x-ui/v3/internal/database"
  10. "github.com/mhsanaei/3x-ui/v3/internal/database/model"
  11. "github.com/mhsanaei/3x-ui/v3/internal/logger"
  12. "github.com/mhsanaei/3x-ui/v3/internal/util/common"
  13. "github.com/mhsanaei/3x-ui/v3/internal/xray"
  14. "gorm.io/gorm"
  15. )
  16. // BulkAttachResult reports the outcome of a bulk attach across target inbounds.
  17. type BulkAttachResult struct {
  18. Attached []string `json:"attached"`
  19. Skipped []string `json:"skipped"`
  20. Errors []string `json:"errors"`
  21. }
  22. // BulkAttach attaches the given existing clients (by email) to each target inbound,
  23. // reusing their identity (email/UUID/password/subId) and a shared traffic row. It adds
  24. // all clients to a target in a single AddInboundClient call, and reports clients already
  25. // present on a target as skipped.
  26. func (s *ClientService) BulkAttach(inboundSvc *InboundService, emails []string, inboundIds []int) (*BulkAttachResult, bool, error) {
  27. result := &BulkAttachResult{}
  28. if len(emails) == 0 || len(inboundIds) == 0 {
  29. return result, false, nil
  30. }
  31. recordErr := func(format string, args ...any) {
  32. msg := fmt.Sprintf(format, args...)
  33. result.Errors = append(result.Errors, msg)
  34. logger.Warningf("[BulkAttach] %s", msg)
  35. }
  36. records := make([]*model.ClientRecord, 0, len(emails))
  37. seenEmail := make(map[string]struct{}, len(emails))
  38. for _, email := range emails {
  39. if email == "" {
  40. continue
  41. }
  42. key := strings.ToLower(email)
  43. if _, ok := seenEmail[key]; ok {
  44. continue
  45. }
  46. seenEmail[key] = struct{}{}
  47. rec, err := s.GetRecordByEmail(nil, email)
  48. if err != nil {
  49. recordErr("%s: %v", email, err)
  50. continue
  51. }
  52. records = append(records, rec)
  53. }
  54. // Same rule as Attach (#4834): clients.flow is unreliable when a non-flow
  55. // inbound synced last, so seed from EffectiveFlow before clientWithInboundFlow.
  56. emailsForFlow := make([]string, 0, len(records))
  57. for _, rec := range records {
  58. emailsForFlow = append(emailsForFlow, rec.Email)
  59. }
  60. flowsByEmail, err := s.EffectiveFlowsByEmails(nil, emailsForFlow)
  61. if err != nil {
  62. return result, false, err
  63. }
  64. needRestart := false
  65. // Prepared in order first, as in Create: fillProtocolDefaults mints the
  66. // shared credentials, so only the node pushes below may overlap.
  67. attachIds := make([]int, 0, len(inboundIds))
  68. attachPayloads := make([]string, 0, len(inboundIds))
  69. attachClients := make([][]model.Client, 0, len(inboundIds))
  70. attachAnyTunnel := false
  71. // A repeated id used to be caught by the second pass seeing the client
  72. // already attached; the applies no longer run before the next prep.
  73. seenInbound := make(map[int]struct{}, len(inboundIds))
  74. for _, ibId := range inboundIds {
  75. if _, dup := seenInbound[ibId]; dup {
  76. continue
  77. }
  78. seenInbound[ibId] = struct{}{}
  79. inbound, err := inboundSvc.GetInbound(ibId)
  80. if err != nil {
  81. recordErr("inbound %d: %v", ibId, err)
  82. continue
  83. }
  84. if inbound.Protocol == model.WireGuard || inbound.Protocol == model.AmneziaWG {
  85. attachAnyTunnel = true
  86. }
  87. existingClients, err := inboundSvc.GetClients(inbound)
  88. if err != nil {
  89. recordErr("inbound %d: %v", ibId, err)
  90. continue
  91. }
  92. have := make(map[string]struct{}, len(existingClients))
  93. for _, c := range existingClients {
  94. have[strings.ToLower(c.Email)] = struct{}{}
  95. }
  96. clientsToAdd := make([]model.Client, 0, len(records))
  97. for _, rec := range records {
  98. if _, attached := have[strings.ToLower(rec.Email)]; attached {
  99. result.Skipped = append(result.Skipped, rec.Email)
  100. continue
  101. }
  102. client := *rec.ToClient()
  103. if flow, ok := flowsByEmail[rec.Email]; ok && flow != "" {
  104. client.Flow = flow
  105. }
  106. client.UpdatedAt = time.Now().UnixMilli()
  107. if err := s.fillProtocolDefaults(&client, inbound); err != nil {
  108. recordErr("%s -> inbound %d: %v", rec.Email, ibId, err)
  109. continue
  110. }
  111. clientsToAdd = append(clientsToAdd, clientWithInboundFlow(client, inbound))
  112. }
  113. if len(clientsToAdd) == 0 {
  114. continue
  115. }
  116. payload, err := json.Marshal(map[string][]model.Client{"clients": clientsToAdd})
  117. if err != nil {
  118. recordErr("inbound %d: %v", ibId, err)
  119. continue
  120. }
  121. attachIds = append(attachIds, ibId)
  122. attachPayloads = append(attachPayloads, string(payload))
  123. attachClients = append(attachClients, clientsToAdd)
  124. }
  125. attachResults, attachPanics := fanoutInboundResults(attachIds, addFanoutLimit(attachAnyTunnel), func(i int) inboundApplyOutcome {
  126. nr, err := s.AddInboundClient(inboundSvc, &model.Inbound{Id: attachIds[i], Settings: attachPayloads[i]})
  127. return inboundApplyOutcome{needRestart: nr, err: err}
  128. })
  129. for i, out := range attachResults {
  130. err := out.err
  131. if attachPanics[i] != nil {
  132. // The apply may already have committed, so ask for the restart the
  133. // lost return value can no longer report.
  134. needRestart = true
  135. err = attachPanics[i]
  136. }
  137. if err != nil {
  138. recordErr("inbound %d: %v", attachIds[i], err)
  139. continue
  140. }
  141. if out.needRestart {
  142. needRestart = true
  143. }
  144. for _, c := range attachClients[i] {
  145. result.Attached = append(result.Attached, c.Email)
  146. }
  147. }
  148. return result, needRestart, nil
  149. }
  150. // BulkDetachResult reports the outcome of a bulk detach across target inbounds.
  151. type BulkDetachResult struct {
  152. Detached []string `json:"detached"`
  153. Skipped []string `json:"skipped"`
  154. Errors []string `json:"errors"`
  155. }
  156. // BulkDetach detaches the given existing clients (by email) from each target inbound.
  157. // (email, inbound) pairs where the client is not currently attached are silently skipped
  158. // at the inbound level; emails that aren't attached to any of the requested inbounds
  159. // are reported under skipped. ClientRecord rows are kept even when they become orphaned
  160. // (matches single-client detach semantics); callers should use bulkDelete for full removal.
  161. func (s *ClientService) BulkDetach(inboundSvc *InboundService, emails []string, inboundIds []int) (*BulkDetachResult, bool, error) {
  162. result := &BulkDetachResult{}
  163. if len(emails) == 0 || len(inboundIds) == 0 {
  164. return result, false, nil
  165. }
  166. recordErr := func(format string, args ...any) {
  167. msg := fmt.Sprintf(format, args...)
  168. result.Errors = append(result.Errors, msg)
  169. logger.Warningf("[BulkDetach] %s", msg)
  170. }
  171. requested := make(map[int]struct{}, len(inboundIds))
  172. for _, id := range inboundIds {
  173. requested[id] = struct{}{}
  174. }
  175. recsByInbound := make(map[int][]*model.ClientRecord)
  176. emailOrder := make([]string, 0, len(emails))
  177. emailRepr := make(map[string]string, len(emails))
  178. emailFailed := make(map[string]bool, len(emails))
  179. seenEmail := make(map[string]struct{}, len(emails))
  180. for _, email := range emails {
  181. if email == "" {
  182. continue
  183. }
  184. key := strings.ToLower(email)
  185. if _, ok := seenEmail[key]; ok {
  186. continue
  187. }
  188. seenEmail[key] = struct{}{}
  189. rec, err := s.GetRecordByEmail(nil, email)
  190. if err != nil {
  191. recordErr("%s: %v", email, err)
  192. continue
  193. }
  194. currentIds, err := s.GetInboundIdsForRecord(rec.Id)
  195. if err != nil {
  196. recordErr("%s: %v", email, err)
  197. continue
  198. }
  199. matched := false
  200. for _, id := range currentIds {
  201. if _, ok := requested[id]; ok {
  202. recsByInbound[id] = append(recsByInbound[id], rec)
  203. matched = true
  204. }
  205. }
  206. if !matched {
  207. result.Skipped = append(result.Skipped, rec.Email)
  208. continue
  209. }
  210. emailOrder = append(emailOrder, key)
  211. emailRepr[key] = rec.Email
  212. }
  213. needRestart := false
  214. // Ordered and de-duplicated up front: the sequential loop dropped each map
  215. // entry as it went, which the concurrent applies can no longer do.
  216. detachIds := make([]int, 0, len(recsByInbound))
  217. detachRecs := make([][]*model.ClientRecord, 0, len(recsByInbound))
  218. for _, ibId := range inboundIds {
  219. recs, ok := recsByInbound[ibId]
  220. if !ok {
  221. continue
  222. }
  223. delete(recsByInbound, ibId)
  224. detachIds = append(detachIds, ibId)
  225. detachRecs = append(detachRecs, recs)
  226. }
  227. detachResults, detachPanics := fanoutInboundResults(detachIds, inboundFanoutConcurrency, func(i int) inboundApplyOutcome {
  228. nr, err := s.delInboundClients(inboundSvc, detachIds[i], detachRecs[i], true)
  229. return inboundApplyOutcome{needRestart: nr, err: err}
  230. })
  231. for i, out := range detachResults {
  232. err := out.err
  233. if detachPanics[i] != nil {
  234. // See BulkAttach: a panicking apply may already have committed.
  235. needRestart = true
  236. err = detachPanics[i]
  237. }
  238. if err != nil {
  239. recordErr("inbound %d: %v", detachIds[i], err)
  240. for _, rec := range detachRecs[i] {
  241. emailFailed[strings.ToLower(rec.Email)] = true
  242. }
  243. continue
  244. }
  245. if out.needRestart {
  246. needRestart = true
  247. }
  248. }
  249. for _, key := range emailOrder {
  250. if emailFailed[key] {
  251. continue
  252. }
  253. result.Detached = append(result.Detached, emailRepr[key])
  254. }
  255. return result, needRestart, nil
  256. }
  257. // inboundApplyOutcome carries one inbound's apply result out of the fanout.
  258. type inboundApplyOutcome struct {
  259. needRestart bool
  260. err error
  261. }
  262. // BulkAdjustResult is returned by BulkAdjust to report how many clients were
  263. // successfully updated and which were skipped (typically because the field
  264. // being adjusted was unlimited for that client) or failed.
  265. type BulkAdjustResult struct {
  266. Adjusted int `json:"adjusted"`
  267. Skipped []BulkAdjustReport `json:"skipped,omitempty"`
  268. }
  269. type BulkAdjustReport struct {
  270. Email string `json:"email"`
  271. Reason string `json:"reason"`
  272. }
  273. type bulkAdjustEntry struct {
  274. record *model.ClientRecord
  275. applyExpiry bool
  276. newExpiry int64
  277. applyTotal bool
  278. newTotal int64
  279. }
  280. // bulkFlowClear is the directive that strips the XTLS flow from every selected
  281. // client. The vision values are the only positive flows xray accepts.
  282. const bulkFlowClear = "none"
  283. // bulkFlowAllowed whitelists the flow directives BulkAdjust accepts. Anything
  284. // outside this set is treated as "" (leave flow untouched) so a malformed or
  285. // hostile value can never be injected into a client's settings. The dropdown in
  286. // ClientBulkAdjustModal.tsx offers the same set ("" / "none" / TLS_FLOW_CONTROL);
  287. // keep the two in sync.
  288. var bulkFlowAllowed = map[string]struct{}{
  289. "": {},
  290. bulkFlowClear: {},
  291. "xtls-rprx-vision": {},
  292. "xtls-rprx-vision-udp443": {},
  293. }
  294. // BulkAdjust shifts ExpiryTime by addDays (days) and TotalGB by addBytes
  295. // for every email in the list. Clients whose corresponding field is
  296. // unlimited (0) are skipped — bulk extend should not accidentally
  297. // limit an unlimited client. addDays and addBytes may be negative.
  298. // flow sets the XTLS flow, limitHwid the max registered devices (0 = unlimited)
  299. // and adTag the MTProto sponsor channel; "none" clears flow or adTag.
  300. //
  301. // Like BulkDelete, the work is grouped by inbound so each inbound's
  302. // settings JSON is parsed and written exactly once regardless of how
  303. // many target emails it contains.
  304. func (s *ClientService) BulkAdjust(inboundSvc *InboundService, emails []string, addDays int, addBytes int64, flow string, limitHwid *int, adTag string) (BulkAdjustResult, bool, error) {
  305. result := BulkAdjustResult{}
  306. if len(emails) == 0 {
  307. return result, false, nil
  308. }
  309. flow = strings.TrimSpace(flow)
  310. if _, ok := bulkFlowAllowed[flow]; !ok {
  311. flow = "" // ignore unknown directives — "" means "leave flow untouched"
  312. }
  313. adTag = strings.TrimSpace(adTag)
  314. if adTag != "" && adTag != bulkFlowClear && !model.ValidMtprotoAdTag(adTag) {
  315. return result, false, common.NewError("mtproto client ad tag must be 32 hex characters")
  316. }
  317. if limitHwid != nil && *limitHwid < 0 {
  318. zero := 0
  319. limitHwid = &zero
  320. }
  321. adjustFlow := flow != ""
  322. adjustHwid := limitHwid != nil
  323. adjustAdTag := adTag != ""
  324. if addDays == 0 && addBytes == 0 && !adjustFlow && !adjustHwid && !adjustAdTag {
  325. return result, false, common.NewError("no adjustment specified")
  326. }
  327. addExpiryMs := int64(addDays) * 24 * 60 * 60 * 1000
  328. cleanEmails := trimmedUniqueEmails(emails)
  329. if len(cleanEmails) == 0 {
  330. return result, false, nil
  331. }
  332. db := database.GetDB()
  333. recordsByEmail, err := clientRecordsByEmail(db, cleanEmails)
  334. if err != nil {
  335. return result, false, err
  336. }
  337. skippedReasons := map[string]string{}
  338. for _, email := range cleanEmails {
  339. if _, ok := recordsByEmail[email]; !ok {
  340. skippedReasons[email] = "client not found"
  341. }
  342. }
  343. plan := map[string]*bulkAdjustEntry{}
  344. for email, rec := range recordsByEmail {
  345. entry := &bulkAdjustEntry{record: rec}
  346. if addDays != 0 {
  347. switch {
  348. case rec.ExpiryTime == 0:
  349. if _, exists := skippedReasons[email]; !exists {
  350. skippedReasons[email] = "unlimited expiry"
  351. }
  352. case rec.ExpiryTime > 0:
  353. next := rec.ExpiryTime + addExpiryMs
  354. if next <= 0 {
  355. if _, exists := skippedReasons[email]; !exists {
  356. skippedReasons[email] = "reduction exceeds remaining time"
  357. }
  358. } else {
  359. entry.applyExpiry = true
  360. entry.newExpiry = next
  361. }
  362. default:
  363. next := rec.ExpiryTime - addExpiryMs
  364. if next >= 0 {
  365. if _, exists := skippedReasons[email]; !exists {
  366. skippedReasons[email] = "reduction exceeds delay window"
  367. }
  368. } else {
  369. entry.applyExpiry = true
  370. entry.newExpiry = next
  371. }
  372. }
  373. }
  374. if addBytes != 0 {
  375. if rec.TotalGB == 0 {
  376. if _, exists := skippedReasons[email]; !exists {
  377. skippedReasons[email] = "unlimited traffic"
  378. }
  379. } else {
  380. next := rec.TotalGB + addBytes
  381. if next <= 0 {
  382. if _, exists := skippedReasons[email]; !exists {
  383. skippedReasons[email] = "reduction exceeds remaining quota"
  384. }
  385. } else {
  386. entry.applyTotal = true
  387. entry.newTotal = next
  388. }
  389. }
  390. }
  391. if entry.applyExpiry || entry.applyTotal || adjustFlow || adjustHwid || adjustAdTag {
  392. plan[email] = entry
  393. }
  394. }
  395. if len(plan) == 0 {
  396. for email, reason := range skippedReasons {
  397. result.Skipped = append(result.Skipped, BulkAdjustReport{Email: email, Reason: reason})
  398. }
  399. return result, false, nil
  400. }
  401. plannedIds := make([]int, 0, len(plan))
  402. recordIdToEmail := make(map[int]string, len(plan))
  403. for email, entry := range plan {
  404. if entry.applyExpiry || entry.applyTotal || adjustFlow || adjustAdTag {
  405. plannedIds = append(plannedIds, entry.record.Id)
  406. recordIdToEmail[entry.record.Id] = email
  407. }
  408. }
  409. var mappings []model.ClientInbound
  410. for _, batch := range chunkInts(plannedIds, sqlInChunk) {
  411. var rows []model.ClientInbound
  412. if err := db.Where("client_id IN ?", batch).Find(&rows).Error; err != nil {
  413. return result, false, err
  414. }
  415. mappings = append(mappings, rows...)
  416. }
  417. emailsByInbound := map[int][]string{}
  418. for _, m := range mappings {
  419. email, ok := recordIdToEmail[m.ClientId]
  420. if !ok {
  421. continue
  422. }
  423. emailsByInbound[m.InboundId] = append(emailsByInbound[m.InboundId], email)
  424. }
  425. needRestart := false
  426. flowHonored := map[string]bool{}
  427. flowIneligible := map[string]bool{}
  428. adTagHonored := map[string]bool{}
  429. adTagIneligible := map[string]bool{}
  430. execFailed := map[string]bool{}
  431. adjustIds := sortedInboundIds(emailsByInbound)
  432. adjustResults, adjustPanics := fanoutInboundResults(adjustIds, inboundFanoutConcurrency, func(i int) bulkInboundAdjustResult {
  433. return s.bulkAdjustInboundClients(inboundSvc, adjustIds[i], emailsByInbound[adjustIds[i]], plan, flow, adTag)
  434. })
  435. for i, ibRes := range adjustResults {
  436. if adjustPanics[i] != nil {
  437. needRestart = true
  438. for _, email := range emailsByInbound[adjustIds[i]] {
  439. execFailed[email] = true
  440. if _, already := skippedReasons[email]; !already {
  441. skippedReasons[email] = adjustPanics[i].Error()
  442. }
  443. }
  444. continue
  445. }
  446. if ibRes.needRestart {
  447. needRestart = true
  448. }
  449. for email := range ibRes.flowHonored {
  450. flowHonored[email] = true
  451. }
  452. for email := range ibRes.flowIneligible {
  453. flowIneligible[email] = true
  454. }
  455. for email := range ibRes.adTagHonored {
  456. adTagHonored[email] = true
  457. }
  458. for email := range ibRes.adTagIneligible {
  459. adTagIneligible[email] = true
  460. }
  461. for email, reason := range ibRes.perEmailSkipped {
  462. execFailed[email] = true
  463. if _, already := skippedReasons[email]; !already {
  464. skippedReasons[email] = reason
  465. }
  466. }
  467. }
  468. cond, condArgs := depletedCond(db)
  469. candidateEmails := make([]string, 0, len(plan))
  470. for email, entry := range plan {
  471. if entry.applyExpiry || entry.applyTotal {
  472. candidateEmails = append(candidateEmails, email)
  473. }
  474. }
  475. wasDisabledDepleted := map[string]struct{}{}
  476. for _, batch := range chunkStrings(candidateEmails, sqlInChunk) {
  477. var rows []string
  478. if err := db.Model(xray.ClientTraffic{}).
  479. Where(cond+" AND enable = ? AND email IN ?", append(append([]any{}, condArgs...), false, batch)...).
  480. Pluck("email", &rows).Error; err != nil {
  481. return result, needRestart, err
  482. }
  483. for _, e := range rows {
  484. wasDisabledDepleted[e] = struct{}{}
  485. }
  486. }
  487. wantAdTag := ""
  488. if adjustAdTag && adTag != bulkFlowClear {
  489. wantAdTag = strings.ToLower(adTag)
  490. }
  491. adjusted := map[string]struct{}{}
  492. for email, entry := range plan {
  493. if execFailed[email] {
  494. continue
  495. }
  496. updates := map[string]any{}
  497. if entry.applyExpiry {
  498. updates["expiry_time"] = entry.newExpiry
  499. }
  500. if entry.applyTotal {
  501. updates["total"] = entry.newTotal
  502. }
  503. if len(updates) > 0 {
  504. if err := db.Model(xray.ClientTraffic{}).Where("email = ?", email).Updates(updates).Error; err != nil {
  505. if _, already := skippedReasons[email]; !already {
  506. skippedReasons[email] = err.Error()
  507. }
  508. continue
  509. }
  510. }
  511. if adjustHwid {
  512. if err := s.setClientLimitHwidByEmail(email, *limitHwid); err != nil {
  513. if _, already := skippedReasons[email]; !already {
  514. skippedReasons[email] = err.Error()
  515. }
  516. continue
  517. }
  518. }
  519. if adjustAdTag && adTagHonored[email] {
  520. if err := db.Model(&model.ClientRecord{}).Where("email = ?", email).UpdateColumn("ad_tag", wantAdTag).Error; err != nil {
  521. if _, already := skippedReasons[email]; !already {
  522. skippedReasons[email] = err.Error()
  523. }
  524. continue
  525. }
  526. }
  527. // Counted when expiry/total changed, flow was honored, adTag was honored, or limitHwid was adjusted.
  528. if len(updates) > 0 || flowHonored[email] || adTagHonored[email] || adjustHwid {
  529. adjusted[email] = struct{}{}
  530. }
  531. }
  532. result.Adjusted = len(adjusted)
  533. for email, reason := range skippedReasons {
  534. result.Skipped = append(result.Skipped, BulkAdjustReport{Email: email, Reason: reason})
  535. }
  536. // Report a flow directive that no inbound could carry — only when it was not
  537. // honored anywhere and the client has no other (expiry/total) skip reason.
  538. // The expiry/total part, if any, has already been applied and counted above.
  539. for email := range flowIneligible {
  540. if flowHonored[email] {
  541. continue
  542. }
  543. if _, already := skippedReasons[email]; already {
  544. continue
  545. }
  546. result.Skipped = append(result.Skipped, BulkAdjustReport{Email: email, Reason: "flow not supported on inbound"})
  547. }
  548. for email := range adTagIneligible {
  549. if adTagHonored[email] {
  550. continue
  551. }
  552. if _, already := skippedReasons[email]; already {
  553. continue
  554. }
  555. result.Skipped = append(result.Skipped, BulkAdjustReport{Email: email, Reason: "adTag not supported on inbound"})
  556. }
  557. if len(wasDisabledDepleted) > 0 {
  558. stillDepleted := map[string]struct{}{}
  559. wasList := make([]string, 0, len(wasDisabledDepleted))
  560. for e := range wasDisabledDepleted {
  561. wasList = append(wasList, e)
  562. }
  563. for _, batch := range chunkStrings(wasList, sqlInChunk) {
  564. var rows []string
  565. if err := db.Model(xray.ClientTraffic{}).
  566. Where(cond+" AND email IN ?", append(append([]any{}, condArgs...), batch)...).
  567. Pluck("email", &rows).Error; err != nil {
  568. return result, needRestart, err
  569. }
  570. for _, e := range rows {
  571. stillDepleted[e] = struct{}{}
  572. }
  573. }
  574. reEnable := make([]string, 0, len(wasDisabledDepleted))
  575. for e := range wasDisabledDepleted {
  576. if _, still := stillDepleted[e]; !still {
  577. reEnable = append(reEnable, e)
  578. }
  579. }
  580. if len(reEnable) > 0 {
  581. _, nr, err := s.BulkSetEnable(inboundSvc, reEnable, true)
  582. if err != nil {
  583. return result, needRestart, err
  584. }
  585. if nr {
  586. needRestart = true
  587. }
  588. }
  589. }
  590. return result, needRestart, nil
  591. }
  592. type bulkInboundAdjustResult struct {
  593. perEmailSkipped map[string]string
  594. flowHonored map[string]bool
  595. // flowIneligible is tracked apart from perEmailSkipped: a flow directive
  596. // that an inbound cannot carry must not suppress the expiry/total write for
  597. // the same client (which would diverge the inbound JSON / ClientRecord from
  598. // ClientTraffic). It only feeds the final Skipped report.
  599. flowIneligible map[string]bool
  600. adTagHonored map[string]bool
  601. adTagIneligible map[string]bool
  602. needRestart bool
  603. }
  604. // bulkAdjustInboundClients applies expiry/total deltas to multiple clients
  605. // inside a single inbound's settings JSON. The xray runtime is updated
  606. // only for remote-node inbounds; local nodes do not need a notification
  607. // because the AddUser payload does not include totalGB/expiryTime —
  608. // changing those fields is identity-preserving and the panel's traffic
  609. // enforcement loop picks up the new limits from ClientTraffic directly.
  610. func (s *ClientService) bulkAdjustInboundClients(
  611. inboundSvc *InboundService,
  612. inboundId int,
  613. emails []string,
  614. plan map[string]*bulkAdjustEntry,
  615. flow string,
  616. adTag string,
  617. ) bulkInboundAdjustResult {
  618. res := bulkInboundAdjustResult{
  619. perEmailSkipped: map[string]string{},
  620. flowHonored: map[string]bool{},
  621. flowIneligible: map[string]bool{},
  622. adTagHonored: map[string]bool{},
  623. adTagIneligible: map[string]bool{},
  624. }
  625. defer lockInbound(inboundId).Unlock()
  626. oldInbound, err := inboundSvc.GetInbound(inboundId)
  627. if err != nil {
  628. logger.Error("Load Old Data Error")
  629. for _, e := range emails {
  630. res.perEmailSkipped[e] = err.Error()
  631. }
  632. return res
  633. }
  634. var settings map[string]any
  635. if err := json.Unmarshal([]byte(oldInbound.Settings), &settings); err != nil {
  636. for _, e := range emails {
  637. res.perEmailSkipped[e] = err.Error()
  638. }
  639. return res
  640. }
  641. // Match by email — the client's stable identity (see Delete). Credentials
  642. // can drift from the inbound JSON, so they are never used for matching.
  643. wantedEmails := make(map[string]struct{}, len(emails))
  644. for _, email := range emails {
  645. if plan[email] == nil {
  646. res.perEmailSkipped[email] = "client not found"
  647. continue
  648. }
  649. wantedEmails[email] = struct{}{}
  650. }
  651. // Flow eligibility is a property of the inbound (protocol + transport), so
  652. // resolve it once. Clearing flow is always allowed; setting a vision flow
  653. // is only honored on an inbound that can carry it.
  654. flowEligible := flow == bulkFlowClear ||
  655. (!oldInbound.DisableFlow &&
  656. inboundCanEnableTlsFlow(string(oldInbound.Protocol), oldInbound.StreamSettings, oldInbound.Settings))
  657. wantAdTag := ""
  658. if adTag != "" && adTag != bulkFlowClear {
  659. wantAdTag = strings.ToLower(adTag)
  660. }
  661. interfaceClients, _ := settings["clients"].([]any)
  662. foundEmails := map[string]bool{}
  663. flowChanged := false
  664. adTagChanged := false
  665. hasInboundChanges := false
  666. nowMs := time.Now().Unix() * 1000
  667. for i, client := range interfaceClients {
  668. c, ok := client.(map[string]any)
  669. if !ok {
  670. continue
  671. }
  672. targetEmail, _ := c["email"].(string)
  673. if _, want := wantedEmails[targetEmail]; !want || targetEmail == "" {
  674. continue
  675. }
  676. clientChanged := false
  677. entry := plan[targetEmail]
  678. if entry.applyExpiry {
  679. c["expiryTime"] = entry.newExpiry
  680. clientChanged = true
  681. }
  682. if entry.applyTotal {
  683. c["totalGB"] = entry.newTotal
  684. clientChanged = true
  685. }
  686. if flow != "" {
  687. if flowEligible {
  688. want := ""
  689. if flow != bulkFlowClear {
  690. want = flow
  691. }
  692. if cur, _ := c["flow"].(string); cur != want {
  693. c["flow"] = want
  694. flowChanged = true
  695. }
  696. res.flowHonored[targetEmail] = true
  697. clientChanged = true
  698. } else {
  699. // Record separately so this never suppresses the expiry/total
  700. // write for the same client (see flowIneligible doc).
  701. res.flowIneligible[targetEmail] = true
  702. }
  703. }
  704. if adTag != "" {
  705. if oldInbound.Protocol == model.MTProto {
  706. if cur, _ := c["adTag"].(string); cur != wantAdTag {
  707. c["adTag"] = wantAdTag
  708. adTagChanged = true
  709. }
  710. res.adTagHonored[targetEmail] = true
  711. clientChanged = true
  712. } else {
  713. res.adTagIneligible[targetEmail] = true
  714. }
  715. }
  716. if clientChanged {
  717. c["updated_at"] = nowMs
  718. hasInboundChanges = true
  719. }
  720. interfaceClients[i] = c
  721. foundEmails[targetEmail] = true
  722. }
  723. for email := range wantedEmails {
  724. if !foundEmails[email] {
  725. res.perEmailSkipped[email] = "Client Not Found In Inbound"
  726. }
  727. }
  728. if len(foundEmails) == 0 || !hasInboundChanges {
  729. return res
  730. }
  731. settings["clients"] = interfaceClients
  732. newSettings, err := json.MarshalIndent(settings, "", " ")
  733. if err != nil {
  734. for email := range foundEmails {
  735. res.perEmailSkipped[email] = err.Error()
  736. }
  737. return res
  738. }
  739. prevSettings := oldInbound.Settings
  740. oldInbound.Settings = string(newSettings)
  741. // A flow change rewrites the user's xray config, which the lightweight
  742. // UpdateUser push below does not carry. Local nodes reload via restart;
  743. // remote nodes get a full reconcile (MarkNodeDirty) instead of a per-user push.
  744. if flowChanged && oldInbound.NodeID == nil {
  745. res.needRestart = true
  746. }
  747. // Serialize against the traffic poll to avoid the cross-transaction
  748. // lock-order deadlock on inbounds/client_records (runSerializedTx).
  749. txErr := runSerializedTx(func(tx *gorm.DB) error {
  750. if err := commitInboundClientSettings(tx, oldInbound, prevSettings); err != nil {
  751. return err
  752. }
  753. finalClients, gcErr := inboundSvc.GetClients(oldInbound)
  754. if gcErr != nil {
  755. return gcErr
  756. }
  757. if err := s.SyncInbound(tx, inboundId, finalClients); err != nil {
  758. return err
  759. }
  760. if oldInbound.NodeID != nil {
  761. return (&NodeService{}).MarkNodeDirtyTx(tx, *oldInbound.NodeID)
  762. }
  763. return nil
  764. })
  765. if txErr != nil {
  766. for email := range foundEmails {
  767. if _, skip := res.perEmailSkipped[email]; !skip {
  768. res.perEmailSkipped[email] = txErr.Error()
  769. }
  770. }
  771. } else {
  772. if adTagChanged && oldInbound.Protocol == model.MTProto && oldInbound.NodeID == nil {
  773. inboundSvc.applyLocalMtproto(oldInbound.Id)
  774. }
  775. if oldInbound.NodeID != nil && !flowChanged && len(foundEmails) <= nodeBulkPushThreshold {
  776. rt, push, _, perr := inboundSvc.nodePushPlan(oldInbound)
  777. if perr != nil {
  778. logger.Warning("BulkAdjust: node runtime lookup after commit failed:", perr)
  779. } else if push {
  780. for email := range foundEmails {
  781. entry := plan[email]
  782. updated := *entry.record.ToClient()
  783. if entry.applyExpiry {
  784. updated.ExpiryTime = entry.newExpiry
  785. }
  786. if entry.applyTotal {
  787. updated.TotalGB = entry.newTotal
  788. }
  789. if adTag != "" && oldInbound.Protocol == model.MTProto {
  790. updated.AdTag = wantAdTag
  791. }
  792. updated.UpdatedAt = nowMs
  793. ctx, cancel := nodePushContext()
  794. err1 := rt.UpdateUser(ctx, oldInbound, email, updated)
  795. cancel()
  796. if err1 != nil {
  797. logger.Warning("Error in updating client on", rt.Name(), ":", err1)
  798. // First failure ends the batch push; the reconcile converges the rest.
  799. break
  800. }
  801. }
  802. }
  803. }
  804. }
  805. return res
  806. }
  807. // BulkDeleteResult mirrors BulkAdjustResult: total deleted plus per-email
  808. // skip reasons when an email could not be processed.
  809. type BulkDeleteResult struct {
  810. Deleted int `json:"deleted"`
  811. Skipped []BulkDeleteReport `json:"skipped,omitempty"`
  812. }
  813. type BulkDeleteReport struct {
  814. Email string `json:"email"`
  815. Reason string `json:"reason"`
  816. }
  817. // BulkDelete removes every client in the list in one optimized pass.
  818. // Instead of running the full single-delete pipeline N times (which would
  819. // re-read, re-parse, and re-write each inbound's settings JSON for every
  820. // email), it groups emails by inbound and performs a single
  821. // read-modify-write per inbound. Per-row DB cleanups are also batched with
  822. // IN-clause queries at the end. Errors on a particular email are recorded
  823. // in the Skipped list and processing continues for the rest.
  824. func (s *ClientService) BulkDelete(inboundSvc *InboundService, emails []string, keepTraffic bool) (BulkDeleteResult, bool, error) {
  825. result := BulkDeleteResult{}
  826. cleanEmails := trimmedUniqueEmails(emails)
  827. if len(cleanEmails) == 0 {
  828. return result, false, nil
  829. }
  830. db := database.GetDB()
  831. recordsByEmail, err := clientRecordsByEmail(db, cleanEmails)
  832. if err != nil {
  833. return result, false, err
  834. }
  835. tombstoneEmails := make([]string, 0, len(recordsByEmail))
  836. for _, email := range cleanEmails {
  837. if recordsByEmail[email] != nil {
  838. tombstoneEmails = append(tombstoneEmails, email)
  839. }
  840. }
  841. tombstoneClientEmails(tombstoneEmails)
  842. skippedReasons := map[string]string{}
  843. for _, email := range cleanEmails {
  844. if _, ok := recordsByEmail[email]; !ok {
  845. skippedReasons[email] = "client not found"
  846. }
  847. }
  848. clientIds := make([]int, 0, len(recordsByEmail))
  849. recordIdToEmail := make(map[int]string, len(recordsByEmail))
  850. for _, r := range recordsByEmail {
  851. clientIds = append(clientIds, r.Id)
  852. recordIdToEmail[r.Id] = r.Email
  853. }
  854. emailsByInbound := map[int][]string{}
  855. if len(clientIds) > 0 {
  856. var mappings []model.ClientInbound
  857. for _, batch := range chunkInts(clientIds, sqlInChunk) {
  858. var rows []model.ClientInbound
  859. if err := db.Where("client_id IN ?", batch).Find(&rows).Error; err != nil {
  860. return result, false, err
  861. }
  862. mappings = append(mappings, rows...)
  863. }
  864. for _, m := range mappings {
  865. email, ok := recordIdToEmail[m.ClientId]
  866. if !ok {
  867. continue
  868. }
  869. emailsByInbound[m.InboundId] = append(emailsByInbound[m.InboundId], email)
  870. }
  871. }
  872. needRestart := false
  873. delIds := sortedInboundIds(emailsByInbound)
  874. delResults, delPanics := fanoutInboundResults(delIds, inboundFanoutConcurrency, func(i int) bulkInboundDeleteResult {
  875. return s.bulkDelInboundClients(inboundSvc, delIds[i], emailsByInbound[delIds[i]], recordsByEmail, keepTraffic)
  876. })
  877. for i, ibResult := range delResults {
  878. if delPanics[i] != nil {
  879. needRestart = true
  880. for _, email := range emailsByInbound[delIds[i]] {
  881. if _, already := skippedReasons[email]; !already {
  882. skippedReasons[email] = delPanics[i].Error()
  883. }
  884. }
  885. continue
  886. }
  887. if ibResult.needRestart {
  888. needRestart = true
  889. }
  890. for email, reason := range ibResult.perEmailSkipped {
  891. if _, already := skippedReasons[email]; !already {
  892. skippedReasons[email] = reason
  893. }
  894. }
  895. }
  896. successEmails := make([]string, 0, len(recordsByEmail))
  897. successIds := make([]int, 0, len(recordsByEmail))
  898. failedEmails := make([]string, 0, len(recordsByEmail))
  899. successSubIDs := make([]string, 0, len(recordsByEmail))
  900. for email, rec := range recordsByEmail {
  901. if _, skipped := skippedReasons[email]; skipped {
  902. failedEmails = append(failedEmails, email)
  903. continue
  904. }
  905. successEmails = append(successEmails, email)
  906. successIds = append(successIds, rec.Id)
  907. successSubIDs = append(successSubIDs, rec.SubID)
  908. }
  909. withdrawClientTombstones(failedEmails...)
  910. if len(successIds) > 0 {
  911. // Serialize the row cleanup against the traffic poll to avoid the
  912. // cross-transaction lock-order deadlock on client_traffics/inbounds.
  913. if err := runSerializedTx(func(tx *gorm.DB) error {
  914. if e := adjustGroupBaselinesForRemovedTraffic(tx, successEmails); e != nil {
  915. return e
  916. }
  917. if e := clearClientHwidsBySubIDTx(tx, successSubIDs...); e != nil {
  918. return e
  919. }
  920. for _, batch := range chunkInts(successIds, sqlInChunk) {
  921. if e := tx.Where("client_id IN ?", batch).Delete(&model.ClientInbound{}).Error; e != nil {
  922. return e
  923. }
  924. if e := tx.Where("client_id IN ?", batch).Delete(&model.ClientExternalLink{}).Error; e != nil {
  925. return e
  926. }
  927. }
  928. if !keepTraffic && len(successEmails) > 0 {
  929. for _, batch := range chunkStrings(successEmails, sqlInChunk) {
  930. if e := tx.Where("email IN ?", batch).Delete(&xray.ClientTraffic{}).Error; e != nil {
  931. return e
  932. }
  933. if e := tx.Where("client_email IN ?", batch).Delete(&model.InboundClientIps{}).Error; e != nil {
  934. return e
  935. }
  936. }
  937. }
  938. for _, batch := range chunkInts(successIds, sqlInChunk) {
  939. if e := tx.Where("id IN ?", batch).Delete(&model.ClientRecord{}).Error; e != nil {
  940. return e
  941. }
  942. }
  943. return nil
  944. }); err != nil {
  945. withdrawClientTombstones(successEmails...)
  946. return result, needRestart, err
  947. }
  948. }
  949. result.Deleted = len(successEmails)
  950. for email, reason := range skippedReasons {
  951. result.Skipped = append(result.Skipped, BulkDeleteReport{Email: email, Reason: reason})
  952. }
  953. return result, needRestart, nil
  954. }
  955. type bulkInboundDeleteResult struct {
  956. perEmailSkipped map[string]string
  957. needRestart bool
  958. }
  959. // bulkDelInboundClients removes multiple clients from a single inbound's
  960. // settings JSON in one read-modify-write cycle, runs the xray runtime
  961. // RemoveUser/DeleteUser calls, and persists the inbound. The returned map
  962. // holds per-email failure reasons; emails not present in the map are
  963. // considered successful for this inbound.
  964. func (s *ClientService) bulkDelInboundClients(
  965. inboundSvc *InboundService,
  966. inboundId int,
  967. emails []string,
  968. records map[string]*model.ClientRecord,
  969. keepTraffic bool,
  970. ) bulkInboundDeleteResult {
  971. res := bulkInboundDeleteResult{perEmailSkipped: map[string]string{}}
  972. defer lockInbound(inboundId).Unlock()
  973. oldInbound, err := inboundSvc.GetInbound(inboundId)
  974. if err != nil {
  975. logger.Error("Load Old Data Error")
  976. for _, e := range emails {
  977. res.perEmailSkipped[e] = err.Error()
  978. }
  979. return res
  980. }
  981. var settings map[string]any
  982. if err := json.Unmarshal([]byte(oldInbound.Settings), &settings); err != nil {
  983. for _, e := range emails {
  984. res.perEmailSkipped[e] = err.Error()
  985. }
  986. return res
  987. }
  988. // Match by email — the client's stable identity (see Delete). The link-derived
  989. // set is deletion intent: an email already absent from settings is successful,
  990. // while foundEmails tracks entries that still need settings-specific cleanup.
  991. wantedEmails := make(map[string]struct{}, len(emails))
  992. for _, email := range emails {
  993. if records[email] == nil {
  994. res.perEmailSkipped[email] = "client not found"
  995. continue
  996. }
  997. wantedEmails[email] = struct{}{}
  998. }
  999. interfaceClients, _ := settings["clients"].([]any)
  1000. newClients := make([]any, 0, len(interfaceClients))
  1001. foundEmails := map[string]bool{}
  1002. enableByEmail := map[string]bool{}
  1003. for _, client := range interfaceClients {
  1004. c, ok := client.(map[string]any)
  1005. if !ok {
  1006. newClients = append(newClients, client)
  1007. continue
  1008. }
  1009. em, _ := c["email"].(string)
  1010. if _, found := wantedEmails[em]; found && em != "" {
  1011. foundEmails[em] = true
  1012. en, _ := c["enable"].(bool)
  1013. enableByEmail[em] = en
  1014. continue
  1015. }
  1016. newClients = append(newClients, client)
  1017. }
  1018. db := database.GetDB()
  1019. newClients = compactOrphans(db, newClients)
  1020. if newClients == nil {
  1021. newClients = []any{}
  1022. }
  1023. settings["clients"] = newClients
  1024. newSettings, err := json.MarshalIndent(settings, "", " ")
  1025. if err != nil {
  1026. for email := range wantedEmails {
  1027. if _, skip := res.perEmailSkipped[email]; !skip {
  1028. res.perEmailSkipped[email] = err.Error()
  1029. }
  1030. }
  1031. return res
  1032. }
  1033. prevSettings := oldInbound.Settings
  1034. oldInbound.Settings = string(newSettings)
  1035. foundList := make([]string, 0, len(foundEmails))
  1036. for email := range foundEmails {
  1037. foundList = append(foundList, email)
  1038. }
  1039. notDepletedByEmail := map[string]bool{}
  1040. if len(foundList) > 0 {
  1041. type trafficRow struct {
  1042. Email string
  1043. Enable bool
  1044. }
  1045. for _, batch := range chunkStrings(foundList, sqlInChunk) {
  1046. var rows []trafficRow
  1047. if err := db.Model(xray.ClientTraffic{}).
  1048. Where("email IN ?", batch).
  1049. Select("email, enable").
  1050. Scan(&rows).Error; err == nil {
  1051. for _, r := range rows {
  1052. notDepletedByEmail[r.Email] = r.Enable
  1053. }
  1054. }
  1055. }
  1056. }
  1057. var sharedSet map[string]bool
  1058. if !keepTraffic {
  1059. var sharedErr error
  1060. sharedSet, sharedErr = inboundSvc.emailsUsedByOtherInbounds(foundList, inboundId)
  1061. if sharedErr != nil {
  1062. for email := range wantedEmails {
  1063. res.perEmailSkipped[email] = sharedErr.Error()
  1064. }
  1065. return res
  1066. }
  1067. }
  1068. if !keepTraffic {
  1069. purge := make([]string, 0, len(foundEmails))
  1070. for email := range foundEmails {
  1071. if !sharedSet[strings.ToLower(strings.TrimSpace(email))] {
  1072. purge = append(purge, email)
  1073. }
  1074. }
  1075. if len(purge) > 0 {
  1076. // Serialize the IP/stat purge against the traffic poll to avoid the
  1077. // cross-transaction lock-order deadlock on client_traffics.
  1078. if delErr := runSerializedTx(func(tx *gorm.DB) error {
  1079. if e := inboundSvc.delClientIPsByEmails(tx, purge); e != nil {
  1080. logger.Error("Error in delete client IPs")
  1081. return e
  1082. }
  1083. if e := inboundSvc.delClientStatsByEmails(tx, purge); e != nil {
  1084. logger.Error("Delete stats Data Error")
  1085. return e
  1086. }
  1087. return nil
  1088. }); delErr != nil {
  1089. for _, email := range purge {
  1090. res.perEmailSkipped[email] = delErr.Error()
  1091. delete(foundEmails, email)
  1092. }
  1093. }
  1094. }
  1095. }
  1096. // Serialize against the traffic poll to avoid the cross-transaction
  1097. // lock-order deadlock on inbounds/client_records (runSerializedTx).
  1098. txErr := runSerializedTx(func(tx *gorm.DB) error {
  1099. if err := commitInboundClientSettings(tx, oldInbound, prevSettings); err != nil {
  1100. return err
  1101. }
  1102. finalClients, err := inboundSvc.GetClients(oldInbound)
  1103. if err != nil {
  1104. return err
  1105. }
  1106. if err := s.SyncInbound(tx, inboundId, finalClients); err != nil {
  1107. return err
  1108. }
  1109. if oldInbound.NodeID != nil {
  1110. return (&NodeService{}).MarkNodeDirtyTx(tx, *oldInbound.NodeID)
  1111. }
  1112. return nil
  1113. })
  1114. if txErr != nil {
  1115. for email := range wantedEmails {
  1116. if _, skip := res.perEmailSkipped[email]; !skip {
  1117. res.perEmailSkipped[email] = txErr.Error()
  1118. }
  1119. }
  1120. } else if oldInbound.NodeID == nil {
  1121. rt, rterr := inboundSvc.runtimeFor(oldInbound)
  1122. if rterr != nil {
  1123. res.needRestart = true
  1124. } else {
  1125. for email := range foundEmails {
  1126. if !enableByEmail[email] || !notDepletedByEmail[email] {
  1127. continue
  1128. }
  1129. err1 := rt.RemoveUser(context.Background(), oldInbound, email)
  1130. if err1 == nil {
  1131. logger.Debug("Client deleted on", rt.Name(), ":", email)
  1132. } else if strings.Contains(err1.Error(), fmt.Sprintf("User %s not found.", email)) {
  1133. logger.Debug("User is already deleted. Nothing to do more...")
  1134. } else {
  1135. logger.Debug("Error in deleting client on", rt.Name(), ":", err1)
  1136. res.needRestart = true
  1137. }
  1138. }
  1139. }
  1140. } else {
  1141. dispatchEmails := make([]string, 0, len(wantedEmails))
  1142. for email := range wantedEmails {
  1143. if _, skip := res.perEmailSkipped[email]; !skip {
  1144. dispatchEmails = append(dispatchEmails, email)
  1145. }
  1146. }
  1147. if len(dispatchEmails) > nodeBulkPushThreshold {
  1148. return res
  1149. }
  1150. rt, push, _, perr := inboundSvc.nodePushPlan(oldInbound)
  1151. if perr != nil {
  1152. logger.Warning("BulkDelete: node runtime lookup after commit failed:", perr)
  1153. } else if push {
  1154. for _, email := range dispatchEmails {
  1155. ctx, cancel := nodePushContext()
  1156. err1 := rt.DeleteClient(ctx, email)
  1157. cancel()
  1158. if err1 != nil {
  1159. logger.Warning("Error in deleting client on", rt.Name(), ":", err1)
  1160. // The node is already dirty, so one reconcile converges the rest of
  1161. // the batch instead of paying another deadline per client.
  1162. break
  1163. }
  1164. }
  1165. }
  1166. }
  1167. return res
  1168. }
  1169. // BulkCreateResult mirrors BulkAdjustResult for the create flow.
  1170. type BulkCreateResult struct {
  1171. Created int `json:"created"`
  1172. Skipped []BulkCreateReport `json:"skipped,omitempty"`
  1173. }
  1174. type BulkCreateReport struct {
  1175. Email string `json:"email"`
  1176. Reason string `json:"reason"`
  1177. }
  1178. func (s *ClientService) BulkCreate(inboundSvc *InboundService, payloads []ClientCreatePayload) (BulkCreateResult, bool, error) {
  1179. result, _, needRestart, err := s.bulkCreate(inboundSvc, payloads)
  1180. return result, needRestart, err
  1181. }
  1182. // bulkCreate also returns the payload indexes that inserted a new client record;
  1183. // a Created payload whose email already existed only reused that client.
  1184. func (s *ClientService) bulkCreate(inboundSvc *InboundService, payloads []ClientCreatePayload) (BulkCreateResult, []int, bool, error) {
  1185. result := BulkCreateResult{}
  1186. if len(payloads) == 0 {
  1187. return result, nil, false, nil
  1188. }
  1189. skip := func(email, reason string) {
  1190. if strings.TrimSpace(email) == "" {
  1191. email = "(missing email)"
  1192. }
  1193. result.Skipped = append(result.Skipped, BulkCreateReport{Email: email, Reason: reason})
  1194. }
  1195. type prepared struct {
  1196. client model.Client
  1197. inboundIds []int
  1198. limitHwid int
  1199. payloadIdx int
  1200. reused bool
  1201. }
  1202. prep := make([]prepared, 0, len(payloads))
  1203. emails := make([]string, 0, len(payloads))
  1204. subIDs := make([]string, 0, len(payloads))
  1205. seenEmail := make(map[string]struct{}, len(payloads))
  1206. seenSubID := make(map[string]string, len(payloads))
  1207. for i := range payloads {
  1208. client := payloads[i].Client
  1209. email := strings.TrimSpace(client.Email)
  1210. if email == "" {
  1211. skip("", "client email is required")
  1212. continue
  1213. }
  1214. if verr := validateClientEmail(email); verr != nil {
  1215. skip(email, verr.Error())
  1216. continue
  1217. }
  1218. if verr := validateClientSubID(client.SubID); verr != nil {
  1219. skip(email, verr.Error())
  1220. continue
  1221. }
  1222. if verr := validateClientRenewal(client); verr != nil {
  1223. skip(email, verr.Error())
  1224. continue
  1225. }
  1226. if verr := validateClientResetMax(client.ResetMax); verr != nil {
  1227. skip(email, verr.Error())
  1228. continue
  1229. }
  1230. if verr := validateClientTrafficReset(client.TrafficReset, client.TrafficResetDay); verr != nil {
  1231. skip(email, verr.Error())
  1232. continue
  1233. }
  1234. if len(payloads[i].InboundIds) == 0 {
  1235. skip(email, "at least one inbound is required")
  1236. continue
  1237. }
  1238. client.Email = email
  1239. if client.SubID == "" {
  1240. client.SubID = uuid.NewString()
  1241. }
  1242. // Preserve enable (omit→true in UnmarshalJSON; explicit false kept) (#6478).
  1243. now := time.Now().UnixMilli()
  1244. if client.CreatedAt == 0 {
  1245. client.CreatedAt = now
  1246. }
  1247. client.UpdatedAt = now
  1248. le := strings.ToLower(email)
  1249. if _, dup := seenEmail[le]; dup {
  1250. skip(email, "email already in use: "+email)
  1251. continue
  1252. }
  1253. if owner, ok := seenSubID[client.SubID]; ok && owner != le {
  1254. skip(email, "subId already in use: "+client.SubID)
  1255. continue
  1256. }
  1257. seenEmail[le] = struct{}{}
  1258. seenSubID[client.SubID] = le
  1259. prep = append(prep, prepared{client: client, inboundIds: payloads[i].InboundIds, limitHwid: payloads[i].LimitHwid, payloadIdx: i})
  1260. emails = append(emails, email)
  1261. subIDs = append(subIDs, client.SubID)
  1262. }
  1263. if len(prep) == 0 {
  1264. return result, nil, false, nil
  1265. }
  1266. db := database.GetDB()
  1267. const lookupChunk = 400
  1268. existingByEmail := make(map[string]model.ClientRecord, len(emails))
  1269. for start := 0; start < len(emails); start += lookupChunk {
  1270. end := min(start+lookupChunk, len(emails))
  1271. var rows []model.ClientRecord
  1272. if e := db.Where("email IN ?", emails[start:end]).Find(&rows).Error; e != nil {
  1273. return result, nil, false, e
  1274. }
  1275. for i := range rows {
  1276. existingByEmail[strings.ToLower(rows[i].Email)] = rows[i]
  1277. }
  1278. }
  1279. existingSubOwner := make(map[string]string, len(subIDs))
  1280. for start := 0; start < len(subIDs); start += lookupChunk {
  1281. end := min(start+lookupChunk, len(subIDs))
  1282. var rows []model.ClientRecord
  1283. if e := db.Where("sub_id IN ?", subIDs[start:end]).Find(&rows).Error; e != nil {
  1284. return result, nil, false, e
  1285. }
  1286. for i := range rows {
  1287. existingSubOwner[rows[i].SubID] = strings.ToLower(rows[i].Email)
  1288. }
  1289. }
  1290. inboundCache := make(map[int]*model.Inbound)
  1291. getIb := func(id int) (*model.Inbound, error) {
  1292. if ib, ok := inboundCache[id]; ok {
  1293. return ib, nil
  1294. }
  1295. ib, e := inboundSvc.GetInbound(id)
  1296. if e != nil {
  1297. return nil, e
  1298. }
  1299. inboundCache[id] = ib
  1300. return ib, nil
  1301. }
  1302. byInbound := make(map[int][]model.Client)
  1303. idxByInbound := make(map[int][]int)
  1304. inboundOrder := make([]int, 0)
  1305. failed := make([]bool, len(prep))
  1306. reason := make([]string, len(prep))
  1307. createAnyTunnel := false
  1308. for idx := range prep {
  1309. le := strings.ToLower(prep[idx].client.Email)
  1310. if rec, ok := existingByEmail[le]; ok {
  1311. if rec.SubID != prep[idx].client.SubID {
  1312. failed[idx] = true
  1313. reason[idx] = "email already in use: " + prep[idx].client.Email
  1314. continue
  1315. }
  1316. prep[idx].reused = true
  1317. if prep[idx].client.ID == "" {
  1318. prep[idx].client.ID = rec.UUID
  1319. }
  1320. if prep[idx].client.Password == "" {
  1321. prep[idx].client.Password = rec.Password
  1322. }
  1323. if prep[idx].client.Auth == "" {
  1324. prep[idx].client.Auth = rec.Auth
  1325. }
  1326. if prep[idx].client.Secret == "" {
  1327. prep[idx].client.Secret = rec.Secret
  1328. }
  1329. }
  1330. if owner, ok := existingSubOwner[prep[idx].client.SubID]; ok && owner != le {
  1331. failed[idx] = true
  1332. reason[idx] = "subId already in use: " + prep[idx].client.SubID
  1333. continue
  1334. }
  1335. ok := true
  1336. for _, ibId := range prep[idx].inboundIds {
  1337. ib, e := getIb(ibId)
  1338. if e != nil {
  1339. failed[idx] = true
  1340. reason[idx] = e.Error()
  1341. ok = false
  1342. break
  1343. }
  1344. if ib.Protocol == model.WireGuard || ib.Protocol == model.AmneziaWG {
  1345. createAnyTunnel = true
  1346. }
  1347. if e := s.fillProtocolDefaults(&prep[idx].client, ib); e != nil {
  1348. failed[idx] = true
  1349. reason[idx] = e.Error()
  1350. ok = false
  1351. break
  1352. }
  1353. }
  1354. if !ok {
  1355. continue
  1356. }
  1357. for _, ibId := range prep[idx].inboundIds {
  1358. ib, _ := getIb(ibId)
  1359. if _, seen := byInbound[ibId]; !seen {
  1360. inboundOrder = append(inboundOrder, ibId)
  1361. }
  1362. byInbound[ibId] = append(byInbound[ibId], clientWithInboundFlow(prep[idx].client, ib))
  1363. idxByInbound[ibId] = append(idxByInbound[ibId], idx)
  1364. }
  1365. }
  1366. needRestart := false
  1367. createResults, createPanics := fanoutInboundResults(inboundOrder, addFanoutLimit(createAnyTunnel), func(i int) inboundApplyOutcome {
  1368. ibId := inboundOrder[i]
  1369. payload, e := json.Marshal(map[string][]model.Client{"clients": byInbound[ibId]})
  1370. if e != nil {
  1371. return inboundApplyOutcome{err: e}
  1372. }
  1373. nr, e := s.AddInboundClient(inboundSvc, &model.Inbound{Id: ibId, Settings: string(payload)})
  1374. return inboundApplyOutcome{needRestart: nr, err: e}
  1375. })
  1376. for i, out := range createResults {
  1377. e := out.err
  1378. if createPanics[i] != nil {
  1379. // See BulkAttach: a panicking apply may already have committed.
  1380. needRestart = true
  1381. e = createPanics[i]
  1382. }
  1383. if e != nil {
  1384. for _, idx := range idxByInbound[inboundOrder[i]] {
  1385. failed[idx] = true
  1386. if reason[idx] == "" {
  1387. reason[idx] = e.Error()
  1388. }
  1389. }
  1390. continue
  1391. }
  1392. if out.needRestart {
  1393. needRestart = true
  1394. }
  1395. }
  1396. inserted := make([]int, 0, len(prep))
  1397. for idx := range prep {
  1398. if failed[idx] {
  1399. skip(prep[idx].client.Email, reason[idx])
  1400. continue
  1401. }
  1402. // The client is already live after fanout; never leave a stale delete
  1403. // tombstone merely because applying its optional HWID limit failed.
  1404. withdrawClientTombstones(prep[idx].client.Email)
  1405. if err := s.setClientLimitHwidByEmail(prep[idx].client.Email, prep[idx].limitHwid); err != nil {
  1406. skip(prep[idx].client.Email, err.Error())
  1407. continue
  1408. }
  1409. result.Created++
  1410. if !prep[idx].reused {
  1411. inserted = append(inserted, prep[idx].payloadIdx)
  1412. }
  1413. }
  1414. // A re-created email is a live identity again: a delete tombstone left
  1415. // standing makes the next node merge prune the new client's inbound links.
  1416. return result, inserted, needRestart, nil
  1417. }
  1418. func (s *ClientService) DelDepleted(inboundSvc *InboundService) (int, bool, error) {
  1419. db := database.GetDB()
  1420. now := time.Now().UnixMilli()
  1421. depletedClause := depletedClientsClause
  1422. var rows []xray.ClientTraffic
  1423. if err := db.Where(depletedClause, now).Find(&rows).Error; err != nil {
  1424. return 0, false, err
  1425. }
  1426. if len(rows) == 0 {
  1427. return 0, false, nil
  1428. }
  1429. seen := make(map[string]struct{}, len(rows))
  1430. emails := make([]string, 0, len(rows))
  1431. for _, r := range rows {
  1432. if r.Email == "" {
  1433. continue
  1434. }
  1435. if _, ok := seen[r.Email]; ok {
  1436. continue
  1437. }
  1438. seen[r.Email] = struct{}{}
  1439. emails = append(emails, r.Email)
  1440. }
  1441. if len(emails) == 0 {
  1442. return 0, false, nil
  1443. }
  1444. res, needRestart, err := s.BulkDelete(inboundSvc, emails, false)
  1445. if err != nil {
  1446. return res.Deleted, needRestart, err
  1447. }
  1448. return res.Deleted, needRestart, nil
  1449. }
  1450. type BulkSetEnableResult struct {
  1451. Changed int `json:"changed"`
  1452. Skipped []BulkSetEnableReport `json:"skipped,omitempty"`
  1453. }
  1454. type BulkSetEnableReport struct {
  1455. Email string `json:"email"`
  1456. Reason string `json:"reason"`
  1457. }
  1458. func (s *ClientService) BulkSetEnable(inboundSvc *InboundService, emails []string, enable bool) (BulkSetEnableResult, bool, error) {
  1459. result := BulkSetEnableResult{}
  1460. cleanEmails := trimmedUniqueEmails(emails)
  1461. if len(cleanEmails) == 0 {
  1462. return result, false, nil
  1463. }
  1464. db := database.GetDB()
  1465. recordsByEmail, err := clientRecordsByEmail(db, cleanEmails)
  1466. if err != nil {
  1467. return result, false, err
  1468. }
  1469. skippedReasons := map[string]string{}
  1470. for _, email := range cleanEmails {
  1471. if _, ok := recordsByEmail[email]; !ok {
  1472. skippedReasons[email] = "client not found"
  1473. }
  1474. }
  1475. clientIds := make([]int, 0, len(recordsByEmail))
  1476. recordIdToEmail := make(map[int]string, len(recordsByEmail))
  1477. for _, r := range recordsByEmail {
  1478. clientIds = append(clientIds, r.Id)
  1479. recordIdToEmail[r.Id] = r.Email
  1480. }
  1481. emailsByInbound := map[int][]string{}
  1482. if len(clientIds) > 0 {
  1483. var mappings []model.ClientInbound
  1484. for _, batch := range chunkInts(clientIds, sqlInChunk) {
  1485. var rows []model.ClientInbound
  1486. if err := db.Where("client_id IN ?", batch).Find(&rows).Error; err != nil {
  1487. return result, false, err
  1488. }
  1489. mappings = append(mappings, rows...)
  1490. }
  1491. for _, m := range mappings {
  1492. email, ok := recordIdToEmail[m.ClientId]
  1493. if !ok {
  1494. continue
  1495. }
  1496. emailsByInbound[m.InboundId] = append(emailsByInbound[m.InboundId], email)
  1497. }
  1498. }
  1499. needRestart := false
  1500. enableIds := sortedInboundIds(emailsByInbound)
  1501. enableResults, enablePanics := fanoutInboundResults(enableIds, inboundFanoutConcurrency, func(i int) bulkSetEnableInboundResult {
  1502. return s.bulkSetEnableInboundClients(inboundSvc, enableIds[i], emailsByInbound[enableIds[i]], enable)
  1503. })
  1504. for i, ibRes := range enableResults {
  1505. if enablePanics[i] != nil {
  1506. needRestart = true
  1507. for _, email := range emailsByInbound[enableIds[i]] {
  1508. if _, already := skippedReasons[email]; !already {
  1509. skippedReasons[email] = enablePanics[i].Error()
  1510. }
  1511. }
  1512. continue
  1513. }
  1514. if ibRes.needRestart {
  1515. needRestart = true
  1516. }
  1517. for email, reason := range ibRes.perEmailSkipped {
  1518. if _, already := skippedReasons[email]; !already {
  1519. skippedReasons[email] = reason
  1520. }
  1521. }
  1522. }
  1523. successEmails := make([]string, 0, len(recordsByEmail))
  1524. for email := range recordsByEmail {
  1525. if _, skipped := skippedReasons[email]; skipped {
  1526. continue
  1527. }
  1528. successEmails = append(successEmails, email)
  1529. }
  1530. if len(successEmails) > 0 {
  1531. now := time.Now().UnixMilli()
  1532. if err := runSerializedTx(func(tx *gorm.DB) error {
  1533. for _, batch := range chunkStrings(successEmails, sqlInChunk) {
  1534. if e := tx.Model(xray.ClientTraffic{}).Where("email IN ?", batch).Update("enable", enable).Error; e != nil {
  1535. return e
  1536. }
  1537. if e := tx.Model(&model.ClientRecord{}).Where("email IN ?", batch).
  1538. Updates(map[string]any{"enable": enable, "updated_at": now}).Error; e != nil {
  1539. return e
  1540. }
  1541. }
  1542. return nil
  1543. }); err != nil {
  1544. return result, needRestart, err
  1545. }
  1546. }
  1547. result.Changed = len(successEmails)
  1548. for email, reason := range skippedReasons {
  1549. result.Skipped = append(result.Skipped, BulkSetEnableReport{Email: email, Reason: reason})
  1550. }
  1551. return result, needRestart, nil
  1552. }
  1553. type bulkSetEnableInboundResult struct {
  1554. perEmailSkipped map[string]string
  1555. needRestart bool
  1556. }
  1557. func (s *ClientService) bulkSetEnableInboundClients(inboundSvc *InboundService, inboundId int, emails []string, enable bool) bulkSetEnableInboundResult {
  1558. res := bulkSetEnableInboundResult{perEmailSkipped: map[string]string{}}
  1559. defer lockInbound(inboundId).Unlock()
  1560. oldInbound, err := inboundSvc.GetInbound(inboundId)
  1561. if err != nil {
  1562. for _, e := range emails {
  1563. res.perEmailSkipped[e] = err.Error()
  1564. }
  1565. return res
  1566. }
  1567. var settings map[string]any
  1568. if err := json.Unmarshal([]byte(oldInbound.Settings), &settings); err != nil {
  1569. for _, e := range emails {
  1570. res.perEmailSkipped[e] = err.Error()
  1571. }
  1572. return res
  1573. }
  1574. wanted := make(map[string]struct{}, len(emails))
  1575. for _, email := range emails {
  1576. wanted[email] = struct{}{}
  1577. }
  1578. cipher := ""
  1579. if oldInbound.Protocol == model.Shadowsocks {
  1580. cipher, _ = settings["method"].(string)
  1581. }
  1582. type changedClient struct {
  1583. email string
  1584. wasEnable bool
  1585. client model.Client
  1586. }
  1587. var changed []changedClient
  1588. found := map[string]bool{}
  1589. nowMs := time.Now().UnixMilli()
  1590. interfaceClients, _ := settings["clients"].([]any)
  1591. for i, c := range interfaceClients {
  1592. entry, ok := c.(map[string]any)
  1593. if !ok {
  1594. continue
  1595. }
  1596. email, _ := entry["email"].(string)
  1597. if _, want := wanted[email]; !want || email == "" {
  1598. continue
  1599. }
  1600. found[email] = true
  1601. prev, _ := entry["enable"].(bool)
  1602. if prev == enable {
  1603. continue
  1604. }
  1605. entry["enable"] = enable
  1606. entry["updated_at"] = nowMs
  1607. interfaceClients[i] = entry
  1608. // Build the pushed client from the inbound JSON (the per-inbound source of
  1609. // truth), so a remote UpdateUser carries every field and never zeroes
  1610. // subId/totalGB/expiry from drifting ClientRecord columns (#4628/#4792).
  1611. var parsed model.Client
  1612. if b, mErr := json.Marshal(entry); mErr == nil {
  1613. _ = json.Unmarshal(b, &parsed)
  1614. }
  1615. parsed.Email = email
  1616. parsed.Enable = enable
  1617. changed = append(changed, changedClient{email: email, wasEnable: prev, client: parsed})
  1618. }
  1619. for email := range wanted {
  1620. if !found[email] {
  1621. res.perEmailSkipped[email] = "Client Not Found In Inbound"
  1622. }
  1623. }
  1624. if len(changed) == 0 {
  1625. return res
  1626. }
  1627. settings["clients"] = interfaceClients
  1628. newSettings, err := json.MarshalIndent(settings, "", " ")
  1629. if err != nil {
  1630. for _, ch := range changed {
  1631. res.perEmailSkipped[ch.email] = err.Error()
  1632. }
  1633. return res
  1634. }
  1635. prevSettings := oldInbound.Settings
  1636. oldInbound.Settings = string(newSettings)
  1637. rt, push, _, perr := inboundSvc.nodePushPlan(oldInbound)
  1638. if perr != nil {
  1639. for _, ch := range changed {
  1640. res.perEmailSkipped[ch.email] = perr.Error()
  1641. }
  1642. return res
  1643. }
  1644. if oldInbound.NodeID != nil && push && len(changed) > nodeBulkPushThreshold {
  1645. push = false
  1646. }
  1647. txErr := runSerializedTx(func(tx *gorm.DB) error {
  1648. if e := commitInboundClientSettings(tx, oldInbound, prevSettings); e != nil {
  1649. return e
  1650. }
  1651. finalClients, gcErr := inboundSvc.GetClients(oldInbound)
  1652. if gcErr != nil {
  1653. return gcErr
  1654. }
  1655. if err := s.SyncInbound(tx, inboundId, finalClients); err != nil {
  1656. return err
  1657. }
  1658. if oldInbound.NodeID != nil {
  1659. return (&NodeService{}).MarkNodeDirtyTx(tx, *oldInbound.NodeID)
  1660. }
  1661. return nil
  1662. })
  1663. if txErr != nil {
  1664. for _, ch := range changed {
  1665. res.perEmailSkipped[ch.email] = txErr.Error()
  1666. }
  1667. return res
  1668. }
  1669. if oldInbound.NodeID == nil {
  1670. if !push {
  1671. res.needRestart = true
  1672. } else {
  1673. for _, ch := range changed {
  1674. if enable {
  1675. err1 := rt.AddUser(context.Background(), oldInbound, map[string]any{
  1676. "email": ch.client.Email,
  1677. "id": ch.client.ID,
  1678. "security": ch.client.Security,
  1679. "flow": ch.client.Flow,
  1680. "auth": ch.client.Auth,
  1681. "password": ch.client.Password,
  1682. "cipher": cipher,
  1683. "reverse": ch.client.Reverse,
  1684. })
  1685. if err1 != nil {
  1686. logger.Debug("Error in adding client on", rt.Name(), ":", err1)
  1687. res.needRestart = true
  1688. }
  1689. } else if ch.wasEnable {
  1690. err1 := rt.RemoveUser(context.Background(), oldInbound, ch.email)
  1691. if err1 != nil && !strings.Contains(err1.Error(), fmt.Sprintf("User %s not found.", ch.email)) {
  1692. logger.Debug("Error in removing client on", rt.Name(), ":", err1)
  1693. res.needRestart = true
  1694. } else if err1 == nil && droppedClientNeedsRestart() {
  1695. // A removed credential does not end the session it was serving.
  1696. res.needRestart = true
  1697. }
  1698. }
  1699. }
  1700. }
  1701. } else if push {
  1702. pushFailed := false
  1703. for _, ch := range changed {
  1704. updated := ch.client
  1705. updated.UpdatedAt = nowMs
  1706. ctx, cancel := nodePushContext()
  1707. err1 := rt.UpdateUser(ctx, oldInbound, ch.email, updated)
  1708. cancel()
  1709. if err1 != nil {
  1710. logger.Warning("Error in updating client on", rt.Name(), ":", err1)
  1711. pushFailed = true
  1712. // First failure ends the batch push; the reconcile converges the rest.
  1713. break
  1714. }
  1715. }
  1716. if !pushFailed {
  1717. advancePushedInbound(rt, prevSettings, string(newSettings), oldInbound)
  1718. }
  1719. }
  1720. return res
  1721. }