mirror of
https://github.com/MHSanaei/3x-ui.git
synced 2026-09-30 19:22:10 +03:00
8f1201553e
* fix(ip-limit): CAS-retry inbound_client_ips merges under Postgres Two writers RMW the same ips JSON blob; on PostgreSQL a lost update drops remote IPs that partitionLiveIps only sees through that blob (#6587). Compare-and-set on the previous blob with re-merge on miss, matching the repo's conditional Where+RowsAffected pattern. Fixes #6587 * ci: retrigger frontend after npm registry maintenance The frontend job failed solely on `npm audit` while registry.npmjs.org returned 503 (Service Under Maintenance). Lint, typecheck, vitest, vite build, and storybook all passed. Local `npm audit --omit=dev --audit-level=high` now reports 0 vulnerabilities. * test(ip-limit): cover the scan's CAS against a mid-scan node sync The job-side compare-and-set had no test of its own. A write injected between the scan's read and its update now has to keep the node's remote IP; main's blind Save drops it. Also keeps the new comments to two lines, as CLAUDE.md requires. --------- Co-authored-by: mrchatam <mrchatam@users.noreply.github.com> Co-authored-by: MHSanaei <ho3ein.sanaei@gmail.com>
336 lines
11 KiB
Go
336 lines
11 KiB
Go
package service
|
|
|
|
import (
|
|
"encoding/json"
|
|
"fmt"
|
|
"sort"
|
|
"time"
|
|
|
|
"github.com/mhsanaei/3x-ui/v3/internal/database"
|
|
"github.com/mhsanaei/3x-ui/v3/internal/database/model"
|
|
|
|
"gorm.io/gorm"
|
|
"gorm.io/gorm/clause"
|
|
)
|
|
|
|
func (s *InboundService) GetAllInboundClientIps() ([]model.InboundClientIps, error) {
|
|
db := database.GetDB()
|
|
var ips []model.InboundClientIps
|
|
err := db.Model(&model.InboundClientIps{}).Find(&ips).Error
|
|
return ips, err
|
|
}
|
|
|
|
// nodeHostedEmails is every client one node serves, its descendants' included.
|
|
// Per-node pushes are scoped to it so their cost tracks the node, not the fleet.
|
|
func nodeHostedEmails(db *gorm.DB, nodeID int) ([]string, error) {
|
|
var emails []string
|
|
err := db.Model(&model.NodeClientTraffic{}).Where("node_id = ?", nodeID).Pluck("email", &emails).Error
|
|
return emails, err
|
|
}
|
|
|
|
// GetNodeInboundClientIps returns the IP rows of the clients nodeID hosts: a node's
|
|
// IP-limit job reads no other row, so pushing the rest only made it echo them back.
|
|
func (s *InboundService) GetNodeInboundClientIps(nodeID int) ([]model.InboundClientIps, error) {
|
|
db := database.GetDB()
|
|
emails, err := nodeHostedEmails(db, nodeID)
|
|
if err != nil || len(emails) == 0 {
|
|
return nil, err
|
|
}
|
|
var ips []model.InboundClientIps
|
|
for _, batch := range chunkStrings(emails, sqlInChunk) {
|
|
var page []model.InboundClientIps
|
|
if err := db.Where("client_email IN ?", batch).Find(&page).Error; err != nil {
|
|
return nil, err
|
|
}
|
|
ips = append(ips, page...)
|
|
}
|
|
return ips, nil
|
|
}
|
|
|
|
// clientIpStaleAfterSeconds mirrors job.ipStaleAfterSeconds: client IPs older than
|
|
// 30 minutes are evicted. Applying the same cutoff inside the cross-node merge keeps
|
|
// the synced blob bounded and stops the master's push-back from resurrecting IPs that
|
|
// a node has already pruned (otherwise the merge defeats the eviction cluster-wide).
|
|
const clientIpStaleAfterSeconds = int64(30 * 60)
|
|
|
|
// clientIpEntry is the on-disk shape of each element of InboundClientIps.Ips. Tags
|
|
// match job.IPWithTimestamp so the blob round-trips with the access.log scanner.
|
|
type clientIpEntry struct {
|
|
IP string `json:"ip"`
|
|
Timestamp int64 `json:"timestamp"`
|
|
}
|
|
|
|
// mergeClientIpEntries unions old and incoming IP observations, dropping anything
|
|
// older than cutoff, keeping the most recent timestamp per IP, and returning the
|
|
// result sorted newest-first.
|
|
func mergeClientIpEntries(old, incoming []clientIpEntry, cutoff int64) []clientIpEntry {
|
|
ipMap := make(map[string]int64, len(old)+len(incoming))
|
|
for _, e := range old {
|
|
if e.Timestamp < cutoff {
|
|
continue
|
|
}
|
|
ipMap[e.IP] = e.Timestamp
|
|
}
|
|
for _, e := range incoming {
|
|
if e.Timestamp < cutoff {
|
|
continue
|
|
}
|
|
if cur, ok := ipMap[e.IP]; !ok || e.Timestamp > cur {
|
|
ipMap[e.IP] = e.Timestamp
|
|
}
|
|
}
|
|
out := make([]clientIpEntry, 0, len(ipMap))
|
|
for ip, ts := range ipMap {
|
|
out = append(out, clientIpEntry{IP: ip, Timestamp: ts})
|
|
}
|
|
sort.Slice(out, func(i, j int) bool { return out[i].Timestamp > out[j].Timestamp })
|
|
return out
|
|
}
|
|
|
|
// MergeInboundClientIps folds client IPs synced from another node into the local
|
|
// inbound_client_ips table without double-counting an IP seen on multiple nodes and
|
|
// without resurrecting stale entries. Existing rows are updated in place; brand-new
|
|
// clients (typically node-only clients with no local row) are created with a fresh
|
|
// local id.
|
|
func (s *InboundService) MergeInboundClientIps(incomingIps []model.InboundClientIps) error {
|
|
db := database.GetDB()
|
|
var currentIps []model.InboundClientIps
|
|
if err := db.Model(&model.InboundClientIps{}).Find(¤tIps).Error; err != nil {
|
|
return err
|
|
}
|
|
|
|
currentMap := make(map[string]*model.InboundClientIps, len(currentIps))
|
|
for i := range currentIps {
|
|
currentMap[currentIps[i].ClientEmail] = ¤tIps[i]
|
|
}
|
|
|
|
now := time.Now().Unix()
|
|
cutoff := now - clientIpStaleAfterSeconds
|
|
|
|
// Node syncs run concurrently (one goroutine per node) and shared clients
|
|
// appear in several nodes' reports. Locking rows in each node's arbitrary
|
|
// report order lets two merges grab the same rows in opposite order, which
|
|
// Postgres aborts as a deadlock (40P01) — take them in one global order.
|
|
sort.Slice(incomingIps, func(i, j int) bool {
|
|
return incomingIps[i].ClientEmail < incomingIps[j].ClientEmail
|
|
})
|
|
|
|
tx := db.Begin()
|
|
defer func() {
|
|
if r := recover(); r != nil {
|
|
tx.Rollback()
|
|
}
|
|
}()
|
|
|
|
for _, incoming := range incomingIps {
|
|
if incoming.ClientEmail == "" || incoming.Ips == "" {
|
|
continue
|
|
}
|
|
|
|
var incomingEntries []clientIpEntry
|
|
_ = json.Unmarshal([]byte(incoming.Ips), &incomingEntries)
|
|
|
|
current, exists := currentMap[incoming.ClientEmail]
|
|
if !exists {
|
|
// New client we've never seen locally. Drop stale entries up front and
|
|
// skip the row entirely if nothing is fresh, so we don't persist a row
|
|
// that is dead on arrival.
|
|
fresh := mergeClientIpEntries(nil, incomingEntries, cutoff)
|
|
if len(fresh) == 0 {
|
|
continue
|
|
}
|
|
b, _ := json.Marshal(fresh)
|
|
incoming.Ips = string(b)
|
|
// Never carry the remote node's primary key into the local table: id
|
|
// spaces are independent across nodes and the remote id would collide
|
|
// with an unrelated local row. OnConflict guards the race where
|
|
// check_client_ip_job creates the same brand-new email between the
|
|
// snapshot above and this insert.
|
|
incoming.Id = 0
|
|
if err := tx.Clauses(clause.OnConflict{
|
|
Columns: []clause.Column{{Name: "client_email"}},
|
|
DoNothing: true,
|
|
}).Create(&incoming).Error; err != nil {
|
|
tx.Rollback()
|
|
return err
|
|
}
|
|
continue
|
|
}
|
|
|
|
// check_client_ip_job rewrites this blob too; an unconditional Update loses
|
|
// whichever writer commits second and its remote IPs with it (#6587).
|
|
if err := mergeExistingClientIps(tx, current.Id, incomingEntries, cutoff); err != nil {
|
|
tx.Rollback()
|
|
return err
|
|
}
|
|
}
|
|
return tx.Commit().Error
|
|
}
|
|
|
|
// ClientIpCasRetries bounds the re-reads after losing a compare-and-set. Running
|
|
// out is an error, so the caller retries on its next schedule instead of dropping IPs.
|
|
const ClientIpCasRetries = 8
|
|
|
|
// CasUpdateInboundClientIps writes newIps only while the row still holds expectedIps;
|
|
// updated=false with a nil error means another writer won and the caller must re-read.
|
|
func CasUpdateInboundClientIps(tx *gorm.DB, id int, expectedIps, newIps string) (updated bool, err error) {
|
|
res := tx.Model(&model.InboundClientIps{}).
|
|
Where("id = ? AND ips = ?", id, expectedIps).
|
|
Update("ips", newIps)
|
|
if res.Error != nil {
|
|
return false, res.Error
|
|
}
|
|
return res.RowsAffected == 1, nil
|
|
}
|
|
|
|
// mergeExistingClientIps folds incoming into the row at id under a CAS loop so
|
|
// a concurrent check_client_ip_job write cannot erase the merge (#6587).
|
|
func mergeExistingClientIps(tx *gorm.DB, id int, incoming []clientIpEntry, cutoff int64) error {
|
|
for attempt := 0; attempt < ClientIpCasRetries; attempt++ {
|
|
var row model.InboundClientIps
|
|
if err := tx.Where("id = ?", id).First(&row).Error; err != nil {
|
|
return err
|
|
}
|
|
var oldEntries []clientIpEntry
|
|
if row.Ips != "" {
|
|
_ = json.Unmarshal([]byte(row.Ips), &oldEntries)
|
|
}
|
|
merged := mergeClientIpEntries(oldEntries, incoming, cutoff)
|
|
b, _ := json.Marshal(merged)
|
|
mergedStr := string(b)
|
|
if row.Ips == mergedStr {
|
|
return nil
|
|
}
|
|
ok, err := CasUpdateInboundClientIps(tx, id, row.Ips, mergedStr)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if ok {
|
|
return nil
|
|
}
|
|
}
|
|
return fmt.Errorf("inbound_client_ips id=%d: exhausted CAS retries merging client IPs", id)
|
|
}
|
|
|
|
func (s *InboundService) UpdateClientIPs(tx *gorm.DB, oldEmail string, newEmail string) error {
|
|
// The caller only renames onto a free identity, so a row already sitting on
|
|
// newEmail is stale tracking data — drop it instead of failing the edit.
|
|
if oldEmail != newEmail {
|
|
if err := tx.Where("client_email = ?", newEmail).Delete(model.InboundClientIps{}).Error; err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return tx.Model(model.InboundClientIps{}).Where("client_email = ?", oldEmail).Update("client_email", newEmail).Error
|
|
}
|
|
|
|
func (s *InboundService) DelClientIPs(tx *gorm.DB, email string) error {
|
|
return tx.Where("client_email = ?", email).Delete(model.InboundClientIps{}).Error
|
|
}
|
|
|
|
func (s *InboundService) delClientIPsByEmails(tx *gorm.DB, emails []string) error {
|
|
const chunk = 400
|
|
for start := 0; start < len(emails); start += chunk {
|
|
end := min(start+chunk, len(emails))
|
|
if err := tx.Where("client_email IN ?", emails[start:end]).Delete(model.InboundClientIps{}).Error; err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (s *InboundService) GetInboundClientIps(clientEmail string) (string, error) {
|
|
db := database.GetDB()
|
|
InboundClientIps := &model.InboundClientIps{}
|
|
err := db.Model(model.InboundClientIps{}).Where("client_email = ?", clientEmail).First(InboundClientIps).Error
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
|
|
if InboundClientIps.Ips == "" {
|
|
return "", nil
|
|
}
|
|
|
|
// Try to parse as new format (with timestamps)
|
|
type IPWithTimestamp struct {
|
|
IP string `json:"ip"`
|
|
Timestamp int64 `json:"timestamp"`
|
|
}
|
|
|
|
var ipsWithTime []IPWithTimestamp
|
|
err = json.Unmarshal([]byte(InboundClientIps.Ips), &ipsWithTime)
|
|
|
|
// If successfully parsed as new format, return with timestamps
|
|
if err == nil && len(ipsWithTime) > 0 {
|
|
return InboundClientIps.Ips, nil
|
|
}
|
|
|
|
// Otherwise, assume it's old format (simple string array)
|
|
// Try to parse as simple array and convert to new format
|
|
var oldIps []string
|
|
err = json.Unmarshal([]byte(InboundClientIps.Ips), &oldIps)
|
|
if err == nil && len(oldIps) > 0 {
|
|
// Convert old format to new format with current timestamp
|
|
newIpsWithTime := make([]IPWithTimestamp, len(oldIps))
|
|
for i, ip := range oldIps {
|
|
newIpsWithTime[i] = IPWithTimestamp{
|
|
IP: ip,
|
|
Timestamp: time.Now().Unix(),
|
|
}
|
|
}
|
|
result, _ := json.Marshal(newIpsWithTime)
|
|
return string(result), nil
|
|
}
|
|
|
|
// Return as-is if parsing fails
|
|
return InboundClientIps.Ips, nil
|
|
}
|
|
|
|
func (s *InboundService) ClearClientIps(clientEmail string) error {
|
|
db := database.GetDB()
|
|
|
|
result := db.Model(model.InboundClientIps{}).
|
|
Where("client_email = ?", clientEmail).
|
|
Update("ips", "")
|
|
err := result.Error
|
|
if err != nil {
|
|
return err
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// PruneStaleClientIps enforces clientIpStaleAfterSeconds for rows the online
|
|
// scan no longer rewrites: an offline client's addresses must still expire.
|
|
func (s *InboundService) PruneStaleClientIps() error {
|
|
db := database.GetDB()
|
|
cutoff := time.Now().Unix() - clientIpStaleAfterSeconds
|
|
|
|
var rows []model.InboundClientIps
|
|
if err := db.Find(&rows).Error; err != nil {
|
|
return err
|
|
}
|
|
for _, row := range rows {
|
|
var entries []clientIpEntry
|
|
if row.Ips != "" {
|
|
// Legacy blobs without timestamps stay untouched; the next scan rewrites them.
|
|
if err := json.Unmarshal([]byte(row.Ips), &entries); err != nil {
|
|
continue
|
|
}
|
|
}
|
|
kept := mergeClientIpEntries(nil, entries, cutoff)
|
|
if len(kept) == 0 {
|
|
if err := db.Delete(&model.InboundClientIps{}, row.Id).Error; err != nil {
|
|
return err
|
|
}
|
|
continue
|
|
}
|
|
if len(kept) == len(entries) {
|
|
continue
|
|
}
|
|
b, _ := json.Marshal(kept)
|
|
if err := db.Model(&model.InboundClientIps{}).Where("id = ?", row.Id).Update("ips", string(b)).Error; err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return pruneStaleNodeClientIps(cutoff)
|
|
}
|