mirror of
https://github.com/Sagit-chu/flvx.git
synced 2026-09-29 07:56:37 +08:00
feat: complete federation sharing backend implementation
This commit is contained in:
@@ -1,9 +0,0 @@
|
||||
---
|
||||
active: true
|
||||
iteration: 1
|
||||
max_iterations: 100
|
||||
completion_promise: "DONE"
|
||||
started_at: "2026-02-08T10:54:42.662Z"
|
||||
session_id: "ses_3c3ac2d4bffe6K5IO2BeAEWe7g"
|
||||
---
|
||||
请进行冒烟测试,确保功能都正常,有问题直接修改,完美后提交全部变更并且push
|
||||
@@ -9,6 +9,7 @@ import (
|
||||
|
||||
"go-backend/internal/http/client"
|
||||
"go-backend/internal/http/response"
|
||||
"go-backend/internal/store/sqlite"
|
||||
)
|
||||
|
||||
type federationTunnelRequest struct {
|
||||
@@ -17,11 +18,129 @@ type federationTunnelRequest struct {
|
||||
Target string `json:"target"`
|
||||
}
|
||||
|
||||
type createPeerShareRequest struct {
|
||||
Name string `json:"name"`
|
||||
NodeID int64 `json:"nodeId"`
|
||||
MaxBandwidth int64 `json:"maxBandwidth"`
|
||||
ExpiryTime int64 `json:"expiryTime"`
|
||||
PortRangeStart int `json:"portRangeStart"`
|
||||
PortRangeEnd int `json:"portRangeEnd"`
|
||||
}
|
||||
|
||||
type deletePeerShareRequest struct {
|
||||
ID int64 `json:"id"`
|
||||
}
|
||||
|
||||
type nodeImportRequest struct {
|
||||
RemoteURL string `json:"remoteUrl"`
|
||||
Token string `json:"token"`
|
||||
}
|
||||
|
||||
func (h *Handler) federationShareList(w http.ResponseWriter, r *http.Request) {
|
||||
if r.Method != http.MethodPost {
|
||||
response.WriteJSON(w, response.ErrDefault("Invalid method"))
|
||||
return
|
||||
}
|
||||
|
||||
shares, err := h.repo.ListPeerShares()
|
||||
if err != nil {
|
||||
response.WriteJSON(w, response.Err(-2, err.Error()))
|
||||
return
|
||||
}
|
||||
response.WriteJSON(w, response.OK(shares))
|
||||
}
|
||||
|
||||
func (h *Handler) federationShareCreate(w http.ResponseWriter, r *http.Request) {
|
||||
if r.Method != http.MethodPost {
|
||||
response.WriteJSON(w, response.ErrDefault("Invalid method"))
|
||||
return
|
||||
}
|
||||
|
||||
var req createPeerShareRequest
|
||||
if err := decodeJSON(r.Body, &req); err != nil {
|
||||
response.WriteJSON(w, response.ErrDefault("Invalid JSON"))
|
||||
return
|
||||
}
|
||||
|
||||
if req.Name == "" || req.NodeID == 0 {
|
||||
response.WriteJSON(w, response.ErrDefault("Name and NodeID are required"))
|
||||
return
|
||||
}
|
||||
|
||||
if req.MaxBandwidth < 0 {
|
||||
response.WriteJSON(w, response.ErrDefault("Max bandwidth cannot be negative"))
|
||||
return
|
||||
}
|
||||
|
||||
if req.ExpiryTime < 0 {
|
||||
response.WriteJSON(w, response.ErrDefault("Expiry time cannot be negative"))
|
||||
return
|
||||
}
|
||||
|
||||
if req.PortRangeStart < 0 || req.PortRangeStart > 65535 || req.PortRangeEnd < 0 || req.PortRangeEnd > 65535 {
|
||||
response.WriteJSON(w, response.ErrDefault("Invalid port range"))
|
||||
return
|
||||
}
|
||||
|
||||
if req.PortRangeStart > req.PortRangeEnd {
|
||||
response.WriteJSON(w, response.ErrDefault("Port range start cannot be greater than end"))
|
||||
return
|
||||
}
|
||||
|
||||
node, err := h.repo.GetNodeByID(req.NodeID)
|
||||
if err != nil {
|
||||
response.WriteJSON(w, response.Err(-2, err.Error()))
|
||||
return
|
||||
}
|
||||
if node == nil {
|
||||
response.WriteJSON(w, response.ErrDefault("Node not found"))
|
||||
return
|
||||
}
|
||||
|
||||
now := time.Now().UnixMilli()
|
||||
token := randomToken(32)
|
||||
|
||||
share := &sqlite.PeerShare{
|
||||
Name: req.Name,
|
||||
NodeID: req.NodeID,
|
||||
Token: token,
|
||||
MaxBandwidth: req.MaxBandwidth,
|
||||
ExpiryTime: req.ExpiryTime,
|
||||
PortRangeStart: req.PortRangeStart,
|
||||
PortRangeEnd: req.PortRangeEnd,
|
||||
IsActive: 1,
|
||||
CreatedTime: now,
|
||||
UpdatedTime: now,
|
||||
}
|
||||
|
||||
if err := h.repo.CreatePeerShare(share); err != nil {
|
||||
response.WriteJSON(w, response.Err(-2, err.Error()))
|
||||
return
|
||||
}
|
||||
|
||||
response.WriteJSON(w, response.OKEmpty())
|
||||
}
|
||||
|
||||
func (h *Handler) federationShareDelete(w http.ResponseWriter, r *http.Request) {
|
||||
if r.Method != http.MethodPost {
|
||||
response.WriteJSON(w, response.ErrDefault("Invalid method"))
|
||||
return
|
||||
}
|
||||
|
||||
var req deletePeerShareRequest
|
||||
if err := decodeJSON(r.Body, &req); err != nil {
|
||||
response.WriteJSON(w, response.ErrDefault("Invalid JSON"))
|
||||
return
|
||||
}
|
||||
|
||||
if err := h.repo.DeletePeerShare(req.ID); err != nil {
|
||||
response.WriteJSON(w, response.Err(-2, err.Error()))
|
||||
return
|
||||
}
|
||||
|
||||
response.WriteJSON(w, response.OKEmpty())
|
||||
}
|
||||
|
||||
func (h *Handler) nodeImport(w http.ResponseWriter, r *http.Request) {
|
||||
if r.Method != http.MethodPost {
|
||||
response.WriteJSON(w, response.ErrDefault("Invalid method"))
|
||||
|
||||
@@ -143,6 +143,9 @@ func (h *Handler) Register(mux *http.ServeMux) {
|
||||
mux.HandleFunc("/api/v1/group/permission/assign", h.groupPermissionAssign)
|
||||
mux.HandleFunc("/api/v1/group/permission/remove", h.groupPermissionRemove)
|
||||
mux.HandleFunc("/api/v1/open_api/sub_store", h.openAPISubStore)
|
||||
mux.HandleFunc("/api/v1/federation/share/list", h.federationShareList)
|
||||
mux.HandleFunc("/api/v1/federation/share/create", h.federationShareCreate)
|
||||
mux.HandleFunc("/api/v1/federation/share/delete", h.federationShareDelete)
|
||||
mux.HandleFunc("/api/v1/federation/connect", h.authPeer(h.federationConnect))
|
||||
mux.HandleFunc("/api/v1/federation/tunnel/create", h.authPeer(h.federationTunnelCreate))
|
||||
mux.HandleFunc("/api/v1/federation/node/import", h.nodeImport)
|
||||
|
||||
@@ -81,6 +81,10 @@ func shouldSkip(path string) bool {
|
||||
return true
|
||||
case path == "/api/v1/user/login":
|
||||
return true
|
||||
case path == "/api/v1/federation/connect":
|
||||
return true
|
||||
case path == "/api/v1/federation/tunnel/create":
|
||||
return true
|
||||
default:
|
||||
return false
|
||||
}
|
||||
@@ -91,6 +95,10 @@ func requiresAdmin(path string) bool {
|
||||
return true
|
||||
}
|
||||
|
||||
if strings.HasPrefix(path, "/api/v1/federation/share/") {
|
||||
return true
|
||||
}
|
||||
|
||||
if strings.HasPrefix(path, "/api/v1/node/") {
|
||||
return true
|
||||
}
|
||||
|
||||
@@ -423,6 +423,22 @@ func (r *Repository) GetNodeBySecret(secret string) (*Node, error) {
|
||||
return &n, nil
|
||||
}
|
||||
|
||||
func (r *Repository) GetNodeByID(id int64) (*Node, error) {
|
||||
if r == nil || r.db == nil {
|
||||
return nil, errors.New("repository not initialized")
|
||||
}
|
||||
|
||||
row := r.db.QueryRow(`SELECT id, secret, version, http, tls, socks, status, is_remote, remote_url, remote_token, remote_config FROM node WHERE id = ? LIMIT 1`, id)
|
||||
var n Node
|
||||
if err := row.Scan(&n.ID, &n.Secret, &n.Version, &n.HTTP, &n.TLS, &n.Socks, &n.Status, &n.IsRemote, &n.RemoteURL, &n.RemoteToken, &n.RemoteConfig); err != nil {
|
||||
if errors.Is(err, sql.ErrNoRows) {
|
||||
return nil, nil
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
return &n, nil
|
||||
}
|
||||
|
||||
func (r *Repository) UpdateNodeOnline(nodeID int64, status int, version string, httpVal, tlsVal, socksVal int) error {
|
||||
if r == nil || r.db == nil {
|
||||
return errors.New("repository not initialized")
|
||||
|
||||
Reference in New Issue
Block a user