Compare commits

...

26 Commits

Author SHA1 Message Date
sagit b32133f81a Merge pull request #92 from Sagit-chu/opencode/lucky-otter
fix(upgrade): stabilize batch node upgrades
2026-02-12 13:37:43 +08:00
sagit ff57bca505 Merge branch 'main' into opencode/lucky-otter 2026-02-12 13:11:45 +08:00
sagit cdb2914dbf fix(upgrade): stabilize batch node upgrades under long-running operations 2026-02-12 04:57:58 +00:00
sagit b56d0a28e7 Merge pull request #91 from Sagit-chu/opencode/lucky-otter
fix(ws): prevent monitor websocket reconnect loop
2026-02-12 12:28:53 +08:00
sagit bdfc704f95 Merge branch 'main' into opencode/lucky-otter 2026-02-12 12:27:23 +08:00
sagit 9223892ca5 fix(ws): prevent monitor websocket reconnect loop 2026-02-12 04:25:47 +00:00
sagit 6f205df37c Merge pull request #90 from Sagit-chu/opencode/lucky-otter
fix(ws): stabilize node connectivity with ping/pong keepalive
2026-02-12 10:29:14 +08:00
sagit 04266165df Merge branch 'main' into opencode/lucky-otter 2026-02-12 10:28:04 +08:00
sagit dbd5773717 fix(ws): stabilize node connectivity with ping/pong keepalive 2026-02-12 02:27:10 +00:00
sagit 6c4d44e7a7 Merge pull request #89 from Sagit-chu/opencode/lucky-otter
fix(store): enable WAL mode and busy timeout to prevent SQLITE_BUSY
2026-02-12 09:09:42 +08:00
sagit 07b8d73956 Merge branch 'main' into opencode/lucky-otter 2026-02-12 09:08:36 +08:00
sagit 71a6a60077 fix(store): enable WAL mode and busy timeout to prevent SQLITE_BUSY errors
Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-opencode)

Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
2026-02-12 01:07:41 +00:00
sagit fe33028934 Merge pull request #87 from Sagit-chu/opencode/lucky-otter
fix: 节点管理按钮改为两行 grid 布局,防止卡片滑动溢出
2026-02-11 16:54:58 +08:00
sagit 01da4bd283 Merge branch 'main' into opencode/lucky-otter 2026-02-11 16:53:15 +08:00
sagit a0b975b62a chore: add go-gost/gost binary to .gitignore 2026-02-11 08:53:03 +00:00
sagit f4e56d091e fix: 节点管理按钮改为两行 grid 布局,防止卡片滑动溢出
- 操作按钮从单行 flex 改为两行 grid (3+2),避免窄屏挤压
- SortableItem 和 Card 添加 overflow-hidden,阻止内容溢出导致页面横滑
2026-02-11 08:45:58 +00:00
sagit 0a5335c1ca Merge pull request #86 from Sagit-chu/opencode/lucky-otter
fix: use systemd-run for agent upgrade/rollback restart
2026-02-11 16:15:55 +08:00
sagit d40e97d73b Merge branch 'main' into opencode/lucky-otter 2026-02-11 16:14:42 +08:00
sagit 73bf672e62 fix: use systemd-run for agent upgrade/rollback restart to avoid cgroup kill 2026-02-11 08:07:04 +00:00
sagit 297f526a92 Merge pull request #85 from Sagit-chu/opencode/lucky-otter
fix: use direct GitHub API for release queries instead of proxy
2026-02-11 15:46:21 +08:00
sagit 2b854d3172 Merge branch 'main' into opencode/lucky-otter 2026-02-11 15:43:18 +08:00
sagit c184a75f22 fix: use direct GitHub API for release queries instead of proxy
The gcode.hostcentral.cc proxy only supports github.com file downloads,
not api.github.com requests. API calls through the proxy returned 404
HTML pages, causing JSON decode error: invalid character '<'.

Also updates repo name from flux-panel to flvx across upgrade and
install URLs.
2026-02-11 07:37:30 +00:00
sagit 8a5bfa5aa8 Merge pull request #84 from Sagit-chu/opencode/lucky-otter
feat: complete agent upgrade system with batch upgrade
2026-02-11 15:09:34 +08:00
sagit 421f18d4da Merge branch 'main' into opencode/lucky-otter 2026-02-11 15:07:07 +08:00
sagit 5133e6f039 fix: remove unused handleUpgradeNode function 2026-02-11 07:04:50 +00:00
sagit 311840b29b feat: complete agent upgrade system with batch upgrade, version selection, progress reporting, and rollback 2026-02-11 07:01:00 +00:00
12 changed files with 922 additions and 37 deletions
+52
View File
@@ -112,6 +112,12 @@ jobs:
upx --best --lzma gost-amd64
upx --best --lzma gost-arm64
- name: Generate SHA256 checksums
working-directory: ./go-gost
run: |
sha256sum gost-amd64 > gost-amd64.sha256
sha256sum gost-arm64 > gost-arm64.sha256
- name: Upload GOST AMD64 artifact
uses: actions/upload-artifact@v4
with:
@@ -124,6 +130,18 @@ jobs:
name: gost-binary-arm64
path: ./go-gost/gost-arm64
- name: Upload GOST AMD64 checksum artifact
uses: actions/upload-artifact@v4
with:
name: gost-checksum-amd64
path: ./go-gost/gost-amd64.sha256
- name: Upload GOST ARM64 checksum artifact
uses: actions/upload-artifact@v4
with:
name: gost-checksum-arm64
path: ./go-gost/gost-arm64.sha256
build-vite:
name: Build & Push Vite Frontend
needs: check-version
@@ -238,7 +256,20 @@ jobs:
name: gost-binary-arm64
path: ./artifacts/arm64
- name: Download GOST AMD64 checksum
uses: actions/download-artifact@v4
with:
name: gost-checksum-amd64
path: ./artifacts/
- name: Download GOST ARM64 checksum
uses: actions/download-artifact@v4
with:
name: gost-checksum-arm64
path: ./artifacts/
- name: Prepare release files
run: |
VERSION="${{ needs.check-version.outputs.version }}"
OWNER="${{ needs.check-version.outputs.image_owner }}"
@@ -331,6 +362,10 @@ jobs:
gh release upload "${VERSION}" ./artifacts/gost-amd64 --clobber
gh release upload "${VERSION}" ./artifacts/gost-arm64 --clobber
echo "📤 上传 GOST 校验文件..."
gh release upload "${VERSION}" ./artifacts/gost-amd64.sha256 --clobber
gh release upload "${VERSION}" ./artifacts/gost-arm64.sha256 --clobber
echo "📤 上传安装脚本..."
gh release upload "${VERSION}" ./artifacts/install.sh --clobber
gh release upload "${VERSION}" ./artifacts/panel_install.sh --clobber
@@ -363,6 +398,18 @@ jobs:
name: gost-binary-arm64
path: ./artifacts/arm64
- name: Download GOST AMD64 checksum
uses: actions/download-artifact@v4
with:
name: gost-checksum-amd64
path: ./artifacts/
- name: Download GOST ARM64 checksum
uses: actions/download-artifact@v4
with:
name: gost-checksum-arm64
path: ./artifacts/
- name: Rename binaries
run: |
mv ./artifacts/amd64/gost-amd64 ./artifacts/gost-amd64
@@ -379,4 +426,9 @@ jobs:
gh release upload "${VERSION}" ./artifacts/gost-amd64 --clobber
gh release upload "${VERSION}" ./artifacts/gost-arm64 --clobber
echo "📤 上传 GOST 校验文件..."
gh release upload "${VERSION}" ./artifacts/gost-amd64.sha256 --clobber
gh release upload "${VERSION}" ./artifacts/gost-arm64.sha256 --clobber
echo "✅ GOST 二进制文件更新完成"
+1
View File
@@ -177,6 +177,7 @@ build/
*.dylib
your_app.exe
go-backend/paneld
go-gost/gost
# Go 测试二进制文件
*.test
@@ -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)
@@ -413,7 +413,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://gcode.hostcentral.cc/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,
}))
}
@@ -165,7 +165,13 @@ func Open(path string) (*Repository, error) {
return nil, err
}
db, 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)"
db, 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
}
}
}
}
+215 -6
View File
@@ -4,10 +4,17 @@ import (
"bytes"
"compress/gzip"
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"fmt"
"io"
"net"
"net/http"
"net/url"
"os"
"os/exec"
"runtime"
"strconv"
"strings"
"sync" // 新增:用于管理连接状态的互斥锁
@@ -21,7 +28,6 @@ import (
"github.com/shirou/gopsutil/v3/host"
"github.com/shirou/gopsutil/v3/mem"
psnet "github.com/shirou/gopsutil/v3/net"
"os"
)
// SystemInfo 系统信息结构体
@@ -85,6 +91,11 @@ type TcpPingResponse struct {
RequestId string `json:"requestId,omitempty"`
}
const (
reporterReadWait = 60 * time.Second
reporterWriteWait = 5 * time.Second
)
type WebSocketReporter struct {
url string
addr string // 保存服务器地址
@@ -237,6 +248,14 @@ func (w *WebSocketReporter) connect() error {
w.conn = conn
w.connected = true
_ = conn.SetReadDeadline(time.Now().Add(reporterReadWait))
conn.SetPingHandler(func(appData string) error {
_ = conn.SetReadDeadline(time.Now().Add(reporterReadWait))
return conn.WriteControl(websocket.PongMessage, []byte(appData), time.Now().Add(reporterWriteWait))
})
conn.SetPongHandler(func(string) error {
return conn.SetReadDeadline(time.Now().Add(reporterReadWait))
})
// 设置关闭处理器来检测连接状态
w.conn.SetCloseHandler(func(code int, text string) error {
@@ -377,7 +396,7 @@ func (w *WebSocketReporter) receiveMessages() {
}
// 设置读取超时
conn.SetReadDeadline(time.Now().Add(30 * time.Second))
conn.SetReadDeadline(time.Now().Add(reporterReadWait))
messageType, message, err := conn.ReadMessage()
if err != nil {
@@ -466,9 +485,8 @@ func (w *WebSocketReporter) handleReceivedMessage(messageType int, message []byt
}
if cmdMsg.Type != "call" {
// TcpPing 诊断命令异步执行,避免阻塞其他命令
// 其他状态变更命令保持同步,确保顺序执行
if cmdMsg.Type == "TcpPing" {
if cmdMsg.Type == "TcpPing" || cmdMsg.Type == "UpgradeAgent" || cmdMsg.Type == "RollbackAgent" {
go w.routeCommand(cmdMsg)
} else {
w.routeCommand(cmdMsg)
@@ -483,9 +501,8 @@ func (w *WebSocketReporter) handleReceivedMessage(messageType int, message []byt
return
}
if cmdMsg.Type != "call" {
// TcpPing 诊断命令异步执行,避免阻塞其他命令
// 其他状态变更命令保持同步,确保顺序执行
if cmdMsg.Type == "TcpPing" {
if cmdMsg.Type == "TcpPing" || cmdMsg.Type == "UpgradeAgent" || cmdMsg.Type == "RollbackAgent" {
go w.routeCommand(cmdMsg)
} else {
w.routeCommand(cmdMsg)
@@ -579,6 +596,18 @@ func (w *WebSocketReporter) routeCommand(cmd CommandMessage) {
response.Type = "SetProtocolResponse"
needSaveConfig = true
// 升级 Agent 命令(异步执行,不需要保存配置)
case "UpgradeAgent":
err = w.handleUpgradeAgent(cmd.Data)
response.Type = "UpgradeAgentResponse"
// needSaveConfig = false (默认值)
// 回退 Agent 到旧版本
case "RollbackAgent":
err = w.handleRollbackAgent(cmd.Data)
response.Type = "RollbackAgentResponse"
// needSaveConfig = false (默认值)
default:
err = fmt.Errorf("未知命令类型: %s", cmd.Type)
response.Type = "UnknownCommandResponse"
@@ -881,6 +910,186 @@ func (w *WebSocketReporter) handleSetProtocol(data interface{}) error {
return nil
}
// sendUpgradeProgress 通过 WS 发送升级进度消息
func (w *WebSocketReporter) sendUpgradeProgress(stage string, percent int, message string) {
response := CommandResponse{
Type: "UpgradeProgress",
Success: true,
Message: message,
Data: map[string]interface{}{
"stage": stage,
"percent": percent,
},
}
w.sendResponse(response)
}
func (w *WebSocketReporter) handleUpgradeAgent(data interface{}) error {
jsonData, err := json.Marshal(data)
if err != nil {
return fmt.Errorf("序列化数据失败: %v", err)
}
var req struct {
DownloadURL string `json:"downloadUrl"`
ChecksumURL string `json:"checksumUrl"`
}
if err := json.Unmarshal(jsonData, &req); err != nil {
return fmt.Errorf("解析升级参数失败: %v", err)
}
if strings.TrimSpace(req.DownloadURL) == "" {
return fmt.Errorf("下载地址不能为空")
}
// 替换架构占位符
downloadURL := strings.ReplaceAll(req.DownloadURL, "{ARCH}", runtime.GOARCH)
checksumURL := strings.ReplaceAll(req.ChecksumURL, "{ARCH}", runtime.GOARCH)
w.sendUpgradeProgress("downloading", 0, "开始下载升级包...")
fmt.Printf("📦 开始下载升级包: %s\n", downloadURL)
// 下载新版本二进制
const binaryPath = "/etc/flux_agent/flux_agent"
tmpPath := binaryPath + ".new"
backupPath := binaryPath + ".old"
resp, err := http.Get(downloadURL)
if err != nil {
return fmt.Errorf("下载升级包失败: %v", err)
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return fmt.Errorf("下载升级包失败, HTTP状态码: %d", resp.StatusCode)
}
outFile, err := os.Create(tmpPath)
if err != nil {
return fmt.Errorf("创建临时文件失败: %v", err)
}
// 带进度的下载
totalSize := resp.ContentLength
var downloaded int64
buf := make([]byte, 32*1024)
lastPercent := 0
hasher := sha256.New()
for {
n, readErr := resp.Body.Read(buf)
if n > 0 {
if _, wErr := outFile.Write(buf[:n]); wErr != nil {
outFile.Close()
os.Remove(tmpPath)
return fmt.Errorf("写入升级包失败: %v", wErr)
}
hasher.Write(buf[:n])
downloaded += int64(n)
if totalSize > 0 {
percent := int(downloaded * 100 / totalSize)
if percent-lastPercent >= 10 {
lastPercent = percent
w.sendUpgradeProgress("downloading", percent, fmt.Sprintf("下载中... %d%%", percent))
}
}
}
if readErr != nil {
if readErr == io.EOF {
break
}
outFile.Close()
os.Remove(tmpPath)
return fmt.Errorf("读取升级包失败: %v", readErr)
}
}
outFile.Close()
if downloaded == 0 {
os.Remove(tmpPath)
return fmt.Errorf("下载的升级包为空")
}
w.sendUpgradeProgress("downloading", 100, fmt.Sprintf("下载完成 (%d bytes)", downloaded))
// Checksum 校验
if checksumURL != "" {
w.sendUpgradeProgress("verifying", 0, "校验文件完整性...")
checksumResp, err := http.Get(checksumURL)
if err == nil {
defer checksumResp.Body.Close()
if checksumResp.StatusCode == http.StatusOK {
checksumBody, err := io.ReadAll(checksumResp.Body)
if err == nil {
// 格式: "<hash> <filename>" 或 "<hash>"
expectedHash := strings.TrimSpace(strings.Split(string(checksumBody), " ")[0])
actualHash := hex.EncodeToString(hasher.Sum(nil))
if !strings.EqualFold(expectedHash, actualHash) {
os.Remove(tmpPath)
return fmt.Errorf("校验失败: 期望 %s, 实际 %s", expectedHash, actualHash)
}
fmt.Printf("✅ Checksum 校验通过: %s\n", actualHash)
}
}
}
w.sendUpgradeProgress("verifying", 100, "校验通过")
}
if err := os.Chmod(tmpPath, 0755); err != nil {
os.Remove(tmpPath)
return fmt.Errorf("设置执行权限失败: %v", err)
}
// 备份旧版本
w.sendUpgradeProgress("installing", 50, "备份旧版本...")
if _, err := os.Stat(binaryPath); err == nil {
// 复制旧文件作为备份(不用 rename,因为可能正在运行)
oldData, err := os.ReadFile(binaryPath)
if err == nil {
_ = os.WriteFile(backupPath, oldData, 0755)
fmt.Println("📦 旧版本已备份到", backupPath)
}
}
w.sendUpgradeProgress("installing", 80, "准备重启...")
fmt.Printf("✅ 升级包下载完成 (%d bytes), 准备重启...\n", downloaded)
// 执行重启脚本
// 使用 systemd-run 在独立的 transient unit 中运行重启脚本,
// 避免 systemctl stop 杀死 flux_agent cgroup 内所有进程(包括此脚本自身)导致 mv 未执行。
script := fmt.Sprintf("sleep 1 && systemctl stop flux_agent && mv %s %s && systemctl start flux_agent", tmpPath, binaryPath)
cmd := exec.Command("systemd-run", "--quiet", "/bin/sh", "-c", script)
if err := cmd.Start(); err != nil {
os.Remove(tmpPath)
return fmt.Errorf("启动重启脚本失败: %v", err)
}
w.sendUpgradeProgress("installing", 100, "重启中...")
fmt.Println("🔄 重启脚本已启动, Agent 将在 1 秒后重启...")
return nil
}
func (w *WebSocketReporter) handleRollbackAgent(data interface{}) error {
const binaryPath = "/etc/flux_agent/flux_agent"
backupPath := binaryPath + ".old"
// 检查备份文件是否存在
if _, err := os.Stat(backupPath); os.IsNotExist(err) {
return fmt.Errorf("没有可用的备份文件,无法回退")
}
fmt.Println("🔄 开始回退到旧版本...")
// 执行回退脚本(同升级逻辑,使用 systemd-run 避免 cgroup 问题)
script := fmt.Sprintf("sleep 1 && systemctl stop flux_agent && cp %s %s && systemctl start flux_agent", backupPath, binaryPath)
cmd := exec.Command("systemd-run", "--quiet", "/bin/sh", "-c", script)
if err := cmd.Start(); err != nil {
return fmt.Errorf("启动回退脚本失败: %v", err)
}
fmt.Println("🔄 回退脚本已启动, Agent 将在 1 秒后重启...")
return nil
}
// updateLocalConfigJSON 将 http/tls/socks 写入工作目录下的 config.json
func updateLocalConfigJSON(httpVal int, tlsVal int, socksVal int) error {
path := "config.json"
+3 -1
View File
@@ -87,6 +87,8 @@ http {
proxy_http_version 1.1;
proxy_set_header Upgrade $http_upgrade;
proxy_set_header Connection "upgrade";
proxy_read_timeout 3600s;
proxy_send_timeout 3600s;
proxy_set_header Host $host;
proxy_set_header X-Real-IP $remote_addr;
@@ -94,4 +96,4 @@ http {
proxy_set_header X-Forwarded-Proto $scheme;
}
}
}
}
+8
View File
@@ -41,6 +41,14 @@ export const checkNodeStatus = (nodeId?: number) => {
return Network.post("/node/check-status", params);
};
export const upgradeNode = (id: number, version?: string) =>
Network.post("/node/upgrade", { id, version: version || "" }, { timeout: 5 * 60 * 1000 });
export const batchUpgradeNodes = (ids: number[], version?: string) =>
Network.post("/node/batch-upgrade", { ids, version: version || "" }, { timeout: 15 * 60 * 1000 });
export const getNodeReleases = () => Network.post("/node/releases");
export const rollbackNode = (id: number) =>
Network.post("/node/rollback", { id });
// 隧道CRUD操作 - 全部使用POST请求
export const createTunnel = (data: any) => Network.post("/tunnel/create", data);
export const getTunnelList = () => Network.post("/tunnel/list");
+8 -2
View File
@@ -43,6 +43,10 @@ interface ApiResponse<T = any> {
data: T;
}
interface RequestOptions {
timeout?: number;
}
// 处理token失效的逻辑
function handleTokenExpired() {
// 清除localStorage中的token
@@ -71,6 +75,7 @@ const Network = {
get: function <T = any>(
path: string = "",
data: any = {},
options: RequestOptions = {},
): Promise<ApiResponse<T>> {
return new Promise(function (resolve) {
// 如果baseURL是默认值且是WebView环境,说明没有设置面板地址
@@ -83,7 +88,7 @@ const Network = {
axios
.get(path, {
params: data,
timeout: 30000,
timeout: options.timeout ?? 30000,
headers: {
Authorization: window.localStorage.getItem("token"),
},
@@ -117,6 +122,7 @@ const Network = {
post: function <T = any>(
path: string = "",
data: any = {},
options: RequestOptions = {},
): Promise<ApiResponse<T>> {
return new Promise(function (resolve) {
// 如果baseURL是默认值且是WebView环境,说明没有设置面板地址
@@ -128,7 +134,7 @@ const Network = {
axios
.post(path, data, {
timeout: 30000,
timeout: options.timeout ?? 30000,
headers: {
Authorization: window.localStorage.getItem("token"),
"Content-Type": "application/json",
+268 -25
View File
@@ -16,6 +16,7 @@ import { Spinner } from "@heroui/spinner";
import { Alert } from "@heroui/alert";
import { Progress } from "@heroui/progress";
import { Accordion, AccordionItem } from "@heroui/accordion";
import { Select, SelectItem } from "@heroui/select";
import { Checkbox } from "@heroui/checkbox";
import toast from "react-hot-toast";
import axios from "axios";
@@ -45,6 +46,10 @@ import {
getNodeInstallCommand,
updateNodeOrder,
batchDeleteNodes,
upgradeNode,
batchUpgradeNodes,
getNodeReleases,
rollbackNode,
} from "@/api";
interface Node {
@@ -77,6 +82,8 @@ interface Node {
uptime: number;
} | null;
copyLoading?: boolean;
upgradeLoading?: boolean;
rollbackLoading?: boolean;
}
interface NodeForm {
@@ -118,7 +125,7 @@ const SortableItem = ({
};
return (
<div ref={setNodeRef} style={style} {...attributes}>
<div ref={setNodeRef} style={style} {...attributes} className="overflow-hidden">
{children(listeners)}
</div>
);
@@ -165,6 +172,16 @@ export default function NodePage() {
const [installCommand, setInstallCommand] = useState("");
const [currentNodeName, setCurrentNodeName] = useState("");
// 升级相关状态
const [upgradeModalOpen, setUpgradeModalOpen] = useState(false);
const [upgradeTarget, setUpgradeTarget] = useState<"single" | "batch">("single");
const [upgradeTargetNodeId, setUpgradeTargetNodeId] = useState<number | null>(null);
const [releases, setReleases] = useState<Array<{ version: string; name: string; publishedAt: string; prerelease: boolean }>>([]);
const [releasesLoading, setReleasesLoading] = useState(false);
const [selectedVersion, setSelectedVersion] = useState("");
const [batchUpgradeLoading, setBatchUpgradeLoading] = useState(false);
const [upgradeProgress, setUpgradeProgress] = useState<Record<number, { stage: string; percent: number; message: string }>>({});
const websocketRef = useRef<WebSocket | null>(null);
const reconnectTimerRef = useRef<NodeJS.Timeout | null>(null);
const reconnectAttemptsRef = useRef(0);
@@ -425,6 +442,22 @@ export default function NodePage() {
return node;
}),
);
} else if (type === "upgrade_progress") {
try {
const progressData = typeof messageData === "string" ? JSON.parse(messageData) : messageData;
if (progressData?.data) {
setUpgradeProgress((prev) => ({
...prev,
[nodeId]: {
stage: progressData.data.stage || "",
percent: progressData.data.percent || 0,
message: progressData.message || "",
},
}));
}
} catch {
// ignore parse errors
}
}
};
@@ -770,6 +803,93 @@ export default function NodePage() {
}
};
// 打开版本选择弹窗
const openUpgradeModal = async (target: "single" | "batch", nodeId?: number) => {
setUpgradeTarget(target);
setUpgradeTargetNodeId(nodeId || null);
setSelectedVersion("");
setUpgradeModalOpen(true);
setReleasesLoading(true);
try {
const res = await getNodeReleases();
if (res.code === 0 && Array.isArray(res.data)) {
setReleases(res.data);
} else {
toast.error(res.msg || "获取版本列表失败");
}
} catch {
toast.error("获取版本列表失败");
} finally {
setReleasesLoading(false);
}
};
// 确认升级(从版本弹窗)
const handleConfirmUpgrade = async () => {
const version = selectedVersion || undefined;
if (upgradeTarget === "single" && upgradeTargetNodeId) {
setUpgradeModalOpen(false);
// Find the node
const node = nodeList.find((n) => n.id === upgradeTargetNodeId);
if (!node) return;
setNodeList((prev) =>
prev.map((n) => (n.id === upgradeTargetNodeId ? { ...n, upgradeLoading: true } : n)),
);
try {
const res = await upgradeNode(upgradeTargetNodeId, version);
if (res.code === 0) {
toast.success(`节点升级命令已发送,节点将自动重启`);
} else {
toast.error(res.msg || "升级失败");
}
} catch {
toast.error("网络错误,请重试");
} finally {
setNodeList((prev) =>
prev.map((n) => (n.id === upgradeTargetNodeId ? { ...n, upgradeLoading: false } : n)),
);
}
} else if (upgradeTarget === "batch") {
setBatchUpgradeLoading(true);
setUpgradeModalOpen(false);
try {
const res = await batchUpgradeNodes(Array.from(selectedIds), version);
if (res.code === 0) {
toast.success(`批量升级命令已发送到 ${selectedIds.size} 个节点`);
} else {
toast.error(res.msg || "批量升级失败");
}
} catch {
toast.error("网络错误,请重试");
} finally {
setBatchUpgradeLoading(false);
}
}
};
// 回退节点
const handleRollbackNode = async (node: Node) => {
setNodeList((prev) =>
prev.map((n) => (n.id === node.id ? { ...n, rollbackLoading: true } : n)),
);
try {
const res = await rollbackNode(node.id);
if (res.code === 0) {
toast.success(`节点 ${node.name} 回退命令已发送,节点将自动重启`);
} else {
toast.error(res.msg || "回退失败");
}
} catch {
toast.error("网络错误,请重试");
} finally {
setNodeList((prev) =>
prev.map((n) => (n.id === node.id ? { ...n, rollbackLoading: false } : n)),
);
}
};
// 提交表单
const handleSubmit = async () => {
if (!validateForm()) return;
@@ -1048,6 +1168,15 @@ export default function NodePage() {
<Button size="sm" variant="flat" onPress={deselectAll}>
清空
</Button>
<Button
color="warning"
isLoading={batchUpgradeLoading}
size="sm"
variant="flat"
onPress={() => openUpgradeModal("batch")}
>
升级
</Button>
<Button
color="danger"
size="sm"
@@ -1124,7 +1253,7 @@ export default function NodePage() {
{(listeners) => (
<Card
key={node.id}
className="group shadow-sm border border-divider hover:shadow-md transition-shadow duration-200"
className="group shadow-sm border border-divider hover:shadow-md transition-shadow duration-200 overflow-hidden"
>
<CardHeader className="pb-2">
<div className="flex justify-between items-start w-full">
@@ -1239,6 +1368,18 @@ export default function NodePage() {
{node.version || "未知"}
</span>
</div>
{upgradeProgress[node.id] && upgradeProgress[node.id].percent < 100 && (
<div className="mt-1">
<Progress
aria-label="升级进度"
color="warning"
label={upgradeProgress[node.id].message}
showValueLabel
size="sm"
value={upgradeProgress[node.id].percent}
/>
</div>
)}
<div className="flex justify-between text-sm">
<span className="text-default-600">开机时间</span>
<span className="text-xs">
@@ -1373,32 +1514,56 @@ export default function NodePage() {
{/* 操作按钮 */}
<div className="space-y-1.5">
<div className="flex gap-1.5">
{!isRemoteNode && (
<div className="grid grid-cols-3 gap-1.5">
<Button
className="min-h-8"
color="success"
isLoading={node.copyLoading}
size="sm"
variant="flat"
onPress={() => handleCopyInstallCommand(node)}
>
安装
</Button>
<Button
className="min-h-8"
color="warning"
isDisabled={node.connectionStatus !== "online"}
isLoading={node.upgradeLoading}
size="sm"
variant="flat"
onPress={() => openUpgradeModal("single", node.id)}
>
升级
</Button>
<Button
className="min-h-8"
color="secondary"
isDisabled={node.connectionStatus !== "online"}
isLoading={node.rollbackLoading}
size="sm"
variant="flat"
onPress={() => handleRollbackNode(node)}
>
回退
</Button>
</div>
)}
<div className={`grid gap-1.5 ${isRemoteNode ? "grid-cols-1" : "grid-cols-2"}`}>
{!isRemoteNode && (
<>
<Button
className="flex-1 min-h-8"
color="success"
isLoading={node.copyLoading}
size="sm"
variant="flat"
onPress={() => handleCopyInstallCommand(node)}
>
安装
</Button>
<Button
className="flex-1 min-h-8"
color="primary"
size="sm"
variant="flat"
onPress={() => handleEdit(node)}
>
编辑
</Button>
</>
<Button
className="min-h-8"
color="primary"
size="sm"
variant="flat"
onPress={() => handleEdit(node)}
>
编辑
</Button>
)}
<Button
className={`min-h-8 ${isRemoteNode ? "w-full" : "flex-1"}`}
className="min-h-8"
color="danger"
size="sm"
variant="flat"
@@ -1844,6 +2009,84 @@ export default function NodePage() {
</ModalContent>
</Modal>
{/* 版本选择升级模态框 */}
<Modal
backdrop="blur"
isOpen={upgradeModalOpen}
placement="center"
scrollBehavior="outside"
size="md"
onOpenChange={setUpgradeModalOpen}
>
<ModalContent>
{(onClose) => (
<>
<ModalHeader className="flex flex-col gap-1">
<h2 className="text-xl font-bold">
{upgradeTarget === "batch"
? `批量升级 (${selectedIds.size} 个节点)`
: "升级节点"}
</h2>
</ModalHeader>
<ModalBody>
{releasesLoading ? (
<div className="flex justify-center py-8">
<Spinner size="lg" />
</div>
) : (
<div className="space-y-4">
<Select
label="选择版本"
placeholder="留空则使用最新版本"
selectedKeys={selectedVersion ? [selectedVersion] : []}
onSelectionChange={(keys) => {
const selected = Array.from(keys)[0] as string;
setSelectedVersion(selected || "");
}}
>
{releases.map((r) => (
<SelectItem key={r.version} textValue={r.version}>
<div className="flex justify-between items-center">
<span>{r.version}</span>
<span className="text-xs text-default-400">
{r.publishedAt
? new Date(r.publishedAt).toLocaleDateString()
: ""}
{r.prerelease && (
<Chip className="ml-1" color="warning" size="sm" variant="flat">
预览
</Chip>
)}
</span>
</div>
</SelectItem>
))}
</Select>
<p className="text-sm text-default-500">
{selectedVersion
? `将升级到版本 ${selectedVersion}`
: "未选择版本,将自动使用最新稳定版"}
</p>
</div>
)}
</ModalBody>
<ModalFooter>
<Button variant="light" onPress={onClose}>
取消
</Button>
<Button
color="warning"
isDisabled={releasesLoading}
onPress={handleConfirmUpgrade}
>
确认升级
</Button>
</ModalFooter>
</>
)}
</ModalContent>
</Modal>
{/* 批量删除确认模态框 */}
<Modal
backdrop="blur"