From 633d683c0e8c3356ce197b643cd68c47d649a082 Mon Sep 17 00:00:00 2001 From: duorameng <2997944583@qq.com> Date: Fri, 12 Jun 2026 18:48:33 +0800 Subject: [PATCH] feat: implement precise worker status tracking and monitoring UI --- internal/controllers/monitor_controller.go | 114 +++++ internal/executor/scheduler.go | 69 ++- internal/router/api_routes.go | 9 + internal/router/register.go | 1 + internal/router/router.go | 1 + web/src/api/index.ts | 47 +++ web/src/layouts/MainLayout.vue | 3 +- web/src/router/index.ts | 1 + web/src/views/monitor/Monitor.vue | 464 +++++++++++++++++++++ 9 files changed, 707 insertions(+), 2 deletions(-) create mode 100644 internal/controllers/monitor_controller.go create mode 100644 web/src/views/monitor/Monitor.vue diff --git a/internal/controllers/monitor_controller.go b/internal/controllers/monitor_controller.go new file mode 100644 index 0000000..baf5656 --- /dev/null +++ b/internal/controllers/monitor_controller.go @@ -0,0 +1,114 @@ +package controllers + +import ( + "net/http" + "runtime" + "time" + + "github.com/engigu/baihu-panel/internal/services/tasks" + "github.com/engigu/baihu-panel/internal/utils" + + "github.com/gin-gonic/gin" + "github.com/gorilla/websocket" +) + +type MonitorController struct { + executorService *tasks.ExecutorService +} + +func NewMonitorController(executorService *tasks.ExecutorService) *MonitorController { + return &MonitorController{ + executorService: executorService, + } +} + +var monitorUpgrader = websocket.Upgrader{ + CheckOrigin: func(r *http.Request) bool { + return true // 开发环境允许所有跨域,生产环境可根据配置限制 + }, +} + +// GetSystemMonitor 获取系统和内存监控信息 (HTTP) +func (mc *MonitorController) GetSystemMonitor(c *gin.Context) { + data := mc.getMonitorData() + utils.Success(c, data) +} + +// MonitorWS WebSocket实时推送系统监控信息 +func (mc *MonitorController) MonitorWS(c *gin.Context) { + ws, err := monitorUpgrader.Upgrade(c.Writer, c.Request, nil) + if err != nil { + return + } + defer ws.Close() + + ticker := time.NewTicker(3 * time.Second) + defer ticker.Stop() + + // 先立即发送一次 + mc.sendMonitorData(ws) + + for { + select { + case <-ticker.C: + if err := mc.sendMonitorData(ws); err != nil { + return // 客户端断开连接或发送失败 + } + case <-c.Request.Context().Done(): + return + } + } +} + +func (mc *MonitorController) sendMonitorData(ws *websocket.Conn) error { + data := mc.getMonitorData() + return ws.WriteJSON(gin.H{ + "code": 200, + "data": data, + "msg": "success", + }) +} + +func (mc *MonitorController) getMonitorData() gin.H { + var m runtime.MemStats + runtime.ReadMemStats(&m) + + return gin.H{ + "env": gin.H{ + "os": runtime.GOOS, + "arch": runtime.GOARCH, + "go_version": runtime.Version(), + "num_cpu": runtime.NumCPU(), + "goroutines": runtime.NumGoroutine(), + }, + "mem": gin.H{ + "alloc": m.Alloc, + "total_alloc": m.TotalAlloc, + "sys": m.Sys, + "lookups": m.Lookups, + "mallocs": m.Mallocs, + "frees": m.Frees, + }, + "heap": gin.H{ + "heap_alloc": m.HeapAlloc, + "heap_sys": m.HeapSys, + "heap_idle": m.HeapIdle, + "heap_inuse": m.HeapInuse, + "heap_released": m.HeapReleased, + "heap_objects": m.HeapObjects, + }, + "gc": gin.H{ + "next_gc": m.NextGC, + "last_gc": m.LastGC, + "pause_total_ns": m.PauseTotalNs, + "num_gc": m.NumGC, + }, + "scheduler": gin.H{ + "scheduled": mc.executorService.GetScheduledCount(), + "running": mc.executorService.GetRunningCount(), + "queue_size": mc.executorService.GetScheduler().GetQueueSize(), + "worker_count": mc.executorService.GetScheduler().GetConfig().WorkerCount, + "workers": mc.executorService.GetScheduler().GetWorkerStatuses(), + }, + } +} diff --git a/internal/executor/scheduler.go b/internal/executor/scheduler.go index 8a06e50..4c7bf3b 100644 --- a/internal/executor/scheduler.go +++ b/internal/executor/scheduler.go @@ -174,6 +174,16 @@ func (h *schedulerHooksAdapter) OnHeartbeat(ctx context.Context, logID string, d // TaskExecutor 定义任务执行函数签名 type TaskExecutor func(ctx context.Context, req *ExecutionRequest, stdout, stderr io.Writer) (*Result, error) +// WorkerStatus 定义并发池中单个 Worker 的状态 +type WorkerStatus struct { + ID int `json:"id"` + Status string `json:"status"` // 状态: "idle" 或 "running" + TaskID string `json:"task_id,omitempty"` + TaskName string `json:"task_name,omitempty"` + StartTime int64 `json:"start_time,omitempty"` // 开始时间戳 (秒) + Duration int64 `json:"duration,omitempty"` // 已运行时长 (秒) +} + // Scheduler 统一调度器(独立组件,可在主服务和 Agent 中复用) // 调度器本身只负责队列管理和任务调度,具体的执行逻辑和事件处理由 Handler 实现 type Scheduler struct { @@ -188,6 +198,9 @@ type Scheduler struct { logger SchedulerLogger runningTasks map[string]context.CancelFunc // 记录运行中的任务,用于停止 (TaskID -> CancelFunc) runningExecs map[string]context.CancelFunc // 记录运行中的执行,用于停止 (LogID -> CancelFunc) + + workers []WorkerStatus + workerMu sync.RWMutex } // NewScheduler 创建调度器 @@ -224,6 +237,14 @@ func NewScheduler(config SchedulerConfig, handler SchedulerEventHandler) *Schedu logger: &DefaultLogger{}, runningTasks: make(map[string]context.CancelFunc), runningExecs: make(map[string]context.CancelFunc), + workers: make([]WorkerStatus, config.WorkerCount), + } + + for i := 0; i < config.WorkerCount; i++ { + s.workers[i] = WorkerStatus{ + ID: i, + Status: "idle", + } } return s @@ -332,7 +353,32 @@ func (s *Scheduler) worker(id int) { }() // 速率限制 <-s.rateLimiter - s.executeTask(req) + + func() { + // 恢复 worker 状态为空闲 + defer func() { + s.workerMu.Lock() + if id >= 0 && id < len(s.workers) { + s.workers[id].Status = "idle" + s.workers[id].TaskID = "" + s.workers[id].TaskName = "" + s.workers[id].StartTime = 0 + } + s.workerMu.Unlock() + }() + + // 更新 worker 状态为运行中 + s.workerMu.Lock() + if id >= 0 && id < len(s.workers) { + s.workers[id].Status = "running" + s.workers[id].TaskID = req.TaskID + s.workers[id].TaskName = req.Name + s.workers[id].StartTime = time.Now().Unix() + } + s.workerMu.Unlock() + + s.executeTask(req) + }() }() } } @@ -615,3 +661,24 @@ func (s *Scheduler) GetConfig() SchedulerConfig { defer s.mu.RUnlock() return s.config } + +// GetWorkerStatuses 获取所有 Worker 的状态 +func (s *Scheduler) GetWorkerStatuses() []WorkerStatus { + s.workerMu.RLock() + defer s.workerMu.RUnlock() + // 返回副本防止外部修改 + statuses := make([]WorkerStatus, len(s.workers)) + now := time.Now().Unix() + for i, w := range s.workers { + statuses[i] = w + // 在服务端计算运行时间,彻底避免客户端与服务端时钟不一致导致的计算偏差 + if w.Status == "running" && w.StartTime > 0 { + duration := now - w.StartTime + if duration < 0 { + duration = 0 + } + statuses[i].Duration = duration + } + } + return statuses +} diff --git a/internal/router/api_routes.go b/internal/router/api_routes.go index f065f6a..784107b 100644 --- a/internal/router/api_routes.go +++ b/internal/router/api_routes.go @@ -65,6 +65,7 @@ func initAuthorizedAPIRoutes(api *gin.RouterGroup, c *Controllers) { registerAppLogRoutes(adminOnly, c) registerSystemWSRoutes(adminOnly, c) registerWebUIRoutes(adminOnly, c) + registerMonitorRoutes(adminOnly, c) } } @@ -268,6 +269,14 @@ func registerSystemWSRoutes(g *gin.RouterGroup, c *Controllers) { g.GET("/ws/events", c.SystemWS.HandleEvents) } +func registerMonitorRoutes(g *gin.RouterGroup, c *Controllers) { + monitor := g.Group("/monitor") + { + monitor.GET("", c.Monitor.GetSystemMonitor) + monitor.GET("/ws", c.Monitor.MonitorWS) + } +} + func initAgentAPIRoutes(root *gin.RouterGroup, c *Controllers) { // Agent API(供远程 Agent 调用,不使用 /v1 版本号) agentAPI := root.Group("/api/agent") diff --git a/internal/router/register.go b/internal/router/register.go index 62e0668..8065f3d 100644 --- a/internal/router/register.go +++ b/internal/router/register.go @@ -64,6 +64,7 @@ func RegisterControllers() *Controllers { AppLog: controllers.NewAppLogController(), SystemWS: controllers.NewSystemWSController(), WebUI: controllers.NewWebUIController(services.NewWebUIService(settingsService)), + Monitor: controllers.NewMonitorController(executorService), } } diff --git a/internal/router/router.go b/internal/router/router.go index db2c83e..0156d44 100644 --- a/internal/router/router.go +++ b/internal/router/router.go @@ -31,6 +31,7 @@ type Controllers struct { AppLog *controllers.AppLogController SystemWS *controllers.SystemWSController WebUI *controllers.WebUIController + Monitor *controllers.MonitorController } func Setup(c *Controllers) *gin.Engine { diff --git a/web/src/api/index.ts b/web/src/api/index.ts index 547e4f3..d9533b0 100644 --- a/web/src/api/index.ts +++ b/web/src/api/index.ts @@ -9,6 +9,52 @@ interface ApiResponse { data: T } +export interface MonitorStats { + env: { + os: string + arch: string + go_version: string + num_cpu: number + goroutines: number + } + mem: { + alloc: number + total_alloc: number + sys: number + lookups: number + mallocs: number + frees: number + } + heap: { + heap_alloc: number + heap_sys: number + heap_idle: number + heap_inuse: number + heap_released: number + heap_objects: number + } + gc: { + next_gc: number + last_gc: number + pause_total_ns: number + num_gc: number + } + scheduler: { + scheduled: number + running: number + queue_size: number + worker_count: number + workers: { + id: number + status: string + task_id?: string + task_name?: string + start_time?: number + duration?: number + }[] + } +} + async function request(url: string, options?: RequestInit): Promise { const res = await fetch(`${API_BASE_URL}${url}`, { ...options, @@ -145,6 +191,7 @@ export const api = { taskStats: (days?: number) => request(`/taskstats${days ? `?days=${days}` : ''}`) }, settings: { + getMonitor: () => request('/monitor'), changePassword: (data: { old_username?: string; username?: string; old_password: string; new_password?: string }) => request('/settings/password', { method: 'POST', body: JSON.stringify(data) }), getSite: () => request('/settings/site'), diff --git a/web/src/layouts/MainLayout.vue b/web/src/layouts/MainLayout.vue index f4a139e..5711668 100644 --- a/web/src/layouts/MainLayout.vue +++ b/web/src/layouts/MainLayout.vue @@ -2,7 +2,7 @@ import { ref, onMounted, computed } from 'vue' import { RouterLink, RouterView, useRoute } from 'vue-router' import { resetAuthCache } from '@/router' -import { LayoutDashboard, ListTodo, FileCode, Settings, LogOut, ScrollText, Terminal, Variable, KeyRound, Menu, X, Server, Globe, Bell } from 'lucide-vue-next' +import { LayoutDashboard, ListTodo, FileCode, Settings, LogOut, ScrollText, Terminal, Variable, KeyRound, Menu, X, Server, Globe, Bell, Activity } from 'lucide-vue-next' import { Button } from '@/components/ui/button' import ThemeToggle from '@/components/ThemeToggle.vue' import SystemNotice from '@/components/SystemNotice.vue' @@ -63,6 +63,7 @@ const navItems = [ { to: '/terminal', icon: Terminal, label: '终端命令', exact: true }, { to: '/notify', icon: Bell, label: '消息推送', exact: true }, { to: '/logs', icon: KeyRound, label: '运行日志', exact: true }, + { to: '/monitor', icon: Activity, label: '系统监控', exact: true }, { to: '/settings', icon: Settings, label: '系统设置', exact: true }, ] diff --git a/web/src/router/index.ts b/web/src/router/index.ts index 604edc3..484f01b 100644 --- a/web/src/router/index.ts +++ b/web/src/router/index.ts @@ -54,6 +54,7 @@ const router = createRouter({ { path: 'logs', name: 'logs', component: () => import('@/views/logs/MessageLogs.vue') }, { path: 'terminal', name: 'terminal', component: () => import('@/views/terminal/Terminal.vue') }, { path: 'notify', name: 'notify', component: () => import('@/views/notify/Notify.vue') }, + { path: 'monitor', name: 'monitor', component: () => import('@/views/monitor/Monitor.vue') }, { path: 'settings', name: 'settings', component: () => import('@/views/settings/Settings.vue') } ] }, diff --git a/web/src/views/monitor/Monitor.vue b/web/src/views/monitor/Monitor.vue new file mode 100644 index 0000000..3dafb55 --- /dev/null +++ b/web/src/views/monitor/Monitor.vue @@ -0,0 +1,464 @@ + + +