Merge branch 'main' into opencode/glowing-orchid

Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-opencode)

Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
This commit is contained in:
sagit
2026-02-12 06:52:06 +00:00
17 changed files with 1340 additions and 66 deletions
@@ -85,6 +85,14 @@ func NewFederationClient() *FederationClient {
}
}
func NewFederationClientWithTimeout(timeout time.Duration) *FederationClient {
return &FederationClient{
client: &http.Client{
Timeout: timeout,
},
}
}
func (c *FederationClient) Connect(url, token, localDomain string) (*RemoteNodeInfo, error) {
url = strings.TrimSuffix(url, "/")
req, err := http.NewRequest("POST", url+"/api/v1/federation/connect", nil)
@@ -8,6 +8,7 @@ import (
"net/http"
"sort"
"strings"
"sync"
"time"
"go-backend/internal/http/client"
@@ -40,6 +41,17 @@ type resetPeerShareFlowRequest struct {
ID int64 `json:"id"`
}
type updatePeerShareRequest struct {
ID int64 `json:"id"`
Name string `json:"name"`
MaxBandwidth int64 `json:"maxBandwidth"`
ExpiryTime int64 `json:"expiryTime"`
PortRangeStart int `json:"portRangeStart"`
PortRangeEnd int `json:"portRangeEnd"`
AllowedDomains string `json:"allowedDomains"`
AllowedIPs string `json:"allowedIps"`
}
type nodeImportRequest struct {
RemoteURL string `json:"remoteUrl"`
Token string `json:"token"`
@@ -325,6 +337,80 @@ func (h *Handler) federationShareResetFlow(w http.ResponseWriter, r *http.Reques
response.WriteJSON(w, response.OKEmpty())
}
func (h *Handler) federationShareUpdate(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
response.WriteJSON(w, response.ErrDefault("Invalid method"))
return
}
var req updatePeerShareRequest
if err := decodeJSON(r.Body, &req); err != nil {
response.WriteJSON(w, response.ErrDefault("Invalid JSON"))
return
}
if req.ID <= 0 {
response.WriteJSON(w, response.ErrDefault("Share ID is required"))
return
}
share, err := h.repo.GetPeerShare(req.ID)
if err != nil {
response.WriteJSON(w, response.Err(-2, err.Error()))
return
}
if share == nil {
response.WriteJSON(w, response.ErrDefault("Share not found"))
return
}
if req.Name == "" {
response.WriteJSON(w, response.ErrDefault("Name is 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
}
allowedIPs, err := normalizePeerShareAllowedIPs(req.AllowedIPs)
if err != nil {
response.WriteJSON(w, response.ErrDefault(err.Error()))
return
}
share.Name = req.Name
share.MaxBandwidth = req.MaxBandwidth
share.ExpiryTime = req.ExpiryTime
share.PortRangeStart = req.PortRangeStart
share.PortRangeEnd = req.PortRangeEnd
share.AllowedDomains = req.AllowedDomains
share.AllowedIPs = allowedIPs
share.UpdatedTime = time.Now().UnixMilli()
if err := h.repo.UpdatePeerShare(share); err != nil {
response.WriteJSON(w, response.Err(-2, err.Error()))
return
}
response.WriteJSON(w, response.OKEmpty())
}
func (h *Handler) federationRemoteUsageList(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
response.WriteJSON(w, response.ErrDefault("Invalid method"))
@@ -1315,6 +1401,74 @@ func isPeerIPAllowed(clientIP net.IP, whitelist string) bool {
return false
}
func (h *Handler) syncRemoteNodeStatuses(items []map[string]interface{}) {
type remoteEntry struct {
index int
remoteURL string
remoteToken string
}
var remotes []remoteEntry
for i, item := range items {
isRemote, _ := item["isRemote"].(int)
if isRemote != 1 {
continue
}
url, _ := item["remoteUrl"].(string)
token, _ := item["remoteToken"].(string)
url = strings.TrimSpace(url)
token = strings.TrimSpace(token)
if url == "" || token == "" {
continue
}
remotes = append(remotes, remoteEntry{index: i, remoteURL: url, remoteToken: token})
}
if len(remotes) == 0 {
return
}
localDomain := h.federationLocalDomain()
fc := client.NewFederationClientWithTimeout(5 * time.Second)
type syncResult struct {
index int
status int
syncError string
}
results := make([]syncResult, len(remotes))
var wg sync.WaitGroup
for i, entry := range remotes {
wg.Add(1)
go func(idx int, e remoteEntry) {
defer wg.Done()
info, err := fc.Connect(e.remoteURL, e.remoteToken, localDomain)
if err != nil {
errMsg := err.Error()
if strings.Contains(errMsg, "401") || strings.Contains(errMsg, "Invalid token") || strings.Contains(errMsg, "Unauthorized") {
results[idx] = syncResult{index: e.index, status: 0, syncError: "provider_share_deleted"}
} else if strings.Contains(errMsg, "403") || strings.Contains(errMsg, "Share is disabled") {
results[idx] = syncResult{index: e.index, status: 0, syncError: "provider_share_disabled"}
} else if strings.Contains(errMsg, "Share expired") {
results[idx] = syncResult{index: e.index, status: 0, syncError: "provider_share_expired"}
} else {
results[idx] = syncResult{index: e.index, status: 0, syncError: errMsg}
}
} else {
results[idx] = syncResult{index: e.index, status: info.Status, syncError: ""}
}
}(i, entry)
}
wg.Wait()
for _, r := range results {
items[r.index]["status"] = r.status
if r.syncError != "" {
items[r.index]["syncError"] = r.syncError
}
}
}
func (h *Handler) cleanupPeerShareRuntimes(shareID int64) {
if h == nil || h.repo == nil || shareID <= 0 {
return
@@ -105,6 +105,10 @@ func (h *Handler) Register(mux *http.ServeMux) {
mux.HandleFunc("/api/v1/node/update-order", h.nodeUpdateOrder)
mux.HandleFunc("/api/v1/node/batch-delete", h.nodeBatchDelete)
mux.HandleFunc("/api/v1/node/check-status", h.nodeCheckStatus)
mux.HandleFunc("/api/v1/node/upgrade", h.nodeUpgrade)
mux.HandleFunc("/api/v1/node/batch-upgrade", h.nodeBatchUpgrade)
mux.HandleFunc("/api/v1/node/rollback", h.nodeRollback)
mux.HandleFunc("/api/v1/node/releases", h.listReleases)
mux.HandleFunc("/api/v1/tunnel/list", h.tunnelList)
mux.HandleFunc("/api/v1/tunnel/create", h.tunnelCreate)
mux.HandleFunc("/api/v1/tunnel/get", h.tunnelGet)
@@ -155,6 +159,7 @@ func (h *Handler) Register(mux *http.ServeMux) {
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/update", h.federationShareUpdate)
mux.HandleFunc("/api/v1/federation/share/delete", h.federationShareDelete)
mux.HandleFunc("/api/v1/federation/share/reset-flow", h.federationShareResetFlow)
mux.HandleFunc("/api/v1/federation/share/remote-usage/list", h.federationRemoteUsageList)
@@ -320,6 +325,9 @@ func (h *Handler) nodeList(w http.ResponseWriter, r *http.Request) {
response.WriteJSON(w, response.Err(-2, err.Error()))
return
}
h.syncRemoteNodeStatuses(items)
response.WriteJSON(w, response.OK(items))
}
@@ -414,7 +414,7 @@ func (h *Handler) nodeInstall(w http.ResponseWriter, r *http.Request) {
response.WriteJSON(w, response.Err(-2, err.Error()))
return
}
cmd := fmt.Sprintf("curl -L https://github.com/Sagit-chu/flux-panel/releases/latest/download/install.sh -o ./install.sh && chmod +x ./install.sh && ./install.sh -a %s -s %s", processServerAddress(panelAddr), secret)
cmd := fmt.Sprintf("curl -L https://gcode.hostcentral.cc/https://github.com/Sagit-chu/flvx/releases/latest/download/install.sh -o ./install.sh && chmod +x ./install.sh && ./install.sh -a %s -s %s", processServerAddress(panelAddr), secret)
response.WriteJSON(w, response.OK(cmd))
}
+291
View File
@@ -0,0 +1,291 @@
package handler
import (
"encoding/json"
"fmt"
"io"
"net/http"
"strings"
"sync"
"time"
"go-backend/internal/http/response"
)
const (
githubRepo = "Sagit-chu/flvx"
githubProxy = "https://gcode.hostcentral.cc"
githubAPIBase = "https://api.github.com"
githubHTMLBase = "https://github.com"
upgradeTimeout = 5 * time.Minute
batchWorkers = 5
)
func (h *Handler) nodeUpgrade(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
response.WriteJSON(w, response.ErrDefault("请求失败"))
return
}
var req struct {
ID int64 `json:"id"`
Version string `json:"version"`
}
if err := decodeJSON(r.Body, &req); err != nil {
response.WriteJSON(w, response.ErrDefault("请求参数错误"))
return
}
if req.ID <= 0 {
response.WriteJSON(w, response.ErrDefault("节点ID无效"))
return
}
version := strings.TrimSpace(req.Version)
if version == "" {
var err error
version, err = resolveLatestRelease()
if err != nil {
response.WriteJSON(w, response.Err(-2, fmt.Sprintf("获取最新版本失败: %v", err)))
return
}
}
downloadURL := fmt.Sprintf(
githubProxy+"/%s/%s/releases/download/%s/gost-{ARCH}",
githubHTMLBase, githubRepo, version,
)
checksumURL := fmt.Sprintf(
githubProxy+"/%s/%s/releases/download/%s/gost-{ARCH}.sha256",
githubHTMLBase, githubRepo, version,
)
result, err := h.wsServer.SendCommand(req.ID, "UpgradeAgent", map[string]interface{}{
"downloadUrl": downloadURL,
"checksumUrl": checksumURL,
}, upgradeTimeout)
if err != nil {
response.WriteJSON(w, response.Err(-2, fmt.Sprintf("升级失败: %v", err)))
return
}
response.WriteJSON(w, response.OK(map[string]interface{}{
"version": version,
"message": result.Message,
}))
}
func resolveLatestRelease() (string, error) {
client := &http.Client{
CheckRedirect: func(req *http.Request, via []*http.Request) error {
return http.ErrUseLastResponse
},
Timeout: 10 * time.Second,
}
resp, err := client.Get(githubProxy + "/" + githubHTMLBase + "/" + githubRepo + "/releases/latest")
if err != nil {
return "", fmt.Errorf("请求GitHub失败: %v", err)
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusFound && resp.StatusCode != http.StatusMovedPermanently {
return resolveLatestReleaseAPI()
}
location := resp.Header.Get("Location")
if location == "" {
return resolveLatestReleaseAPI()
}
parts := strings.Split(location, "/")
tag := parts[len(parts)-1]
if tag == "" || tag == "latest" {
return resolveLatestReleaseAPI()
}
return tag, nil
}
func resolveLatestReleaseAPI() (string, error) {
client := &http.Client{Timeout: 10 * time.Second}
resp, err := client.Get(githubAPIBase + "/repos/" + githubRepo + "/releases/latest")
if err != nil {
return "", fmt.Errorf("请求GitHub API失败: %v", err)
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
body, _ := io.ReadAll(io.LimitReader(resp.Body, 512))
return "", fmt.Errorf("GitHub API返回 %d: %s", resp.StatusCode, string(body))
}
var release struct {
TagName string `json:"tag_name"`
}
if err := json.NewDecoder(resp.Body).Decode(&release); err != nil {
return "", fmt.Errorf("解析GitHub API响应失败: %v", err)
}
if strings.TrimSpace(release.TagName) == "" {
return "", fmt.Errorf("无法从GitHub获取最新版本号")
}
return release.TagName, nil
}
func (h *Handler) nodeBatchUpgrade(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
response.WriteJSON(w, response.ErrDefault("请求失败"))
return
}
var req struct {
IDs []int64 `json:"ids"`
Version string `json:"version"`
}
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
response.WriteJSON(w, response.ErrDefault("请求参数错误"))
return
}
if len(req.IDs) == 0 {
response.WriteJSON(w, response.ErrDefault("ids不能为空"))
return
}
version := strings.TrimSpace(req.Version)
if version == "" {
var err error
version, err = resolveLatestRelease()
if err != nil {
response.WriteJSON(w, response.Err(-2, fmt.Sprintf("获取最新版本失败: %v", err)))
return
}
}
downloadURL := fmt.Sprintf(
githubProxy+"/%s/%s/releases/download/%s/gost-{ARCH}",
githubHTMLBase, githubRepo, version,
)
checksumURL := fmt.Sprintf(
githubProxy+"/%s/%s/releases/download/%s/gost-{ARCH}.sha256",
githubHTMLBase, githubRepo, version,
)
type upgradeResult struct {
ID int64 `json:"id"`
Success bool `json:"success"`
Message string `json:"message"`
}
results := make([]upgradeResult, len(req.IDs))
sem := make(chan struct{}, batchWorkers)
var wg sync.WaitGroup
for i, id := range req.IDs {
wg.Add(1)
go func(index int, nodeID int64) {
defer wg.Done()
sem <- struct{}{}
defer func() { <-sem }()
result, err := h.wsServer.SendCommand(nodeID, "UpgradeAgent", map[string]interface{}{
"downloadUrl": downloadURL,
"checksumUrl": checksumURL,
}, upgradeTimeout)
if err != nil {
results[index] = upgradeResult{ID: nodeID, Success: false, Message: err.Error()}
return
}
results[index] = upgradeResult{ID: nodeID, Success: true, Message: result.Message}
}(i, id)
}
wg.Wait()
response.WriteJSON(w, response.OK(map[string]interface{}{
"version": version,
"results": results,
}))
}
func (h *Handler) listReleases(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
response.WriteJSON(w, response.ErrDefault("请求失败"))
return
}
client := &http.Client{Timeout: 15 * time.Second}
resp, err := client.Get(githubAPIBase + "/repos/" + githubRepo + "/releases?per_page=20")
if err != nil {
response.WriteJSON(w, response.Err(-2, fmt.Sprintf("获取版本列表失败: %v", err)))
return
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
body, _ := io.ReadAll(io.LimitReader(resp.Body, 512))
response.WriteJSON(w, response.Err(-2, fmt.Sprintf("获取版本列表失败: GitHub API返回 %d: %s", resp.StatusCode, string(body))))
return
}
var releases []struct {
TagName string `json:"tag_name"`
Name string `json:"name"`
PublishedAt string `json:"published_at"`
Prerelease bool `json:"prerelease"`
Draft bool `json:"draft"`
}
if err := json.NewDecoder(resp.Body).Decode(&releases); err != nil {
response.WriteJSON(w, response.Err(-2, fmt.Sprintf("解析版本列表失败: %v", err)))
return
}
type releaseItem struct {
Version string `json:"version"`
Name string `json:"name"`
PublishedAt string `json:"publishedAt"`
Prerelease bool `json:"prerelease"`
}
items := make([]releaseItem, 0, len(releases))
for _, r := range releases {
if r.Draft {
continue
}
items = append(items, releaseItem{
Version: r.TagName,
Name: r.Name,
PublishedAt: r.PublishedAt,
Prerelease: r.Prerelease,
})
}
response.WriteJSON(w, response.OK(items))
}
func (h *Handler) nodeRollback(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
response.WriteJSON(w, response.ErrDefault("请求失败"))
return
}
var req struct {
ID int64 `json:"id"`
}
if err := decodeJSON(r.Body, &req); err != nil {
response.WriteJSON(w, response.ErrDefault("请求参数错误"))
return
}
if req.ID <= 0 {
response.WriteJSON(w, response.ErrDefault("节点ID无效"))
return
}
result, err := h.wsServer.SendCommand(req.ID, "RollbackAgent", map[string]interface{}{}, 30*time.Second)
if err != nil {
response.WriteJSON(w, response.Err(-2, fmt.Sprintf("回退失败: %v", err)))
return
}
response.WriteJSON(w, response.OK(map[string]interface{}{
"message": result.Message,
}))
}
@@ -168,7 +168,13 @@ func Open(path string) (*Repository, error) {
return nil, err
}
raw, err := sql.Open("sqlite", path)
// Use _pragma DSN parameters so every connection from the pool gets
// the same settings (busy_timeout and synchronous are per-connection).
dsn := "file:" + path +
"?_pragma=busy_timeout(5000)" +
"&_pragma=journal_mode(WAL)" +
"&_pragma=synchronous(NORMAL)"
raw, err := sql.Open("sqlite", dsn)
if err != nil {
return nil, err
}
+64 -1
View File
@@ -54,6 +54,12 @@ type pendingRequest struct {
ch chan CommandResult
}
const (
wsPingPeriod = 15 * time.Second
wsPongWait = 45 * time.Second
wsWriteWait = 5 * time.Second
)
type CommandResult struct {
Type string `json:"type"`
Success bool `json:"success"`
@@ -120,12 +126,19 @@ func (s *Server) handleAdmin(w http.ResponseWriter, r *http.Request) {
return
}
cw := &connWrap{conn: conn}
_ = conn.SetReadDeadline(time.Now().Add(wsPongWait))
conn.SetPongHandler(func(string) error {
return conn.SetReadDeadline(time.Now().Add(wsPongWait))
})
done := make(chan struct{})
go startKeepalive(cw, done)
s.mu.Lock()
s.admins[cw] = struct{}{}
s.mu.Unlock()
defer func() {
close(done)
s.mu.Lock()
delete(s.admins, cw)
s.mu.Unlock()
@@ -145,6 +158,12 @@ func (s *Server) handleNode(w http.ResponseWriter, r *http.Request, nodeID int64
return
}
cw := &connWrap{conn: conn}
_ = conn.SetReadDeadline(time.Now().Add(wsPongWait))
conn.SetPongHandler(func(string) error {
return conn.SetReadDeadline(time.Now().Add(wsPongWait))
})
done := make(chan struct{})
go startKeepalive(cw, done)
version := r.URL.Query().Get("version")
httpVal := parseIntDefault(r.URL.Query().Get("http"), 0)
@@ -165,6 +184,7 @@ func (s *Server) handleNode(w http.ResponseWriter, r *http.Request, nodeID int64
s.broadcastStatus(nodeID, 1)
defer func() {
close(done)
needOfflineBroadcast := false
s.mu.Lock()
current, ok := s.nodes[nodeID]
@@ -190,7 +210,15 @@ func (s *Server) handleNode(w http.ResponseWriter, r *http.Request, nodeID int64
msg := decryptIfNeeded(payload, secret)
s.tryResolvePending(nodeID, msg)
s.broadcastInfo(nodeID, msg)
var parsed struct {
Type string `json:"type"`
}
if json.Unmarshal([]byte(msg), &parsed) == nil && parsed.Type == "UpgradeProgress" {
s.broadcastTyped(nodeID, "upgrade_progress", msg)
} else {
s.broadcastInfo(nodeID, msg)
}
}
}
@@ -264,7 +292,9 @@ func (s *Server) SendCommand(nodeID int64, cmdType string, data interface{}, tim
}
ns.conn.mu.Lock()
_ = ns.conn.conn.SetWriteDeadline(time.Now().Add(wsWriteWait))
err = ns.conn.conn.WriteMessage(websocket.TextMessage, messageData)
_ = ns.conn.conn.SetWriteDeadline(time.Time{})
ns.conn.mu.Unlock()
if err != nil {
cleanup()
@@ -385,6 +415,12 @@ func (s *Server) broadcastInfo(nodeID int64, data string) {
s.broadcastToAdmins(string(raw))
}
func (s *Server) broadcastTyped(nodeID int64, msgType string, data string) {
payload := broadcastMessage{ID: nodeID, Type: msgType, Data: data}
raw, _ := json.Marshal(payload)
s.broadcastToAdmins(string(raw))
}
func (s *Server) broadcastToAdmins(message string) {
s.mu.RLock()
admins := make([]*connWrap, 0, len(s.admins))
@@ -395,7 +431,9 @@ func (s *Server) broadcastToAdmins(message string) {
for _, c := range admins {
c.mu.Lock()
_ = c.conn.SetWriteDeadline(time.Now().Add(wsWriteWait))
err := c.conn.WriteMessage(websocket.TextMessage, []byte(message))
_ = c.conn.SetWriteDeadline(time.Time{})
c.mu.Unlock()
if err != nil {
log.Printf("websocket broadcast failed: %v", err)
@@ -428,3 +466,28 @@ func parseIntDefault(v string, fallback int) int {
}
return x
}
func startKeepalive(cw *connWrap, done <-chan struct{}) {
if cw == nil || cw.conn == nil {
return
}
ticker := time.NewTicker(wsPingPeriod)
defer ticker.Stop()
for {
select {
case <-done:
return
case <-ticker.C:
cw.mu.Lock()
_ = cw.conn.SetWriteDeadline(time.Now().Add(wsWriteWait))
err := cw.conn.WriteMessage(websocket.PingMessage, nil)
_ = cw.conn.SetWriteDeadline(time.Time{})
cw.mu.Unlock()
if err != nil {
_ = cw.conn.Close()
return
}
}
}
}