mirror of
https://github.com/MHSanaei/3x-ui.git
synced 2026-10-05 05:32:07 +03:00
0054e671f8
* feat(tuic): implement native in-process Go TUIC v5 server - Implement native TUIC v5 protocol server on pure Go using quic-go - Bridge decrypted TCP/UDP traffic into Xray-core via loopback SOCKS5 inbound - Support full Xray routing rules (geosite/geoip) and cascading outbounds - Implement atomic per-client traffic accounting with TotalGB and ExpiryTime - Add automatic legacy cleanup for older Rust tuic-server binaries, configs, and orphaned processes - Eliminate external Rust tuic-server downloads from install/CI scripts * fix(tuic): address traffic accounting, client reload, and socket lifecycle issues * fix(service): update checkTuicSocksReverseConflict to use bindAddr for listenOverlaps * fix(tuic): resolve traffic double-accounting, UDP fragmentation, and socket lifecycle issues * feat(tuic): complete native Go integration and address audit findings - Integrate an isolated QUIC fork pinned to a specific commit - Preserve original QUIC dependencies for Xray, Hysteria and Gin - Apply BBR, CUBIC and Reno to server connections and exported client profiles - Bridge Xray BBR with correct monotonic time and congestion type conversions - Handle congestion sender recreation after PMTU changes - Update congestion control for new connections without restarting the listener - Preserve existing connections and their selected congestion controller - Apply per-inbound log levels through the shared panel logger - Add lifecycle, authentication and TCP/UDP relay events without exposing secrets - Rate-limit repeated authentication and relay warnings - Support native and QUIC UDP relay modes on the same listener - Recover UDP associations after relay worker failures - Fix TCP relay cancellation, idle shutdown and half-close handling - Close active sessions when client credentials are revoked or disabled - Track traffic by immutable client statistics IDs across email and UUID changes - Prevent ambiguous accounting and duplicate UUIDs within TUIC inbounds - Persist pending traffic in a durable shutdown journal - Replay journal batches transactionally without duplicate accounting - Report server shutdown failures through the shared logger - Preserve legacy flat and nested TUIC settings compatibility - Normalize congestion controller values consistently across backend and frontend - Preserve controller, UDP mode and SNI in client links and subscriptions - Separate client profile options from server settings in the TUIC form - Keep certificate path autofill explicit when changing client SNI - Align UDP packet size validation with protocol limits - Simplify and localize TUIC field hints and certificate autofill messages - Add controller, TCP/UDP, logging and live settings update tests - Add accounting identity, journal replay and shutdown regression tests - Add relay recovery, session revocation and legacy frontend form tests * fix(service): alias the TUIC duplicate-UUID subquery for PostgreSQL < 16 syncInboundClients runs a COUNT(*) FROM (subquery) for every client sync, whatever the protocol. PostgreSQL before 16 rejects a FROM subquery with no alias, so on the distro PostgreSQL install.sh provisions (14 on Ubuntu 22.04, 15 on Debian 12) every client add or edit failed with SQLSTATE 42601. Reproduced against postgres:15 with the new env-gated test. * fix(database): create tuic_traffic_receipts through the model migration AddTuicTrafficBatch issued CREATE TABLE IF NOT EXISTS at runtime, a schema change outside db.go. The table was invisible to allModels and migrationModels, so x-ui migrate-db dropped the receipts and a retained journal could be counted twice after a SQLite to PostgreSQL move. It is now a GORM model in both lists, and the insert uses OnConflict DoNothing. * chore(tuic): skip the ICMP-dependent relay test on Windows, drop dead collectors Go disables SIO_UDP_CONNRESET on Windows, so a dead UDP bridge never fails a read there and TestAudit3UDPAssociationMustRecoverAfterBridgeReadFailure was red on every Windows run. Server.CollectTotalTraffic and Manager.CollectTraffic had no caller. * refactor(tuic): serve TUIC on apernet/quic-go instead of a personal fork The native server depended on github.com/poise52/quic-go, a personal fork of apernet/quic-go patched only to pick the congestion controller before the handshake. That put a second QUIC/TLS stack in the binary that no upstream security fix reaches. apernet/quic-go is already in the graph through xray-core and exposes SetCongestionControl, so BBR is now installed on each accepted connection with Xray's own congestion.UseBBR; the cross-module BBR adapter is gone. apernet ships New Reno as its only built-in sender, so a cubic setting is served as new_reno server-side (clients still get cubic in their profile). The test inspectors now read the sender under congestionMutex, which the post-handshake install writes under. Linux loopback, single stream through Xray: 2428 -> 3383 Mbit/s (bbr). --------- Co-authored-by: Sanaei <ho3ein.sanaei@gmail.com>
1312 lines
42 KiB
Go
1312 lines
42 KiB
Go
package service
|
||
|
||
import (
|
||
"context"
|
||
"encoding/json"
|
||
"errors"
|
||
"fmt"
|
||
"maps"
|
||
"slices"
|
||
"strconv"
|
||
"strings"
|
||
"time"
|
||
|
||
"github.com/mhsanaei/3x-ui/v3/internal/database"
|
||
"github.com/mhsanaei/3x-ui/v3/internal/database/model"
|
||
"github.com/mhsanaei/3x-ui/v3/internal/logger"
|
||
"github.com/mhsanaei/3x-ui/v3/internal/web/runtime"
|
||
"github.com/mhsanaei/3x-ui/v3/internal/xray"
|
||
|
||
"gorm.io/gorm"
|
||
"gorm.io/gorm/clause"
|
||
)
|
||
|
||
// A client with a renewal day set auto-renews too, so it must not read as
|
||
// depleted — otherwise the operator's purge deletes it between cycles (#6239).
|
||
const depletedClientsClause = "reset = 0 and reset_day = 0 and reset_weekday = 0 and ((total > 0 and up + down >= total) or (expiry_time > 0 and expiry_time <= ?))"
|
||
|
||
func (s *InboundService) AddTraffic(inboundTraffics []*xray.Traffic, clientTraffics []*xray.ClientTraffic) (needRestart bool, clientsDisabled bool, err error) {
|
||
var disabledNodeIDs []int
|
||
var remotePlans []trafficInboundUpdatePlan
|
||
var renewed []string
|
||
err = submitTrafficWrite(func() error {
|
||
var inner error
|
||
needRestart, clientsDisabled, disabledNodeIDs, remotePlans, renewed, inner = s.addTrafficLocked(inboundTraffics, clientTraffics)
|
||
return inner
|
||
})
|
||
if err != nil {
|
||
return
|
||
}
|
||
s.resetMtprotoClientQuotas(renewed)
|
||
// Off the serial writer: a hanging node must not stall traffic accounting.
|
||
needRestart = s.applyTrafficRemotePlans(remotePlans) || needRestart
|
||
if len(disabledNodeIDs) > 0 {
|
||
s.restartRemoteNodesOnDisable(disabledNodeIDs)
|
||
}
|
||
return
|
||
}
|
||
|
||
func (s *InboundService) addTrafficLocked(inboundTraffics []*xray.Traffic, clientTraffics []*xray.ClientTraffic) (bool, bool, []int, []trafficInboundUpdatePlan, []string, error) {
|
||
db := database.GetDB()
|
||
// Commit durable traffic before best-effort lifecycle maintenance so helper
|
||
// failures cannot discard usage already reported by Xray.
|
||
if err := db.Transaction(func(tx *gorm.DB) error {
|
||
if err := s.addInboundTraffic(tx, inboundTraffics); err != nil {
|
||
return err
|
||
}
|
||
return s.addClientTraffic(tx, clientTraffics)
|
||
}); err != nil {
|
||
return false, false, nil, nil, nil, err
|
||
}
|
||
|
||
var (
|
||
needRestart bool
|
||
clientsDisabled bool
|
||
disabledNodeIDs []int
|
||
disabledClientsCount int64
|
||
)
|
||
batch := newTrafficMutationBatch()
|
||
err := db.Transaction(func(tx *gorm.DB) error {
|
||
needRestart0, count, err := s.autoRenewClients(tx, batch)
|
||
if err != nil {
|
||
return fmt.Errorf("renew clients: %w", err)
|
||
}
|
||
if count > 0 {
|
||
logger.Debugf("%v clients renewed", count)
|
||
}
|
||
|
||
needRestart1, count, nodeIDs, err := s.disableInvalidClients(tx, batch)
|
||
if err != nil {
|
||
return fmt.Errorf("disable invalid clients: %w", err)
|
||
}
|
||
if count > 0 {
|
||
logger.Debugf("%v clients disabled", count)
|
||
disabledClientsCount = count
|
||
}
|
||
|
||
needRestart2, count, err := s.disableInvalidInbounds(tx, batch)
|
||
if err != nil {
|
||
return fmt.Errorf("disable invalid inbounds: %w", err)
|
||
}
|
||
if count > 0 {
|
||
logger.Debugf("%v inbounds disabled", count)
|
||
}
|
||
if err := batch.markNodesTx(tx); err != nil {
|
||
return err
|
||
}
|
||
needRestart = needRestart0 || needRestart1 || needRestart2
|
||
clientsDisabled = disabledClientsCount > 0
|
||
disabledNodeIDs = nodeIDs
|
||
return nil
|
||
})
|
||
if err != nil {
|
||
logger.Warning("traffic lifecycle maintenance failed after traffic commit:", err)
|
||
return false, false, nil, nil, nil, nil
|
||
}
|
||
needRestart = needRestart || s.applyTrafficMutationBatch(batch)
|
||
return needRestart, clientsDisabled, disabledNodeIDs, batch.remotePlans, batch.renewedEmails, nil
|
||
}
|
||
|
||
func (s *InboundService) addInboundTraffic(tx *gorm.DB, traffics []*xray.Traffic) error {
|
||
if len(traffics) == 0 {
|
||
return nil
|
||
}
|
||
|
||
var err error
|
||
|
||
for _, traffic := range traffics {
|
||
if traffic.IsInbound {
|
||
err = tx.Model(&model.Inbound{}).Where("tag = ? AND node_id IS NULL", traffic.Tag).
|
||
Updates(map[string]any{
|
||
"up": gorm.Expr(database.ClampedAddExpr("up"), traffic.Up),
|
||
"down": gorm.Expr(database.ClampedAddExpr("down"), traffic.Down),
|
||
}).Error
|
||
if err != nil {
|
||
return err
|
||
}
|
||
}
|
||
}
|
||
return nil
|
||
}
|
||
|
||
func (s *InboundService) addClientTraffic(tx *gorm.DB, traffics []*xray.ClientTraffic) (err error) {
|
||
if len(traffics) == 0 {
|
||
return nil
|
||
}
|
||
traffics, err = canonicalizeClientTraffic(tx, traffics)
|
||
if err != nil {
|
||
return fmt.Errorf("resolve client traffic identities: %w", err)
|
||
}
|
||
|
||
emails := make([]string, 0, len(traffics))
|
||
for _, traffic := range traffics {
|
||
emails = append(emails, traffic.Email)
|
||
}
|
||
dbClientTraffics := make([]*xray.ClientTraffic, 0, len(traffics))
|
||
// Match purely by email. client_traffics is email-keyed (one shared row per
|
||
// email regardless of how many inbounds the client is attached to), and these
|
||
// emails come from the local xray's report, so they always belong to a client
|
||
// attached to a local inbound. The old `inbound_id NOT IN (node inbounds)`
|
||
// filter dropped the local traffic of a client attached to both a node and the
|
||
// mother inbound whenever the node inbound happened to be attached first — its
|
||
// shared row then carried the node inbound's id (AddClientStat used to use
|
||
// OnConflict DoNothing and never refreshed it; it now refreshes inbound_id on
|
||
// conflict, but this filter was removed rather than relying on that ordering).
|
||
err = tx.Model(xray.ClientTraffic{}).
|
||
Where("email IN (?)", emails).
|
||
Order("id").
|
||
Find(&dbClientTraffics).Error
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
// Avoid empty slice error
|
||
if len(dbClientTraffics) == 0 {
|
||
return nil
|
||
}
|
||
|
||
dbClientTraffics, convertedExpiryByEmail, err := s.adjustTraffics(tx, dbClientTraffics)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
// Index by email for O(N) merge.
|
||
trafficByEmail := make(map[string]*xray.ClientTraffic, len(traffics))
|
||
for i := range traffics {
|
||
if traffics[i] != nil {
|
||
trafficByEmail[traffics[i].Email] = traffics[i]
|
||
}
|
||
}
|
||
now := time.Now().UnixMilli()
|
||
// Use atomic per-row UPDATE instead of read-modify-write Save. tx.Save
|
||
// issues UPDATEs in slice order, which varies between concurrent callers;
|
||
// on PostgreSQL two transactions locking the same rows in opposite order
|
||
// deadlock. An atomic "SET up = up + ?" never holds a row lock across a
|
||
// subsequent lock acquisition, so concurrent writers cannot deadlock.
|
||
for _, ct := range dbClientTraffics {
|
||
t, ok := trafficByEmail[ct.Email]
|
||
if !ok || (t.Up == 0 && t.Down == 0) {
|
||
continue
|
||
}
|
||
if err = tx.Exec(
|
||
fmt.Sprintf(
|
||
`UPDATE client_traffics SET up = %s, down = %s, last_online = %s WHERE email = ?`,
|
||
database.ClampedAddExpr("up"),
|
||
database.ClampedAddExpr("down"),
|
||
database.GreatestExpr("last_online", "?"),
|
||
),
|
||
t.Up, t.Down, now, ct.Email,
|
||
).Error; err != nil {
|
||
return fmt.Errorf("update traffic for %s: %w", ct.Email, err)
|
||
}
|
||
}
|
||
|
||
// adjustTraffics converts delayed-start rows (negative ExpiryTime → absolute
|
||
// deadline) in-memory. Persist that conversion now since the traffic UPDATE
|
||
// above only touches up/down/last_online. Only converted emails are written:
|
||
// updating every polled row issued one no-op UPDATE per active client per
|
||
// poll. Sorted order keeps concurrent writers lock-compatible on Postgres.
|
||
for _, email := range slices.Sorted(maps.Keys(convertedExpiryByEmail)) {
|
||
if err = tx.Exec(
|
||
`UPDATE client_traffics SET expiry_time = ? WHERE email = ? AND expiry_time < 0`,
|
||
convertedExpiryByEmail[email], email,
|
||
).Error; err != nil {
|
||
return fmt.Errorf("update expiry time for %s: %w", email, err)
|
||
}
|
||
}
|
||
|
||
return nil
|
||
}
|
||
|
||
func canonicalizeClientTraffic(tx *gorm.DB, traffics []*xray.ClientTraffic) ([]*xray.ClientTraffic, error) {
|
||
byEmail := make(map[string]*xray.ClientTraffic, len(traffics))
|
||
for _, traffic := range traffics {
|
||
if traffic == nil {
|
||
continue
|
||
}
|
||
email := traffic.Email
|
||
if traffic.TuicTrafficID > 0 {
|
||
var owner xray.ClientTraffic
|
||
err := tx.Select("email").Where("id = ?", traffic.TuicTrafficID).Take(&owner).Error
|
||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||
continue
|
||
}
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
email = owner.Email
|
||
}
|
||
current := byEmail[email]
|
||
if current == nil {
|
||
copy := *traffic
|
||
copy.Email = email
|
||
copy.TuicUUID = ""
|
||
copy.TuicInboundId = 0
|
||
byEmail[email] = ©
|
||
continue
|
||
}
|
||
current.Up += traffic.Up
|
||
current.Down += traffic.Down
|
||
current.Enable = current.Enable || traffic.Enable
|
||
}
|
||
result := make([]*xray.ClientTraffic, 0, len(byEmail))
|
||
for _, traffic := range byEmail {
|
||
result = append(result, traffic)
|
||
}
|
||
return result, nil
|
||
}
|
||
|
||
func (s *InboundService) adjustTraffics(tx *gorm.DB, dbClientTraffics []*xray.ClientTraffic) ([]*xray.ClientTraffic, map[string]int64, error) {
|
||
now := time.Now().UnixMilli()
|
||
|
||
// "Start After First Use" stores a negative expiry (the duration). On the
|
||
// first traffic tick it becomes an absolute deadline of now+duration. Compute
|
||
// it once per email so every inbound the client is attached to lands on the
|
||
// same value (recomputing per inbound would skip all but the first one).
|
||
newExpiryByEmail := make(map[string]int64, len(dbClientTraffics))
|
||
for traffic_index := range dbClientTraffics {
|
||
if dbClientTraffics[traffic_index].ExpiryTime < 0 {
|
||
newExpiryByEmail[dbClientTraffics[traffic_index].Email] = now - dbClientTraffics[traffic_index].ExpiryTime
|
||
}
|
||
}
|
||
if len(newExpiryByEmail) == 0 {
|
||
return dbClientTraffics, nil, nil
|
||
}
|
||
|
||
delayedEmails := make([]string, 0, len(newExpiryByEmail))
|
||
for email := range newExpiryByEmail {
|
||
delayedEmails = append(delayedEmails, email)
|
||
}
|
||
|
||
// Resolve the owning inbounds through the client_inbounds link, which is
|
||
// authoritative. client_traffics.inbound_id goes stale when an inbound is
|
||
// deleted and recreated, which would leave the negative expiry unconverted.
|
||
var inboundIds []int
|
||
err := tx.Table("client_inbounds").
|
||
Joins("JOIN clients ON clients.id = client_inbounds.client_id").
|
||
Where("clients.email IN (?)", delayedEmails).
|
||
Distinct().
|
||
Pluck("client_inbounds.inbound_id", &inboundIds).Error
|
||
if err != nil {
|
||
return nil, nil, err
|
||
}
|
||
if len(inboundIds) == 0 {
|
||
return dbClientTraffics, nil, nil
|
||
}
|
||
|
||
var inbounds []*model.Inbound
|
||
err = tx.Model(model.Inbound{}).Where("id IN (?)", inboundIds).Find(&inbounds).Error
|
||
if err != nil {
|
||
return nil, nil, err
|
||
}
|
||
for inbound_index := range inbounds {
|
||
settings := map[string]any{}
|
||
_ = json.Unmarshal([]byte(inbounds[inbound_index].Settings), &settings)
|
||
clients, ok := settings["clients"].([]any)
|
||
if ok {
|
||
var newClients []any
|
||
for client_index := range clients {
|
||
c := clients[client_index].(map[string]any)
|
||
email, _ := c["email"].(string)
|
||
if newExpiry, ok := newExpiryByEmail[email]; ok {
|
||
c["expiryTime"] = newExpiry
|
||
c["updated_at"] = now
|
||
}
|
||
if _, ok := c["created_at"]; !ok {
|
||
c["created_at"] = now
|
||
}
|
||
if _, ok := c["updated_at"]; !ok {
|
||
c["updated_at"] = now
|
||
}
|
||
newClients = append(newClients, any(c))
|
||
}
|
||
settings["clients"] = newClients
|
||
modifiedSettings, err := json.MarshalIndent(settings, "", " ")
|
||
if err != nil {
|
||
return nil, nil, err
|
||
}
|
||
|
||
inbounds[inbound_index].Settings = string(modifiedSettings)
|
||
}
|
||
}
|
||
|
||
for traffic_index := range dbClientTraffics {
|
||
if newExpiry, ok := newExpiryByEmail[dbClientTraffics[traffic_index].Email]; ok {
|
||
dbClientTraffics[traffic_index].ExpiryTime = newExpiry
|
||
}
|
||
}
|
||
|
||
err = tx.Save(inbounds).Error
|
||
if err != nil {
|
||
logger.Warning("AddClientTraffic update inbounds ", err)
|
||
logger.Error(inbounds)
|
||
} else {
|
||
for _, ib := range inbounds {
|
||
if ib == nil {
|
||
continue
|
||
}
|
||
cs, gcErr := s.GetClients(ib)
|
||
if gcErr != nil {
|
||
logger.Warning("AddClientTraffic sync clients: GetClients failed", gcErr)
|
||
continue
|
||
}
|
||
if syncErr := s.clientService.SyncInbound(tx, ib.Id, cs); syncErr != nil {
|
||
logger.Warning("AddClientTraffic sync clients: SyncInbound failed", syncErr)
|
||
}
|
||
}
|
||
}
|
||
|
||
return dbClientTraffics, newExpiryByEmail, nil
|
||
}
|
||
|
||
// apiUserFromClient prepares a stored client object for the runtime AddUser
|
||
// call. The copy matters twice over: the stored object keeps being mutated and
|
||
// marshalled back into the inbound's settings, which must not gain an API-only
|
||
// key, and shadowsocks clients carry no cipher of their own — it lives on the
|
||
// inbound, and without it the API cannot tell which of xray's two shadowsocks
|
||
// account types the running inbound expects.
|
||
func apiUserFromClient(client map[string]any, cipher string) map[string]any {
|
||
user := maps.Clone(client)
|
||
if user == nil {
|
||
user = map[string]any{}
|
||
}
|
||
if cipher != "" {
|
||
user["cipher"] = cipher
|
||
}
|
||
return user
|
||
}
|
||
|
||
// Candidates and renewals are not the same set: a skipped candidate keeps its
|
||
// counters, so only the clients actually reset may lose their cross-panel rows.
|
||
func (s *InboundService) autoRenewClients(tx *gorm.DB, mutationBatch *trafficMutationBatch) (bool, int64, error) {
|
||
// check for time expired
|
||
var traffics []*xray.ClientTraffic
|
||
now := time.Now().Unix() * 1000
|
||
var err error
|
||
|
||
// Filter to clients that have at least one local inbound. Using
|
||
// client_traffics.inbound_id is wrong: it goes stale after an inbound is
|
||
// deleted/recreated and always points to the first inbound the client was
|
||
// attached to, so it could be a node inbound even when the client also has
|
||
// local inbounds. The email-based join through client_inbounds is authoritative.
|
||
err = tx.Model(xray.ClientTraffic{}).
|
||
Where("(reset > 0 or reset_day > 0 or reset_weekday > 0) and expiry_time > 0 and expiry_time <= ?", now).
|
||
// A prepaid plan stops itself: once as many renewals have fired as the
|
||
// operator allowed, the client is left to expire like any other.
|
||
Where("reset_max <= 0 or reset_count < reset_max").
|
||
Where("email IN (?)", tx.Table("client_inbounds ci").
|
||
Select("c.email").
|
||
Joins("JOIN clients c ON c.id = ci.client_id").
|
||
Joins("JOIN inbounds i ON i.id = ci.inbound_id").
|
||
Where("i.node_id IS NULL")).
|
||
Find(&traffics).Error
|
||
if err != nil {
|
||
return false, 0, err
|
||
}
|
||
// return if there is no client to renew
|
||
if len(traffics) == 0 {
|
||
return false, 0, nil
|
||
}
|
||
|
||
renewLocation, locErr := (&SettingService{}).GetTimeLocation()
|
||
if locErr != nil || renewLocation == nil {
|
||
// Falling back to UTC keeps renewals happening; the alternative is
|
||
// skipping them entirely because a setting could not be read.
|
||
logger.Warning("autoRenewClients: could not read the panel time zone, using UTC:", locErr)
|
||
renewLocation = time.UTC
|
||
}
|
||
|
||
var inbound_ids []int
|
||
var inbounds []*model.Inbound
|
||
needRestart := false
|
||
type inboundClientKey struct {
|
||
inboundID int
|
||
email string
|
||
}
|
||
var clientsToAdd []struct {
|
||
inbound model.Inbound
|
||
client map[string]any
|
||
}
|
||
clientsToAddSet := make(map[inboundClientKey]struct{})
|
||
|
||
// Resolve the inbounds to renew through the client_inbounds link rather than
|
||
// client_traffics.inbound_id, which goes stale after an inbound is deleted and
|
||
// recreated and would otherwise skip the renew entirely.
|
||
renewEmails := make([]string, 0, len(traffics))
|
||
for _, traffic := range traffics {
|
||
renewEmails = append(renewEmails, traffic.Email)
|
||
}
|
||
for _, batch := range chunkStrings(renewEmails, sqliteMaxVars) {
|
||
var ids []int
|
||
if err = tx.Table("client_inbounds").
|
||
Joins("JOIN clients ON clients.id = client_inbounds.client_id").
|
||
Where("clients.email IN ?", batch).
|
||
Distinct().
|
||
Pluck("client_inbounds.inbound_id", &ids).Error; err != nil {
|
||
return false, 0, err
|
||
}
|
||
inbound_ids = append(inbound_ids, ids...)
|
||
}
|
||
// Dedupe so an inbound hosting N expired clients is fetched and saved once
|
||
// per tick instead of N times across chunk boundaries.
|
||
inbound_ids = uniqueInts(inbound_ids)
|
||
// Chunked to stay under SQLite's bind-variable limit when many inbounds
|
||
// are touched in a single tick.
|
||
for _, batch := range chunkInts(inbound_ids, sqliteMaxVars) {
|
||
var page []*model.Inbound
|
||
if err = tx.Model(model.Inbound{}).Where("id IN ?", batch).Find(&page).Error; err != nil {
|
||
return false, 0, err
|
||
}
|
||
inbounds = append(inbounds, page...)
|
||
}
|
||
// Index the expired traffics by email so each client is an O(1) lookup
|
||
// instead of a linear scan of every expired row (O(clients × expired) per
|
||
// inbound, quadratic at scale). Pointers keep the in-place mutation below.
|
||
trafficByEmail := make(map[string]*xray.ClientTraffic, len(traffics))
|
||
// Keep the pre-renewal quota state: the shared pointer becomes enabled while
|
||
// processing the first inbound, while an already-enabled row paired with
|
||
// disabled settings represents an operator-disabled client we must preserve.
|
||
trafficWasEnabled := make(map[string]bool, len(traffics))
|
||
for i := range traffics {
|
||
trafficByEmail[traffics[i].Email] = traffics[i]
|
||
trafficWasEnabled[traffics[i].Email] = traffics[i].Enable
|
||
}
|
||
renewedEmails := make([]string, 0, len(traffics))
|
||
for inbound_index := range inbounds {
|
||
settings := map[string]any{}
|
||
_ = json.Unmarshal([]byte(inbounds[inbound_index].Settings), &settings)
|
||
clients, _ := settings["clients"].([]any)
|
||
if len(clients) == 0 {
|
||
continue
|
||
}
|
||
cipher := ""
|
||
if inbounds[inbound_index].Protocol == model.Shadowsocks {
|
||
cipher, _ = settings["method"].(string)
|
||
}
|
||
for client_index := range clients {
|
||
c := clients[client_index].(map[string]any)
|
||
email, _ := c["email"].(string)
|
||
traffic, ok := trafficByEmail[email]
|
||
if !ok {
|
||
continue
|
||
}
|
||
// One allowance per period, not per tick: a client away for three
|
||
// cycles must not catch up three of them against a prepaid cap.
|
||
newExpiryTime, renewals := catchUpClientRenewal(traffic, now, renewLocation)
|
||
if renewals > 0 {
|
||
traffic.ExpiryTime = newExpiryTime
|
||
traffic.ResetCount += renewals
|
||
}
|
||
c["expiryTime"] = traffic.ExpiryTime
|
||
if traffic.ExpiryTime <= now {
|
||
// Cap ran out mid-catch-up and the client is still expired: enabling it
|
||
// for disableInvalidClients to undo adds and removes an xray user for nothing.
|
||
clients[client_index] = any(c)
|
||
continue
|
||
}
|
||
if renewals > 0 {
|
||
traffic.Down = 0
|
||
traffic.Up = 0
|
||
renewedEmails = append(renewedEmails, email)
|
||
}
|
||
if !trafficWasEnabled[email] {
|
||
traffic.Enable = true
|
||
c["enable"] = true
|
||
key := inboundClientKey{inboundID: inbounds[inbound_index].Id, email: email}
|
||
if _, planned := clientsToAddSet[key]; !planned {
|
||
clientsToAddSet[key] = struct{}{}
|
||
clientsToAdd = append(clientsToAdd,
|
||
struct {
|
||
inbound model.Inbound
|
||
client map[string]any
|
||
}{
|
||
inbound: *inbounds[inbound_index],
|
||
client: apiUserFromClient(c, cipher),
|
||
})
|
||
}
|
||
}
|
||
clients[client_index] = any(c)
|
||
}
|
||
settings["clients"] = clients
|
||
newSettings, err := json.MarshalIndent(settings, "", " ")
|
||
if err != nil {
|
||
return false, 0, err
|
||
}
|
||
inbounds[inbound_index].Settings = string(newSettings)
|
||
}
|
||
err = tx.Save(inbounds).Error
|
||
if err != nil {
|
||
return false, 0, err
|
||
}
|
||
for _, ib := range inbounds {
|
||
if ib == nil {
|
||
continue
|
||
}
|
||
cs, gcErr := s.GetClients(ib)
|
||
if gcErr != nil {
|
||
logger.Warning("autoRenewClients sync clients: GetClients failed", gcErr)
|
||
continue
|
||
}
|
||
if syncErr := s.clientService.SyncInbound(tx, ib.Id, cs); syncErr != nil {
|
||
logger.Warning("autoRenewClients sync clients: SyncInbound failed", syncErr)
|
||
}
|
||
}
|
||
err = tx.Save(traffics).Error
|
||
if err != nil {
|
||
return false, 0, err
|
||
}
|
||
// A renewed client starts a fresh quota window: drop the cross-panel rows
|
||
// too, or the stale pushed totals would re-deplete it immediately.
|
||
if err = clearGlobalTraffic(tx, renewedEmails...); err != nil {
|
||
return false, 0, err
|
||
}
|
||
mutationBatch.renewedEmails = append(mutationBatch.renewedEmails, renewedEmails...)
|
||
for _, clientToAdd := range clientsToAdd {
|
||
if clientToAdd.inbound.NodeID != nil {
|
||
mutationBatch.addNode(*clientToAdd.inbound.NodeID)
|
||
continue
|
||
}
|
||
mutationBatch.localPlans = append(mutationBatch.localPlans, trafficLocalApplyPlan{
|
||
action: trafficAddUser, inbound: clientToAdd.inbound, client: clientToAdd.client,
|
||
})
|
||
}
|
||
return needRestart, int64(len(renewedEmails)), nil
|
||
}
|
||
|
||
// AddClientStat inserts a per-client accounting row, or refreshes the
|
||
// config-derived columns on an email conflict. Xray reports traffic per
|
||
// email, so the surviving row also acts as the shared accumulator for
|
||
// inbounds that re-use the same identity — every call for that identity
|
||
// (one per attached inbound) carries the same enable/expiry/reset/total,
|
||
// so re-asserting them here is idempotent for that legitimate case.
|
||
//
|
||
// The conflict path matters on its own for a second reason: an inbound
|
||
// delete detaches its clients (InboundService.DelInbound) without deleting
|
||
// their client_traffics row, by design — mirroring ClientService.Detach,
|
||
// which intentionally leaves a fully-detached client's row in place so a
|
||
// later Attach can resume it with its accumulated traffic intact. If that
|
||
// same email is instead reused for a freshly (re)created client, the new
|
||
// config's enable/expiry/reset/total must win over whatever the orphaned
|
||
// row still holds; DoNothing left them stale indefinitely (#5958).
|
||
//
|
||
// up/down are deliberately excluded from the refresh: they are the
|
||
// accumulated traffic totals, and zeroing them here would erase real usage
|
||
// every time an existing, actively-used client is attached to one more
|
||
// inbound. One tradeoff this does not resolve: a genuinely new client that
|
||
// happens to reuse an orphaned email still inherits that row's leftover
|
||
// up/down, since nothing at this call site can tell the two cases apart.
|
||
func (s *InboundService) AddClientStat(tx *gorm.DB, inboundId int, client *model.Client) error {
|
||
if err := validateClientRenewal(*client); err != nil {
|
||
return err
|
||
}
|
||
clientTraffic := xray.ClientTraffic{
|
||
InboundId: inboundId,
|
||
Email: client.Email,
|
||
Total: client.TotalGB,
|
||
ExpiryTime: client.ExpiryTime,
|
||
Enable: client.Enable,
|
||
Reset: client.Reset,
|
||
ResetDay: client.ResetDay,
|
||
ResetWeekday: client.ResetWeekday,
|
||
ResetMax: client.ResetMax,
|
||
}
|
||
return tx.Clauses(clause.OnConflict{
|
||
Columns: []clause.Column{{Name: "email"}},
|
||
DoUpdates: clause.AssignmentColumns([]string{"inbound_id", "total", "expiry_time", "enable", "reset", "reset_day", "reset_weekday", "reset_max"}),
|
||
}).Create(&clientTraffic).Error
|
||
}
|
||
|
||
func (s *InboundService) UpdateClientStat(tx *gorm.DB, email string, client *model.Client) error {
|
||
if err := validateClientRenewal(*client); err != nil {
|
||
return err
|
||
}
|
||
result := tx.Model(xray.ClientTraffic{}).
|
||
Where("email = ?", email).
|
||
Updates(map[string]any{
|
||
"enable": client.Enable,
|
||
"email": client.Email,
|
||
"total": client.TotalGB,
|
||
"expiry_time": client.ExpiryTime,
|
||
"reset": client.Reset,
|
||
"reset_day": client.ResetDay,
|
||
"reset_weekday": client.ResetWeekday,
|
||
"reset_max": client.ResetMax,
|
||
})
|
||
err := result.Error
|
||
return err
|
||
}
|
||
|
||
func (s *InboundService) DelClientStat(tx *gorm.DB, email string) error {
|
||
if err := adjustGroupBaselinesForRemovedTraffic(tx, []string{email}); err != nil {
|
||
return err
|
||
}
|
||
if err := tx.Where("email = ?", email).Delete(xray.ClientTraffic{}).Error; err != nil {
|
||
return err
|
||
}
|
||
if err := clearGlobalTraffic(tx, email); err != nil {
|
||
return err
|
||
}
|
||
return tx.Where("email = ?", email).Delete(&model.NodeClientTraffic{}).Error
|
||
}
|
||
|
||
func (s *InboundService) delClientStatsByEmails(tx *gorm.DB, emails []string) error {
|
||
if err := adjustGroupBaselinesForRemovedTraffic(tx, emails); err != nil {
|
||
return err
|
||
}
|
||
const chunk = 400
|
||
for start := 0; start < len(emails); start += chunk {
|
||
end := min(start+chunk, len(emails))
|
||
batch := emails[start:end]
|
||
if err := tx.Where("email IN ?", batch).Delete(xray.ClientTraffic{}).Error; err != nil {
|
||
return err
|
||
}
|
||
if err := tx.Where("email IN ?", batch).Delete(&model.ClientGlobalTraffic{}).Error; err != nil {
|
||
return err
|
||
}
|
||
if err := tx.Where("email IN ?", batch).Delete(&model.NodeClientTraffic{}).Error; err != nil {
|
||
return err
|
||
}
|
||
}
|
||
return nil
|
||
}
|
||
|
||
func (s *InboundService) ResetClientTrafficByEmail(clientEmail string) error {
|
||
err := submitTrafficWrite(func() error {
|
||
return database.GetDB().Transaction(func(tx *gorm.DB) error {
|
||
if err := adjustGroupBaselinesForRemovedTraffic(tx, []string{clientEmail}); err != nil {
|
||
return err
|
||
}
|
||
if err := clearGlobalTraffic(tx, clientEmail); err != nil {
|
||
return err
|
||
}
|
||
if err := tx.Model(xray.ClientTraffic{}).
|
||
Where("email = ?", clientEmail).
|
||
Updates(map[string]any{"enable": true, "up": 0, "down": 0}).Error; err != nil {
|
||
return err
|
||
}
|
||
return tx.Where("email = ?", clientEmail).Delete(&model.NodeClientTraffic{}).Error
|
||
})
|
||
})
|
||
if err == nil {
|
||
s.resetMtprotoClientQuota(clientEmail)
|
||
}
|
||
return err
|
||
}
|
||
|
||
func (s *InboundService) ResetClientTraffic(id int, clientEmail string) (needRestart bool, err error) {
|
||
var ownNode *int
|
||
err = submitTrafficWrite(func() error {
|
||
var inner error
|
||
needRestart, ownNode, inner = s.resetClientTrafficLocked(id, clientEmail)
|
||
return inner
|
||
})
|
||
if err == nil {
|
||
s.resetMtprotoClientQuota(clientEmail)
|
||
// Siblings on other nodes are delivered by their own inbound's reset.
|
||
if ownNode != nil {
|
||
s.deliverNodeResetsNow([]int{*ownNode})
|
||
}
|
||
}
|
||
return
|
||
}
|
||
|
||
func (s *InboundService) resetClientTrafficLocked(id int, clientEmail string) (bool, *int, error) {
|
||
needRestart := false
|
||
var reenablePlan *trafficLocalApplyPlan
|
||
var reenableNodeID *int
|
||
|
||
traffic, err := s.GetClientTrafficByEmail(clientEmail)
|
||
if err != nil {
|
||
return false, nil, err
|
||
}
|
||
|
||
if !traffic.Enable {
|
||
inbound, err := s.GetInbound(id)
|
||
if err != nil {
|
||
return false, nil, err
|
||
}
|
||
clients, err := s.GetClients(inbound)
|
||
if err != nil {
|
||
return false, nil, err
|
||
}
|
||
for _, client := range clients {
|
||
if client.Email == clientEmail && client.Enable {
|
||
cipher := ""
|
||
if string(inbound.Protocol) == "shadowsocks" {
|
||
var oldSettings map[string]any
|
||
err = json.Unmarshal([]byte(inbound.Settings), &oldSettings)
|
||
if err != nil {
|
||
return false, nil, err
|
||
}
|
||
cipher, _ = oldSettings["method"].(string)
|
||
}
|
||
clientMap := map[string]any{
|
||
"email": client.Email,
|
||
"id": client.ID,
|
||
"auth": client.Auth,
|
||
"security": client.Security,
|
||
"flow": client.Flow,
|
||
"password": client.Password,
|
||
"cipher": cipher,
|
||
"reverse": client.Reverse,
|
||
}
|
||
if inbound.NodeID != nil {
|
||
reenableNodeID = inbound.NodeID
|
||
} else {
|
||
reenablePlan = &trafficLocalApplyPlan{action: trafficAddUser, inbound: *inbound, client: clientMap}
|
||
}
|
||
break
|
||
}
|
||
}
|
||
}
|
||
|
||
traffic.Up = 0
|
||
traffic.Down = 0
|
||
traffic.Enable = true
|
||
|
||
db := database.GetDB()
|
||
now := time.Now().UnixMilli()
|
||
inbound, err := s.GetInbound(id)
|
||
if err != nil {
|
||
return false, nil, err
|
||
}
|
||
if err := db.Transaction(func(tx *gorm.DB) error {
|
||
if err := adjustGroupBaselinesForRemovedTraffic(tx, []string{clientEmail}); err != nil {
|
||
return err
|
||
}
|
||
if err := tx.Save(traffic).Error; err != nil {
|
||
return err
|
||
}
|
||
if err := clearGlobalTraffic(tx, clientEmail); err != nil {
|
||
return err
|
||
}
|
||
if err := tx.Where("email = ?", clientEmail).Delete(&model.NodeClientTraffic{}).Error; err != nil {
|
||
return err
|
||
}
|
||
if _, err := queueNodeResets(tx, []string{clientEmail}); err != nil {
|
||
return err
|
||
}
|
||
if err := tx.Model(model.Inbound{}).
|
||
Where("id = ?", id).
|
||
Update("last_traffic_reset_time", now).Error; err != nil {
|
||
return err
|
||
}
|
||
if reenableNodeID != nil {
|
||
return (&NodeService{}).MarkNodeDirtyTx(tx, *reenableNodeID)
|
||
}
|
||
if inbound != nil && inbound.NodeID != nil {
|
||
return (&NodeService{}).MarkNodeDirtyTx(tx, *inbound.NodeID)
|
||
}
|
||
return nil
|
||
}); err != nil {
|
||
return false, nil, err
|
||
}
|
||
|
||
if reenablePlan != nil {
|
||
rt, err := s.runtimeFor(&reenablePlan.inbound)
|
||
if err != nil {
|
||
needRestart = true
|
||
} else if err := rt.AddUser(context.Background(), &reenablePlan.inbound, reenablePlan.client); err != nil {
|
||
logger.Debug("Error in enabling client on", rt.Name(), ":", err)
|
||
needRestart = true
|
||
} else {
|
||
logger.Debug("Client enabled on", rt.Name(), "due to reset traffic:", clientEmail)
|
||
}
|
||
}
|
||
|
||
if inbound != nil {
|
||
return needRestart, inbound.NodeID, nil
|
||
}
|
||
return needRestart, nil, nil
|
||
}
|
||
|
||
func (s *InboundService) ResetAllTraffics() error {
|
||
err := submitTrafficWrite(func() error {
|
||
return s.resetAllTrafficsLocked()
|
||
})
|
||
if err == nil {
|
||
s.propagateResetAllTrafficsToNodes()
|
||
}
|
||
return err
|
||
}
|
||
|
||
func (s *InboundService) resetAllTrafficsLocked() error {
|
||
db := database.GetDB()
|
||
now := time.Now().UnixMilli()
|
||
|
||
return db.Model(model.Inbound{}).
|
||
Where("user_id > ?", 0).
|
||
Updates(map[string]any{
|
||
"up": 0,
|
||
"down": 0,
|
||
"last_traffic_reset_time": now,
|
||
}).Error
|
||
}
|
||
|
||
// propagateResetAllTrafficsToNodes tells every node to zero its own counters.
|
||
// Kept OUT of the traffic-writer transaction: each remote call can block up to
|
||
// remoteHTTPTimeout, and holding the single serial writer across N such calls
|
||
// stalls traffic accounting and drops the deltas of every concurrent poll.
|
||
func (s *InboundService) propagateResetAllTrafficsToNodes() {
|
||
nodes, err := (&NodeService{}).GetAll()
|
||
if err != nil {
|
||
return
|
||
}
|
||
ids := make([]int, len(nodes))
|
||
for i, node := range nodes {
|
||
ids[i] = node.Id
|
||
}
|
||
fanoutInboundResults(ids, nodeFanoutConcurrency, func(i int) struct{} {
|
||
if rt, err := runtime.GetManager().RuntimeFor(&ids[i]); err == nil {
|
||
if e := rt.ResetAllTraffics(context.Background()); e != nil {
|
||
logger.Warning("ResetAllTraffics: remote propagation to", rt.Name(), "failed:", e)
|
||
}
|
||
}
|
||
return struct{}{}
|
||
})
|
||
}
|
||
|
||
func (s *InboundService) ResetInboundTraffic(id int) error {
|
||
var inbound *model.Inbound
|
||
if err := submitTrafficWrite(func() error {
|
||
db := database.GetDB()
|
||
if err := db.Model(model.Inbound{}).
|
||
Where("id = ?", id).
|
||
Updates(map[string]any{"up": 0, "down": 0}).Error; err != nil {
|
||
return err
|
||
}
|
||
var err error
|
||
inbound, err = s.GetInbound(id)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
return nil
|
||
}); err != nil {
|
||
return err
|
||
}
|
||
if inbound != nil && inbound.NodeID != nil {
|
||
if rt, rterr := s.runtimeFor(inbound); rterr == nil {
|
||
if e := rt.ResetInboundTraffic(context.Background(), inbound); e != nil {
|
||
logger.Warning("ResetInboundTraffic: remote propagation to", rt.Name(), "failed:", e)
|
||
}
|
||
} else {
|
||
logger.Warning("ResetInboundTraffic: runtime lookup failed:", rterr)
|
||
}
|
||
}
|
||
return nil
|
||
}
|
||
|
||
func (s *InboundService) DelDepletedClients(id int) (err error) {
|
||
db := database.GetDB()
|
||
var deletedInbounds []model.Inbound
|
||
err = db.Transaction(func(tx *gorm.DB) error {
|
||
// Collect depleted emails globally — a shared-email row owned by one
|
||
// inbound depletes every sibling that lists the email.
|
||
now := time.Now().Unix() * 1000
|
||
depletedClause := depletedClientsClause
|
||
var depletedRows []xray.ClientTraffic
|
||
if err := tx.Model(xray.ClientTraffic{}).
|
||
Where(depletedClause, now).
|
||
Find(&depletedRows).Error; err != nil {
|
||
return err
|
||
}
|
||
if len(depletedRows) == 0 {
|
||
return nil
|
||
}
|
||
|
||
depletedEmails := make(map[string]struct{}, len(depletedRows))
|
||
for _, r := range depletedRows {
|
||
if r.Email == "" {
|
||
continue
|
||
}
|
||
depletedEmails[strings.ToLower(r.Email)] = struct{}{}
|
||
}
|
||
if len(depletedEmails) == 0 {
|
||
return nil
|
||
}
|
||
|
||
var inbounds []*model.Inbound
|
||
inboundQuery := tx.Model(model.Inbound{})
|
||
if id >= 0 {
|
||
inboundQuery = inboundQuery.Where("id = ?", id)
|
||
}
|
||
if err := inboundQuery.Find(&inbounds).Error; err != nil {
|
||
return err
|
||
}
|
||
|
||
for _, inbound := range inbounds {
|
||
var settings map[string]any
|
||
if err := json.Unmarshal([]byte(inbound.Settings), &settings); err != nil {
|
||
return err
|
||
}
|
||
rawClients, ok := settings["clients"].([]any)
|
||
if !ok {
|
||
continue
|
||
}
|
||
newClients := make([]any, 0, len(rawClients))
|
||
removed := 0
|
||
for _, client := range rawClients {
|
||
c, ok := client.(map[string]any)
|
||
if !ok {
|
||
newClients = append(newClients, client)
|
||
continue
|
||
}
|
||
email, _ := c["email"].(string)
|
||
if _, isDepleted := depletedEmails[strings.ToLower(email)]; isDepleted {
|
||
removed++
|
||
continue
|
||
}
|
||
newClients = append(newClients, client)
|
||
}
|
||
if removed == 0 {
|
||
continue
|
||
}
|
||
if len(newClients) == 0 {
|
||
deletedInbounds = append(deletedInbounds, *inbound)
|
||
if err := s.clientService.DetachInbound(tx, inbound.Id); err != nil {
|
||
return err
|
||
}
|
||
if err := tx.Where("inbound_id = ?", inbound.Id).Delete(&model.Host{}).Error; err != nil {
|
||
return err
|
||
}
|
||
if err := tx.Delete(model.Inbound{}, inbound.Id).Error; err != nil {
|
||
return err
|
||
}
|
||
if inbound.NodeID != nil {
|
||
if err := (&NodeService{}).MarkNodeDirtyTx(tx, *inbound.NodeID); err != nil {
|
||
return err
|
||
}
|
||
}
|
||
continue
|
||
}
|
||
settings["clients"] = newClients
|
||
ns, mErr := json.MarshalIndent(settings, "", " ")
|
||
if mErr != nil {
|
||
return mErr
|
||
}
|
||
inbound.Settings = string(ns)
|
||
if err := tx.Save(inbound).Error; err != nil {
|
||
return err
|
||
}
|
||
survivingClients, gcErr := s.GetClients(inbound)
|
||
if gcErr != nil {
|
||
return gcErr
|
||
}
|
||
if err := s.clientService.SyncInbound(tx, inbound.Id, survivingClients); err != nil {
|
||
return err
|
||
}
|
||
if inbound.NodeID != nil {
|
||
if err := (&NodeService{}).MarkNodeDirtyTx(tx, *inbound.NodeID); err != nil {
|
||
return err
|
||
}
|
||
}
|
||
}
|
||
|
||
// Drop now-orphaned rows. With id >= 0, a row is safe to drop only when
|
||
// no out-of-scope inbound still references the email.
|
||
if id < 0 {
|
||
return tx.Where(depletedClause, now).Delete(xray.ClientTraffic{}).Error
|
||
}
|
||
emails := make([]string, 0, len(depletedEmails))
|
||
for e := range depletedEmails {
|
||
emails = append(emails, e)
|
||
}
|
||
var stillReferenced []string
|
||
emailExpr := database.JSONFieldText("client.value", "email")
|
||
stillQuery := fmt.Sprintf(
|
||
"SELECT DISTINCT LOWER(%s) %s WHERE LOWER(%s) IN ?",
|
||
emailExpr,
|
||
database.JSONClientsFromInbound(),
|
||
emailExpr,
|
||
)
|
||
if err := tx.Raw(stillQuery, emails).Scan(&stillReferenced).Error; err != nil {
|
||
return err
|
||
}
|
||
stillSet := make(map[string]struct{}, len(stillReferenced))
|
||
for _, e := range stillReferenced {
|
||
stillSet[e] = struct{}{}
|
||
}
|
||
toDelete := make([]string, 0, len(emails))
|
||
for _, e := range emails {
|
||
if _, kept := stillSet[e]; !kept {
|
||
toDelete = append(toDelete, e)
|
||
}
|
||
}
|
||
if len(toDelete) > 0 {
|
||
if err := tx.Where("LOWER(email) IN ?", toDelete).Delete(xray.ClientTraffic{}).Error; err != nil {
|
||
return err
|
||
}
|
||
}
|
||
return nil
|
||
})
|
||
if err != nil {
|
||
return err
|
||
}
|
||
for i := range deletedInbounds {
|
||
inbound := &deletedInbounds[i]
|
||
if rt, rtErr := s.runtimeFor(inbound); rtErr != nil {
|
||
logger.Warning("DelDepletedClients: runtime lookup failed after commit:", rtErr)
|
||
} else if rtErr = rt.DelInbound(context.Background(), inbound); rtErr != nil && !xray.IsMissingHandlerErr(rtErr) {
|
||
logger.Warning("DelDepletedClients: runtime cleanup failed after commit:", rtErr)
|
||
}
|
||
if inbound.Tag != "" {
|
||
if _, syncErr := (&XraySettingService{}).RemoveInboundTagReferences(inbound.Tag); syncErr != nil {
|
||
logger.Warning("DelDepletedClients: routing cleanup failed after commit:", syncErr)
|
||
}
|
||
}
|
||
}
|
||
return nil
|
||
}
|
||
|
||
func (s *InboundService) GetClientTrafficTgBot(tgId int64) ([]*xray.ClientTraffic, error) {
|
||
db := database.GetDB()
|
||
|
||
idQuery := fmt.Sprintf(
|
||
"SELECT DISTINCT inbounds.id %s WHERE %s = ?",
|
||
database.JSONClientsFromInbound(),
|
||
database.JSONFieldText("client.value", "tgId"),
|
||
)
|
||
var inboundIds []int
|
||
if err := db.Raw(idQuery, strconv.FormatInt(tgId, 10)).Scan(&inboundIds).Error; err != nil {
|
||
logger.Errorf("Error retrieving inbounds with tgId %d: %v", tgId, err)
|
||
return nil, err
|
||
}
|
||
|
||
var inbounds []*model.Inbound
|
||
if len(inboundIds) > 0 {
|
||
err := db.Model(model.Inbound{}).Where("id IN ?", inboundIds).Find(&inbounds).Error
|
||
if err != nil && !errors.Is(err, gorm.ErrRecordNotFound) {
|
||
logger.Errorf("Error retrieving inbounds with tgId %d: %v", tgId, err)
|
||
return nil, err
|
||
}
|
||
}
|
||
|
||
var emails []string
|
||
for _, inbound := range inbounds {
|
||
clients, err := s.GetClients(inbound)
|
||
if err != nil {
|
||
logger.Errorf("Error retrieving clients for inbound %d: %v", inbound.Id, err)
|
||
continue
|
||
}
|
||
for _, client := range clients {
|
||
if client.TgID == tgId {
|
||
emails = append(emails, client.Email)
|
||
}
|
||
}
|
||
}
|
||
|
||
// Chunked to stay under SQLite's bind-variable limit when a single Telegram
|
||
// account owns thousands of clients across inbounds.
|
||
uniqEmails := uniqueNonEmptyStrings(emails)
|
||
traffics := make([]*xray.ClientTraffic, 0, len(uniqEmails))
|
||
for _, batch := range chunkStrings(uniqEmails, sqliteMaxVars) {
|
||
var page []*xray.ClientTraffic
|
||
if err := db.Model(xray.ClientTraffic{}).Where("email IN ?", batch).Find(&page).Error; err != nil {
|
||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||
continue
|
||
}
|
||
logger.Errorf("Error retrieving ClientTraffic for emails %v: %v", batch, err)
|
||
return nil, err
|
||
}
|
||
traffics = append(traffics, page...)
|
||
}
|
||
if len(traffics) == 0 {
|
||
logger.Warning("No ClientTraffic records found for emails:", emails)
|
||
return nil, nil
|
||
}
|
||
|
||
// Populate UUID and other client data for each traffic record
|
||
for i := range traffics {
|
||
if ct, client, e := s.GetClientByEmail(traffics[i].Email); e == nil && ct != nil && client != nil {
|
||
traffics[i].Enable = client.Enable
|
||
traffics[i].UUID = client.ID
|
||
traffics[i].SubId = client.SubID
|
||
}
|
||
}
|
||
|
||
return traffics, nil
|
||
}
|
||
|
||
// BumpClientsLastOnline sets client_traffics.last_online to now for the given
|
||
// emails. Used in online-API mode for clients that hold a live connection but
|
||
// moved no bytes this poll — the traffic path (addClientTraffic) only bumps
|
||
// last_online on a non-zero delta, so idle-but-connected clients would
|
||
// otherwise show a stale "last online" while being reported online.
|
||
func (s *InboundService) BumpClientsLastOnline(emails []string) error {
|
||
uniq := uniqueNonEmptyStrings(emails)
|
||
if len(uniq) == 0 {
|
||
return nil
|
||
}
|
||
now := time.Now().UnixMilli()
|
||
return submitTrafficWrite(func() error {
|
||
db := database.GetDB()
|
||
for _, batch := range chunkStrings(uniq, sqliteMaxVars) {
|
||
if err := db.Model(xray.ClientTraffic{}).Where("email IN ?", batch).Update("last_online", now).Error; err != nil {
|
||
return err
|
||
}
|
||
}
|
||
return nil
|
||
})
|
||
}
|
||
|
||
func (s *InboundService) GetActiveClientTraffics(emails []string) ([]*xray.ClientTraffic, error) {
|
||
uniq := uniqueNonEmptyStrings(emails)
|
||
if len(uniq) == 0 {
|
||
return nil, nil
|
||
}
|
||
db := database.GetDB()
|
||
traffics := make([]*xray.ClientTraffic, 0, len(uniq))
|
||
for _, batch := range chunkStrings(uniq, sqliteMaxVars) {
|
||
var page []*xray.ClientTraffic
|
||
if err := db.Model(xray.ClientTraffic{}).Where("email IN ?", batch).Find(&page).Error; err != nil {
|
||
return nil, err
|
||
}
|
||
traffics = append(traffics, page...)
|
||
}
|
||
overlayGlobalTraffic(db, traffics)
|
||
return traffics, nil
|
||
}
|
||
|
||
// GetAllClientTraffics returns the full set of client_traffics rows so the
|
||
// websocket broadcasters can ship a complete snapshot every cycle. A pure
|
||
// delta path silently dropped the per-client section whenever no client moved
|
||
// bytes in the cycle or a node sync failed, leaving client rows in the UI
|
||
// stuck at stale numbers — so small installs broadcast this snapshot, and only
|
||
// above the traffic job's snapshot threshold (where the marshaled snapshot
|
||
// would exceed the hub's payload cap and be dropped wholesale) does the job
|
||
// fall back to active-row deltas.
|
||
func (s *InboundService) GetAllClientTraffics() ([]*xray.ClientTraffic, error) {
|
||
db := database.GetDB()
|
||
var traffics []*xray.ClientTraffic
|
||
if err := db.Model(xray.ClientTraffic{}).Find(&traffics).Error; err != nil {
|
||
return nil, err
|
||
}
|
||
overlayGlobalTraffic(db, traffics)
|
||
return traffics, nil
|
||
}
|
||
|
||
func (s *InboundService) CountClientTraffics() (int64, error) {
|
||
db := database.GetDB()
|
||
var count int64
|
||
err := db.Model(xray.ClientTraffic{}).Count(&count).Error
|
||
return count, err
|
||
}
|
||
|
||
type InboundTrafficSummary struct {
|
||
Id int `json:"id" example:"1"`
|
||
Up int64 `json:"up" example:"1048576"`
|
||
Down int64 `json:"down" example:"2097152"`
|
||
Total int64 `json:"total" example:"10737418240"`
|
||
Enable bool `json:"enable" example:"true"`
|
||
}
|
||
|
||
func (s *InboundService) GetInboundsTrafficSummary() ([]InboundTrafficSummary, error) {
|
||
db := database.GetDB()
|
||
var summaries []InboundTrafficSummary
|
||
if err := db.Model(&model.Inbound{}).
|
||
Select("id, up, down, total, enable").
|
||
Find(&summaries).Error; err != nil {
|
||
return nil, err
|
||
}
|
||
return summaries, nil
|
||
}
|
||
|
||
func (s *InboundService) GetClientTrafficByEmail(email string) (traffic *xray.ClientTraffic, err error) {
|
||
db := database.GetDB()
|
||
var traffics []*xray.ClientTraffic
|
||
if err := db.Model(xray.ClientTraffic{}).Where("email = ?", email).Find(&traffics).Error; err != nil {
|
||
logger.Warningf("Error retrieving ClientTraffic with email %s: %v", email, err)
|
||
return nil, err
|
||
}
|
||
if len(traffics) == 0 {
|
||
return nil, nil
|
||
}
|
||
overlayGlobalTraffic(db, traffics)
|
||
t := traffics[0]
|
||
|
||
if rec, rErr := s.clientService.GetRecordByEmail(db, email); rErr == nil && rec != nil {
|
||
c := rec.ToClient()
|
||
t.UUID = c.ID
|
||
t.SubId = c.SubID
|
||
return t, nil
|
||
}
|
||
|
||
t2, client, err := s.GetClientByEmail(email)
|
||
if err != nil {
|
||
logger.Warningf("Error retrieving ClientTraffic with email %s: %v", email, err)
|
||
return nil, err
|
||
}
|
||
if t2 != nil && client != nil {
|
||
t2.UUID = client.ID
|
||
t2.SubId = client.SubID
|
||
return t2, nil
|
||
}
|
||
return nil, nil
|
||
}
|
||
|
||
func (s *InboundService) UpdateClientTrafficByEmail(email string, upload int64, download int64) error {
|
||
return submitTrafficWrite(func() error {
|
||
db := database.GetDB()
|
||
err := db.Model(xray.ClientTraffic{}).
|
||
Where("email = ?", email).
|
||
Updates(map[string]any{
|
||
"up": upload,
|
||
"down": download,
|
||
}).Error
|
||
if err != nil {
|
||
logger.Warningf("Error updating ClientTraffic with email %s: %v", email, err)
|
||
}
|
||
return err
|
||
})
|
||
}
|
||
|
||
func (s *InboundService) SearchClientTraffic(query string) (traffic *xray.ClientTraffic, err error) {
|
||
db := database.GetDB()
|
||
inbound := &model.Inbound{}
|
||
traffic = &xray.ClientTraffic{}
|
||
|
||
// Search for inbound settings that contain the query
|
||
err = db.Model(model.Inbound{}).Where("settings LIKE ?", "%\""+query+"\"%").First(inbound).Error
|
||
if err != nil {
|
||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||
logger.Warningf("Inbound settings containing query %s not found: %v", query, err)
|
||
return nil, err
|
||
}
|
||
logger.Errorf("Error searching for inbound settings with query %s: %v", query, err)
|
||
return nil, err
|
||
}
|
||
|
||
traffic.InboundId = inbound.Id
|
||
|
||
clients, err := ParseInboundSettingsClients(inbound.Settings)
|
||
if err != nil {
|
||
logger.Errorf("Error unmarshalling inbound settings for inbound ID %d: %v", inbound.Id, err)
|
||
return nil, err
|
||
}
|
||
|
||
for _, client := range clients {
|
||
if (client.ID == query || client.Password == query) && client.Email != "" {
|
||
traffic.Email = client.Email
|
||
break
|
||
}
|
||
}
|
||
|
||
if traffic.Email == "" {
|
||
logger.Warningf("No client found with query %s in inbound ID %d", query, inbound.Id)
|
||
return nil, gorm.ErrRecordNotFound
|
||
}
|
||
|
||
// Retrieve ClientTraffic based on the found email
|
||
err = db.Model(xray.ClientTraffic{}).Where("email = ?", traffic.Email).First(traffic).Error
|
||
if err != nil {
|
||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||
logger.Warningf("ClientTraffic for email %s not found: %v", traffic.Email, err)
|
||
return nil, err
|
||
}
|
||
logger.Errorf("Error retrieving ClientTraffic for email %s: %v", traffic.Email, err)
|
||
return nil, err
|
||
}
|
||
|
||
return traffic, nil
|
||
}
|