Files
3x-ui/internal/web/runtime/local.go
T
Kuzz007 f78dfa6f67 feat(amneziawg): swap the app's integration points to the embedded manager
Hard cutover, part 2: every real call site that used to drive
internal/amneziawg's kernel-module Manager now drives
internal/amneziawgnet's instead --

- internal/web/job/amneziawg_job.go: the reconcile cron job. Traffic/
  online-status accounting is dropped entirely (not ported) -- once a
  peer's traffic is relayed through Xray's own SOCKS5 inbound, it's an
  ordinary Xray user and XrayTrafficJob's existing generic stats polling
  already handles it, with zero AmneziaWG-specific code.
- internal/web/runtime/local.go: the immediate-apply CRUD path
  (AddInbound/DelInbound/updateAmneziaWGInbound).
- internal/web/web.go: panel shutdown's StopAll.

internal/amneziawgnet.Manager gains Remove(id) to match the kernel-module
Manager's shape at these call sites (Reconcile alone doesn't cover a
single-inbound removal outside a full reconcile pass).

internal/web/service/inbound_amneziawg.go's applyLocalAmneziaWG needed no
change: it already goes through runtime.Runtime.UpdateInbound, which now
resolves to the updated local.go path.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-08-02 15:04:51 +03:00

263 lines
7.0 KiB
Go

package runtime
import (
"context"
"encoding/json"
"errors"
"strconv"
"strings"
"sync"
"github.com/mhsanaei/3x-ui/v3/internal/amneziawg"
"github.com/mhsanaei/3x-ui/v3/internal/amneziawgnet"
"github.com/mhsanaei/3x-ui/v3/internal/database/model"
"github.com/mhsanaei/3x-ui/v3/internal/mtproto"
"github.com/mhsanaei/3x-ui/v3/internal/xray"
)
type LocalDeps struct {
APIPort func() int
SetNeedRestart func()
}
type Local struct {
deps LocalDeps
mu sync.Mutex
}
func NewLocal(deps LocalDeps) *Local {
return &Local{deps: deps}
}
func (l *Local) Name() string { return "local" }
func (l *Local) withAPI(fn func(api *xray.XrayAPI) error) error {
l.mu.Lock()
defer l.mu.Unlock()
port := l.deps.APIPort()
if port <= 0 {
return errors.New("local xray is not running")
}
var api xray.XrayAPI
if err := api.Init(port); err != nil {
return err
}
defer api.Close()
return fn(&api)
}
func (l *Local) AddInbound(_ context.Context, ib *model.Inbound) error {
if ib.Protocol == model.MTProto {
inst, ok := mtproto.InstanceFromInbound(ib)
if !ok {
return nil
}
return mtproto.GetManager().Ensure(inst)
}
if ib.Protocol == model.AmneziaWG {
inst, ok := amneziawg.InstanceFromInbound(ib)
if !ok {
return nil
}
return amneziawgnet.GetManager().Ensure(amneziawgnet.Desired{Instance: inst})
}
body, err := json.MarshalIndent(ib.GenXrayInboundConfig(), "", " ")
if err != nil {
return err
}
return l.withAPI(func(api *xray.XrayAPI) error {
return api.AddInbound(body)
})
}
func (l *Local) DelInbound(_ context.Context, ib *model.Inbound) error {
if ib.Protocol == model.MTProto {
mtproto.GetManager().Remove(ib.Id)
return nil
}
if ib.Protocol == model.AmneziaWG {
amneziawgnet.GetManager().Remove(ib.Id)
return nil
}
return l.withAPI(func(api *xray.XrayAPI) error {
return api.DelInbound(ib.Tag)
})
}
func (l *Local) UpdateInbound(ctx context.Context, oldIb, newIb *model.Inbound) error {
if oldIb.Protocol == model.MTProto || newIb.Protocol == model.MTProto {
return l.updateMtprotoInbound(ctx, oldIb, newIb)
}
if oldIb.Protocol == model.AmneziaWG || newIb.Protocol == model.AmneziaWG {
return l.updateAmneziaWGInbound(ctx, oldIb, newIb)
}
_ = l.DelInbound(ctx, oldIb)
if !newIb.Enable {
return nil
}
return l.AddInbound(ctx, newIb)
}
// updateMtprotoInbound applies an inbound update without the Del+Add sequence
// the xray path uses: Remove would drop the manager's fingerprint state, which
// is what lets Ensure keep the running mtg process (and its live connections)
// when nothing in the generated config changed. The sidecar is only stopped
// when the inbound is disabled, loses its last active secret, or moves to a
// different protocol.
func (l *Local) updateMtprotoInbound(ctx context.Context, oldIb, newIb *model.Inbound) error {
if oldIb.Protocol == model.MTProto && newIb.Protocol != model.MTProto {
mtproto.GetManager().Remove(oldIb.Id)
if !newIb.Enable {
return nil
}
return l.AddInbound(ctx, newIb)
}
if oldIb.Protocol != model.MTProto {
_ = l.DelInbound(ctx, oldIb)
}
if !newIb.Enable {
mtproto.GetManager().Remove(newIb.Id)
return nil
}
inst, ok := mtproto.InstanceFromInbound(newIb)
if !ok {
mtproto.GetManager().Remove(newIb.Id)
return nil
}
return mtproto.GetManager().Ensure(inst)
}
// updateAmneziaWGInbound mirrors updateMtprotoInbound: it skips the
// Remove+Ensure sequence a plain Del+Add would force so that, on an
// AmneziaWG-to-AmneziaWG edit, Manager.Ensure's own fingerprint comparison
// can reconfigure the running embedded Device in place via IpcSet instead
// of always rebuilding it (see internal/amneziawgnet.Manager.ensureLocked --
// only an address/MTU change forces a rebuild there, not a peer edit).
func (l *Local) updateAmneziaWGInbound(ctx context.Context, oldIb, newIb *model.Inbound) error {
if oldIb.Protocol == model.AmneziaWG && newIb.Protocol != model.AmneziaWG {
amneziawgnet.GetManager().Remove(oldIb.Id)
if !newIb.Enable {
return nil
}
return l.AddInbound(ctx, newIb)
}
if oldIb.Protocol != model.AmneziaWG {
_ = l.DelInbound(ctx, oldIb)
}
if !newIb.Enable {
amneziawgnet.GetManager().Remove(newIb.Id)
return nil
}
inst, ok := amneziawg.InstanceFromInbound(newIb)
if !ok {
amneziawgnet.GetManager().Remove(newIb.Id)
return nil
}
return amneziawgnet.GetManager().Ensure(amneziawgnet.Desired{Instance: inst})
}
func (l *Local) AddUser(_ context.Context, ib *model.Inbound, userMap map[string]any) error {
if ib.Protocol == model.MTProto || ib.Protocol == model.AmneziaWG {
return nil
}
return l.withAPI(func(api *xray.XrayAPI) error {
return api.AddUser(string(ib.Protocol), ib.Tag, userMap)
})
}
func (l *Local) RemoveUser(_ context.Context, ib *model.Inbound, email string) error {
if ib.Protocol == model.MTProto || ib.Protocol == model.AmneziaWG {
return nil
}
return l.withAPI(func(api *xray.XrayAPI) error {
return api.RemoveUser(ib.Tag, email)
})
}
func (l *Local) AddClient(ctx context.Context, ib *model.Inbound, client model.Client) error {
if !client.Enable {
return nil
}
user := map[string]any{
"email": client.Email,
"id": client.ID,
"security": client.Security,
"flow": client.Flow,
"auth": client.Auth,
"password": client.Password,
"publicKey": client.PublicKey,
"allowedIPs": client.AllowedIPs,
"preSharedKey": client.PreSharedKey,
"keepAlive": wgKeepAlive(client.KeepAlive),
}
return l.AddUser(ctx, ib, user)
}
func (l *Local) DeleteUser(ctx context.Context, ib *model.Inbound, email string) error {
if email == "" {
return nil
}
if err := l.RemoveUser(ctx, ib, email); err != nil {
if strings.Contains(err.Error(), "not found") {
return nil
}
return err
}
return nil
}
func (l *Local) DeleteClient(context.Context, string) error {
return nil
}
func (l *Local) UpdateUser(ctx context.Context, ib *model.Inbound, oldEmail string, payload model.Client) error {
if oldEmail != "" {
if err := l.RemoveUser(ctx, ib, oldEmail); err != nil && !strings.Contains(err.Error(), "not found") {
return err
}
}
if !payload.Enable {
return nil
}
user := map[string]any{
"email": payload.Email,
"id": payload.ID,
"security": payload.Security,
"flow": payload.Flow,
"auth": payload.Auth,
"password": payload.Password,
"publicKey": payload.PublicKey,
"allowedIPs": payload.AllowedIPs,
"preSharedKey": payload.PreSharedKey,
"keepAlive": wgKeepAlive(payload.KeepAlive),
}
return l.AddUser(ctx, ib, user)
}
func wgKeepAlive(seconds int) string {
if seconds <= 0 {
return ""
}
return strconv.Itoa(seconds)
}
func (l *Local) RestartXray(_ context.Context) error {
if l.deps.SetNeedRestart != nil {
l.deps.SetNeedRestart()
}
return nil
}
func (l *Local) ResetClientTraffic(_ context.Context, _ *model.Inbound, _ string) error {
return nil
}
func (l *Local) ResetAllTraffics(_ context.Context) error {
return nil
}
func (l *Local) ResetInboundTraffic(_ context.Context, _ *model.Inbound) error {
return nil
}