mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-01 22:46:38 +08:00
[优化] 更新节点类型支持,添加隧道客户端,重构相关路由和配置
This commit is contained in:
@@ -89,6 +89,23 @@ func AgentHeartbeat(c *gin.Context) {
|
||||
// @Success 200 {object} map[string]interface{}
|
||||
// @Router /api/agent/config-versions/active [get]
|
||||
func AgentGetActiveConfig(c *gin.Context) {
|
||||
authNode, ok := c.Get("node")
|
||||
if !ok {
|
||||
respondUnauthorized(c, "Node object missing from context")
|
||||
return
|
||||
}
|
||||
node := authNode.(*model.Node)
|
||||
|
||||
if node.NodeType == "tunnel_client" {
|
||||
config, err := service.GetFlaredTunnelConfig(node)
|
||||
if err != nil {
|
||||
respondFailure(c, "无法生成隧道配置: "+err.Error())
|
||||
return
|
||||
}
|
||||
respondSuccess(c, config)
|
||||
return
|
||||
}
|
||||
|
||||
config, err := service.GetActiveConfigForAgent()
|
||||
if err != nil {
|
||||
respondFailure(c, "当前没有激活版本")
|
||||
|
||||
@@ -1,171 +0,0 @@
|
||||
package controller
|
||||
|
||||
import (
|
||||
"log/slog"
|
||||
"net"
|
||||
"openflare/model"
|
||||
"openflare/service"
|
||||
"time"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
"golang.org/x/net/websocket"
|
||||
)
|
||||
|
||||
// FlaredHeartbeat godoc
|
||||
// @Summary Report OpenFlared client heartbeat
|
||||
// @Tags Flared
|
||||
// @Accept json
|
||||
// @Produce json
|
||||
// @Security TunnelTokenAuth
|
||||
// @Param payload body service.FlaredHeartbeatPayload true "Flared heartbeat payload"
|
||||
// @Success 200 {object} map[string]interface{}
|
||||
// @Failure 400 {object} map[string]interface{}
|
||||
// @Router /api/flared/heartbeat [post]
|
||||
func FlaredHeartbeat(c *gin.Context) {
|
||||
var payload service.FlaredHeartbeatPayload
|
||||
if !bindJSON(c, &payload) {
|
||||
return
|
||||
}
|
||||
authTunnel, ok := c.Get("tunnel")
|
||||
if !ok {
|
||||
respondUnauthorized(c, "无权进行此操作")
|
||||
return
|
||||
}
|
||||
tunnel := authTunnel.(*model.Tunnel)
|
||||
result, err := service.HeartbeatFlared(tunnel, payload)
|
||||
if err != nil {
|
||||
respondFailure(c, err.Error())
|
||||
return
|
||||
}
|
||||
respondSuccess(c, result)
|
||||
}
|
||||
|
||||
// FlaredGetActiveConfig godoc
|
||||
// @Summary Get active tunnel config for OpenFlared client
|
||||
// @Tags Flared
|
||||
// @Produce json
|
||||
// @Security TunnelTokenAuth
|
||||
// @Success 200 {object} map[string]interface{}
|
||||
// @Router /api/flared/config/active [get]
|
||||
func FlaredGetActiveConfig(c *gin.Context) {
|
||||
authTunnel, ok := c.Get("tunnel")
|
||||
if !ok {
|
||||
respondUnauthorized(c, "无权进行此操作")
|
||||
return
|
||||
}
|
||||
tunnel := authTunnel.(*model.Tunnel)
|
||||
config, err := service.GetFlaredTunnelConfig(tunnel)
|
||||
if err != nil {
|
||||
respondFailure(c, err.Error())
|
||||
return
|
||||
}
|
||||
respondSuccess(c, config)
|
||||
}
|
||||
|
||||
// FlaredReportApplyLog godoc
|
||||
// @Summary Report apply log for OpenFlared client
|
||||
// @Tags Flared
|
||||
// @Accept json
|
||||
// @Produce json
|
||||
// @Security TunnelTokenAuth
|
||||
// @Param payload body service.ApplyLogPayload true "Apply log payload"
|
||||
// @Success 200 {object} map[string]interface{}
|
||||
// @Failure 400 {object} map[string]interface{}
|
||||
// @Router /api/flared/apply-log [post]
|
||||
func FlaredReportApplyLog(c *gin.Context) {
|
||||
var payload service.ApplyLogPayload
|
||||
if !bindJSON(c, &payload) {
|
||||
return
|
||||
}
|
||||
authTunnel, ok := c.Get("tunnel")
|
||||
if !ok {
|
||||
respondUnauthorized(c, "无权进行此操作")
|
||||
return
|
||||
}
|
||||
tunnel := authTunnel.(*model.Tunnel)
|
||||
payload.NodeID = tunnel.TunnelID
|
||||
log, err := service.ReportApplyLog(payload)
|
||||
if err != nil {
|
||||
respondFailure(c, err.Error())
|
||||
return
|
||||
}
|
||||
respondSuccess(c, log)
|
||||
}
|
||||
|
||||
// FlaredWebSocket godoc
|
||||
// @Summary Upgrade OpenFlared connection to websocket
|
||||
// @Tags Flared
|
||||
// @Security TunnelTokenAuth
|
||||
// @Router /api/flared/ws [get]
|
||||
func FlaredWebSocket(c *gin.Context) {
|
||||
authTunnel, ok := c.Get("tunnel")
|
||||
if !ok {
|
||||
respondUnauthorized(c, "无权进行此操作")
|
||||
return
|
||||
}
|
||||
tunnel := authTunnel.(*model.Tunnel)
|
||||
slog.Debug("flared ws upgrade requested", "tunnel_id", tunnel.TunnelID, "remote", c.Request.RemoteAddr)
|
||||
websocket.Handler(func(conn *websocket.Conn) {
|
||||
client := service.RegisterFlaredWSClient(tunnel.TunnelID)
|
||||
defer service.UnregisterFlaredWSClient(client)
|
||||
defer func() {
|
||||
_ = conn.Close()
|
||||
slog.Debug("flared ws connection closed", "tunnel_id", tunnel.TunnelID)
|
||||
}()
|
||||
|
||||
slog.Debug("flared ws upgrade succeeded", "tunnel_id", tunnel.TunnelID, "remote", c.Request.RemoteAddr)
|
||||
|
||||
go func() {
|
||||
<-client.Done()
|
||||
_ = conn.Close()
|
||||
}()
|
||||
|
||||
go streamFlaredWSMessages(c, conn, client)
|
||||
|
||||
for {
|
||||
var message service.WSMessage
|
||||
_ = conn.SetReadDeadline(time.Now().Add(agentWSReadTimeout()))
|
||||
if err := websocket.JSON.Receive(conn, &message); err != nil {
|
||||
if netErr, ok := err.(net.Error); ok && netErr.Timeout() {
|
||||
slog.Debug("flared ws receive timeout", "tunnel_id", tunnel.TunnelID)
|
||||
return
|
||||
}
|
||||
slog.Debug("flared ws receive failed", "tunnel_id", tunnel.TunnelID, "error", err)
|
||||
return
|
||||
}
|
||||
slog.Debug("flared ws message received", "tunnel_id", tunnel.TunnelID, "type", message.Type)
|
||||
switch message.Type {
|
||||
case "status":
|
||||
// Handle status if needed for flared
|
||||
case "ping":
|
||||
if !service.SendFlaredWSPong(tunnel.TunnelID) {
|
||||
slog.Debug("flared ws pong enqueue failed", "tunnel_id", tunnel.TunnelID)
|
||||
}
|
||||
case "pong":
|
||||
slog.Debug("flared ws pong received", "tunnel_id", tunnel.TunnelID)
|
||||
default:
|
||||
slog.Debug("flared ws unsupported message type", "tunnel_id", tunnel.TunnelID, "type", message.Type)
|
||||
}
|
||||
}
|
||||
}).ServeHTTP(c.Writer, c.Request)
|
||||
}
|
||||
|
||||
func streamFlaredWSMessages(c *gin.Context, conn *websocket.Conn, client *service.WSClient) {
|
||||
for {
|
||||
select {
|
||||
case <-c.Request.Context().Done():
|
||||
return
|
||||
case <-client.Done():
|
||||
return
|
||||
case message, ok := <-client.Messages():
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
_ = conn.SetWriteDeadline(time.Now().Add(agentWSWriteTimeout()))
|
||||
if err := websocket.JSON.Send(conn, message); err != nil {
|
||||
slog.Debug("flared ws send failed", "tunnel_id", client.ID(), "error", err)
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,84 +0,0 @@
|
||||
package controller
|
||||
|
||||
import (
|
||||
"openflare/service"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
)
|
||||
|
||||
func GetTunnels(c *gin.Context) {
|
||||
tunnels, err := service.ListTunnels()
|
||||
if err != nil {
|
||||
respondFailure(c, err.Error())
|
||||
return
|
||||
}
|
||||
respondSuccess(c, tunnels)
|
||||
}
|
||||
|
||||
func GetTunnel(c *gin.Context) {
|
||||
id, ok := parseIDParam(c)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
tunnel, err := service.GetTunnel(id)
|
||||
if err != nil {
|
||||
respondFailure(c, err.Error())
|
||||
return
|
||||
}
|
||||
respondSuccess(c, tunnel)
|
||||
}
|
||||
|
||||
func CreateTunnel(c *gin.Context) {
|
||||
var input service.TunnelInput
|
||||
if !bindJSON(c, &input) {
|
||||
return
|
||||
}
|
||||
tunnel, err := service.CreateTunnel(input)
|
||||
if err != nil {
|
||||
respondFailure(c, err.Error())
|
||||
return
|
||||
}
|
||||
respondSuccess(c, tunnel)
|
||||
}
|
||||
|
||||
func UpdateTunnel(c *gin.Context) {
|
||||
id, ok := parseIDParam(c)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
var input service.TunnelInput
|
||||
if !bindJSON(c, &input) {
|
||||
return
|
||||
}
|
||||
tunnel, err := service.UpdateTunnel(id, input)
|
||||
if err != nil {
|
||||
respondFailure(c, err.Error())
|
||||
return
|
||||
}
|
||||
respondSuccess(c, tunnel)
|
||||
}
|
||||
|
||||
func DeleteTunnel(c *gin.Context) {
|
||||
id, ok := parseIDParam(c)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
if err := service.DeleteTunnel(id); err != nil {
|
||||
respondFailure(c, err.Error())
|
||||
return
|
||||
}
|
||||
respondSuccess(c, nil)
|
||||
}
|
||||
|
||||
func RotateTunnelToken(c *gin.Context) {
|
||||
id, ok := parseIDParam(c)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
tunnel, err := service.RotateTunnelToken(id)
|
||||
if err != nil {
|
||||
respondFailure(c, err.Error())
|
||||
return
|
||||
}
|
||||
respondSuccess(c, tunnel)
|
||||
}
|
||||
Reference in New Issue
Block a user