xray.go 55 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688
  1. package service
  2. import (
  3. "encoding/json"
  4. "errors"
  5. "fmt"
  6. "path"
  7. "path/filepath"
  8. "runtime"
  9. "slices"
  10. "strings"
  11. "sync"
  12. "github.com/mhsanaei/3x-ui/v3/internal/amneziawg"
  13. "github.com/mhsanaei/3x-ui/v3/internal/amneziawgnet"
  14. "github.com/mhsanaei/3x-ui/v3/internal/config"
  15. "github.com/mhsanaei/3x-ui/v3/internal/database/model"
  16. "github.com/mhsanaei/3x-ui/v3/internal/logger"
  17. "github.com/mhsanaei/3x-ui/v3/internal/util/json_util"
  18. "github.com/mhsanaei/3x-ui/v3/internal/util/maskcompat"
  19. "github.com/mhsanaei/3x-ui/v3/internal/xray"
  20. "go.uber.org/atomic"
  21. )
  22. var (
  23. lock sync.Mutex
  24. isNeedXrayRestart atomic.Bool // Indicates that restart was requested for Xray
  25. isManuallyStopped atomic.Bool // Indicates that Xray was stopped manually from the panel
  26. xrayState xrayLifecycle
  27. )
  28. type xrayLifecycle struct {
  29. mu sync.RWMutex
  30. process *xray.Process
  31. result string
  32. // heldBack is why the running core still serves the previous config.
  33. heldBack string
  34. }
  35. func (s *xrayLifecycle) snapshot() (*xray.Process, string) {
  36. s.mu.RLock()
  37. defer s.mu.RUnlock()
  38. return s.process, s.result
  39. }
  40. func (s *xrayLifecycle) replace(process *xray.Process) {
  41. s.mu.Lock()
  42. s.process = process
  43. s.result = ""
  44. s.heldBack = ""
  45. s.mu.Unlock()
  46. }
  47. func (s *xrayLifecycle) holdBack(reason string) {
  48. s.mu.Lock()
  49. s.heldBack = reason
  50. s.mu.Unlock()
  51. }
  52. func (s *xrayLifecycle) heldBackReason() string {
  53. s.mu.RLock()
  54. defer s.mu.RUnlock()
  55. return s.heldBack
  56. }
  57. func (s *xrayLifecycle) storeResult(process *xray.Process, result string) {
  58. s.mu.Lock()
  59. if s.process == process && s.result == "" {
  60. s.result = result
  61. }
  62. s.mu.Unlock()
  63. }
  64. func currentXrayProcess() *xray.Process {
  65. process, _ := xrayState.snapshot()
  66. return process
  67. }
  68. // XrayService provides business logic for Xray process management.
  69. // It handles starting, stopping, restarting Xray, and managing its configuration.
  70. type XrayService struct {
  71. inboundService InboundService
  72. settingService SettingService
  73. nodeService NodeService
  74. xrayAPI xray.XrayAPI
  75. }
  76. // IsXrayRunning checks if the Xray process is currently running.
  77. func (s *XrayService) IsXrayRunning() bool {
  78. process := currentXrayProcess()
  79. return process != nil && process.IsRunning()
  80. }
  81. // XrayProcess returns the current Xray process instance (may be nil when Xray
  82. // is not running). It exposes the lifecycle snapshot to callers outside this
  83. // package (e.g. the tgbot subpackage).
  84. func XrayProcess() *xray.Process {
  85. return currentXrayProcess()
  86. }
  87. // GetXrayErr returns the error from the Xray process, if any.
  88. func (s *XrayService) GetXrayErr() error {
  89. process := currentXrayProcess()
  90. if process == nil {
  91. return nil
  92. }
  93. err := process.GetErr()
  94. if err == nil {
  95. return nil
  96. }
  97. if runtime.GOOS == "windows" && err.Error() == "exit status 1" {
  98. // exit status 1 on Windows means that Xray process was killed
  99. // as we kill process to stop in on Windows, this is not an error
  100. return nil
  101. }
  102. return err
  103. }
  104. // GetHeldBackConfig returns why the running core still serves its previous
  105. // config, or "" when the pending config was applied.
  106. func (s *XrayService) GetHeldBackConfig() string {
  107. return xrayState.heldBackReason()
  108. }
  109. // GetXrayResult returns the result string from the Xray process.
  110. func (s *XrayService) GetXrayResult() string {
  111. process, cachedResult := xrayState.snapshot()
  112. if cachedResult != "" {
  113. return cachedResult
  114. }
  115. if process == nil || process.IsRunning() {
  116. return ""
  117. }
  118. result := process.GetResult()
  119. if runtime.GOOS == "windows" && result == "exit status 1" {
  120. // exit status 1 on Windows means that Xray process was killed
  121. // as we kill process to stop in on Windows, this is not an error
  122. return ""
  123. }
  124. xrayState.storeResult(process, result)
  125. return result
  126. }
  127. // GetXrayVersion returns the version of the running Xray process.
  128. func (s *XrayService) GetXrayVersion() string {
  129. process := currentXrayProcess()
  130. if process == nil {
  131. return "Unknown"
  132. }
  133. return process.GetXrayVersion()
  134. }
  135. // RemoveIndex removes an element at the specified index from a slice.
  136. // Returns a new slice with the element removed.
  137. func RemoveIndex(s []any, index int) []any {
  138. return append(s[:index], s[index+1:]...)
  139. }
  140. // GetXrayConfig retrieves and builds the Xray configuration from settings and inbounds.
  141. func (s *XrayService) GetXrayConfig() (*xray.Config, error) {
  142. templateConfig, err := s.settingService.GetXrayConfigTemplate()
  143. if err != nil {
  144. return nil, err
  145. }
  146. xrayConfig := &xray.Config{}
  147. err = json.Unmarshal([]byte(templateConfig), xrayConfig)
  148. if err != nil {
  149. return nil, err
  150. }
  151. xrayConfig.LogConfig = resolveXrayLogPaths(xrayConfig.LogConfig)
  152. xrayConfig.API = ensureAPIServices(xrayConfig.API)
  153. xrayConfig.Policy = ensureStatsPolicy(xrayConfig.Policy)
  154. xrayConfig.RouterConfig = stripDisabledRules(xrayConfig.RouterConfig)
  155. // Template outbounds authored before the xray-core #6258 XHTTP rename may
  156. // still carry sessionPlacement/sessionKey; lift them too (same reason as
  157. // the per-inbound lift below).
  158. xrayConfig.OutboundConfigs = liftOutboundsXhttpSessionIDKeys(xrayConfig.OutboundConfigs)
  159. // Bridge amneziawg outbounds before anything else reads OutboundConfigs;
  160. // the core has no amneziawg proxy and would reject the raw entry.
  161. if err := transformAmneziaWGOutbounds(xrayConfig); err != nil {
  162. return nil, err
  163. }
  164. _, _, _ = s.inboundService.AddTraffic(nil, nil)
  165. inbounds, err := s.inboundService.GetAllInbounds()
  166. if err != nil {
  167. return nil, err
  168. }
  169. for _, inbound := range inbounds {
  170. if !inbound.Enable {
  171. continue
  172. }
  173. if inbound.NodeID != nil {
  174. continue
  175. }
  176. if inbound.Protocol == model.MTProto || inbound.Protocol == model.AmneziaWG || inbound.Protocol == model.TUIC {
  177. continue
  178. }
  179. settings := map[string]any{}
  180. _ = json.Unmarshal([]byte(inbound.Settings), &settings)
  181. var wireguardClientsByEmail map[string]model.Client
  182. if inbound.Protocol == model.WireGuard {
  183. inboundClients, _ := ParseInboundSettingsClients(inbound.Settings)
  184. if len(inboundClients) > 0 {
  185. wireguardClientsByEmail = make(map[string]model.Client, len(inboundClients))
  186. for _, client := range inboundClients {
  187. wireguardClientsByEmail[strings.ToLower(strings.TrimSpace(client.Email))] = client
  188. }
  189. }
  190. }
  191. dbClients, listErr := s.inboundService.clientService.ListForInbound(nil, inbound.Id)
  192. if listErr != nil {
  193. return nil, listErr
  194. }
  195. clientStats := inbound.ClientStats
  196. enableMap := make(map[string]bool, len(clientStats))
  197. for _, clientTraffic := range clientStats {
  198. enableMap[clientTraffic.Email] = clientTraffic.Enable
  199. }
  200. finalClients := make([]any, 0, len(dbClients))
  201. var wgPeers []any
  202. for i := range dbClients {
  203. c := dbClients[i]
  204. if enable, exists := enableMap[c.Email]; exists && !enable {
  205. logger.Infof("Remove Inbound User %s due to expiration or traffic limit", c.Email)
  206. continue
  207. }
  208. if !c.Enable {
  209. continue
  210. }
  211. flow := c.Flow
  212. if flow == "xtls-rprx-vision-udp443" {
  213. flow = "xtls-rprx-vision"
  214. }
  215. if inbound.DisableFlow {
  216. flow = ""
  217. }
  218. entry := map[string]any{"email": c.Email}
  219. switch inbound.Protocol {
  220. case model.VLESS:
  221. if c.ID != "" {
  222. entry["id"] = c.ID
  223. }
  224. if flow != "" {
  225. entry["flow"] = flow
  226. }
  227. if c.Reverse != nil {
  228. entry["reverse"] = c.Reverse
  229. }
  230. case model.VMESS:
  231. if c.ID != "" {
  232. entry["id"] = c.ID
  233. }
  234. if c.Security != "" {
  235. entry["security"] = c.Security
  236. }
  237. case model.Trojan:
  238. if c.Password != "" {
  239. entry["password"] = c.Password
  240. }
  241. if flow != "" {
  242. entry["flow"] = flow
  243. }
  244. case model.Shadowsocks:
  245. if c.Password != "" {
  246. entry["password"] = c.Password
  247. }
  248. case model.Hysteria:
  249. if c.Auth != "" {
  250. entry["auth"] = c.Auth
  251. }
  252. case model.WireGuard:
  253. if inboundClient, ok := wireguardClientsByEmail[strings.ToLower(strings.TrimSpace(c.Email))]; ok {
  254. c.AllowedIPs = inboundClient.AllowedIPs
  255. c.PreSharedKey = inboundClient.PreSharedKey
  256. }
  257. wgPeers = append(wgPeers, model.WireguardPeerFromClient(c))
  258. continue
  259. }
  260. finalClients = append(finalClients, entry)
  261. }
  262. var mutated bool
  263. if inbound.Protocol == model.WireGuard {
  264. delete(settings, "clients")
  265. if wgPeers == nil {
  266. wgPeers = []any{}
  267. }
  268. settings["peers"] = wgPeers
  269. mutated = true
  270. } else {
  271. _, hadClients := settings["clients"]
  272. mutated = hadClients || len(finalClients) > 0
  273. if mutated {
  274. settings["clients"] = finalClients
  275. }
  276. }
  277. if inboundCanHostFallbacks(inbound) {
  278. fallbacks, fbErr := s.inboundService.fallbackService.BuildFallbacksJSON(nil, inbound.Id)
  279. if fbErr != nil {
  280. return nil, fbErr
  281. }
  282. if len(fallbacks) > 0 {
  283. generic := make([]any, 0, len(fallbacks))
  284. for _, f := range fallbacks {
  285. generic = append(generic, f)
  286. }
  287. settings["fallbacks"] = generic
  288. mutated = true
  289. }
  290. }
  291. if mutated {
  292. modifiedSettings, err := json.MarshalIndent(settings, "", " ")
  293. if err != nil {
  294. return nil, err
  295. }
  296. inbound.Settings = string(modifiedSettings)
  297. }
  298. if len(inbound.StreamSettings) > 0 {
  299. // Unmarshal stream JSON
  300. var stream map[string]any
  301. _ = json.Unmarshal([]byte(inbound.StreamSettings), &stream)
  302. // Remove the "settings" field under "tlsSettings" and "realitySettings"
  303. tlsSettings, ok1 := stream["tlsSettings"].(map[string]any)
  304. realitySettings, ok2 := stream["realitySettings"].(map[string]any)
  305. if ok1 || ok2 {
  306. if ok1 {
  307. delete(tlsSettings, "settings")
  308. } else if ok2 {
  309. delete(realitySettings, "settings")
  310. }
  311. }
  312. delete(stream, "externalProxy")
  313. // finalmask.tcp + REALITY panics Xray-core on the first connection
  314. // (XTLS/Xray-core#6453). AddInbound/UpdateInbound reject this
  315. // combination at save time, but a row saved before that guard
  316. // existed (upgrade, node sync, restored backup, direct DB edit)
  317. // would still crash Xray on the next restart without this — drop
  318. // it here too, the same way liftXhttpSessionIDKeys and
  319. // HealShadowsocksClientMethods heal other legacy data in place.
  320. if len(finalMaskRealityTcpMasks(stream)) > 0 {
  321. logger.Warningf("Inbound %q: dropping finalmask, incompatible with REALITY security (crashes Xray-core, see XTLS/Xray-core#6453)", inbound.Tag)
  322. delete(stream, "finalmask")
  323. }
  324. dropEmptyRandPackets(stream["finalmask"])
  325. if dropped := stripIncompleteXmcMasks(stream); dropped > 0 {
  326. logger.Warningf("Inbound %q: dropping %d XMC finalmask mask(s) without complete Minecraft profiles — reconfigure them to restore the obfuscation (see XTLS/Xray-core#6487)", inbound.Tag, dropped)
  327. }
  328. // A row that skipped the save path can still carry the pre-26.9.30 xdns lists.
  329. maskcompat.UpgradeLegacyXdns(stream["finalmask"])
  330. // xray-core v26.6.22 (#6258) renamed the XHTTP session keys and
  331. // kept no fallback. Lift legacy sessionPlacement/sessionKey onto the
  332. // new names here so inbounds stored before the rename keep working
  333. // without the admin re-saving them.
  334. liftXhttpSessionIDKeys(stream)
  335. newStream, err := json.MarshalIndent(stream, "", " ")
  336. if err != nil {
  337. return nil, err
  338. }
  339. inbound.StreamSettings = string(newStream)
  340. }
  341. if inbound.Protocol == model.Shadowsocks {
  342. if healed, ok := model.HealShadowsocksClientMethods(inbound.Settings); ok {
  343. inbound.Settings = healed
  344. }
  345. }
  346. inboundConfig := inbound.GenXrayInboundConfig()
  347. xrayConfig.InboundConfigs = append(xrayConfig.InboundConfigs, *inboundConfig)
  348. }
  349. // Merge subscription-derived outbounds (if any) into the final outbounds array.
  350. // These are additive: each subscription is placed before or after the template
  351. // outbounds based on its Prepend flag, ordered by Priority. Tags assigned by the
  352. // subscription service are kept stable across refreshes so that balancers and
  353. // routing rules continue to work.
  354. subSvc := &OutboundSubscriptionService{}
  355. if prepend, appendList, err := subSvc.activeOutboundsSplit(); err == nil && (len(prepend) > 0 || len(appendList) > 0) {
  356. mergeSubscriptionOutbounds(xrayConfig, prepend, appendList)
  357. }
  358. // Route opted-in local mtproto inbounds through the core's router. Each one
  359. // gets a loopback SOCKS bridge — tagged with the inbound's own tag so it is
  360. // matchable in routing rules — that its mtg sidecar dials Telegram through.
  361. // Done after the subscription merge so a selected subscription outbound (or
  362. // balancer) is a valid rule target.
  363. for i := range inbounds {
  364. inbound := inbounds[i]
  365. if inbound.Protocol != model.MTProto || !inbound.Enable || inbound.NodeID != nil {
  366. continue
  367. }
  368. injectMtprotoEgress(xrayConfig, inbound)
  369. }
  370. // Every AmneziaWG inbound is embedded (internal/amneziawgnet: amneziawg-go
  371. // over a gVisor netstack, no kernel module) and relays every peer's
  372. // decapsulated traffic into its own loopback SOCKS5 inbound, always on —
  373. // unlike mtproto's bridge above, there's no opt-in gate here: once
  374. // traffic is decapsulated in gVisor, Xray's own freedom outbound is the
  375. // only way it reaches the real internet at all, not an optional extra
  376. // hop. Whether it goes anywhere beyond Xray's default routing is up to
  377. // whatever rules the admin adds through the stock Routing page, exactly
  378. // like routing any other protocol.
  379. injectAmneziawgnetSocks(xrayConfig, inbounds)
  380. // Restores each opted-in peer's own distinct public IPv6 source identity
  381. // for its outbound connections — a peer that has an IPv6 address in its
  382. // AllowedIPs, on an inbound with IPv6Enabled, gets its own freedom
  383. // outbound bound to that exact address via sendThrough.
  384. // internal/amneziawgnet's own Manager is responsible for actually
  385. // aliasing that address onto the host (see v6alias.go) so the kernel
  386. // lets Xray bind an egress socket to it at all; this call only builds
  387. // the Xray-side outbound/routing-rule half.
  388. injectAmneziawgV6Egress(xrayConfig, inbounds)
  389. // Wire the panel's own HTTP traffic through the configured outbound, after
  390. // the subscription merge so subscription outbound tags are valid targets.
  391. if egressTag, err := s.settingService.GetPanelOutbound(); err != nil {
  392. logger.Warning("read panelOutbound setting failed:", err)
  393. } else if egressTag != "" {
  394. injectPanelEgress(xrayConfig, egressTag)
  395. }
  396. nodes, err := s.nodeService.GetAll()
  397. if err != nil {
  398. logger.Warning("read nodes for egress injection failed:", err)
  399. } else {
  400. injectNodeEgresses(xrayConfig, nodes)
  401. }
  402. return xrayConfig, nil
  403. }
  404. // PanelEgressInboundTag is the tag of the loopback SOCKS inbound injected into
  405. // the generated config when a panel outbound is configured. The panel's own
  406. // HTTP clients dial through it to egress via the chosen outbound.
  407. const PanelEgressInboundTag = "panel-egress"
  408. // panelEgressBasePort is the first port tried for the egress bridge; ports
  409. // already taken by other inbounds in the generated config are skipped.
  410. const panelEgressBasePort = 62790
  411. // injectPanelEgress appends a loopback SOCKS inbound and routing rule only when
  412. // outboundTag resolves in the final outbound or balancer set. Otherwise the
  413. // entire injection is skipped. Generated state is hot-appliable and never
  414. // modifies the stored template or restarts the core.
  415. func injectPanelEgress(cfg *xray.Config, outboundTag string) {
  416. for i := range cfg.InboundConfigs {
  417. if cfg.InboundConfigs[i].Tag == PanelEgressInboundTag {
  418. logger.Warning("panel egress: inbound tag [", PanelEgressInboundTag, "] already exists, skipping injection")
  419. return
  420. }
  421. }
  422. // The rule must exist before the inbound takes traffic, otherwise the
  423. // bridge would silently egress through the default outbound instead.
  424. routing := map[string]any{}
  425. if len(cfg.RouterConfig) > 0 {
  426. if err := json.Unmarshal(cfg.RouterConfig, &routing); err != nil {
  427. logger.Warning("panel egress: routing section is unparsable, skipping injection:", err)
  428. return
  429. }
  430. }
  431. if !routingTargetExists(routing, cfg.OutboundConfigs, outboundTag) {
  432. logger.Warning("panel egress: target tag [", outboundTag, "] not found, skipping injection")
  433. return
  434. }
  435. rules, _ := routing["rules"].([]any)
  436. rule := map[string]any{
  437. "type": "field",
  438. "inboundTag": []any{PanelEgressInboundTag},
  439. }
  440. // The configured tag may name a routing balancer instead of a concrete
  441. // outbound. A field rule can target either, so emit the matching key —
  442. // balancerTag load-balances the panel's own traffic across the balancer's
  443. // outbounds, while a plain outbound tag keeps the original behavior.
  444. if routingTagIsBalancer(routing, outboundTag) {
  445. rule["balancerTag"] = outboundTag
  446. } else {
  447. rule["outboundTag"] = outboundTag
  448. }
  449. routing["rules"] = append([]any{rule}, rules...)
  450. newRouting, err := json.Marshal(routing)
  451. if err != nil {
  452. logger.Warning("panel egress: failed to rebuild routing section, skipping injection:", err)
  453. return
  454. }
  455. cfg.RouterConfig = json_util.RawMessage(newRouting)
  456. used := make(map[int]struct{}, len(cfg.InboundConfigs))
  457. for i := range cfg.InboundConfigs {
  458. used[cfg.InboundConfigs[i].Port] = struct{}{}
  459. }
  460. port := panelEgressBasePort
  461. for {
  462. if _, taken := used[port]; !taken {
  463. break
  464. }
  465. port++
  466. }
  467. cfg.InboundConfigs = append(cfg.InboundConfigs, xray.InboundConfig{
  468. Listen: json_util.RawMessage(`"127.0.0.1"`),
  469. Port: port,
  470. Protocol: "socks",
  471. Settings: json_util.RawMessage(`{"auth":"noauth","udp":false}`),
  472. Tag: PanelEgressInboundTag,
  473. })
  474. }
  475. func outboundTagExists(outbounds json_util.RawMessage, tag string) bool {
  476. var parsed []struct {
  477. Tag string `json:"tag"`
  478. }
  479. if tag == "" || json.Unmarshal(outbounds, &parsed) != nil {
  480. return false
  481. }
  482. for _, outbound := range parsed {
  483. if outbound.Tag == tag {
  484. return true
  485. }
  486. }
  487. return false
  488. }
  489. func routingTargetExists(routing map[string]any, outbounds json_util.RawMessage, tag string) bool {
  490. return routingTagIsBalancer(routing, tag) || outboundTagExists(outbounds, tag)
  491. }
  492. // NodeEgressInboundTag returns the loopback SOCKS inbound tag for a given node.
  493. func NodeEgressInboundTag(nodeID int) string {
  494. return fmt.Sprintf("node-egress-%d", nodeID)
  495. }
  496. // nodeEgressBasePort is the first port tried for node egress bridges.
  497. const nodeEgressBasePort = 62800
  498. // injectNodeEgresses appends a loopback SOCKS inbound per enabled node that has
  499. // an OutboundTag, and prepends a routing rule sending that inbound's traffic to
  500. // the selected outbound tag. These bridges are hot-appliable.
  501. func injectNodeEgresses(cfg *xray.Config, nodes []*model.Node) {
  502. routing := map[string]any{}
  503. if len(cfg.RouterConfig) > 0 {
  504. if err := json.Unmarshal(cfg.RouterConfig, &routing); err != nil {
  505. logger.Warning("node egress: routing section is unparsable, skipping injection:", err)
  506. return
  507. }
  508. }
  509. used := make(map[int]struct{}, len(cfg.InboundConfigs))
  510. usedTags := make(map[string]struct{}, len(cfg.InboundConfigs))
  511. for i := range cfg.InboundConfigs {
  512. used[cfg.InboundConfigs[i].Port] = struct{}{}
  513. usedTags[cfg.InboundConfigs[i].Tag] = struct{}{}
  514. }
  515. rules, _ := routing["rules"].([]any)
  516. newRules := make([]any, 0)
  517. for _, n := range nodes {
  518. if !n.Enable || n.OutboundTag == "" {
  519. continue
  520. }
  521. if !routingTargetExists(routing, cfg.OutboundConfigs, n.OutboundTag) {
  522. logger.Warning("node egress: target tag [", n.OutboundTag, "] not found, skipping node [", n.Id, "]")
  523. continue
  524. }
  525. tag := NodeEgressInboundTag(n.Id)
  526. if _, exists := usedTags[tag]; exists {
  527. logger.Warning("node egress: inbound tag [", tag, "] already exists, skipping")
  528. continue
  529. }
  530. usedTags[tag] = struct{}{}
  531. rule := map[string]any{
  532. "type": "field",
  533. "inboundTag": []any{tag},
  534. }
  535. if routingTagIsBalancer(routing, n.OutboundTag) {
  536. rule["balancerTag"] = n.OutboundTag
  537. } else {
  538. rule["outboundTag"] = n.OutboundTag
  539. }
  540. newRules = append(newRules, rule)
  541. port := nodeEgressBasePort + n.Id
  542. for {
  543. if _, taken := used[port]; !taken {
  544. break
  545. }
  546. port++
  547. }
  548. used[port] = struct{}{}
  549. cfg.InboundConfigs = append(cfg.InboundConfigs, xray.InboundConfig{
  550. Listen: json_util.RawMessage(`"127.0.0.1"`),
  551. Port: port,
  552. Protocol: "socks",
  553. Settings: json_util.RawMessage(`{"auth":"noauth","udp":false}`),
  554. Tag: tag,
  555. })
  556. }
  557. if len(newRules) == 0 {
  558. return
  559. }
  560. routing["rules"] = append(newRules, rules...)
  561. newRouting, err := json.Marshal(routing)
  562. if err != nil {
  563. logger.Warning("node egress: failed to rebuild routing section, skipping injection:", err)
  564. return
  565. }
  566. cfg.RouterConfig = json_util.RawMessage(newRouting)
  567. }
  568. // routingTagIsBalancer reports whether tag names a balancer in the parsed
  569. // routing section. The panel-egress rule targets a balancer via balancerTag and
  570. // a concrete outbound via outboundTag, so the caller picks the key from this.
  571. func routingTagIsBalancer(routing map[string]any, tag string) bool {
  572. if tag == "" {
  573. return false
  574. }
  575. balancers, ok := routing["balancers"].([]any)
  576. if !ok {
  577. return false
  578. }
  579. for _, b := range balancers {
  580. bm, ok := b.(map[string]any)
  581. if !ok {
  582. continue
  583. }
  584. if t, ok := bm["tag"].(string); ok && t == tag {
  585. return true
  586. }
  587. }
  588. return false
  589. }
  590. // mtprotoEgressSocksSettings is the loopback SOCKS server a routed mtproto
  591. // inbound exposes for its mtg sidecar to dial Telegram through. mtg makes plain
  592. // TCP connections, so UDP is left off (matching the panel egress bridge).
  593. const mtprotoEgressSocksSettings = `{"auth":"noauth","udp":false}`
  594. // injectMtprotoEgress wires one routed mtproto inbound into the generated
  595. // config after any selected outbound resolves in the final target set. Invalid
  596. // selected targets or routing data skip the entire injection; without a selected
  597. // outbound, the bridge retains default-route behavior. Generated state remains
  598. // hot-appliable, leaves the stored template untouched, and never forces a full
  599. // Xray restart. Mirrors injectPanelEgress.
  600. func injectMtprotoEgress(cfg *xray.Config, inbound *model.Inbound) {
  601. var parsed struct {
  602. RouteThroughXray bool `json:"routeThroughXray"`
  603. RouteXrayPort int `json:"routeXrayPort"`
  604. OutboundTag string `json:"outboundTag"`
  605. }
  606. if err := json.Unmarshal([]byte(inbound.Settings), &parsed); err != nil {
  607. return
  608. }
  609. if !parsed.RouteThroughXray || parsed.RouteXrayPort <= 0 || inbound.Tag == "" {
  610. return
  611. }
  612. tag := inbound.Tag
  613. for i := range cfg.InboundConfigs {
  614. if cfg.InboundConfigs[i].Tag == tag {
  615. logger.Warning("mtproto egress: inbound tag [", tag, "] already present in generated config, skipping bridge")
  616. return
  617. }
  618. }
  619. if parsed.OutboundTag != "" {
  620. routing := map[string]any{}
  621. if len(cfg.RouterConfig) > 0 {
  622. if err := json.Unmarshal(cfg.RouterConfig, &routing); err != nil {
  623. logger.Warning("mtproto egress: routing section is unparsable, skipping injection:", err)
  624. return
  625. }
  626. }
  627. if !routingTargetExists(routing, cfg.OutboundConfigs, parsed.OutboundTag) {
  628. logger.Warning("mtproto egress: target tag [", parsed.OutboundTag, "] not found, skipping injection")
  629. return
  630. }
  631. rules, _ := routing["rules"].([]any)
  632. rule := map[string]any{
  633. "type": "field",
  634. "inboundTag": []any{tag},
  635. }
  636. if routingTagIsBalancer(routing, parsed.OutboundTag) {
  637. rule["balancerTag"] = parsed.OutboundTag
  638. } else {
  639. rule["outboundTag"] = parsed.OutboundTag
  640. }
  641. routing["rules"] = append([]any{rule}, rules...)
  642. newRouting, err := json.Marshal(routing)
  643. if err != nil {
  644. logger.Warning("mtproto egress: failed to rebuild routing section, skipping injection:", err)
  645. return
  646. }
  647. cfg.RouterConfig = json_util.RawMessage(newRouting)
  648. }
  649. cfg.InboundConfigs = append(cfg.InboundConfigs, xray.InboundConfig{
  650. Listen: json_util.RawMessage(`"127.0.0.1"`),
  651. Port: parsed.RouteXrayPort,
  652. Protocol: "socks",
  653. Settings: json_util.RawMessage(mtprotoEgressSocksSettings),
  654. Tag: tag,
  655. })
  656. }
  657. // Peers resolve DNS inside the tunnel, so domain rules match only via sniffing; routeOnly
  658. // keeps the dial on the peer's IP, else Telegram's FakeTLS (IP + foreign SNI) breaks.
  659. const amneziawgEgressSniffingSettings = `{"enabled":true,"destOverride":["http","tls","quic","fakedns"],"routeOnly":true}`
  660. // injectAmneziawgnetSocks gives every enabled AmneziaWG inbound with at
  661. // least one qualifying peer its own loopback SOCKS5 inbound for the
  662. // embedded (amneziawg-go) relay path (internal/amneziawgnet) -- always on,
  663. // since there is no alternative datapath once traffic is decapsulated in
  664. // gVisor: Xray's own freedom outbound is how it reaches the real internet at
  665. // all (see internal/amneziawgnet/relay.go's doc comment, Finding 3 of the
  666. // migration plan). Tagged with the inbound's own real tag: it's already
  667. // selectable in the panel's stock Routing page (InboundService.GetInboundTags
  668. // is protocol-blind), and per-inbound traffic totals
  669. // (internal/web/service/inbound_traffic.go's addClientTraffic) match by
  670. // exact tag -- reusing it isn't a style choice.
  671. func injectAmneziawgnetSocks(cfg *xray.Config, inbounds []*model.Inbound) {
  672. existingTags := make(map[string]struct{}, len(cfg.InboundConfigs))
  673. for i := range cfg.InboundConfigs {
  674. existingTags[cfg.InboundConfigs[i].Tag] = struct{}{}
  675. }
  676. for _, inbound := range inbounds {
  677. if inbound.Protocol != model.AmneziaWG || !inbound.Enable || inbound.NodeID != nil {
  678. continue
  679. }
  680. inst, ok := amneziawg.InstanceFromInbound(inbound)
  681. if !ok {
  682. continue
  683. }
  684. if _, taken := existingTags[inbound.Tag]; taken {
  685. logger.Warning("amneziawgnet socks: inbound tag [", inbound.Tag, "] already present in generated config, skipping its relay inbound")
  686. continue
  687. }
  688. emails := make([]string, 0, len(inst.Peers))
  689. for _, p := range inst.Peers {
  690. if p.Email != "" {
  691. emails = append(emails, p.Email)
  692. }
  693. }
  694. if len(emails) == 0 {
  695. continue
  696. }
  697. settings, err := amneziawgnet.SocksInboundSettings(emails, amneziawgnet.SocksPassword())
  698. if err != nil {
  699. logger.Warning("amneziawgnet socks: building settings for inbound [", inbound.Tag, "]: ", err)
  700. continue
  701. }
  702. existingTags[inbound.Tag] = struct{}{}
  703. cfg.InboundConfigs = append(cfg.InboundConfigs, xray.InboundConfig{
  704. Listen: json_util.RawMessage(`"127.0.0.1"`),
  705. Port: amneziawgnet.SOCKSPortForInbound(inbound.Id),
  706. Protocol: "socks",
  707. Settings: json_util.RawMessage(settings),
  708. Sniffing: json_util.RawMessage(amneziawgEgressSniffingSettings),
  709. Tag: inbound.Tag,
  710. })
  711. }
  712. }
  713. // amneziawgV6EgressTag returns the stable, globally-unique freedom outbound
  714. // tag for one peer's IPv6 source-identity egress. Stable across config
  715. // regenerations (a pure function of two stable identifiers), so
  716. // internal/xray/hot_diff.go's tag-keyed outbound/routing diffing recognizes
  717. // "unchanged" rather than remove+recreate on every poll. The inbound.Id
  718. // prefix is defense in depth, not load-bearing on its own: email is already
  719. // enforced globally unique across the whole panel's client table
  720. // (model.ClientRecord.Email has a gorm uniqueIndex) — kept anyway since it
  721. // costs nothing and makes the tag self-describing, matching
  722. // NodeEgressInboundTag's own style.
  723. func amneziawgV6EgressTag(inboundID int, email string) string {
  724. return fmt.Sprintf("amneziawg-v6-%d-%s", inboundID, email)
  725. }
  726. // injectAmneziawgV6Egress gives every enabled, non-node-hosted AmneziaWG
  727. // peer with an IPv6 AllowedIPs entry its own single-purpose freedom
  728. // outbound, bound via sendThrough to that exact address, plus a routing
  729. // rule sending only that peer's IPv6-destined traffic through it — restoring the
  730. // per-client public IPv6 identity the hard cutover temporarily dropped
  731. // (Phase 3.5 of the migration plan). Scoped to outbound source identity
  732. // only: it depends on internal/amneziawgnet's own alias mechanism actually
  733. // giving the host that address at the OS level (see v6alias.go's
  734. // V6AliasesActive, the exact same gate this function uses below) — without
  735. // that, sendThrough fails to bind and every connection through it errors
  736. // outright (freedom.go's dial failure); there is no fallback outbound.
  737. //
  738. // The routing rule matches both inboundTag and user: SocksInboundSettings
  739. // (used by injectAmneziawgnetSocks above) already authenticates each
  740. // connection as the peer's own email via stock SOCKS5 auth, and a stock
  741. // Xray SOCKS5 inbound sets that connection's stats/routing identity from
  742. // the authenticated username — so "user" reliably isolates exactly one
  743. // peer's traffic, the same building block Finding 3 of the migration plan
  744. // already established for per-client stats.
  745. //
  746. // Modeled on injectNodeEgresses (the established N-per-slice inbound+rule
  747. // precedent, not injectAmneziawgnetSocks itself, which only ever emits a
  748. // single inbound and never touches outbounds/routing) and
  749. // mergeSubscriptionOutbounds's unmarshal-append-remarshal pattern for
  750. // cfg.OutboundConfigs. Synthetic rules are prepended ahead of whatever's
  751. // already in the routing rules array, the same pattern injectNodeEgresses/
  752. // injectMtprotoEgress already use for their own always-must-win infra
  753. // rules — this never touches the admin's own saved Routing-page rule
  754. // order.
  755. func injectAmneziawgV6Egress(cfg *xray.Config, inbounds []*model.Inbound) {
  756. // Protocol is checked alongside Tag, not just Tag alone: a tag collision
  757. // with some unrelated (non-socks) inbound must not be mistaken for this
  758. // instance's own relay having been created.
  759. liveInboundTags := make(map[string]struct{}, len(cfg.InboundConfigs))
  760. for i := range cfg.InboundConfigs {
  761. if cfg.InboundConfigs[i].Protocol == "socks" {
  762. liveInboundTags[cfg.InboundConfigs[i].Tag] = struct{}{}
  763. }
  764. }
  765. var existingOutbounds []any
  766. if len(cfg.OutboundConfigs) > 0 {
  767. if err := json.Unmarshal(cfg.OutboundConfigs, &existingOutbounds); err != nil {
  768. logger.Warning("amneziawg v6 egress: outbounds section is unparsable, skipping injection:", err)
  769. return
  770. }
  771. }
  772. usedOutboundTags := make(map[string]struct{}, len(existingOutbounds))
  773. for _, o := range existingOutbounds {
  774. if m, ok := o.(map[string]any); ok {
  775. if t, ok := m["tag"].(string); ok {
  776. usedOutboundTags[t] = struct{}{}
  777. }
  778. }
  779. }
  780. routing := map[string]any{}
  781. if len(cfg.RouterConfig) > 0 {
  782. if err := json.Unmarshal(cfg.RouterConfig, &routing); err != nil {
  783. logger.Warning("amneziawg v6 egress: routing section is unparsable, skipping injection:", err)
  784. return
  785. }
  786. }
  787. rules, _ := routing["rules"].([]any)
  788. newRules := make([]any, 0)
  789. newOutbounds := make([]any, 0)
  790. for _, inbound := range inbounds {
  791. if inbound.Protocol != model.AmneziaWG || !inbound.Enable || inbound.NodeID != nil {
  792. continue
  793. }
  794. if _, live := liveInboundTags[inbound.Tag]; !live {
  795. // The relay inbound itself wasn't created this pass (e.g. a tag
  796. // collision inside injectAmneziawgnetSocks) -- no SOCKS5 inbound
  797. // exists for hot_diff.go's inboundTag match to ever fire against.
  798. continue
  799. }
  800. inst, ok := amneziawg.InstanceFromInbound(inbound)
  801. if !ok || !amneziawgnet.V6AliasesActive(inst) {
  802. continue
  803. }
  804. for _, p := range inst.Peers {
  805. if p.Email == "" {
  806. continue
  807. }
  808. v6 := amneziawg.FirstIPv6(p.AllowedIPs)
  809. if v6 == "" {
  810. continue
  811. }
  812. tag := amneziawgV6EgressTag(inbound.Id, p.Email)
  813. if _, taken := usedOutboundTags[tag]; taken {
  814. logger.Warning("amneziawg v6 egress: outbound tag [", tag, "] already exists, skipping peer [", p.Email, "]")
  815. continue
  816. }
  817. usedOutboundTags[tag] = struct{}{}
  818. newOutbounds = append(newOutbounds, map[string]any{
  819. "tag": tag,
  820. "protocol": "freedom",
  821. "sendThrough": v6,
  822. "settings": map[string]any{},
  823. })
  824. newRules = append(newRules, map[string]any{
  825. "type": "field",
  826. "inboundTag": []any{inbound.Tag},
  827. "user": []any{p.Email},
  828. "ip": []any{"::/0"},
  829. "outboundTag": tag,
  830. })
  831. }
  832. }
  833. if len(newOutbounds) == 0 {
  834. return
  835. }
  836. merged := make([]any, 0, len(existingOutbounds))
  837. merged = append(merged, existingOutbounds...)
  838. merged = append(merged, newOutbounds...)
  839. combined, err := json.MarshalIndent(merged, "", " ")
  840. if err != nil {
  841. logger.Warning("amneziawg v6 egress: failed to rebuild outbounds section, skipping injection:", err)
  842. return
  843. }
  844. cfg.OutboundConfigs = json_util.RawMessage(combined)
  845. routing["rules"] = append(newRules, rules...)
  846. newRouting, err := json.Marshal(routing)
  847. if err != nil {
  848. logger.Warning("amneziawg v6 egress: failed to rebuild routing section, skipping injection:", err)
  849. return
  850. }
  851. cfg.RouterConfig = json_util.RawMessage(newRouting)
  852. }
  853. // mergeSubscriptionOutbounds appends the subscription outbounds to the
  854. // OutboundConfigs array of the xray config. It works on the already-unmarshaled
  855. // template so that manually configured outbounds are never overwritten.
  856. //
  857. // Safety: if we cannot parse the template's outbounds array, we leave
  858. // OutboundConfigs exactly as it came from the template (we do not inject
  859. // subscription outbounds). This prevents us from accidentally dropping the
  860. // user's manually configured outbounds when the template is in a weird state.
  861. func mergeSubscriptionOutbounds(cfg *xray.Config, prepend, appendList []any) {
  862. if len(prepend) == 0 && len(appendList) == 0 {
  863. return
  864. }
  865. var templateOutbounds []any
  866. if len(cfg.OutboundConfigs) > 0 {
  867. if err := json.Unmarshal(cfg.OutboundConfigs, &templateOutbounds); err != nil {
  868. // Corrupt template outbounds — do not touch the field at all.
  869. // The user will see problems on Xray start / next save.
  870. return
  871. }
  872. }
  873. var merged []any
  874. merged = append(merged, prepend...)
  875. merged = append(merged, templateOutbounds...)
  876. merged = append(merged, appendList...)
  877. combined, err := json.MarshalIndent(merged, "", " ")
  878. if err != nil {
  879. return
  880. }
  881. cfg.OutboundConfigs = json_util.RawMessage(combined)
  882. }
  883. // ensureAPIServices guarantees the gRPC services the panel depends on are
  884. // listed in the generated config's api block: HandlerService and StatsService
  885. // have always been required for inbound/user management and traffic polling,
  886. // and RoutingService enables hot routing reload on templates saved before it
  887. // was added to the default template. The stored template itself is not
  888. // modified — only the generated runtime config.
  889. func ensureAPIServices(api json_util.RawMessage) json_util.RawMessage {
  890. if len(api) == 0 {
  891. // No api block means the panel's API integration is deliberately
  892. // disabled; don't resurrect it behind the user's back.
  893. return api
  894. }
  895. var parsed map[string]any
  896. if err := json.Unmarshal(api, &parsed); err != nil {
  897. return api
  898. }
  899. services, _ := parsed["services"].([]any)
  900. have := make(map[string]bool, len(services))
  901. for _, svc := range services {
  902. if name, ok := svc.(string); ok {
  903. have[name] = true
  904. }
  905. }
  906. added := false
  907. for _, name := range []string{"HandlerService", "StatsService", "RoutingService"} {
  908. if !have[name] {
  909. services = append(services, name)
  910. added = true
  911. }
  912. }
  913. if !added {
  914. return api
  915. }
  916. parsed["services"] = services
  917. out, err := json.Marshal(parsed)
  918. if err != nil {
  919. return api
  920. }
  921. return out
  922. }
  923. // ensureStatsPolicy guarantees every policy level in the generated config has
  924. // statsUserOnline enabled, so the core tracks per-email online IPs for the
  925. // panel's online view and access-log-free IP limiting. Generated clients carry
  926. // no explicit level, so level "0" is created when absent. The flag is panel
  927. // infrastructure and is forced on even over an explicit false in the template,
  928. // same as the api services above. An entirely missing or unparsable policy
  929. // block is left alone; the stored template itself is never modified — only the
  930. // generated runtime config.
  931. func ensureStatsPolicy(policy json_util.RawMessage) json_util.RawMessage {
  932. if len(policy) == 0 {
  933. return policy
  934. }
  935. var parsed map[string]any
  936. if err := json.Unmarshal(policy, &parsed); err != nil {
  937. return policy
  938. }
  939. levels, _ := parsed["levels"].(map[string]any)
  940. if levels == nil {
  941. levels = make(map[string]any)
  942. }
  943. if _, ok := levels["0"]; !ok {
  944. levels["0"] = map[string]any{}
  945. }
  946. changed := false
  947. for _, raw := range levels {
  948. level, ok := raw.(map[string]any)
  949. if !ok {
  950. continue
  951. }
  952. if enabled, ok := level["statsUserOnline"].(bool); !ok || !enabled {
  953. level["statsUserOnline"] = true
  954. changed = true
  955. }
  956. }
  957. if !changed {
  958. return policy
  959. }
  960. parsed["levels"] = levels
  961. out, err := json.Marshal(parsed)
  962. if err != nil {
  963. return policy
  964. }
  965. return out
  966. }
  967. // caseVariantKeys returns every key of parsed that equals want ignoring case,
  968. // lowest first so the fold is deterministic when several variants are present.
  969. func caseVariantKeys(parsed map[string]any, want string) []string {
  970. var keys []string
  971. for key := range parsed {
  972. if strings.EqualFold(key, want) {
  973. keys = append(keys, key)
  974. }
  975. }
  976. slices.Sort(keys)
  977. return keys
  978. }
  979. func resolveXrayLogPaths(logCfg json_util.RawMessage) json_util.RawMessage {
  980. if len(logCfg) == 0 {
  981. return logCfg
  982. }
  983. var parsed map[string]any
  984. if err := json.Unmarshal(logCfg, &parsed); err != nil {
  985. return logCfg
  986. }
  987. changed := false
  988. for _, key := range []string{"access", "error"} {
  989. // xray-core decodes this object with encoding/json, whose case-insensitive
  990. // field match makes "Access" reach AccessLog too — fold every variant.
  991. variants := caseVariantKeys(parsed, key)
  992. value, hasValue := parsed[key]
  993. for _, variant := range variants {
  994. if variant == key {
  995. continue
  996. }
  997. if !hasValue {
  998. value, hasValue = parsed[variant], true
  999. }
  1000. delete(parsed, variant)
  1001. changed = true
  1002. }
  1003. v, ok := value.(string)
  1004. if !ok {
  1005. continue
  1006. }
  1007. trimmed := strings.TrimSpace(v)
  1008. if trimmed == "" || strings.EqualFold(trimmed, "none") {
  1009. if changed {
  1010. parsed[key] = v
  1011. }
  1012. continue
  1013. }
  1014. base := path.Base(filepath.ToSlash(trimmed))
  1015. if base == "" || base == "." || base == ".." || base == "/" {
  1016. continue
  1017. }
  1018. confined := filepath.Join(config.GetLogFolder(), base)
  1019. if confined == trimmed {
  1020. continue
  1021. }
  1022. parsed[key] = confined
  1023. changed = true
  1024. }
  1025. if !changed {
  1026. return logCfg
  1027. }
  1028. out, err := json.Marshal(parsed)
  1029. if err != nil {
  1030. return logCfg
  1031. }
  1032. return out
  1033. }
  1034. // stripDisabledRules removes routing rules marked `enabled: false` from the
  1035. // generated runtime config and strips panel-only keys (`enabled`, `comment`)
  1036. // from the rest, since xray-core has no such fields. The internal api rule is
  1037. // always kept (see isApiRule) so traffic stats can't be toggled off. The stored
  1038. // template is untouched — only the generated config is filtered.
  1039. func stripDisabledRules(routerCfg json_util.RawMessage) json_util.RawMessage {
  1040. if len(routerCfg) == 0 {
  1041. return routerCfg
  1042. }
  1043. var parsed map[string]any
  1044. if err := json.Unmarshal(routerCfg, &parsed); err != nil {
  1045. return routerCfg
  1046. }
  1047. rules, ok := parsed["rules"].([]any)
  1048. if !ok || len(rules) == 0 {
  1049. return routerCfg
  1050. }
  1051. var activeRules []any
  1052. changed := false
  1053. for _, rawRule := range rules {
  1054. rule, ok := rawRule.(map[string]any)
  1055. if !ok {
  1056. activeRules = append(activeRules, rawRule)
  1057. continue
  1058. }
  1059. if enabledRaw, exists := rule["enabled"]; exists {
  1060. // The internal api rule carries traffic stats and must never be
  1061. // dropped, even if it was somehow marked disabled.
  1062. enabled, ok := enabledRaw.(bool)
  1063. if ok && !enabled && !isApiRule(rule) {
  1064. changed = true
  1065. continue
  1066. }
  1067. delete(rule, "enabled")
  1068. changed = true
  1069. }
  1070. if _, exists := rule["comment"]; exists {
  1071. delete(rule, "comment")
  1072. changed = true
  1073. }
  1074. activeRules = append(activeRules, rule)
  1075. }
  1076. if !changed {
  1077. return routerCfg
  1078. }
  1079. parsed["rules"] = activeRules
  1080. out, err := json.Marshal(parsed)
  1081. if err != nil {
  1082. return routerCfg
  1083. }
  1084. return out
  1085. }
  1086. // GetXrayTraffic fetches the current traffic statistics from the running Xray process.
  1087. func (s *XrayService) GetXrayTraffic() ([]*xray.Traffic, []*xray.ClientTraffic, error) {
  1088. process := currentXrayProcess()
  1089. if process == nil || !process.IsRunning() {
  1090. err := errors.New("xray is not running")
  1091. logger.Debug("Attempted to fetch Xray traffic, but Xray is not running:", err)
  1092. return nil, nil, err
  1093. }
  1094. apiPort := process.GetAPIPort()
  1095. if err := s.xrayAPI.Init(apiPort); err != nil {
  1096. logger.Debug("Failed to initialize Xray API:", err)
  1097. return nil, nil, err
  1098. }
  1099. defer s.xrayAPI.Close()
  1100. traffic, clientTraffic, err := s.xrayAPI.GetTraffic()
  1101. if err != nil {
  1102. logger.Debug("Failed to fetch Xray traffic:", err)
  1103. return nil, nil, err
  1104. }
  1105. return traffic, clientTraffic, nil
  1106. }
  1107. // GetOnlineUsers returns connection-based online users (email + source IPs)
  1108. // from the running core's online-stats API. ok=false means the API is not
  1109. // available — xray isn't running or the core predates the online-stats RPCs —
  1110. // and callers must use the legacy traffic-delta / access-log paths. The
  1111. // capability is probed lazily per process: an Unimplemented answer pins this
  1112. // core as unsupported until the next restart, while transient errors leave the
  1113. // capability undecided so a flaky poll can't lock in legacy mode.
  1114. func (s *XrayService) GetOnlineUsers() ([]xray.OnlineUser, bool, error) {
  1115. process := currentXrayProcess()
  1116. if process == nil || !process.IsRunning() {
  1117. return nil, false, nil
  1118. }
  1119. if process.OnlineAPISupport() == xray.OnlineAPIUnsupported {
  1120. return nil, false, nil
  1121. }
  1122. if err := s.xrayAPI.Init(process.GetAPIPort()); err != nil {
  1123. logger.Debug("Failed to initialize Xray API:", err)
  1124. return nil, false, err
  1125. }
  1126. defer s.xrayAPI.Close()
  1127. users, err := s.xrayAPI.GetOnlineUsers()
  1128. if err != nil {
  1129. if xray.IsUnimplementedErr(err) {
  1130. process.SetOnlineAPISupport(xray.OnlineAPIUnsupported)
  1131. logger.Info("xray core does not support the online-stats API; falling back to traffic-delta onlines and access-log IP limit")
  1132. return nil, false, nil
  1133. }
  1134. logger.Debug("Failed to fetch Xray online users:", err)
  1135. return nil, false, err
  1136. }
  1137. if process.OnlineAPISupport() == xray.OnlineAPIUnknown {
  1138. process.SetOnlineAPISupport(xray.OnlineAPISupported)
  1139. logger.Info("xray core supports the online-stats API; using connection-based onlines and access-log-free IP limit")
  1140. }
  1141. return users, true, nil
  1142. }
  1143. // BalancerStatus is the live view of one balancer for the panel UI. Running
  1144. // is false when the balancer isn't present in the running core (e.g. xray is
  1145. // stopped or the balancer hasn't been saved/applied yet).
  1146. type BalancerStatus struct {
  1147. Tag string `json:"tag"`
  1148. Running bool `json:"running"`
  1149. Override string `json:"override"`
  1150. Selected []string `json:"selected"`
  1151. }
  1152. // GetBalancersStatus queries the running core for the live state of the
  1153. // given balancer tags. Per-tag failures are reported as Running=false rather
  1154. // than failing the whole call, so the UI can render saved-but-not-applied
  1155. // balancers alongside live ones.
  1156. func (s *XrayService) GetBalancersStatus(tags []string) ([]BalancerStatus, error) {
  1157. statuses := make([]BalancerStatus, 0, len(tags))
  1158. process := currentXrayProcess()
  1159. if process == nil || !process.IsRunning() {
  1160. for _, tag := range tags {
  1161. statuses = append(statuses, BalancerStatus{Tag: tag})
  1162. }
  1163. return statuses, nil
  1164. }
  1165. if err := s.xrayAPI.Init(process.GetAPIPort()); err != nil {
  1166. return nil, err
  1167. }
  1168. defer s.xrayAPI.Close()
  1169. for _, tag := range tags {
  1170. info, err := s.xrayAPI.GetBalancerInfo(tag)
  1171. if err != nil {
  1172. logger.Debug("get balancer info [", tag, "] failed:", err)
  1173. statuses = append(statuses, BalancerStatus{Tag: tag})
  1174. continue
  1175. }
  1176. statuses = append(statuses, BalancerStatus{
  1177. Tag: tag,
  1178. Running: true,
  1179. Override: info.Override,
  1180. Selected: info.Selected,
  1181. })
  1182. }
  1183. return statuses, nil
  1184. }
  1185. // OverrideBalancer forces a balancer in the running core to use the given
  1186. // outbound tag; an empty target clears the override. When target names
  1187. // another balancer, the override resolves to the loopback outbound that
  1188. // routes traffic through the target balancer via the routing rules.
  1189. func (s *XrayService) OverrideBalancer(tag, target string) error {
  1190. process := currentXrayProcess()
  1191. if process == nil || !process.IsRunning() {
  1192. return errors.New("xray is not running")
  1193. }
  1194. if target != "" {
  1195. resolved, err := s.resolveOverrideTarget(target)
  1196. if err != nil {
  1197. return err
  1198. }
  1199. if resolved != "" {
  1200. target = resolved
  1201. }
  1202. }
  1203. if err := s.xrayAPI.Init(process.GetAPIPort()); err != nil {
  1204. return err
  1205. }
  1206. defer s.xrayAPI.Close()
  1207. return s.xrayAPI.SetBalancerTarget(tag, target)
  1208. }
  1209. // resolveOverrideTarget checks if target names a balancer and, if so,
  1210. // returns the loopback outbound tag that routes to it through the
  1211. // routing rules. Returns empty if target is already a concrete outbound.
  1212. func (s *XrayService) resolveOverrideTarget(target string) (string, error) {
  1213. template, err := s.settingService.GetXrayConfigTemplate()
  1214. if err != nil {
  1215. return "", err
  1216. }
  1217. var cfg map[string]any
  1218. if err := json.Unmarshal([]byte(template), &cfg); err != nil {
  1219. return "", err
  1220. }
  1221. routing, _ := cfg["routing"].(map[string]any)
  1222. if routing == nil {
  1223. return "", nil
  1224. }
  1225. rules, _ := routing["rules"].([]any)
  1226. for _, r := range rules {
  1227. rule, ok := r.(map[string]any)
  1228. if !ok {
  1229. continue
  1230. }
  1231. if rule["balancerTag"] != target {
  1232. continue
  1233. }
  1234. inboundTags, ok := rule["inboundTag"].([]any)
  1235. if !ok || len(inboundTags) == 0 {
  1236. continue
  1237. }
  1238. if lbTag, ok := inboundTags[0].(string); ok && strings.HasPrefix(lbTag, "_bl_") {
  1239. return lbTag, nil
  1240. }
  1241. }
  1242. return "", nil
  1243. }
  1244. // TestRoute asks the running core which outbound its router picks for the
  1245. // described connection.
  1246. func (s *XrayService) TestRoute(req xray.RouteTestRequest) (*xray.RouteTestResult, error) {
  1247. process := currentXrayProcess()
  1248. if process == nil || !process.IsRunning() {
  1249. return nil, errors.New("xray is not running")
  1250. }
  1251. if err := s.xrayAPI.Init(process.GetAPIPort()); err != nil {
  1252. return nil, err
  1253. }
  1254. defer s.xrayAPI.Close()
  1255. return s.xrayAPI.TestRoute(req)
  1256. }
  1257. // RestartXray reconciles the running Xray process with the current desired
  1258. // config. When isForce is false it first tries to apply the changes through
  1259. // the Xray gRPC API without restarting the process (inbounds, outbounds and
  1260. // routing rules/balancers are hot-reloadable); only changes the core cannot
  1261. // take at runtime — or a force request — stop and restart the process.
  1262. func (s *XrayService) RestartXray(isForce bool) error {
  1263. lock.Lock()
  1264. defer lock.Unlock()
  1265. logger.Debug("restart Xray, force:", isForce)
  1266. if !isForce && isManuallyStopped.Load() {
  1267. return nil
  1268. }
  1269. isManuallyStopped.Store(false)
  1270. xrayConfig, err := s.GetXrayConfig()
  1271. if err != nil {
  1272. return err
  1273. }
  1274. process := currentXrayProcess()
  1275. if process != nil && process.IsRunning() {
  1276. configUnchanged := process.GetConfig().Equals(xrayConfig)
  1277. if !isForce && configUnchanged && !isNeedXrayRestart.Load() {
  1278. logger.Debug("It does not need to restart Xray")
  1279. return nil
  1280. }
  1281. // A config the core cannot bind never replaces one that works: its failed
  1282. // start exits the core, and the watchdog would then loop on it forever.
  1283. if conflicts := bindConflicts(xrayConfig, process.GetConfig()); len(conflicts) > 0 {
  1284. refused := fmt.Sprintf("config refused: %s", conflicts[0])
  1285. for _, conflict := range conflicts {
  1286. logger.Error("xray config refused:", conflict.String())
  1287. }
  1288. // The refusal is otherwise invisible: the operator's request
  1289. // succeeded, so the status page has to carry the stale state.
  1290. xrayState.holdBack(refused)
  1291. return fmt.Errorf("xray %s", refused)
  1292. }
  1293. if !isForce && !configUnchanged && s.tryHotApply(process, xrayConfig) {
  1294. logger.Info("Xray config changes applied through the core API, no restart needed")
  1295. return nil
  1296. }
  1297. _ = process.Stop()
  1298. } else if conflicts := bindConflicts(xrayConfig, nil); len(conflicts) > 0 {
  1299. // Nothing is running to protect and the core is the authority on what it
  1300. // can bind: start it and let its own error name the port it lost.
  1301. logger.Warning("xray config may not start:", conflicts[0].String())
  1302. }
  1303. process = xray.NewProcess(xrayConfig)
  1304. xrayState.replace(process)
  1305. s.xrayAPI.StatsLastValues = nil
  1306. err = process.Start()
  1307. if err != nil {
  1308. return err
  1309. }
  1310. return nil
  1311. }
  1312. // restartToDropClients reports whether a diff that strands clients must be
  1313. // applied by restarting instead of through the API.
  1314. func (s *XrayService) restartToDropClients(diff *xray.HotDiff) bool {
  1315. if diff == nil || !diff.DropsUsers() {
  1316. return false
  1317. }
  1318. restart, err := s.settingService.GetRestartXrayOnClientDisable()
  1319. if err != nil {
  1320. logger.Warning("get RestartXrayOnClientDisable failed:", err)
  1321. return false
  1322. }
  1323. return restart
  1324. }
  1325. // tryHotApply attempts to reconcile the running Xray instance with newCfg
  1326. // through the core gRPC API (HandlerService for inbounds/outbounds,
  1327. // RoutingService for rules/balancers). It returns true when the running
  1328. // instance now matches newCfg; on any failure it returns false and the
  1329. // caller falls back to a full process restart, which cleans up whatever was
  1330. // partially applied. Callers must hold the package-level lock.
  1331. func (s *XrayService) tryHotApply(process *xray.Process, newCfg *xray.Config) bool {
  1332. oldCfg := process.GetConfig()
  1333. diff, ok := xray.ComputeHotDiff(oldCfg, newCfg)
  1334. if !ok {
  1335. logger.Debug("hot apply: config change is not API-applicable, falling back to restart")
  1336. return false
  1337. }
  1338. if diff.Empty() {
  1339. process.SetConfig(newCfg)
  1340. return true
  1341. }
  1342. // The core's RemoveUser drops the credential only, so a disabled or deleted
  1343. // client needs the restart this setting asks for.
  1344. if s.restartToDropClients(diff) {
  1345. logger.Info("hot apply: clients left the config, restarting to drop their live sessions")
  1346. return false
  1347. }
  1348. apiPort := process.GetAPIPort()
  1349. if apiPort <= 0 {
  1350. return false
  1351. }
  1352. // A dedicated client: s.xrayAPI may be in use by traffic polling on other
  1353. // service instances and is reset around restarts.
  1354. hotAPI := xray.XrayAPI{}
  1355. if err := hotAPI.Init(apiPort); err != nil {
  1356. logger.Debug("hot apply: failed to init xray api:", err)
  1357. return false
  1358. }
  1359. defer hotAPI.Close()
  1360. // Removals first so changed handlers and port swaps never collide with
  1361. // the additions that follow.
  1362. for _, u := range diff.RemovedUsers {
  1363. if err := hotAPI.RemoveUser(u.Tag, u.Email); err != nil && !xray.IsMissingHandlerErr(err) {
  1364. logger.Info("hot apply: remove user [", u.Email, "] from [", u.Tag, "] failed:", err)
  1365. return false
  1366. }
  1367. }
  1368. for _, tag := range diff.RemovedInboundTags {
  1369. if err := hotAPI.DelInbound(tag); err != nil && !xray.IsMissingHandlerErr(err) {
  1370. logger.Info("hot apply: remove inbound [", tag, "] failed:", err)
  1371. return false
  1372. }
  1373. }
  1374. for _, tag := range diff.RemovedOutboundTags {
  1375. if err := hotAPI.DelOutbound(tag); err != nil && !xray.IsMissingHandlerErr(err) {
  1376. logger.Info("hot apply: remove outbound [", tag, "] failed:", err)
  1377. return false
  1378. }
  1379. }
  1380. for _, ob := range diff.AddedOutbounds {
  1381. if err := addOutboundReconciling(&hotAPI, ob); err != nil {
  1382. logger.Info("hot apply: add outbound failed:", err)
  1383. return false
  1384. }
  1385. }
  1386. for _, ib := range diff.AddedInbounds {
  1387. if err := addInboundReconciling(&hotAPI, ib); err != nil {
  1388. logger.Info("hot apply: add inbound failed:", err)
  1389. return false
  1390. }
  1391. }
  1392. for _, u := range diff.AddedUsers {
  1393. if err := addUserReconciling(&hotAPI, u); err != nil {
  1394. logger.Info("hot apply: add user [", u.Email, "] to [", u.Tag, "] failed:", err)
  1395. return false
  1396. }
  1397. }
  1398. if diff.RoutingConfig != nil {
  1399. if err := hotAPI.ApplyRoutingConfig(diff.RoutingConfig); err != nil {
  1400. logger.Info("hot apply: apply routing config failed:", err)
  1401. return false
  1402. }
  1403. }
  1404. process.SetConfig(newCfg)
  1405. return true
  1406. }
  1407. // addUserReconciling adds a user, and on an email conflict (the user was
  1408. // already applied through the runtime API) replaces the existing user instead.
  1409. func addUserReconciling(api *xray.XrayAPI, u xray.UserOp) error {
  1410. err := api.AddUser(u.Protocol, u.Tag, u.User)
  1411. if err == nil || !xray.IsUserExistsErr(err) {
  1412. return err
  1413. }
  1414. if delErr := api.RemoveUser(u.Tag, u.Email); delErr != nil && !xray.IsMissingHandlerErr(delErr) {
  1415. return delErr
  1416. }
  1417. return api.AddUser(u.Protocol, u.Tag, u.User)
  1418. }
  1419. // addInboundReconciling adds an inbound, and on a tag conflict (the handler
  1420. // was already created through the runtime API while the stored snapshot was
  1421. // stale) replaces the existing handler instead.
  1422. func addInboundReconciling(api *xray.XrayAPI, inbound []byte) error {
  1423. err := api.AddInbound(inbound)
  1424. if err == nil || !xray.IsExistingTagErr(err) {
  1425. return err
  1426. }
  1427. var meta struct {
  1428. Tag string `json:"tag"`
  1429. }
  1430. if jsonErr := json.Unmarshal(inbound, &meta); jsonErr != nil || meta.Tag == "" {
  1431. return err
  1432. }
  1433. if delErr := api.DelInbound(meta.Tag); delErr != nil && !xray.IsMissingHandlerErr(delErr) {
  1434. return delErr
  1435. }
  1436. return api.AddInbound(inbound)
  1437. }
  1438. // addOutboundReconciling mirrors addInboundReconciling for outbounds.
  1439. func addOutboundReconciling(api *xray.XrayAPI, outbound []byte) error {
  1440. err := api.AddOutbound(outbound)
  1441. if err == nil || !xray.IsExistingTagErr(err) {
  1442. return err
  1443. }
  1444. var meta struct {
  1445. Tag string `json:"tag"`
  1446. }
  1447. if jsonErr := json.Unmarshal(outbound, &meta); jsonErr != nil || meta.Tag == "" {
  1448. return err
  1449. }
  1450. if delErr := api.DelOutbound(meta.Tag); delErr != nil && !xray.IsMissingHandlerErr(delErr) {
  1451. return delErr
  1452. }
  1453. return api.AddOutbound(outbound)
  1454. }
  1455. // StopXray stops the running Xray process.
  1456. func (s *XrayService) StopXray() error {
  1457. lock.Lock()
  1458. defer lock.Unlock()
  1459. isManuallyStopped.Store(true)
  1460. logger.Debug("Attempting to stop Xray...")
  1461. process := currentXrayProcess()
  1462. if process != nil && process.IsRunning() {
  1463. return process.Stop()
  1464. }
  1465. return errors.New("xray is not running")
  1466. }
  1467. // SetToNeedRestart marks that Xray needs to be restarted.
  1468. func (s *XrayService) SetToNeedRestart() {
  1469. isNeedXrayRestart.Store(true)
  1470. }
  1471. // GetXrayAPIPort returns the port the local xray process is listening on
  1472. // for its gRPC HandlerService, or 0 when xray isn't currently running.
  1473. // Exposed for the runtime package's LocalRuntime adapter without a
  1474. // service-package import cycle.
  1475. func (s *XrayService) GetXrayAPIPort() int {
  1476. process := currentXrayProcess()
  1477. if process == nil || !process.IsRunning() {
  1478. return 0
  1479. }
  1480. return process.GetAPIPort()
  1481. }
  1482. // IsNeedRestartAndSetFalse checks if restart is needed and resets the flag to false.
  1483. func (s *XrayService) IsNeedRestartAndSetFalse() bool {
  1484. return isNeedXrayRestart.CompareAndSwap(true, false)
  1485. }
  1486. // ApplyPendingRestart consumes the need-restart flag and restarts Xray. If the
  1487. // restart fails (for example GetXrayConfig hits a transient DB error and leaves
  1488. // the old process running), it re-arms the flag so the next tick retries instead
  1489. // of silently dropping the pending config change.
  1490. func (s *XrayService) ApplyPendingRestart() {
  1491. if !s.IsNeedRestartAndSetFalse() {
  1492. return
  1493. }
  1494. if err := s.RestartXray(false); err != nil {
  1495. logger.Error("restart xray failed:", err)
  1496. s.SetToNeedRestart()
  1497. }
  1498. }
  1499. // DidXrayCrash checks if Xray crashed by verifying it's not running and wasn't manually stopped.
  1500. func (s *XrayService) DidXrayCrash() bool {
  1501. return !s.IsXrayRunning() && !isManuallyStopped.Load()
  1502. }
  1503. // liftXhttpSessionIDKeys renames the legacy XHTTP session keys
  1504. // (sessionPlacement/sessionKey) to the v26.6.22 #6258 names
  1505. // (sessionIDPlacement/sessionIDKey) inside a streamSettings map. xray-core kept
  1506. // no fallback for the old names, so a config stored before the rename would be
  1507. // silently ignored by the engine. Returns true if it changed anything.
  1508. func liftXhttpSessionIDKeys(stream map[string]any) bool {
  1509. xhttp, ok := stream["xhttpSettings"].(map[string]any)
  1510. if !ok {
  1511. return false
  1512. }
  1513. changed := false
  1514. for legacy, renamed := range map[string]string{
  1515. "sessionPlacement": "sessionIDPlacement",
  1516. "sessionKey": "sessionIDKey",
  1517. } {
  1518. v, has := xhttp[legacy]
  1519. if !has {
  1520. continue
  1521. }
  1522. if _, exists := xhttp[renamed]; !exists {
  1523. xhttp[renamed] = v
  1524. }
  1525. delete(xhttp, legacy)
  1526. changed = true
  1527. }
  1528. return changed
  1529. }
  1530. // liftOutboundsXhttpSessionIDKeys applies liftXhttpSessionIDKeys to every
  1531. // outbound's streamSettings in the raw outbounds array. The original bytes are
  1532. // returned untouched when nothing needs lifting, so an unchanged config never
  1533. // looks modified to the hot-reload diff.
  1534. func liftOutboundsXhttpSessionIDKeys(raw json_util.RawMessage) json_util.RawMessage {
  1535. if len(raw) == 0 {
  1536. return raw
  1537. }
  1538. var outbounds []map[string]any
  1539. if err := json.Unmarshal(raw, &outbounds); err != nil {
  1540. return raw
  1541. }
  1542. changed := false
  1543. for _, ob := range outbounds {
  1544. if stream, ok := ob["streamSettings"].(map[string]any); ok {
  1545. if liftXhttpSessionIDKeys(stream) {
  1546. changed = true
  1547. }
  1548. }
  1549. }
  1550. if !changed {
  1551. return raw
  1552. }
  1553. if rewritten, err := json.Marshal(outbounds); err == nil {
  1554. return rewritten
  1555. }
  1556. return raw
  1557. }