From eb9a2a88149e307c364d12a2bc7eff0345eeec15 Mon Sep 17 00:00:00 2001 From: ryan Date: Sun, 15 Mar 2026 11:57:52 +0800 Subject: [PATCH] =?UTF-8?q?[=E4=BC=98=E5=8C=96]=20=E6=B7=BB=E5=8A=A0?= =?UTF-8?q?=E6=9C=8D=E5=8A=A1=E5=99=A8=E5=8D=87=E7=BA=A7=E6=97=A5=E5=BF=97?= =?UTF-8?q?=E7=9A=84WebSocket=E6=B5=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- atsf_server/Dockerfile | 2 +- atsf_server/controller/update.go | 38 ++++++++ atsf_server/router/api-router.go | 1 + atsf_server/service/update.go | 89 +++++++++++++++++++ .../components/layout/dashboard-topbar.tsx | 73 ++++++++++++++- atsf_server/web/features/update/api/update.ts | 42 ++++++++- atsf_server/web/features/update/types.ts | 6 ++ docs/deployment.md | 2 +- docs/development-guidelines.md | 7 +- 9 files changed, 255 insertions(+), 5 deletions(-) diff --git a/atsf_server/Dockerfile b/atsf_server/Dockerfile index 926b41c8..7dc67fd3 100644 --- a/atsf_server/Dockerfile +++ b/atsf_server/Dockerfile @@ -11,7 +11,7 @@ RUN corepack enable && pnpm install --frozen-lockfile COPY ./web ./ RUN NEXT_PUBLIC_APP_VERSION="$VERSION" pnpm build -FROM golang:1.23 AS builder2 +FROM golang:1.24 AS builder2 ARG VERSION diff --git a/atsf_server/controller/update.go b/atsf_server/controller/update.go index 56f6f237..a0755e42 100644 --- a/atsf_server/controller/update.go +++ b/atsf_server/controller/update.go @@ -6,8 +6,10 @@ import ( "io" "net/http" "strings" + "time" "github.com/gin-gonic/gin" + "golang.org/x/net/websocket" ) type confirmManualUpgradeRequest struct { @@ -74,6 +76,42 @@ func UpgradeServer(c *gin.Context) { }) } +// StreamServerUpgradeLogs godoc +// @Summary Stream server upgrade logs over websocket +// @Tags Update +// @Router /api/update/logs/ws [get] +func StreamServerUpgradeLogs(c *gin.Context) { + websocket.Handler(func(conn *websocket.Conn) { + defer func() { + _ = conn.Close() + }() + + updates, unsubscribe := service.SubscribeServerUpgradeStream() + defer unsubscribe() + + heartbeatTicker := time.NewTicker(15 * time.Second) + defer heartbeatTicker.Stop() + + for { + select { + case snapshot, ok := <-updates: + if !ok { + return + } + if err := websocket.JSON.Send(conn, snapshot); err != nil { + return + } + case <-heartbeatTicker.C: + if err := websocket.JSON.Send(conn, service.ServerUpgradeStreamSnapshot{}); err != nil { + return + } + case <-c.Request.Context().Done(): + return + } + } + }).ServeHTTP(c.Writer, c.Request) +} + // UploadManualServerBinary godoc // @Summary Upload server binary and inspect version before upgrade // @Tags Update diff --git a/atsf_server/router/api-router.go b/atsf_server/router/api-router.go index f561f94f..ea5b39e8 100644 --- a/atsf_server/router/api-router.go +++ b/atsf_server/router/api-router.go @@ -59,6 +59,7 @@ func SetApiRouter(router *gin.Engine) { updateRoute.Use(middleware.RootAuth(), middleware.NoTokenAuth()) { updateRoute.GET("/latest-release", controller.GetLatestRelease) + updateRoute.GET("/logs/ws", controller.StreamServerUpgradeLogs) updateRoute.POST("/manual-upload", controller.UploadManualServerBinary) updateRoute.POST("/manual-upgrade", controller.ConfirmManualServerUpgrade) updateRoute.POST("/upgrade", controller.UpgradeServer) diff --git a/atsf_server/service/update.go b/atsf_server/service/update.go index 749a38d9..b373e1df 100644 --- a/atsf_server/service/update.go +++ b/atsf_server/service/update.go @@ -44,6 +44,12 @@ var serverUpgradeState struct { logs []ServerUpgradeLogRecord } +var serverUpgradeSubscribers struct { + sync.Mutex + nextID int + listeners map[int]chan ServerUpgradeStreamSnapshot +} + var manualServerBinaryState struct { sync.Mutex candidate *manualServerBinaryCandidate @@ -74,6 +80,12 @@ type ServerUpgradeLogRecord struct { CreatedAt time.Time `json:"created_at"` } +type ServerUpgradeStreamSnapshot struct { + InProgress bool `json:"in_progress"` + UpgradeStatus string `json:"upgrade_status"` + UpgradeLogs []ServerUpgradeLogRecord `json:"upgrade_logs"` +} + type githubReleaseResponse struct { TagName string `json:"tag_name"` Body string `json:"body"` @@ -137,17 +149,23 @@ func ScheduleServerUpgrade(channel string) (*LatestServerRelease, error) { resetServerUpgradeLogsLocked() serverUpgradeState.status = "running" appendServerUpgradeLogLocked("info", fmt.Sprintf("Automatic upgrade scheduled for channel: %s.", normalizedChannel.String())) + serverUpgradeState.Unlock() + broadcastServerUpgradeSnapshot() prepared, err := prepareServerUpgrade(context.Background(), normalizedChannel) if err != nil { + serverUpgradeState.Lock() serverUpgradeState.status = "failed" appendServerUpgradeLogLocked("error", err.Error()) serverUpgradeState.Unlock() + broadcastServerUpgradeSnapshot() return nil, err } + serverUpgradeState.Lock() serverUpgradeState.inProgress = true serverUpgradeState.Unlock() + broadcastServerUpgradeSnapshot() prepared.release.InProgress = true @@ -262,6 +280,7 @@ func ConfirmManualServerUpgrade(uploadToken string) (*UploadedServerBinary, erro serverUpgradeState.status = "running" appendServerUpgradeLogLocked("info", fmt.Sprintf("Manual upgrade confirmed for version: %s.", strings.TrimSpace(candidate.DetectedVersion))) serverUpgradeState.Unlock() + broadcastServerUpgradeSnapshot() go func(task *manualServerBinaryCandidate) { time.Sleep(serverUpgradeDispatchDelay) @@ -852,6 +871,15 @@ func snapshotServerUpgradeState() (bool, string, []ServerUpgradeLogRecord) { return serverUpgradeState.inProgress, status, logs } +func snapshotServerUpgradeStream() ServerUpgradeStreamSnapshot { + inProgress, status, logs := snapshotServerUpgradeState() + return ServerUpgradeStreamSnapshot{ + InProgress: inProgress, + UpgradeStatus: status, + UpgradeLogs: logs, + } +} + func resetServerUpgradeLogsLocked() { serverUpgradeState.logs = nil } @@ -871,6 +899,7 @@ func recordServerUpgradeLog(level string, message string) { serverUpgradeState.Lock() appendServerUpgradeLogLocked(level, message) serverUpgradeState.Unlock() + broadcastServerUpgradeSnapshot() } func markServerUpgradeSucceeded() { @@ -879,6 +908,7 @@ func markServerUpgradeSucceeded() { serverUpgradeState.status = "succeeded" appendServerUpgradeLogLocked("info", "Upgrade binary is ready; server restart will begin.") serverUpgradeState.Unlock() + broadcastServerUpgradeSnapshot() } func recordServerUpgradeFailure(err error) { @@ -889,6 +919,65 @@ func recordServerUpgradeFailure(err error) { appendServerUpgradeLogLocked("error", err.Error()) } serverUpgradeState.Unlock() + broadcastServerUpgradeSnapshot() +} + +func SubscribeServerUpgradeStream() (<-chan ServerUpgradeStreamSnapshot, func()) { + serverUpgradeSubscribers.Lock() + if serverUpgradeSubscribers.listeners == nil { + serverUpgradeSubscribers.listeners = make(map[int]chan ServerUpgradeStreamSnapshot) + } + serverUpgradeSubscribers.nextID++ + listenerID := serverUpgradeSubscribers.nextID + listener := make(chan ServerUpgradeStreamSnapshot, 8) + serverUpgradeSubscribers.listeners[listenerID] = listener + serverUpgradeSubscribers.Unlock() + + listener <- snapshotServerUpgradeStream() + + unsubscribe := func() { + serverUpgradeSubscribers.Lock() + ch, ok := serverUpgradeSubscribers.listeners[listenerID] + if ok { + delete(serverUpgradeSubscribers.listeners, listenerID) + } + serverUpgradeSubscribers.Unlock() + if ok { + close(ch) + } + } + + return listener, unsubscribe +} + +func broadcastServerUpgradeSnapshot() { + snapshot := snapshotServerUpgradeStream() + + serverUpgradeSubscribers.Lock() + if len(serverUpgradeSubscribers.listeners) == 0 { + serverUpgradeSubscribers.Unlock() + return + } + listeners := make([]chan ServerUpgradeStreamSnapshot, 0, len(serverUpgradeSubscribers.listeners)) + for _, listener := range serverUpgradeSubscribers.listeners { + listeners = append(listeners, listener) + } + serverUpgradeSubscribers.Unlock() + + for _, listener := range listeners { + select { + case listener <- snapshot: + default: + select { + case <-listener: + default: + } + select { + case listener <- snapshot: + default: + } + } + } } func UpdateHTTPClientForTest() *http.Client { diff --git a/atsf_server/web/components/layout/dashboard-topbar.tsx b/atsf_server/web/components/layout/dashboard-topbar.tsx index 7148dcfb..8987308f 100644 --- a/atsf_server/web/components/layout/dashboard-topbar.tsx +++ b/atsf_server/web/components/layout/dashboard-topbar.tsx @@ -8,14 +8,18 @@ import { useAuth } from '@/components/providers/auth-provider'; import { ThemeToggle } from '@/components/ui/theme-toggle'; import { getPublicStatus } from '@/features/auth/api/public'; import { + createUpgradeLogsWebSocket, confirmManualServerUpgrade, getLatestRelease, + parseUpgradeStreamSnapshot, upgradeServer, uploadServerBinary, } from '@/features/update/api/update'; import { VersionUpgradeModal } from '@/features/update/components/version-upgrade-modal'; import type { + LatestReleaseInfo, ReleaseChannel, + UpgradeStreamSnapshot, UploadedServerBinaryInfo, } from '@/features/update/types'; import { publicEnv } from '@/lib/env/public-env'; @@ -46,6 +50,8 @@ export function DashboardTopbar() { const [uploadedBinary, setUploadedBinary] = useState(null); const [uploadProgress, setUploadProgress] = useState(0); + const [upgradeStream, setUpgradeStream] = + useState(null); const menuRef = useRef(null); const isRoot = (user?.role ?? 0) >= 100; const upgradeStatusPollInterval = 3000; @@ -146,6 +152,51 @@ export function DashboardTopbar() { }, }); + useEffect(() => { + if (!isVersionModalOpen || !isRoot) { + setUpgradeStream(null); + return; + } + + let closed = false; + let reconnectTimer: number | null = null; + let socket: WebSocket | null = null; + + const connect = () => { + if (closed) { + return; + } + + socket = createUpgradeLogsWebSocket(); + if (!socket) { + return; + } + + socket.onmessage = (event) => { + const snapshot = parseUpgradeStreamSnapshot(String(event.data)); + if (snapshot) { + setUpgradeStream(snapshot); + } + }; + + socket.onclose = () => { + if (!closed) { + reconnectTimer = window.setTimeout(connect, 1500); + } + }; + }; + + connect(); + + return () => { + closed = true; + if (reconnectTimer !== null) { + window.clearTimeout(reconnectTimer); + } + socket?.close(); + }; + }, [isRoot, isVersionModalOpen]); + useEffect(() => { if (!isUserMenuOpen) { return; @@ -245,6 +296,10 @@ export function DashboardTopbar() { selectedReleaseChannel === 'preview' ? previewReleaseQuery.data : stableReleaseQuery.data; + const releaseWithStream = mergeReleaseWithUpgradeStream( + selectedRelease, + upgradeStream, + ); const selectedReleaseError = selectedReleaseChannel === 'preview' ? previewReleaseQuery.error @@ -345,7 +400,7 @@ export function DashboardTopbar() { isOpen={isVersionModalOpen} onClose={() => setIsVersionModalOpen(false)} currentVersion={currentVersion} - release={selectedRelease} + release={releaseWithStream} selectedChannel={selectedReleaseChannel} uploadedBinary={uploadedBinary} isLoading={ @@ -376,3 +431,19 @@ export function DashboardTopbar() { ); } + +function mergeReleaseWithUpgradeStream( + release: LatestReleaseInfo | null | undefined, + stream: UpgradeStreamSnapshot | null, +) { + if (!release || !stream) { + return release; + } + + return { + ...release, + in_progress: stream.in_progress, + upgrade_status: stream.upgrade_status, + upgrade_logs: stream.upgrade_logs, + }; +} diff --git a/atsf_server/web/features/update/api/update.ts b/atsf_server/web/features/update/api/update.ts index 7b5e8dad..88ce39a7 100644 --- a/atsf_server/web/features/update/api/update.ts +++ b/atsf_server/web/features/update/api/update.ts @@ -4,6 +4,7 @@ import type { ApiEnvelope } from '@/types/api'; import type { LatestReleaseInfo, ReleaseChannel, + UpgradeStreamSnapshot, UploadedServerBinaryInfo, } from '@/features/update/types'; @@ -48,7 +49,9 @@ export function uploadServerBinary( xhr.addEventListener('load', () => { let payload: ApiEnvelope | null = null; try { - payload = JSON.parse(xhr.responseText) as ApiEnvelope; + payload = JSON.parse( + xhr.responseText, + ) as ApiEnvelope; } catch { payload = null; } @@ -86,3 +89,40 @@ export function confirmManualServerUpgrade(uploadToken: string) { body: JSON.stringify({ upload_token: uploadToken }), }); } + +export function createUpgradeLogsWebSocket() { + if (typeof window === 'undefined') { + return null; + } + + const apiUrl = getApiUrl('/update/logs/ws'); + const resolvedUrl = apiUrl.startsWith('http://') + ? `ws://${apiUrl.slice('http://'.length)}` + : apiUrl.startsWith('https://') + ? `wss://${apiUrl.slice('https://'.length)}` + : `${window.location.protocol === 'https:' ? 'wss:' : 'ws:'}//${window.location.host}${apiUrl}`; + + return new WebSocket(resolvedUrl); +} + +export function parseUpgradeStreamSnapshot( + rawMessage: string, +): UpgradeStreamSnapshot | null { + try { + const parsed = JSON.parse(rawMessage) as Partial; + if ( + typeof parsed.in_progress !== 'boolean' || + typeof parsed.upgrade_status !== 'string' || + !Array.isArray(parsed.upgrade_logs) + ) { + return null; + } + return { + in_progress: parsed.in_progress, + upgrade_status: parsed.upgrade_status, + upgrade_logs: parsed.upgrade_logs, + }; + } catch { + return null; + } +} diff --git a/atsf_server/web/features/update/types.ts b/atsf_server/web/features/update/types.ts index 8fef0f28..1a2ced8c 100644 --- a/atsf_server/web/features/update/types.ts +++ b/atsf_server/web/features/update/types.ts @@ -21,6 +21,12 @@ export interface LatestReleaseInfo { upgrade_logs: UpgradeLogItem[]; } +export interface UpgradeStreamSnapshot { + in_progress: boolean; + upgrade_status: 'idle' | 'running' | 'succeeded' | 'failed' | string; + upgrade_logs: UpgradeLogItem[]; +} + export interface UploadedServerBinaryInfo { upload_token: string; file_name: string; diff --git a/docs/deployment.md b/docs/deployment.md index 829b763c..a2e31c71 100644 --- a/docs/deployment.md +++ b/docs/deployment.md @@ -8,7 +8,7 @@ ### 1.1 Server -* Go 1.23+ +* Go 1.24+ * Node.js 18+ * 可写 SQLite 文件目录 diff --git a/docs/development-guidelines.md b/docs/development-guidelines.md index a1a5c4ed..f0c9c051 100644 --- a/docs/development-guidelines.md +++ b/docs/development-guidelines.md @@ -22,7 +22,7 @@ `atsf_server` 继续作为单体控制面: -* Go 1.23+ +* Go 1.24+ * Gin * GORM * SQLite @@ -216,6 +216,11 @@ Agent: * 写入 `config_versions` * 通过切换 `is_active` 激活版本 +Go 版本基线约束: + +* 当 `go.mod` 中的 Go 主版本或次版本发生变化时,必须同步检查并更新所有相关构建入口,至少包括 Docker 构建使用的基础镜像版本以及 GitHub Actions 中的 release / docker 发布工作流 +* 版本升级后必须确保本地构建、Docker 构建与发布工作流使用一致的 Go 版本,避免因为 `go.mod`、`Dockerfile` 与 CI 工作流版本漂移导致发布失败 + 第五版新增要求: * “完整 OpenResty 配置”至少包括主配置文件与路由配置文件