diff --git a/.gitignore b/.gitignore index 0f6fc6d..1a26f43 100644 --- a/.gitignore +++ b/.gitignore @@ -174,6 +174,7 @@ build/ *.dll *.so *.dylib +go-gost/gost your_app.exe # Go 测试二进制文件 diff --git a/README.md b/README.md index 21bf555..264e596 100644 --- a/README.md +++ b/README.md @@ -1,6 +1,6 @@ # 转发面板 -本项目基于 [go-gost/gost](https://github.com/go-gost/gost) 和 [go-gost/x](https://github.com/go-gost/x) 两个开源库,实现了转发面板。 +本项目基于 [go-gost/gost](https://github.com/go-gost/gost) 和 [go-gost/x](https://.com/go-gost/x) 两个开源库,实现了转发面板。 --- ## 文档地址 - [文档地址](https://tes.cc) @@ -73,9 +73,8 @@ ### Docker Compose部署 #### 快速部署 -不在提供gitee 自行解决github访问问题 ```bash -curl -L https://raw.githubusercontent.com/bqlpfy/forward-panel/refs/heads/main/panel_install.sh -o panel_install.sh && chmod +x panel_install.sh && ./panel_install.sh +curl -L https://file.tes.cc/panel_install.sh -o panel_install.sh && chmod +x panel_install.sh && ./panel_install.sh ``` diff --git a/beego-backend/beego-backend b/beego-backend/beego-backend deleted file mode 100755 index fa2e484..0000000 Binary files a/beego-backend/beego-backend and /dev/null differ diff --git a/beego-backend/common/database.go b/beego-backend/common/database.go deleted file mode 100644 index 3f7fe19..0000000 --- a/beego-backend/common/database.go +++ /dev/null @@ -1,96 +0,0 @@ -package common - -import ( - "database/sql" - "fmt" - "log" - - _ "github.com/mattn/go-sqlite3" -) - -var DB *sql.DB - -// InitDatabase 初始化sqlite3数据库 -func InitDatabase() error { - var err error - DB, err = sql.Open("sqlite3", "data.db") - if err != nil { - return fmt.Errorf("连接数据库失败: %v", err) - } - - if err = DB.Ping(); err != nil { - return fmt.Errorf("数据库连接测试失败: %v", err) - } - - // 创建表结构 - if err = createTables(); err != nil { - return fmt.Errorf("创建表结构失败: %v", err) - } - - return nil -} - -// createTables 创建表结构 -func createTables() error { - tables := []string{ - `CREATE TABLE IF NOT EXISTS users ( - id INTEGER PRIMARY KEY AUTOINCREMENT, - username VARCHAR(100) NOT NULL UNIQUE, - password VARCHAR(100) NOT NULL, - email VARCHAR(100), - is_active BOOLEAN NOT NULL DEFAULT 1 - )`, - } - - for _, table := range tables { - if _, err := DB.Exec(table); err != nil { - return err - } - } - - // 创建默认管理员用户 - //createDefaultUser() - - fmt.Println("✅ 数据库表创建完成") - return nil -} - -// createDefaultUser 创建默认管理员用户 -func createDefaultUser() { - var count int - err := DB.QueryRow("SELECT COUNT(*) FROM users").Scan(&count) - if err != nil || count > 0 { - return - } - - // 插入默认管理员用户 - _, err = DB.Exec(`INSERT INTO users (username, password, email, is_active) - VALUES (?, ?, ?, ?)`, - "admin", "123456", "admin@example.com", true) - - if err != nil { - log.Printf("创建默认用户失败: %v", err) - } else { - fmt.Println("✅ 创建默认管理员用户成功") - } -} - -// Insert 插入数据 -func Insert(sql string, args ...interface{}) (sql.Result, error) { - return DB.Exec(sql, args...) -} - -// Update 更新数据 -func Update(sql string, args ...interface{}) (sql.Result, error) { - return DB.Exec(sql, args...) -} - -// Delete 删除数据 -func Delete(sql string, args ...interface{}) (sql.Result, error) { - return DB.Exec(sql, args...) -} - -// Select 查询数据 -func Select(sql string, args ...interface{}) (*sql.Rows, error) { - return DB.Query(sql, args...) -} diff --git a/beego-backend/common/gost_util.go b/beego-backend/common/gost_util.go deleted file mode 100644 index 805d0c7..0000000 --- a/beego-backend/common/gost_util.go +++ /dev/null @@ -1 +0,0 @@ -package common diff --git a/beego-backend/common/res_util.go b/beego-backend/common/res_util.go deleted file mode 100644 index 3c63d65..0000000 --- a/beego-backend/common/res_util.go +++ /dev/null @@ -1,8 +0,0 @@ -package common - -type ResponseResult struct { - Code int64 `json:"code"` - Success bool `json:"success"` - Message string `json:"message"` - Data interface{} `json:"data"` -} diff --git a/beego-backend/common/web_socket_util.go b/beego-backend/common/web_socket_util.go deleted file mode 100644 index 805d0c7..0000000 --- a/beego-backend/common/web_socket_util.go +++ /dev/null @@ -1 +0,0 @@ -package common diff --git a/beego-backend/conf/app.conf b/beego-backend/conf/app.conf deleted file mode 100644 index 03748fc..0000000 --- a/beego-backend/conf/app.conf +++ /dev/null @@ -1,11 +0,0 @@ -appname = beego-backend -httpport = 8080 -runmode = dev -autorender = false -copyrequestbody = true -EnableDocs = true - -# 数据库配置 -db_type = sqlite3 -db_conn = data.db -db_debug = true diff --git a/beego-backend/controllers/base_controller.go b/beego-backend/controllers/base_controller.go deleted file mode 100644 index 3cb0481..0000000 --- a/beego-backend/controllers/base_controller.go +++ /dev/null @@ -1,30 +0,0 @@ -package controllers - -import ( - "encoding/json" - "github.com/beego/beego/v2/server/web" - "net/http" -) - -type BaseController struct { - web.Controller -} - -func (c *BaseController) ServeJSON() { - c.Ctx.Output.Header("Content-Type", "application/json; charset=utf-8") - - var err error - - data := c.Data["json"] - - encoder := json.NewEncoder(c.Ctx.Output.Context.ResponseWriter) - encoder.SetEscapeHTML(false) - - encoder.SetIndent("", " ") - err = encoder.Encode(data) - - if err != nil { - http.Error(c.Ctx.Output.Context.ResponseWriter, err.Error(), http.StatusInternalServerError) - return - } -} diff --git a/beego-backend/controllers/login_controller.go b/beego-backend/controllers/login_controller.go deleted file mode 100644 index ae8f9c6..0000000 --- a/beego-backend/controllers/login_controller.go +++ /dev/null @@ -1,32 +0,0 @@ -package controllers - -import ( - "beego-backend/common" - "beego-backend/models/users" - "encoding/json" - "fmt" -) - -type LoginController struct { - BaseController -} - -// @router /login_by_pwd [post] -func (c *LoginController) LoginByPwd() { - var Data users.LoginModel - err := json.Unmarshal(c.Ctx.Input.RequestBody, &Data) - if err != nil { - Result := common.ResponseResult{ - Code: -1, - Success: false, - Message: fmt.Sprintf("系统异常:%v", err.Error()), - Data: nil, - } - c.Data["json"] = &Result - c.ServeJSON() - return - } - result := users.LoginByPassword(Data) - c.Data["json"] = &result - c.ServeJSON() -} diff --git a/beego-backend/go.mod b/beego-backend/go.mod deleted file mode 100644 index 9578fa3..0000000 --- a/beego-backend/go.mod +++ /dev/null @@ -1,28 +0,0 @@ -module beego-backend - -go 1.24 - -require ( - github.com/beego/beego/v2 v2.3.8 - github.com/mattn/go-sqlite3 v1.14.28 -) - -require ( - github.com/beorn7/perks v1.0.1 // indirect - github.com/cespare/xxhash/v2 v2.2.0 // indirect - github.com/hashicorp/golang-lru v0.5.4 // indirect - github.com/kr/text v0.2.0 // indirect - github.com/mitchellh/mapstructure v1.5.0 // indirect - github.com/prometheus/client_golang v1.19.0 // indirect - github.com/prometheus/client_model v0.5.0 // indirect - github.com/prometheus/common v0.48.0 // indirect - github.com/prometheus/procfs v0.12.0 // indirect - github.com/shiena/ansicolor v0.0.0-20200904210342-c7312218db18 // indirect - github.com/valyala/bytebufferpool v1.0.0 // indirect - golang.org/x/crypto v0.24.0 // indirect - golang.org/x/net v0.23.0 // indirect - golang.org/x/sys v0.21.0 // indirect - golang.org/x/text v0.16.0 // indirect - google.golang.org/protobuf v1.34.2 // indirect - gopkg.in/yaml.v3 v3.0.1 // indirect -) diff --git a/beego-backend/go.sum b/beego-backend/go.sum deleted file mode 100644 index 7ddb0c5..0000000 --- a/beego-backend/go.sum +++ /dev/null @@ -1,59 +0,0 @@ -filippo.io/edwards25519 v1.1.0 h1:FNf4tywRC1HmFuKW5xopWpigGjJKiJSV0Cqo0cJWDaA= -filippo.io/edwards25519 v1.1.0/go.mod h1:BxyFTGdWcka3PhytdK4V28tE5sGfRvvvRV7EaN4VDT4= -github.com/beego/beego/v2 v2.3.8 h1:wplhB1pF4TxR+2SS4PUej8eDoH4xGfxuHfS7wAk9VBc= -github.com/beego/beego/v2 v2.3.8/go.mod h1:8vl9+RrXqvodrl9C8yivX1e6le6deCK6RWeq8R7gTTg= -github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= -github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= -github.com/cespare/xxhash/v2 v2.2.0 h1:DC2CZ1Ep5Y4k3ZQ899DldepgrayRUGE6BBZ/cd9Cj44= -github.com/cespare/xxhash/v2 v2.2.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= -github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E= -github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= -github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= -github.com/elazarl/go-bindata-assetfs v1.0.1 h1:m0kkaHRKEu7tUIUFVwhGGGYClXvyl4RE03qmvRTNfbw= -github.com/elazarl/go-bindata-assetfs v1.0.1/go.mod h1:v+YaWX3bdea5J/mo8dSETolEo7R71Vk1u8bnjau5yw4= -github.com/go-sql-driver/mysql v1.8.1 h1:LedoTUt/eveggdHS9qUFC1EFSa8bU2+1pZjSRpvNJ1Y= -github.com/go-sql-driver/mysql v1.8.1/go.mod h1:wEBSXgmK//2ZFJyE+qWnIsVGmvmEKlqwuVSjsCm7DZg= -github.com/google/go-cmp v0.6.0 h1:ofyhxvXcZhMsU5ulbFiLKl/XBFqE1GSq7atu8tAmTRI= -github.com/google/go-cmp v0.6.0/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= -github.com/hashicorp/golang-lru v0.5.4 h1:YDjusn29QI/Das2iO9M0BHnIbxPeyuCHsjMW+lJfyTc= -github.com/hashicorp/golang-lru v0.5.4/go.mod h1:iADmTwqILo4mZ8BN3D2Q6+9jd8WM5uGBxy+E8yxSoD4= -github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= -github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= -github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= -github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= -github.com/mattn/go-sqlite3 v1.14.28 h1:ThEiQrnbtumT+QMknw63Befp/ce/nUPgBPMlRFEum7A= -github.com/mattn/go-sqlite3 v1.14.28/go.mod h1:Uh1q+B4BYcTPb+yiD3kU8Ct7aC0hY9fxUwlHK0RXw+Y= -github.com/mitchellh/mapstructure v1.5.0 h1:jeMsZIYE/09sWLaz43PL7Gy6RuMjD2eJVyuac5Z2hdY= -github.com/mitchellh/mapstructure v1.5.0/go.mod h1:bFUtVrKA4DC2yAKiSyO/QUcy7e+RRV2QTWOzhPopBRo= -github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= -github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= -github.com/prometheus/client_golang v1.19.0 h1:ygXvpU1AoN1MhdzckN+PyD9QJOSD4x7kmXYlnfbA6JU= -github.com/prometheus/client_golang v1.19.0/go.mod h1:ZRM9uEAypZakd+q/x7+gmsvXdURP+DABIEIjnmDdp+k= -github.com/prometheus/client_model v0.5.0 h1:VQw1hfvPvk3Uv6Qf29VrPF32JB6rtbgI6cYPYQjL0Qw= -github.com/prometheus/client_model v0.5.0/go.mod h1:dTiFglRmd66nLR9Pv9f0mZi7B7fk5Pm3gvsjB5tr+kI= -github.com/prometheus/common v0.48.0 h1:QO8U2CdOzSn1BBsmXJXduaaW+dY/5QLjfB8svtSzKKE= -github.com/prometheus/common v0.48.0/go.mod h1:0/KsvlIEfPQCQ5I2iNSAWKPZziNCvRs5EC6ILDTlAPc= -github.com/prometheus/procfs v0.12.0 h1:jluTpSng7V9hY0O2R9DzzJHYb2xULk9VTR1V1R/k6Bo= -github.com/prometheus/procfs v0.12.0/go.mod h1:pcuDEFsWDnvcgNzo4EEweacyhjeA9Zk3cnaOZAZEfOo= -github.com/rogpeppe/go-internal v1.10.0 h1:TMyTOH3F/DB16zRVcYyreMH6GnZZrwQVAoYjRBZyWFQ= -github.com/rogpeppe/go-internal v1.10.0/go.mod h1:UQnix2H7Ngw/k4C5ijL5+65zddjncjaFoBhdsK/akog= -github.com/shiena/ansicolor v0.0.0-20200904210342-c7312218db18 h1:DAYUYH5869yV94zvCES9F51oYtN5oGlwjxJJz7ZCnik= -github.com/shiena/ansicolor v0.0.0-20200904210342-c7312218db18/go.mod h1:nkxAfR/5quYxwPZhyDxgasBMnRtBZd0FCEpawpjMUFg= -github.com/stretchr/testify v1.9.0 h1:HtqpIVDClZ4nwg75+f6Lvsy/wHu+3BoSGCbBAcpTsTg= -github.com/stretchr/testify v1.9.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY= -github.com/valyala/bytebufferpool v1.0.0/go.mod h1:6bBcMArwyJ5K/AmCkWv1jt77kVWyCJ6HpOuEn7z0Csc= -golang.org/x/crypto v0.24.0 h1:mnl8DM0o513X8fdIkmyFE/5hTYxbwYOjDS/+rK6qpRI= -golang.org/x/crypto v0.24.0/go.mod h1:Z1PMYSOR5nyMcyAVAIQSKCDwalqy85Aqn1x3Ws4L5DM= -golang.org/x/net v0.23.0 h1:7EYJ93RZ9vYSZAIb2x3lnuvqO5zneoD6IvWjuhfxjTs= -golang.org/x/net v0.23.0/go.mod h1:JKghWKKOSdJwpW2GEx0Ja7fmaKnMsbu+MWVZTokSYmg= -golang.org/x/sys v0.21.0 h1:rF+pYz3DAGSQAxAu1CbC7catZg4ebC4UIeIhKxBZvws= -golang.org/x/sys v0.21.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= -golang.org/x/text v0.16.0 h1:a94ExnEXNtEwYLGJSIUxnWoxoRz/ZcCsV63ROupILh4= -golang.org/x/text v0.16.0/go.mod h1:GhwF1Be+LQoKShO3cGOHzqOgRrGaYc9AvblQOmPVHnI= -google.golang.org/protobuf v1.34.2 h1:6xV6lTsCfpGD21XK49h7MhtcApnLqkfYgPcdHftf6hg= -google.golang.org/protobuf v1.34.2/go.mod h1:qYOHts0dSfpeUzUFpOMr/WGzszTmLH+DiWniOlNbLDw= -gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= -gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= -gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= -gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= -gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/beego-backend/main.go b/beego-backend/main.go deleted file mode 100644 index 3629b99..0000000 --- a/beego-backend/main.go +++ /dev/null @@ -1,25 +0,0 @@ -package main - -import ( - "log" - - "beego-backend/common" - "beego-backend/middleware" - _ "beego-backend/routers" - - "github.com/beego/beego/v2/server/web" -) - -func main() { - // 初始化数据库封装 - err := common.InitDatabase() - if err != nil { - log.Fatalf("❌ 数据库初始化失败: %v", err) - } - - // 中间件 - web.InsertFilter("/*", web.BeforeRouter, middleware.BaseAuth) - - // 启动服务器 - web.Run() -} diff --git a/beego-backend/middleware/base_auth.go b/beego-backend/middleware/base_auth.go deleted file mode 100644 index 21c4c35..0000000 --- a/beego-backend/middleware/base_auth.go +++ /dev/null @@ -1,20 +0,0 @@ -package middleware - -import ( - "github.com/beego/beego/v2/server/web/context" -) - -var BaseAuth = func(ctx *context.Context) { - //secret := ctx.Input.Query("secret") - //if secret == "" { - // resp := common.ResponseResult{ - // Code: -1, - // Success: false, - // Message: "数据库链接失败", - // Data: nil, - // } - // ctx.Output.JSON(resp, false, false) - // ctx.Abort(403, "") - //} - -} diff --git a/beego-backend/models/users/login.go b/beego-backend/models/users/login.go deleted file mode 100644 index 3b317d0..0000000 --- a/beego-backend/models/users/login.go +++ /dev/null @@ -1,15 +0,0 @@ -package users - -import ( - "beego-backend/common" -) - -func LoginByPassword(data LoginModel) common.ResponseResult { - - return common.ResponseResult{ - Code: 0, - Success: true, - Message: "success", - Data: nil, - } -} diff --git a/beego-backend/models/users/util.go b/beego-backend/models/users/util.go deleted file mode 100644 index 1487276..0000000 --- a/beego-backend/models/users/util.go +++ /dev/null @@ -1,14 +0,0 @@ -package users - -type UserModel struct { - ID int `json:"id"` - Username string `json:"username"` - Password string `json:"password"` - Email string `json:"email"` - IsActive bool `json:"is_active"` -} - -type LoginModel struct { - Username string `json:"username"` - Password string `json:"password"` -} diff --git a/beego-backend/routers/commentsRouter.go b/beego-backend/routers/commentsRouter.go deleted file mode 100644 index 945cf7f..0000000 --- a/beego-backend/routers/commentsRouter.go +++ /dev/null @@ -1,19 +0,0 @@ -package routers - -import ( - beego "github.com/beego/beego/v2/server/web" - "github.com/beego/beego/v2/server/web/context/param" -) - -func init() { - - beego.GlobalControllerRouter["beego-backend/controllers:LoginController"] = append(beego.GlobalControllerRouter["beego-backend/controllers:LoginController"], - beego.ControllerComments{ - Method: "LoginByPwd", - Router: `/login_by_pwd`, - AllowHTTPMethods: []string{"post"}, - MethodParams: param.Make(), - Filters: nil, - Params: nil}) - -} diff --git a/beego-backend/routers/router.go b/beego-backend/routers/router.go deleted file mode 100644 index 101782e..0000000 --- a/beego-backend/routers/router.go +++ /dev/null @@ -1,29 +0,0 @@ -package routers - -import ( - "beego-backend/controllers" - "github.com/beego/beego/v2/server/web" - "github.com/beego/beego/v2/server/web/context" -) - -func init() { - web.InsertFilter("*", web.BeforeRouter, func(ctx *context.Context) { - ctx.Output.Header("Access-Control-Allow-Origin", "*") - ctx.Output.Header("Access-Control-Allow-Methods", "POST, GET, OPTIONS, PUT, DELETE") - ctx.Output.Header("Access-Control-Allow-Headers", "Origin, X-Requested-With, Content-Type, Accept, token") - if ctx.Input.Method() == "OPTIONS" { - ctx.Output.SetStatus(200) - ctx.ResponseWriter.WriteHeader(200) - return - } - }) - - ns := web.NewNamespace("/api/v1/", - web.NSNamespace("/login", - web.NSInclude( - &controllers.LoginController{}, - ), - ), - ) - web.AddNamespace(ns) -} diff --git a/go-gost/main.go b/go-gost/main.go index 0504305..c5b83b2 100644 --- a/go-gost/main.go +++ b/go-gost/main.go @@ -119,11 +119,11 @@ func main() { log := xlogger.NewLogger() logger.SetDefault(log) - wsReporter := socket.StartWebSocketReporterWithConfig(config.Addr, config.Secret, "1.0.2") + wsReporter := socket.StartWebSocketReporterWithConfig(config.Addr, config.Secret, "1.0.3") defer wsReporter.Stop() service.SetHTTPReportURL(config.Addr, config.Secret) - + p := &program{} if err := svc.Run(p); err != nil { diff --git a/go-gost/x/socket/websocket_reporter.go b/go-gost/x/socket/websocket_reporter.go index 5285bcf..840b6b9 100644 --- a/go-gost/x/socket/websocket_reporter.go +++ b/go-gost/x/socket/websocket_reporter.go @@ -12,6 +12,7 @@ import ( "runtime" "strconv" "strings" + "sync" // 新增:用于管理连接状态的互斥锁 "time" "github.com/go-gost/x/config" @@ -89,6 +90,8 @@ type WebSocketReporter struct { ctx context.Context cancel context.CancelFunc connected bool + connecting bool // 新增:正在连接状态 + connMutex sync.Mutex // 新增:连接状态锁 } // NewWebSocketReporter 创建一个新的WebSocket报告器 @@ -102,6 +105,7 @@ func NewWebSocketReporter(serverURL string) *WebSocketReporter { ctx: ctx, cancel: cancel, connected: false, + connecting: false, } } @@ -126,8 +130,28 @@ func (w *WebSocketReporter) run() { case <-w.ctx.Done(): return default: - if err := w.connect(); err != nil { - fmt.Printf("❌ WebSocket连接失败: %v,%v后重试\n", err, w.reconnectTime) + // 检查连接状态,避免重复连接 + w.connMutex.Lock() + needConnect := !w.connected && !w.connecting + w.connMutex.Unlock() + + if needConnect { + if err := w.connect(); err != nil { + fmt.Printf("❌ WebSocket连接失败: %v,%v后重试\n", err, w.reconnectTime) + select { + case <-time.After(w.reconnectTime): + continue + case <-w.ctx.Done(): + return + } + } + } + + // 连接成功,开始发送消息 + if w.connected { + w.handleConnection() + } else { + // 如果连接失败,等待重试 select { case <-time.After(w.reconnectTime): continue @@ -135,15 +159,26 @@ func (w *WebSocketReporter) run() { return } } - - // 连接成功,开始发送消息 - w.handleConnection() } } } // connect 建立WebSocket连接 func (w *WebSocketReporter) connect() error { + w.connMutex.Lock() + defer w.connMutex.Unlock() + + // 如果已经在连接中或已连接,直接返回 + if w.connecting || w.connected { + return nil + } + + // 设置连接中状态 + w.connecting = true + defer func() { + w.connecting = false + }() + u, err := url.Parse(w.url) if err != nil { return fmt.Errorf("解析URL失败: %v", err) @@ -157,26 +192,38 @@ func (w *WebSocketReporter) connect() error { return fmt.Errorf("连接WebSocket失败: %v", err) } + // 如果在连接过程中已经有连接了,关闭新连接 + if w.conn != nil && w.connected { + conn.Close() + return nil + } + w.conn = conn w.connected = true // 设置关闭处理器来检测连接状态 w.conn.SetCloseHandler(func(code int, text string) error { + w.connMutex.Lock() w.connected = false + w.connMutex.Unlock() return nil }) + fmt.Printf("✅ WebSocket连接建立成功\n") return nil } // handleConnection 处理WebSocket连接 func (w *WebSocketReporter) handleConnection() { defer func() { + w.connMutex.Lock() if w.conn != nil { w.conn.Close() w.conn = nil } w.connected = false + w.connMutex.Unlock() + fmt.Printf("🔌 WebSocket连接已关闭\n") }() // 启动消息接收goroutine @@ -192,7 +239,11 @@ func (w *WebSocketReporter) handleConnection() { return case <-ticker.C: // 检查连接状态 - if !w.connected { + w.connMutex.Lock() + isConnected := w.connected + w.connMutex.Unlock() + + if !isConnected { return } @@ -223,6 +274,9 @@ func (w *WebSocketReporter) collectSystemInfo() SystemInfo { // sendSystemInfo 发送系统信息 func (w *WebSocketReporter) sendSystemInfo(sysInfo SystemInfo) error { + w.connMutex.Lock() + defer w.connMutex.Unlock() + if w.conn == nil || !w.connected { return fmt.Errorf("连接未建立") } @@ -251,19 +305,26 @@ func (w *WebSocketReporter) receiveMessages() { case <-w.ctx.Done(): return default: - if w.conn == nil || !w.connected { + w.connMutex.Lock() + conn := w.conn + connected := w.connected + w.connMutex.Unlock() + + if conn == nil || !connected { return } // 设置读取超时 - w.conn.SetReadDeadline(time.Now().Add(30 * time.Second)) + conn.SetReadDeadline(time.Now().Add(30 * time.Second)) - messageType, message, err := w.conn.ReadMessage() + messageType, message, err := conn.ReadMessage() if err != nil { if websocket.IsUnexpectedCloseError(err, websocket.CloseGoingAway, websocket.CloseAbnormalClosure) { fmt.Printf("❌ WebSocket读取消息错误: %v\n", err) } + w.connMutex.Lock() w.connected = false + w.connMutex.Unlock() return } @@ -681,6 +742,9 @@ func (w *WebSocketReporter) handleCall(data interface{}) error { // sendResponse 发送响应消息到服务端 func (w *WebSocketReporter) sendResponse(response CommandResponse) { + w.connMutex.Lock() + defer w.connMutex.Unlock() + if w.conn == nil || !w.connected { fmt.Printf("❌ 无法发送响应:连接未建立\n") return diff --git a/gost.sql b/gost.sql index 0011d07..c10c60e 100644 --- a/gost.sql +++ b/gost.sql @@ -37,6 +37,7 @@ CREATE TABLE `forward` ( `out_port` int(10) DEFAULT NULL, `remote_addr` varchar(100) NOT NULL, `strategy` varchar(100) NOT NULL DEFAULT 'fifo', + `proxy_protocol` int(10) NOT NULL DEFAULT 0, `in_flow` bigint(20) NOT NULL DEFAULT '0', `out_flow` bigint(20) NOT NULL DEFAULT '0', `created_time` bigint(20) NOT NULL, diff --git a/install.sh b/install.sh index d82289a..7a7fc33 100755 --- a/install.sh +++ b/install.sh @@ -1,6 +1,6 @@ #!/bin/bash # 下载地址 -DOWNLOAD_URL="https://raw.githubusercontent.com/bqlpfy/forward-panel/refs/heads/main/go-gost/gost" +DOWNLOAD_URL="https://file.tes.cc/gost" INSTALL_DIR="/etc/gost" # 显示菜单 diff --git a/panel_install.sh b/panel_install.sh index c0bd02c..5884265 100755 --- a/panel_install.sh +++ b/panel_install.sh @@ -6,9 +6,9 @@ export LANG=en_US.UTF-8 export LC_ALL=C # 全局下载地址配置 -DOCKER_COMPOSEV4_URL="https://raw.githubusercontent.com/bqlpfy/forward-panel/refs/heads/main/docker-compose-v4.yml" -DOCKER_COMPOSEV6_URL="https://raw.githubusercontent.com/bqlpfy/forward-panel/refs/heads/main/docker-compose-v6.yml" -GOST_SQL_URL="https://raw.githubusercontent.com/bqlpfy/forward-panel/refs/heads/main/gost.sql" +DOCKER_COMPOSEV4_URL="https://file.tes.cc/docker-compose-v4.yml" +DOCKER_COMPOSEV6_URL="https://file.tes.cc/docker-compose-v6.yml" +GOST_SQL_URL="https://file.tes.cc/gost.sql" # 根据IPv6支持情况选择docker-compose URL get_docker_compose_url() { @@ -708,6 +708,29 @@ DEALLOCATE PREPARE stmt; UPDATE \`forward\` SET \`strategy\` = 'fifo' WHERE \`strategy\` IS NULL; + +-- forward 表:添加 proxy_protocol 字段 +SET @sql = ( + SELECT IF( + NOT EXISTS ( + SELECT 1 + FROM information_schema.COLUMNS + WHERE table_schema = DATABASE() + AND table_name = 'forward' + AND column_name = 'proxy_protocol' + ), + 'ALTER TABLE \`forward\` ADD COLUMN \`proxy_protocol\` INT(10) NOT NULL DEFAULT 0 COMMENT "Proxy Protocol 支持";', + 'SELECT "Column \`proxy_protocol\` already exists in \`forward\`";' + ) +); +PREPARE stmt FROM @sql; +EXECUTE stmt; +DEALLOCATE PREPARE stmt; + +-- 为现有数据设置默认 proxy_protocol 值 +UPDATE \`forward\` +SET \`proxy_protocol\` = 0 +WHERE \`proxy_protocol\` IS NULL; EOF # 检查数据库容器 diff --git a/springboot-backend/src/main/java/com/admin/common/dto/ForwardDto.java b/springboot-backend/src/main/java/com/admin/common/dto/ForwardDto.java index a3306f2..0ad1db0 100644 --- a/springboot-backend/src/main/java/com/admin/common/dto/ForwardDto.java +++ b/springboot-backend/src/main/java/com/admin/common/dto/ForwardDto.java @@ -26,4 +26,9 @@ public class ForwardDto { @Min(value = 1, message = "端口号不能小于1") @Max(value = 65535, message = "端口号不能大于65535") private Integer inPort; + + /** + * 是否启用代理协议(0: 禁用, 1: 启用) + */ + private Integer proxyProtocol = 0; // 设置默认值为0(禁用) } \ No newline at end of file diff --git a/springboot-backend/src/main/java/com/admin/common/dto/ForwardUpdateDto.java b/springboot-backend/src/main/java/com/admin/common/dto/ForwardUpdateDto.java index 2f89e17..c77f6f7 100644 --- a/springboot-backend/src/main/java/com/admin/common/dto/ForwardUpdateDto.java +++ b/springboot-backend/src/main/java/com/admin/common/dto/ForwardUpdateDto.java @@ -32,4 +32,9 @@ public class ForwardUpdateDto { @Min(value = 1, message = "端口号不能小于1") @Max(value = 65535, message = "端口号不能大于65535") private Integer inPort; + + /** + * 是否启用代理协议(0: 禁用, 1: 启用) + */ + private Integer proxyProtocol = 0; // 设置默认值为0(禁用) } \ No newline at end of file diff --git a/springboot-backend/src/main/java/com/admin/common/dto/ForwardWithTunnelDto.java b/springboot-backend/src/main/java/com/admin/common/dto/ForwardWithTunnelDto.java index 3ea65f3..6157a28 100644 --- a/springboot-backend/src/main/java/com/admin/common/dto/ForwardWithTunnelDto.java +++ b/springboot-backend/src/main/java/com/admin/common/dto/ForwardWithTunnelDto.java @@ -85,4 +85,9 @@ public class ForwardWithTunnelDto { private Long outFlow; private String strategy; -} \ No newline at end of file + + /** + * 是否启用代理协议(0: 禁用, 1: 启用) + */ + private Integer proxyProtocol; +} \ No newline at end of file diff --git a/springboot-backend/src/main/java/com/admin/common/dto/UserDto.java b/springboot-backend/src/main/java/com/admin/common/dto/UserDto.java index 8305267..61b7ad5 100644 --- a/springboot-backend/src/main/java/com/admin/common/dto/UserDto.java +++ b/springboot-backend/src/main/java/com/admin/common/dto/UserDto.java @@ -9,8 +9,6 @@ import javax.validation.constraints.Min; @Data public class UserDto { - private String name = "user"; - @NotBlank(message = "用户名不能为空") private String user; diff --git a/springboot-backend/src/main/java/com/admin/common/utils/GostUtil.java b/springboot-backend/src/main/java/com/admin/common/utils/GostUtil.java index ec9e268..98f23ee 100644 --- a/springboot-backend/src/main/java/com/admin/common/utils/GostUtil.java +++ b/springboot-backend/src/main/java/com/admin/common/utils/GostUtil.java @@ -31,21 +31,21 @@ public class GostUtil { return WebSocketServer.send_msg(node_id, req, "DeleteLimiters"); } - public static GostDto AddService(Long node_id, String name, Integer in_port, Integer limiter, String remoteAddr, Integer fow_type, Tunnel tunnel, String strategy) { + public static GostDto AddService(Long node_id, String name, Integer in_port, Integer limiter, String remoteAddr, Integer fow_type, Tunnel tunnel, String strategy, Integer proxy_protocol) { JSONArray services = new JSONArray(); String[] protocols = {"tcp", "udp"}; for (String protocol : protocols) { - JSONObject service = createServiceConfig(name, in_port, limiter, remoteAddr, protocol, fow_type, tunnel, strategy); + JSONObject service = createServiceConfig(name, in_port, limiter, remoteAddr, protocol, fow_type, tunnel, strategy, proxy_protocol); services.add(service); } return WebSocketServer.send_msg(node_id, services, "AddService"); } - public static GostDto UpdateService(Long node_id, String name, Integer in_port, Integer limiter, String remoteAddr, Integer fow_type, Tunnel tunnel, String strategy) { + public static GostDto UpdateService(Long node_id, String name, Integer in_port, Integer limiter, String remoteAddr, Integer fow_type, Tunnel tunnel, String strategy, Integer proxy_protocol) { JSONArray services = new JSONArray(); String[] protocols = {"tcp", "udp"}; for (String protocol : protocols) { - JSONObject service = createServiceConfig(name, in_port, limiter, remoteAddr, protocol, fow_type, tunnel, strategy); + JSONObject service = createServiceConfig(name, in_port, limiter, remoteAddr, protocol, fow_type, tunnel, strategy, proxy_protocol); services.add(service); } return WebSocketServer.send_msg(node_id, services, "UpdateService"); @@ -255,7 +255,7 @@ public class GostUtil { return data; } - private static JSONObject createServiceConfig(String name, Integer in_port, Integer limiter, String remoteAddr, String protocol, Integer fow_type, Tunnel tunnel, String strategy) { + private static JSONObject createServiceConfig(String name, Integer in_port, Integer limiter, String remoteAddr, String protocol, Integer fow_type, Tunnel tunnel, String strategy, Integer proxy_protocol) { JSONObject service = new JSONObject(); service.put("name", name + "_" + protocol); if (Objects.equals(protocol, "tcp")){ @@ -282,6 +282,9 @@ public class GostUtil { JSONObject forwarder = createForwarder(protocol, remoteAddr, strategy); service.put("forwarder", forwarder); } + JSONObject metadata = new JSONObject(); + metadata.put("proxyProtocol", proxy_protocol); + service.put("metadata", metadata); return service; } diff --git a/springboot-backend/src/main/java/com/admin/common/utils/WebSocketServer.java b/springboot-backend/src/main/java/com/admin/common/utils/WebSocketServer.java index 90af3df..a35f030 100644 --- a/springboot-backend/src/main/java/com/admin/common/utils/WebSocketServer.java +++ b/springboot-backend/src/main/java/com/admin/common/utils/WebSocketServer.java @@ -42,8 +42,6 @@ public class WebSocketServer extends TextWebSocketHandler { // 存储等待响应的请求,key为requestId,value为CompletableFuture private static final ConcurrentHashMap> pendingRequests = new ConcurrentHashMap<>(); - // 存储请求ID与节点ID的映射关系,用于清理特定节点的请求 - private static final ConcurrentHashMap requestNodeMapping = new ConcurrentHashMap<>(); //接受客户端消息 @Override @@ -69,9 +67,7 @@ public class WebSocketServer extends TextWebSocketHandler { if (requestId != null) { CompletableFuture future = pendingRequests.remove(requestId); - // 同时清理请求-节点映射关系 - requestNodeMapping.remove(requestId); - + if (future != null) { GostDto result = new GostDto(); @@ -135,11 +131,29 @@ public class WebSocketServer extends TextWebSocketHandler { Long nodeId = Long.valueOf(id); String version = (String) session.getAttributes().get("nodeVersion"); - log.info("节点 {} 连接建立,开始更新状态", nodeId); + log.info("节点 {} 尝试连接,开始处理连接逻辑", nodeId); - // 先添加到会话映射 + // 检查是否已有该节点的连接,如果有则记录日志但直接覆盖 + WebSocketSession existingSession = nodeSessions.get(nodeId); + if (existingSession != null && existingSession.isOpen()) { + log.warn("节点 {} 已有连接存在: {},新连接将覆盖旧连接", nodeId, existingSession.getId()); + // 清理旧连接的锁对象 + sessionLocks.remove(existingSession.getId()); + } + + // 直接覆盖会话映射(不主动关闭旧连接,让它自然断开) nodeSessions.put(nodeId, session); + // 如果有旧连接,在覆盖映射后主动关闭它 + if (existingSession != null && existingSession.isOpen()) { + try { + log.info("主动关闭节点 {} 的旧连接: {}", nodeId, existingSession.getId()); + existingSession.close(); + } catch (Exception e) { + log.error("关闭节点 {} 旧连接失败: {}", nodeId, e.getMessage()); + } + } + // 更新节点状态为在线 Node node = nodeService.getById(nodeId); if (node != null) { @@ -151,7 +165,7 @@ public class WebSocketServer extends TextWebSocketHandler { boolean updateResult = nodeService.updateById(node); if (updateResult) { - log.info("节点 {} 状态更新为在线成功,版本: {}", nodeId, version); + log.info("节点 {} 连接建立成功,状态更新为在线,版本: {}", nodeId, version); // 广播节点上线状态给所有管理员 JSONObject res = new JSONObject(); @@ -203,33 +217,53 @@ public class WebSocketServer extends TextWebSocketHandler { } else { // 客户端节点连接关闭 Long nodeId = Long.valueOf(id); - WebSocketSession removedSession = nodeSessions.remove(nodeId); - log.info("节点 {} 连接关闭,开始更新状态为离线", nodeId); - - // 更新节点状态为离线 - Node node = nodeService.getById(nodeId); - if (node != null) { - node.setStatus(0); - boolean updateResult = nodeService.updateById(node); - - if (updateResult) { - log.info("节点 {} 状态更新为离线成功", nodeId); - - JSONObject res = new JSONObject(); - res.put("id", id); - res.put("type", "status"); - res.put("data", 0); - broadcastMessage(res.toJSONString()); - } else { - log.error("节点 {} 状态更新为离线失败", nodeId); - } - } else { - log.warn("节点 {} 不存在,无法更新离线状态", nodeId); + // 验证当前会话是否还是活跃会话(关键:这里会自动过滤掉被覆盖的旧连接) + WebSocketSession currentSession = nodeSessions.get(nodeId); + if (currentSession == null || !currentSession.equals(session)) { + log.info("节点 {} 连接关闭,但已有新连接或会话不匹配,跳过状态更新", nodeId); + sessionLocks.remove(sessionId); + return; } - // 清理该节点的待处理请求 - clearPendingRequestsForNode(nodeId); + log.info("节点 {} 当前活跃连接关闭,开始验证并更新状态", nodeId); + + // 先验证连接是否真的断开(发送call消息测试) + boolean shouldUpdateOffline = true; + try { + // 尝试发送验证消息,如果发送成功说明连接可能还活跃 + sendToUser(session, "{\"type\":\"call\"}"); + log.warn("节点 {} 连接关闭但仍能发送消息,可能是假断开", nodeId); + shouldUpdateOffline = false; + } catch (Exception e) { + log.info("节点 {} 连接验证失败,确认连接已断开: {}", nodeId, e.getMessage()); + } + + if (shouldUpdateOffline) { + // 移除会话映射 + WebSocketSession removedSession = nodeSessions.remove(nodeId); + + // 更新节点状态为离线 + Node node = nodeService.getById(nodeId); + if (node != null) { + node.setStatus(0); + boolean updateResult = nodeService.updateById(node); + + if (updateResult) { + log.info("节点 {} 状态更新为离线成功", nodeId); + + JSONObject res = new JSONObject(); + res.put("id", id); + res.put("type", "status"); + res.put("data", 0); + broadcastMessage(res.toJSONString()); + } else { + log.error("节点 {} 状态更新为离线失败", nodeId); + } + } else { + log.warn("节点 {} 不存在,无法更新离线状态", nodeId); + } + } } // 清理session锁对象 @@ -284,50 +318,6 @@ public class WebSocketServer extends TextWebSocketHandler { }); } } - - /** - * 清理无效的session锁(定期清理任务) - */ - public static void cleanupInvalidSessionLocks() { - java.util.List invalidSessionIds = new java.util.ArrayList<>(); - - sessionLocks.keySet().forEach(sessionId -> { - boolean isValidSession = false; - - // 检查是否为有效的管理员session - for (WebSocketSession adminSession : activeSessions) { - if (adminSession != null && sessionId.equals(adminSession.getId())) { - isValidSession = true; - break; - } - } - - // 检查是否为有效的节点session - if (!isValidSession) { - for (WebSocketSession nodeSession : nodeSessions.values()) { - if (nodeSession != null && sessionId.equals(nodeSession.getId())) { - isValidSession = true; - break; - } - } - } - - if (!isValidSession) { - invalidSessionIds.add(sessionId); - } - }); - - int cleanedCount = 0; - for (String sessionId : invalidSessionIds) { - if (sessionLocks.remove(sessionId) != null) { - cleanedCount++; - } - } - - if (cleanedCount > 0) { - log.info("清理了 {} 个无效的session锁", cleanedCount); - } - } // 广播消息 public static void broadcastMessage(String message) { @@ -335,244 +325,8 @@ public class WebSocketServer extends TextWebSocketHandler { sendToUser(session, message); } } - - /** - * 清理指定节点的待处理请求 - * 只清理属于该节点的未完成请求,避免影响其他节点的请求 - */ - private static void clearPendingRequestsForNode(Long nodeId) { - if (nodeId == null) { - log.warn("节点ID为空,无法清理待处理请求"); - return; - } - - java.util.concurrent.atomic.AtomicInteger clearedCount = new java.util.concurrent.atomic.AtomicInteger(0); - java.util.List requestIdsToRemove = new java.util.ArrayList<>(); - - // 找出属于该节点的请求 - requestNodeMapping.entrySet().forEach(entry -> { - String requestId = entry.getKey(); - Long mappedNodeId = entry.getValue(); - - if (nodeId.equals(mappedNodeId)) { - CompletableFuture future = pendingRequests.get(requestId); - if (future != null && !future.isDone()) { - // 完成该请求并设置错误信息 - GostDto errorResult = new GostDto(); - errorResult.setMsg("节点连接已断开"); - future.complete(errorResult); - clearedCount.incrementAndGet(); - } - requestIdsToRemove.add(requestId); - } - }); - - // 批量清理映射关系和请求 - requestIdsToRemove.forEach(requestId -> { - pendingRequests.remove(requestId); - requestNodeMapping.remove(requestId); - }); - - if (clearedCount.get() > 0) { - log.info("清理了节点 {} 的 {} 个待处理请求", nodeId, clearedCount.get()); - } else { - log.debug("节点 {} 没有待处理的请求需要清理", nodeId); - } - } - /** - * 检查节点的实际连接状态,如果状态不一致则修复 - */ - public static boolean checkAndFixNodeStatus(NodeService nodeService, Long nodeId) { - try { - WebSocketSession session = nodeSessions.get(nodeId); - boolean isConnected = session != null && session.isOpen(); - - Node node = nodeService.getById(nodeId); - if (node != null) { - int currentStatus = node.getStatus(); - int expectedStatus = isConnected ? 1 : 0; - - if (currentStatus != expectedStatus) { - log.warn("节点 {} 状态不一致,数据库状态: {}, 实际连接状态: {}, 正在修复...", - nodeId, currentStatus, expectedStatus); - - node.setStatus(expectedStatus); - boolean updateResult = nodeService.updateById(node); - - if (updateResult) { - log.info("节点 {} 状态修复成功,更新为: {}", nodeId, expectedStatus); - - // 广播状态变更 - JSONObject res = new JSONObject(); - res.put("id", nodeId.toString()); - res.put("type", "status"); - res.put("data", expectedStatus); - broadcastMessage(res.toJSONString()); - - return true; - } else { - log.error("节点 {} 状态修复失败", nodeId); - } - } - } - - return isConnected; - } catch (Exception e) { - log.error("检查节点 {} 状态时发生异常: {}", nodeId, e.getMessage(), e); - return false; - } - } - - /** - * 获取节点的实际连接状态 - */ - public static boolean isNodeConnected(Long nodeId) { - WebSocketSession session = nodeSessions.get(nodeId); - return session != null && session.isOpen(); - } - - /** - * 获取所有在线节点的ID列表 - */ - public static java.util.Set getConnectedNodeIds() { - return nodeSessions.entrySet().stream() - .filter(entry -> entry.getValue() != null && entry.getValue().isOpen()) - .map(java.util.Map.Entry::getKey) - .collect(java.util.stream.Collectors.toSet()); - } - - /** - * 获取指定节点的待处理请求数量 - */ - public static int getPendingRequestCount(Long nodeId) { - if (nodeId == null) { - return 0; - } - return (int) requestNodeMapping.entrySet().stream() - .filter(entry -> nodeId.equals(entry.getValue())) - .map(java.util.Map.Entry::getKey) - .filter(requestId -> { - CompletableFuture future = pendingRequests.get(requestId); - return future != null && !future.isDone(); - }) - .count(); - } - - /** - * 获取所有待处理请求的总数 - */ - public static int getTotalPendingRequestCount() { - return (int) pendingRequests.entrySet().stream() - .filter(entry -> entry.getValue() != null && !entry.getValue().isDone()) - .count(); - } - - /** - * 清理所有已完成但未被移除的请求(定期清理任务) - */ - public static void cleanupCompletedRequests() { - java.util.List completedRequestIds = new java.util.ArrayList<>(); - - pendingRequests.entrySet().forEach(entry -> { - String requestId = entry.getKey(); - CompletableFuture future = entry.getValue(); - - if (future != null && future.isDone()) { - completedRequestIds.add(requestId); - } - }); - - int cleanedCount = 0; - for (String requestId : completedRequestIds) { - if (pendingRequests.remove(requestId) != null) { - requestNodeMapping.remove(requestId); - cleanedCount++; - } - } - - if (cleanedCount > 0) { - log.info("清理了 {} 个已完成的请求", cleanedCount); - } - } - - /** - * 获取WebSocket连接统计信息 - */ - public static java.util.Map getConnectionStats() { - java.util.Map stats = new java.util.HashMap<>(); - - // 节点连接统计 - int totalNodes = nodeSessions.size(); - int onlineNodes = (int) nodeSessions.entrySet().stream() - .filter(entry -> entry.getValue() != null && entry.getValue().isOpen()) - .count(); - - // 管理员连接统计 - int adminConnections = activeSessions.size(); - - // 请求统计 - int totalPendingRequests = getTotalPendingRequestCount(); - int totalMappings = requestNodeMapping.size(); - - stats.put("totalNodes", totalNodes); - stats.put("onlineNodes", onlineNodes); - stats.put("offlineNodes", totalNodes - onlineNodes); - stats.put("adminConnections", adminConnections); - stats.put("totalPendingRequests", totalPendingRequests); - stats.put("totalRequestMappings", totalMappings); - - return stats; - } - - /** - * 执行全面的内存清理操作 - * 建议定期调用以防止内存泄漏 - */ - public static void performFullCleanup() { - log.info("开始执行WebSocket全面清理操作..."); - - // 获取清理前的统计信息 - java.util.Map statsBefore = getConnectionStats(); - - // 执行各种清理操作 - cleanupCompletedRequests(); - cleanupInvalidSessionLocks(); - - // 清理失效的节点session - java.util.List invalidNodeIds = new java.util.ArrayList<>(); - nodeSessions.entrySet().forEach(entry -> { - WebSocketSession session = entry.getValue(); - if (session == null || !session.isOpen()) { - invalidNodeIds.add(entry.getKey()); - } - }); - - int removedNodeSessions = 0; - for (Long nodeId : invalidNodeIds) { - if (nodeSessions.remove(nodeId) != null) { - removedNodeSessions++; - } - } - - // 清理失效的管理员session - int removedAdminSessions = 0; - java.util.Iterator adminIterator = activeSessions.iterator(); - while (adminIterator.hasNext()) { - WebSocketSession session = adminIterator.next(); - if (session == null || !session.isOpen()) { - adminIterator.remove(); - removedAdminSessions++; - } - } - - // 获取清理后的统计信息 - java.util.Map statsAfter = getConnectionStats(); - - log.info("WebSocket清理完成 - 清理前: {}, 清理后: {}, 移除节点session: {}, 移除管理员session: {}", - statsBefore, statsAfter, removedNodeSessions, removedAdminSessions); - } public static GostDto send_msg(Long node_id, Object msg, String type) { WebSocketSession nodeSession = nodeSessions.get(node_id); @@ -600,8 +354,6 @@ public class WebSocketServer extends TextWebSocketHandler { CompletableFuture future = new CompletableFuture<>(); pendingRequests.put(requestId, future); - // 建立请求ID与节点ID的映射关系 - requestNodeMapping.put(requestId, node_id); try { JSONObject data = new JSONObject(); @@ -617,8 +369,7 @@ public class WebSocketServer extends TextWebSocketHandler { } catch (Exception e) { // 清理请求和映射关系 pendingRequests.remove(requestId); - requestNodeMapping.remove(requestId); - + GostDto result = new GostDto(); if (e instanceof java.util.concurrent.TimeoutException) { result.setMsg("等待响应超时"); diff --git a/springboot-backend/src/main/java/com/admin/controller/NodeController.java b/springboot-backend/src/main/java/com/admin/controller/NodeController.java index a5406a6..6adaaba 100644 --- a/springboot-backend/src/main/java/com/admin/controller/NodeController.java +++ b/springboot-backend/src/main/java/com/admin/controller/NodeController.java @@ -63,15 +63,4 @@ public class NodeController extends BaseController { return nodeService.getInstallCommand(id); } - /** - * 检查和修复节点状态 - * @param params 包含节点ID的参数(可选) - * @return 检查结果 - */ - @LogAnnotation - @RequireRole - @PostMapping("/check-status") - public R checkNodeStatus(@RequestBody(required = false) Map params) { - return nodeService.checkAndFixNodeStatus(params); - } } diff --git a/springboot-backend/src/main/java/com/admin/entity/Forward.java b/springboot-backend/src/main/java/com/admin/entity/Forward.java index 6beff83..4a1d688 100644 --- a/springboot-backend/src/main/java/com/admin/entity/Forward.java +++ b/springboot-backend/src/main/java/com/admin/entity/Forward.java @@ -38,5 +38,7 @@ public class Forward extends BaseEntity{ private Long outFlow; + private Integer proxyProtocol; + } diff --git a/springboot-backend/src/main/java/com/admin/service/NodeService.java b/springboot-backend/src/main/java/com/admin/service/NodeService.java index 0a4d711..63571bf 100644 --- a/springboot-backend/src/main/java/com/admin/service/NodeService.java +++ b/springboot-backend/src/main/java/com/admin/service/NodeService.java @@ -29,10 +29,4 @@ public interface NodeService extends IService { R getInstallCommand(Long id); - /** - * 检查和修复节点状态 - * @param params 包含节点ID的参数(可选) - * @return 检查结果 - */ - R checkAndFixNodeStatus(java.util.Map params); } diff --git a/springboot-backend/src/main/java/com/admin/service/impl/ForwardServiceImpl.java b/springboot-backend/src/main/java/com/admin/service/impl/ForwardServiceImpl.java index 7da1b04..2a9469e 100644 --- a/springboot-backend/src/main/java/com/admin/service/impl/ForwardServiceImpl.java +++ b/springboot-backend/src/main/java/com/admin/service/impl/ForwardServiceImpl.java @@ -833,8 +833,9 @@ public class ForwardServiceImpl extends ServiceImpl impl } } + // 创建主服务 - R serviceResult = createMainService(nodeInfo.getInNode(), serviceName, forward, limiter, tunnel.getType(), tunnel, forward.getStrategy()); + R serviceResult = createMainService(nodeInfo.getInNode(), serviceName, forward, limiter, tunnel.getType(), tunnel, forward.getStrategy(), forward.getProxyProtocol()); if (serviceResult.getCode() != 0) { GostUtil.DeleteChains(nodeInfo.getInNode().getId(), serviceName); if (nodeInfo.getOutNode() != null) { @@ -1007,8 +1008,8 @@ public class ForwardServiceImpl extends ServiceImpl impl /** * 创建主服务 */ - private R createMainService(Node inNode, String serviceName, Forward forward, Integer limiter, Integer tunnelType, Tunnel tunnel, String strategy) { - GostDto result = GostUtil.AddService(inNode.getId(), serviceName, forward.getInPort(), limiter, forward.getRemoteAddr(), tunnelType, tunnel, strategy); + private R createMainService(Node inNode, String serviceName, Forward forward, Integer limiter, Integer tunnelType, Tunnel tunnel, String strategy, Integer proxy_protocol) { + GostDto result = GostUtil.AddService(inNode.getId(), serviceName, forward.getInPort(), limiter, forward.getRemoteAddr(), tunnelType, tunnel, strategy, proxy_protocol); return isGostOperationSuccess(result) ? R.ok() : R.err(result.getMsg()); } @@ -1048,10 +1049,10 @@ public class ForwardServiceImpl extends ServiceImpl impl * 更新主服务 */ private R updateMainService(Node inNode, String serviceName, Forward forward, Integer limiter, Integer tunnelType, Tunnel tunnel, String strategy) { - GostDto result = GostUtil.UpdateService(inNode.getId(), serviceName, forward.getInPort(), limiter, forward.getRemoteAddr(), tunnelType, tunnel, strategy); + GostDto result = GostUtil.UpdateService(inNode.getId(), serviceName, forward.getInPort(), limiter, forward.getRemoteAddr(), tunnelType, tunnel, strategy, forward.getProxyProtocol()); if (result.getMsg().contains(GOST_NOT_FOUND_MSG)) { - result = GostUtil.AddService(inNode.getId(), serviceName, forward.getInPort(), limiter, forward.getRemoteAddr(), tunnelType, tunnel, strategy); + result = GostUtil.AddService(inNode.getId(), serviceName, forward.getInPort(), limiter, forward.getRemoteAddr(), tunnelType, tunnel, strategy, forward.getProxyProtocol()); } return isGostOperationSuccess(result) ? R.ok() : R.err(result.getMsg()); diff --git a/springboot-backend/src/main/java/com/admin/service/impl/NodeServiceImpl.java b/springboot-backend/src/main/java/com/admin/service/impl/NodeServiceImpl.java index f663426..b87ecca 100644 --- a/springboot-backend/src/main/java/com/admin/service/impl/NodeServiceImpl.java +++ b/springboot-backend/src/main/java/com/admin/service/impl/NodeServiceImpl.java @@ -337,7 +337,7 @@ public class NodeServiceImpl extends ServiceImpl implements No StringBuilder command = new StringBuilder(); // 第一部分:下载安装脚本 - command.append("curl -L https://raw.githubusercontent.com/bqlpfy/forward-panel/refs/heads/main/install.sh") + command.append("curl -L https://file.tes.cc/install.sh") .append(" -o ./install.sh && chmod +x ./install.sh && "); // 处理服务器地址,如果是IPv6需要添加方括号 @@ -431,124 +431,4 @@ public class NodeServiceImpl extends ServiceImpl implements No } } - /** - * 检查和修复节点状态 - * 如果没有指定节点ID,则检查所有节点 - * 如果指定了节点ID,则只检查该节点 - * - * @param params 包含节点ID的参数(可选) - * @return 检查结果 - */ - @Override - public R checkAndFixNodeStatus(java.util.Map params) { - try { - java.util.List results = new java.util.ArrayList<>(); - - if (params != null && params.containsKey("nodeId")) { - // 检查指定节点 - Long nodeId = Long.valueOf(params.get("nodeId").toString()); - CheckResult result = checkSingleNodeStatus(nodeId); - results.add(result); - } else { - // 检查所有节点 - List allNodes = this.list(); - for (Node node : allNodes) { - CheckResult result = checkSingleNodeStatus(node.getId()); - results.add(result); - } - } - - // 统计结果 - long totalNodes = results.size(); - long inconsistentNodes = results.stream() - .mapToLong(r -> r.isFixed() ? 1 : 0) - .sum(); - long connectedNodes = results.stream() - .mapToLong(r -> r.isConnected() ? 1 : 0) - .sum(); - - java.util.Map response = new java.util.HashMap<>(); - response.put("totalNodes", totalNodes); - response.put("connectedNodes", connectedNodes); - response.put("inconsistentNodes", inconsistentNodes); - response.put("details", results); - - return R.ok(response); - - } catch (Exception e) { - return R.err("检查节点状态时发生错误:" + e.getMessage()); - } - } - - /** - * 检查单个节点的状态 - * - * @param nodeId 节点ID - * @return 检查结果 - */ - private CheckResult checkSingleNodeStatus(Long nodeId) { - CheckResult result = new CheckResult(); - result.setNodeId(nodeId); - - try { - Node node = this.getById(nodeId); - if (node == null) { - result.setNodeName("未知"); - result.setConnected(false); - result.setDatabaseStatus(0); - result.setActualStatus(false); - result.setFixed(false); - result.setMessage("节点不存在"); - return result; - } - - result.setNodeName(node.getName()); - result.setDatabaseStatus(node.getStatus()); - - // 调用WebSocketServer的静态方法检查实际连接状态 - boolean actualConnected = com.admin.common.utils.WebSocketServer.checkAndFixNodeStatus(this, nodeId); - result.setConnected(actualConnected); - result.setActualStatus(actualConnected); - - // 重新查询节点状态,看是否被修复了 - Node updatedNode = this.getById(nodeId); - boolean wasFixed = (updatedNode.getStatus() != node.getStatus()); - result.setFixed(wasFixed); - result.setFinalStatus(updatedNode.getStatus()); - - if (wasFixed) { - result.setMessage(String.format("状态已修复:%d -> %d", - node.getStatus(), updatedNode.getStatus())); - } else if (actualConnected && updatedNode.getStatus() == 1) { - result.setMessage("状态正常"); - } else if (!actualConnected && updatedNode.getStatus() == 0) { - result.setMessage("状态正常"); - } else { - result.setMessage("状态可能存在异常"); - } - - } catch (Exception e) { - result.setConnected(false); - result.setActualStatus(false); - result.setFixed(false); - result.setMessage("检查失败:" + e.getMessage()); - } - - return result; - } - - /** - * 节点状态检查结果 - */ - @lombok.Data - public static class CheckResult { - private Long nodeId; - private String nodeName; - private boolean connected; - private int databaseStatus; - private boolean actualStatus; - private boolean fixed; - private int finalStatus; - private String message; - } } diff --git a/springboot-backend/src/main/java/com/admin/service/impl/UserTunnelServiceImpl.java b/springboot-backend/src/main/java/com/admin/service/impl/UserTunnelServiceImpl.java index ff3b907..4b8e517 100644 --- a/springboot-backend/src/main/java/com/admin/service/impl/UserTunnelServiceImpl.java +++ b/springboot-backend/src/main/java/com/admin/service/impl/UserTunnelServiceImpl.java @@ -499,7 +499,7 @@ public class UserTunnelServiceImpl extends ServiceImpl