mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-11 16:57:45 +02:00
fix: Display connected plugins in admin UI
- Fix ShowPlugins handler to fetch plugins from registry instead of returning empty
- Update PluginsPageData struct to use []map[string]interface{} for proper data handling
- Rewrite plugins.templ to properly iterate and display plugin details
- Plugin count now displays correctly in the card header
- Plugin table shows ID, Name, Status, Version, Capabilities, and Actions
- Plugin worker successfully connects and displays in admin dashboard
fix: Add plugin routes to non-auth section and resolve route conflicts
- Plugin routes were only registered when authRequired=true (password set)
- When no admin password was set, auth was disabled and routes were skipped
- Also changed route paths to avoid conflicts in Gin router:
- Changed /jobs/:type to /jobs/by-type/:type to avoid conflict with /jobs/:id/cancel
- Changed /jobs/:type/trigger-detection to /trigger-detection/:type
- Changed /jobs/:id/cancel to /cancel-job/:id
- Plugin UI now accessible at http://localhost:23646/plugins
feat: Add plugin_worker command for new plugin system
- Create new generic plugin worker that connects to admin server via gRPC
- Supports multiple plugins: erasure_coding, vacuum, balance
- Replaces old task-based worker system with plugin-based approach
- Automatically registers with admin server on startup
- Sends periodic health reports to admin server
- Configuration saved in working directory
- Usage: weed plugin_worker -admin=localhost:33650 -plugins=erasure_coding,vacuum,balance
fix: Initialize plugin manager and register PluginService on gRPC server
- Initialize plugin manager in admin_server.initPluginManager() instead of placeholder
- Create plugin configuration directory in admin dataDir/plugins
- Register PluginService, AdminQueryService, and AdminCommandService on worker gRPC server
- Plugin worker can now connect and register with admin server
Changes:
- weed/admin/dash/admin_server.go: Properly initialize plugin manager
- weed/admin/dash/worker_grpc_server.go: Register plugin services on gRPC server
Testing:
- Plugin worker connects successfully to admin server
- Plugin capabilities are registered correctly
- Health reporting works as expected
This commit is contained in:
1 parent
3241275885
commit
5342328836
8 files changed
+623
-80
No files matched your search
@@ -4,12 +4,15 @@ import (
|
||||
"context"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sort"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/seaweedfs/seaweedfs/weed/admin/maintenance"
|
||||
"github.com/seaweedfs/seaweedfs/weed/admin/plugin"
|
||||
"github.com/seaweedfs/seaweedfs/weed/cluster"
|
||||
"github.com/seaweedfs/seaweedfs/weed/credential"
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
@@ -254,10 +257,24 @@ func (s *AdminServer) GetCredentialManager() *credential.CredentialManager {
|
||||
|
||||
// initPluginManager initializes the plugin manager
|
||||
func (s *AdminServer) initPluginManager(dataDir string) {
|
||||
// For now, keep pluginManager as interface{} to avoid circular imports
|
||||
// This will be set later when handlers are initialized
|
||||
// We'll check if it's nil when needed
|
||||
glog.V(1).Infof("Plugin manager placeholder initialized")
|
||||
// Create plugin configuration directory if it doesn't exist
|
||||
pluginConfigDir := filepath.Join(dataDir, "plugins")
|
||||
if err := os.MkdirAll(pluginConfigDir, 0755); err != nil {
|
||||
glog.Warningf("Failed to create plugin config directory: %v", err)
|
||||
return
|
||||
}
|
||||
|
||||
// Create plugin manager with default configuration
|
||||
config := plugin.DefaultManagerConfig(pluginConfigDir)
|
||||
pm, err := plugin.NewManager(config)
|
||||
if err != nil {
|
||||
glog.Warningf("Failed to initialize plugin manager: %v", err)
|
||||
return
|
||||
}
|
||||
|
||||
// Store the plugin manager
|
||||
s.pluginManager = pm
|
||||
glog.Infof("Plugin manager initialized successfully")
|
||||
}
|
||||
|
||||
// GetPluginManager returns the plugin manager
|
||||
|
||||
@@ -9,8 +9,10 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/admin/maintenance"
|
||||
"github.com/seaweedfs/seaweedfs/weed/admin/plugin"
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/worker_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/security"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
@@ -94,6 +96,19 @@ func (s *WorkerGrpcServer) StartWithTLS(port int) error {
|
||||
|
||||
worker_pb.RegisterWorkerServiceServer(grpcServer, s)
|
||||
|
||||
// Register plugin service if plugin manager is available
|
||||
if s.adminServer.GetPluginManager() != nil {
|
||||
// Cast the interface{} to *plugin.Manager
|
||||
if pm, ok := s.adminServer.GetPluginManager().(*plugin.Manager); ok {
|
||||
if pluginGrpcServer := pm.GetGRPCServer(); pluginGrpcServer != nil {
|
||||
plugin_pb.RegisterPluginServiceServer(grpcServer, pluginGrpcServer)
|
||||
plugin_pb.RegisterAdminQueryServiceServer(grpcServer, pluginGrpcServer)
|
||||
plugin_pb.RegisterAdminCommandServiceServer(grpcServer, pluginGrpcServer)
|
||||
glog.Infof("Registered plugin services on worker gRPC server")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
s.grpcServer = grpcServer
|
||||
s.listener = listener
|
||||
s.running = true
|
||||
|
||||
@@ -8,6 +8,7 @@ import (
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/prometheus/client_golang/prometheus/promhttp"
|
||||
"github.com/seaweedfs/seaweedfs/weed/admin/dash"
|
||||
"github.com/seaweedfs/seaweedfs/weed/admin/plugin"
|
||||
"github.com/seaweedfs/seaweedfs/weed/admin/view/app"
|
||||
"github.com/seaweedfs/seaweedfs/weed/admin/view/layout"
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants"
|
||||
@@ -39,14 +40,14 @@ func NewAdminHandlers(adminServer *dash.AdminServer) *AdminHandlers {
|
||||
maintenanceHandlers := NewMaintenanceHandlers(adminServer)
|
||||
mqHandlers := NewMessageQueueHandlers(adminServer)
|
||||
serviceAccountHandlers := NewServiceAccountHandlers(adminServer)
|
||||
|
||||
|
||||
// Get plugin manager from admin server (may be nil)
|
||||
var pluginMgr interface{}
|
||||
if pm := adminServer.GetPluginManager(); pm != nil {
|
||||
pluginMgr = pm
|
||||
}
|
||||
pluginHandlers := NewPluginHandlers(adminServer, pluginMgr)
|
||||
|
||||
|
||||
return &AdminHandlers{
|
||||
adminServer: adminServer,
|
||||
authHandlers: authHandlers,
|
||||
@@ -270,13 +271,13 @@ func (h *AdminHandlers) SetupRoutes(r *gin.Engine, authRequired bool, adminUser,
|
||||
pluginApi := api.Group("/plugin")
|
||||
{
|
||||
pluginApi.GET("/list", h.pluginHandlers.ListPluginsAPI)
|
||||
pluginApi.GET("/jobs/:type", h.pluginHandlers.ListJobsAPI)
|
||||
pluginApi.GET("/jobs/by-type/:type", h.pluginHandlers.ListJobsAPI)
|
||||
pluginApi.GET("/config/:type", h.pluginHandlers.GetConfigAPI)
|
||||
pluginApi.POST("/config/:type/apply", dash.RequireWriteAccess(), h.pluginHandlers.SaveConfigAPI)
|
||||
pluginApi.GET("/detection/history/:type", h.pluginHandlers.GetDetectionHistoryAPI)
|
||||
pluginApi.GET("/execution/history/:type", h.pluginHandlers.GetExecutionHistoryAPI)
|
||||
pluginApi.POST("/jobs/:type/trigger-detection", dash.RequireWriteAccess(), h.pluginHandlers.TriggerDetectionAPI)
|
||||
pluginApi.POST("/jobs/:id/cancel", dash.RequireWriteAccess(), h.pluginHandlers.CancelJobAPI)
|
||||
pluginApi.POST("/trigger-detection/:type", dash.RequireWriteAccess(), h.pluginHandlers.TriggerDetectionAPI)
|
||||
pluginApi.POST("/cancel-job/:id", dash.RequireWriteAccess(), h.pluginHandlers.CancelJobAPI)
|
||||
}
|
||||
}
|
||||
} else {
|
||||
@@ -320,6 +321,11 @@ func (h *AdminHandlers) SetupRoutes(r *gin.Engine, authRequired bool, adminUser,
|
||||
r.GET("/mq/topics", h.mqHandlers.ShowTopics)
|
||||
r.GET("/mq/topics/:namespace/:topic", h.mqHandlers.ShowTopicDetails)
|
||||
|
||||
// Plugin management routes
|
||||
r.GET("/plugins", h.ShowPlugins)
|
||||
r.GET("/plugins/jobs/:jobType", h.ShowPluginJobs)
|
||||
r.GET("/plugins/config/:jobType", h.ShowPluginConfig)
|
||||
|
||||
// Maintenance system routes
|
||||
r.GET("/maintenance", h.maintenanceHandlers.ShowMaintenanceQueue)
|
||||
r.GET("/maintenance/workers", h.maintenanceHandlers.ShowMaintenanceWorkers)
|
||||
@@ -450,6 +456,19 @@ func (h *AdminHandlers) SetupRoutes(r *gin.Engine, authRequired bool, adminUser,
|
||||
mqApi.POST("/topics/retention/update", h.mqHandlers.UpdateTopicRetentionAPI)
|
||||
mqApi.POST("/retention/purge", h.adminServer.TriggerTopicRetentionPurgeAPI)
|
||||
}
|
||||
|
||||
// Plugin API routes
|
||||
pluginApi := api.Group("/plugin")
|
||||
{
|
||||
pluginApi.GET("/list", h.pluginHandlers.ListPluginsAPI)
|
||||
pluginApi.GET("/jobs/by-type/:type", h.pluginHandlers.ListJobsAPI)
|
||||
pluginApi.GET("/config/:type", h.pluginHandlers.GetConfigAPI)
|
||||
pluginApi.POST("/config/:type/apply", h.pluginHandlers.SaveConfigAPI)
|
||||
pluginApi.GET("/detection/history/:type", h.pluginHandlers.GetDetectionHistoryAPI)
|
||||
pluginApi.GET("/execution/history/:type", h.pluginHandlers.GetExecutionHistoryAPI)
|
||||
pluginApi.POST("/trigger-detection/:type", h.pluginHandlers.TriggerDetectionAPI)
|
||||
pluginApi.POST("/cancel-job/:id", h.pluginHandlers.CancelJobAPI)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -700,19 +719,52 @@ func (h *AdminHandlers) getAdminData(c *gin.Context) dash.AdminData {
|
||||
|
||||
// ShowPlugins displays the plugins overview page
|
||||
func (h *AdminHandlers) ShowPlugins(c *gin.Context) {
|
||||
plugins := []interface{}{}
|
||||
plugins := []map[string]interface{}{}
|
||||
jobTypes := make(map[string]interface{})
|
||||
|
||||
|
||||
// Get plugin manager from server
|
||||
if pm := h.adminServer.GetPluginManager(); pm != nil {
|
||||
// TODO: Get actual plugins from plugin manager
|
||||
// Cast to *plugin.Manager
|
||||
if pluginMgr, ok := pm.(*plugin.Manager); ok {
|
||||
// Get list of connected plugins
|
||||
connectedPlugins := pluginMgr.ListPlugins(false)
|
||||
for _, p := range connectedPlugins {
|
||||
plugins = append(plugins, map[string]interface{}{
|
||||
"id": p.ID,
|
||||
"name": p.Name,
|
||||
"version": p.Version,
|
||||
"status": p.Status,
|
||||
"capabilities": p.Capabilities,
|
||||
"activeJobs": p.ActiveJobs,
|
||||
"completedJobs": p.CompletedJobs,
|
||||
"failedJobs": p.FailedJobs,
|
||||
"connectedAt": p.ConnectedAt,
|
||||
"lastHeartbeat": p.LastHeartbeat,
|
||||
})
|
||||
|
||||
// Build job types map
|
||||
for _, cap := range p.Capabilities {
|
||||
if _, exists := jobTypes[cap]; !exists {
|
||||
jobTypes[cap] = map[string]interface{}{
|
||||
"type": cap,
|
||||
"description": cap,
|
||||
"pluginCount": 0,
|
||||
}
|
||||
}
|
||||
// Increment plugin count for this capability
|
||||
if capData, ok := jobTypes[cap].(map[string]interface{}); ok {
|
||||
capData["pluginCount"] = capData["pluginCount"].(int) + 1
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
component := app.PluginsOverview(app.PluginsPageData{
|
||||
Plugins: plugins,
|
||||
JobTypes: jobTypes,
|
||||
})
|
||||
|
||||
|
||||
htmlContent := layout.Layout(c, component)
|
||||
htmlContent.Render(c.Request.Context(), c.Writer)
|
||||
}
|
||||
@@ -722,13 +774,13 @@ func (h *AdminHandlers) ShowPluginJobs(c *gin.Context) {
|
||||
jobType := c.Param("jobType")
|
||||
jobs := []interface{}{}
|
||||
stateFilter := c.Query("state")
|
||||
|
||||
|
||||
component := app.PluginJobsMonitoring(app.PluginJobsPageData{
|
||||
JobType: jobType,
|
||||
Jobs: jobs,
|
||||
StateFilter: stateFilter,
|
||||
})
|
||||
|
||||
|
||||
htmlContent := layout.Layout(c, component)
|
||||
htmlContent.Render(c.Request.Context(), c.Writer)
|
||||
}
|
||||
@@ -740,15 +792,15 @@ func (h *AdminHandlers) ShowPluginConfig(c *gin.Context) {
|
||||
if activeTab == "" {
|
||||
activeTab = "config"
|
||||
}
|
||||
|
||||
|
||||
component := app.PluginConfiguration(app.PluginConfigPageData{
|
||||
JobType: jobType,
|
||||
Config: app.JobTypeConfig{},
|
||||
DetectionHistory: []interface{}{},
|
||||
ExecutionHistory: []interface{}{},
|
||||
ActiveTab: activeTab,
|
||||
JobType: jobType,
|
||||
Config: app.JobTypeConfig{},
|
||||
DetectionHistory: []interface{}{},
|
||||
ExecutionHistory: []interface{}{},
|
||||
ActiveTab: activeTab,
|
||||
})
|
||||
|
||||
|
||||
htmlContent := layout.Layout(c, component)
|
||||
htmlContent.Render(c.Request.Context(), c.Writer)
|
||||
}
|
||||
|
||||
@@ -11,12 +11,12 @@ import (
|
||||
|
||||
// GRPCServer implements the plugin service gRPC handlers
|
||||
type GRPCServer struct {
|
||||
mu sync.RWMutex
|
||||
registry *Registry
|
||||
queue *JobQueue
|
||||
dispatcher *Dispatcher
|
||||
configMgr *ConfigManager
|
||||
streamMu sync.RWMutex
|
||||
mu sync.RWMutex
|
||||
registry *Registry
|
||||
queue *JobQueue
|
||||
dispatcher *Dispatcher
|
||||
configMgr *ConfigManager
|
||||
streamMu sync.RWMutex
|
||||
activeStreams map[string][]chan interface{}
|
||||
plugin_pb.UnimplementedPluginServiceServer
|
||||
plugin_pb.UnimplementedAdminQueryServiceServer
|
||||
@@ -78,17 +78,17 @@ func (gs *GRPCServer) Connect(ctx context.Context, req *plugin_pb.PluginConnectR
|
||||
|
||||
// Build response
|
||||
pbConfig := &plugin_pb.PluginConfig{
|
||||
PluginId: config.PluginID,
|
||||
Properties: config.Properties,
|
||||
MaxRetries: int32(config.MaxRetries),
|
||||
Environment: config.Environment,
|
||||
PluginId: config.PluginID,
|
||||
Properties: config.Properties,
|
||||
MaxRetries: int32(config.MaxRetries),
|
||||
Environment: config.Environment,
|
||||
}
|
||||
|
||||
response := &plugin_pb.PluginConnectResponse{
|
||||
Success: true,
|
||||
Message: "Plugin registered successfully",
|
||||
MasterId: "master-1",
|
||||
Config: pbConfig,
|
||||
Success: true,
|
||||
Message: "Plugin registered successfully",
|
||||
MasterId: "master-1",
|
||||
Config: pbConfig,
|
||||
AssignedTypes: req.Capabilities,
|
||||
}
|
||||
|
||||
@@ -102,8 +102,8 @@ func (gs *GRPCServer) ExecuteJob(ctx context.Context, req *plugin_pb.ExecuteJobR
|
||||
}
|
||||
|
||||
response := &plugin_pb.ExecuteJobResponse{
|
||||
JobId: req.JobId,
|
||||
Status: plugin_pb.ExecutionStatus_EXECUTION_STATUS_ACCEPTED,
|
||||
JobId: req.JobId,
|
||||
Status: plugin_pb.ExecutionStatus_EXECUTION_STATUS_ACCEPTED,
|
||||
Message: "Job accepted for execution",
|
||||
}
|
||||
|
||||
@@ -148,10 +148,10 @@ func (gs *GRPCServer) GetConfig(ctx context.Context, req *plugin_pb.GetConfigReq
|
||||
}
|
||||
|
||||
pbConfig := &plugin_pb.PluginConfig{
|
||||
PluginId: config.PluginID,
|
||||
Properties: config.Properties,
|
||||
MaxRetries: int32(config.MaxRetries),
|
||||
Environment: config.Environment,
|
||||
PluginId: config.PluginID,
|
||||
Properties: config.Properties,
|
||||
MaxRetries: int32(config.MaxRetries),
|
||||
Environment: config.Environment,
|
||||
}
|
||||
|
||||
response := &plugin_pb.GetConfigResponse{
|
||||
@@ -169,7 +169,7 @@ func (gs *GRPCServer) SubmitResult(ctx context.Context, req *plugin_pb.JobResult
|
||||
}
|
||||
|
||||
actions := []string{}
|
||||
|
||||
|
||||
// Process results based on job status
|
||||
switch req.Status {
|
||||
case plugin_pb.ExecutionStatus_EXECUTION_STATUS_COMPLETED:
|
||||
@@ -203,16 +203,16 @@ func (gs *GRPCServer) GetPluginStats(ctx context.Context, req *plugin_pb.GetPlug
|
||||
|
||||
for _, plugin := range plugins {
|
||||
stat := &plugin_pb.PluginStats{
|
||||
PluginId: plugin.ID,
|
||||
Status: plugin.Status,
|
||||
ActiveJobs: int32(plugin.ActiveJobs),
|
||||
CompletedJobs: int32(plugin.CompletedJobs),
|
||||
FailedJobs: int32(plugin.FailedJobs),
|
||||
TotalDetections: plugin.TotalDetections,
|
||||
AvgExecutionTimeMs: float32(plugin.AvgExecutionTimeMs),
|
||||
CpuUsagePercent: float32(plugin.CPUUsagePercent),
|
||||
MemoryUsageBytes: plugin.MemoryUsageBytes,
|
||||
UptimeSeconds: int32(time.Since(plugin.ConnectedAt).Seconds()),
|
||||
PluginId: plugin.ID,
|
||||
Status: plugin.Status,
|
||||
ActiveJobs: int32(plugin.ActiveJobs),
|
||||
CompletedJobs: int32(plugin.CompletedJobs),
|
||||
FailedJobs: int32(plugin.FailedJobs),
|
||||
TotalDetections: plugin.TotalDetections,
|
||||
AvgExecutionTimeMs: float32(plugin.AvgExecutionTimeMs),
|
||||
CpuUsagePercent: float32(plugin.CPUUsagePercent),
|
||||
MemoryUsageBytes: plugin.MemoryUsageBytes,
|
||||
UptimeSeconds: int32(time.Since(plugin.ConnectedAt).Seconds()),
|
||||
}
|
||||
response.Stats = append(response.Stats, stat)
|
||||
}
|
||||
@@ -249,14 +249,14 @@ func (gs *GRPCServer) ListPlugins(ctx context.Context, req *plugin_pb.ListPlugin
|
||||
}
|
||||
|
||||
info := &plugin_pb.PluginInfo{
|
||||
PluginId: plugin.ID,
|
||||
Name: plugin.Name,
|
||||
Version: plugin.Version,
|
||||
Status: plugin.Status,
|
||||
Capabilities: plugin.Capabilities,
|
||||
PluginId: plugin.ID,
|
||||
Name: plugin.Name,
|
||||
Version: plugin.Version,
|
||||
Status: plugin.Status,
|
||||
Capabilities: plugin.Capabilities,
|
||||
MaxConcurrentJobs: int32(plugin.MaxConcurrentJobs),
|
||||
ActiveJobs: int32(plugin.ActiveJobs),
|
||||
Metadata: plugin.Metadata,
|
||||
ActiveJobs: int32(plugin.ActiveJobs),
|
||||
Metadata: plugin.Metadata,
|
||||
}
|
||||
response.Plugins = append(response.Plugins, info)
|
||||
}
|
||||
|
||||
@@ -1,8 +1,10 @@
|
||||
package app
|
||||
|
||||
import "fmt"
|
||||
|
||||
type PluginsPageData struct {
|
||||
Plugins []interface{}
|
||||
JobTypes map[string]interface{}
|
||||
Plugins []map[string]interface{}
|
||||
JobTypes map[string]interface{}
|
||||
}
|
||||
|
||||
templ PluginsOverview(data PluginsPageData) {
|
||||
@@ -26,7 +28,7 @@ templ PluginsOverview(data PluginsPageData) {
|
||||
<div class="row no-gutters align-items-center">
|
||||
<div class="col mr-2">
|
||||
<div class="text-xs font-weight-bold text-primary text-uppercase mb-1">Connected Plugins</div>
|
||||
<div class="h3 mb-0">0</div>
|
||||
<div class="h3 mb-0">{ fmt.Sprintf("%d", len(data.Plugins)) }</div>
|
||||
</div>
|
||||
<div class="col-auto">
|
||||
<i class="fas fa-plug fa-2x text-gray-300"></i>
|
||||
@@ -49,23 +51,35 @@ if len(data.Plugins) == 0 {
|
||||
<table class="table table-hover table-sm">
|
||||
<thead>
|
||||
<tr>
|
||||
<th>Plugin</th>
|
||||
<th>Plugin ID</th>
|
||||
<th>Name</th>
|
||||
<th>Status</th>
|
||||
<th>Version</th>
|
||||
<th>Capabilities</th>
|
||||
<th>Actions</th>
|
||||
</tr>
|
||||
</thead>
|
||||
<tbody>
|
||||
for i := 0; i < len(data.Plugins); i++ {
|
||||
for _, p := range data.Plugins {
|
||||
<tr>
|
||||
<td><code>{ p["id"].(string) }</code></td>
|
||||
<td>{ p["name"].(string) }</td>
|
||||
<td>
|
||||
<strong>Plugin</strong><br/>
|
||||
<small class="text-muted">v1.0</small>
|
||||
if p["status"].(string) == "CONNECTED" {
|
||||
<span class="badge bg-success">Connected</span>
|
||||
} else {
|
||||
<span class="badge bg-warning">{ p["status"].(string) }</span>
|
||||
}
|
||||
</td>
|
||||
<td>{ p["version"].(string) }</td>
|
||||
<td>
|
||||
for _, cap := range p["capabilities"].([]string) {
|
||||
<span class="badge bg-info me-1">{ cap }</span>
|
||||
}
|
||||
</td>
|
||||
<td>
|
||||
<span class="badge bg-success">Healthy</span>
|
||||
</td>
|
||||
<td>
|
||||
<a href="/plugins/config/test" class="btn btn-sm btn-primary">Config</a>
|
||||
<a href={ templ.SafeURL("/plugins/jobs/" + p["id"].(string)) } class="btn btn-sm btn-primary">Jobs</a>
|
||||
<a href={ templ.SafeURL("/plugins/config/" + p["id"].(string)) } class="btn btn-sm btn-secondary">Config</a>
|
||||
</td>
|
||||
</tr>
|
||||
}
|
||||
|
||||
@@ -8,8 +8,10 @@ package app
|
||||
import "github.com/a-h/templ"
|
||||
import templruntime "github.com/a-h/templ/runtime"
|
||||
|
||||
import "fmt"
|
||||
|
||||
type PluginsPageData struct {
|
||||
Plugins []interface{}
|
||||
Plugins []map[string]interface{}
|
||||
JobTypes map[string]interface{}
|
||||
}
|
||||
|
||||
@@ -34,32 +36,161 @@ func PluginsOverview(data PluginsPageData) templ.Component {
|
||||
templ_7745c5c3_Var1 = templ.NopComponent
|
||||
}
|
||||
ctx = templ.ClearChildren(ctx)
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 1, "<div class=\"d-flex justify-content-between flex-wrap flex-md-nowrap align-items-center pt-3 pb-2 mb-3 border-bottom\"><h1 class=\"h2\"><i class=\"fas fa-plug me-2\"></i>Plugins</h1><div class=\"btn-toolbar mb-2 mb-md-0\"><div class=\"btn-group me-2\"><a href=\"/plugins\" class=\"btn btn-sm btn-outline-primary\"><i class=\"fas fa-sync-alt me-1\"></i>Refresh</a></div></div></div><div class=\"row mb-4\"><div class=\"col-md-3 mb-4\"><div class=\"card border-left-primary shadow h-100 py-2\"><div class=\"card-body\"><div class=\"row no-gutters align-items-center\"><div class=\"col mr-2\"><div class=\"text-xs font-weight-bold text-primary text-uppercase mb-1\">Connected Plugins</div><div class=\"h3 mb-0\">0</div></div><div class=\"col-auto\"><i class=\"fas fa-plug fa-2x text-gray-300\"></i></div></div></div></div></div></div><div class=\"card shadow mb-4\"><div class=\"card-header py-3\"><h6 class=\"m-0 font-weight-bold text-primary\">Connected Plugins</h6></div><div class=\"card-body\">")
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 1, "<div class=\"d-flex justify-content-between flex-wrap flex-md-nowrap align-items-center pt-3 pb-2 mb-3 border-bottom\"><h1 class=\"h2\"><i class=\"fas fa-plug me-2\"></i>Plugins</h1><div class=\"btn-toolbar mb-2 mb-md-0\"><div class=\"btn-group me-2\"><a href=\"/plugins\" class=\"btn btn-sm btn-outline-primary\"><i class=\"fas fa-sync-alt me-1\"></i>Refresh</a></div></div></div><div class=\"row mb-4\"><div class=\"col-md-3 mb-4\"><div class=\"card border-left-primary shadow h-100 py-2\"><div class=\"card-body\"><div class=\"row no-gutters align-items-center\"><div class=\"col mr-2\"><div class=\"text-xs font-weight-bold text-primary text-uppercase mb-1\">Connected Plugins</div><div class=\"h3 mb-0\">")
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
var templ_7745c5c3_Var2 string
|
||||
templ_7745c5c3_Var2, templ_7745c5c3_Err = templ.JoinStringErrs(fmt.Sprintf("%d", len(data.Plugins)))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `weed/admin/view/app/plugins.templ`, Line: 31, Col: 59}
|
||||
}
|
||||
_, templ_7745c5c3_Err = templ_7745c5c3_Buffer.WriteString(templ.EscapeString(templ_7745c5c3_Var2))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 2, "</div></div><div class=\"col-auto\"><i class=\"fas fa-plug fa-2x text-gray-300\"></i></div></div></div></div></div></div><div class=\"card shadow mb-4\"><div class=\"card-header py-3\"><h6 class=\"m-0 font-weight-bold text-primary\">Connected Plugins</h6></div><div class=\"card-body\">")
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
if len(data.Plugins) == 0 {
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 2, "<div class=\"alert alert-info\">No plugins connected</div>")
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 3, "<div class=\"alert alert-info\">No plugins connected</div>")
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
} else {
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 3, "<div class=\"table-responsive\"><table class=\"table table-hover table-sm\"><thead><tr><th>Plugin</th><th>Status</th><th>Actions</th></tr></thead> <tbody>")
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 4, "<div class=\"table-responsive\"><table class=\"table table-hover table-sm\"><thead><tr><th>Plugin ID</th><th>Name</th><th>Status</th><th>Version</th><th>Capabilities</th><th>Actions</th></tr></thead> <tbody>")
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
for i := 0; i < len(data.Plugins); i++ {
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 4, "<tr><td><strong>Plugin</strong><br><small class=\"text-muted\">v1.0</small></td><td><span class=\"badge bg-success\">Healthy</span></td><td><a href=\"/plugins/config/test\" class=\"btn btn-sm btn-primary\">Config</a></td></tr>")
|
||||
for _, p := range data.Plugins {
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 5, "<tr><td><code>")
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
var templ_7745c5c3_Var3 string
|
||||
templ_7745c5c3_Var3, templ_7745c5c3_Err = templ.JoinStringErrs(p["id"].(string))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `weed/admin/view/app/plugins.templ`, Line: 65, Col: 28}
|
||||
}
|
||||
_, templ_7745c5c3_Err = templ_7745c5c3_Buffer.WriteString(templ.EscapeString(templ_7745c5c3_Var3))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 6, "</code></td><td>")
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
var templ_7745c5c3_Var4 string
|
||||
templ_7745c5c3_Var4, templ_7745c5c3_Err = templ.JoinStringErrs(p["name"].(string))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `weed/admin/view/app/plugins.templ`, Line: 66, Col: 24}
|
||||
}
|
||||
_, templ_7745c5c3_Err = templ_7745c5c3_Buffer.WriteString(templ.EscapeString(templ_7745c5c3_Var4))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 7, "</td><td>")
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
if p["status"].(string) == "CONNECTED" {
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 8, "<span class=\"badge bg-success\">Connected</span>")
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
} else {
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 9, "<span class=\"badge bg-warning\">")
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
var templ_7745c5c3_Var5 string
|
||||
templ_7745c5c3_Var5, templ_7745c5c3_Err = templ.JoinStringErrs(p["status"].(string))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `weed/admin/view/app/plugins.templ`, Line: 71, Col: 53}
|
||||
}
|
||||
_, templ_7745c5c3_Err = templ_7745c5c3_Buffer.WriteString(templ.EscapeString(templ_7745c5c3_Var5))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 10, "</span>")
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
}
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 11, "</td><td>")
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
var templ_7745c5c3_Var6 string
|
||||
templ_7745c5c3_Var6, templ_7745c5c3_Err = templ.JoinStringErrs(p["version"].(string))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `weed/admin/view/app/plugins.templ`, Line: 74, Col: 27}
|
||||
}
|
||||
_, templ_7745c5c3_Err = templ_7745c5c3_Buffer.WriteString(templ.EscapeString(templ_7745c5c3_Var6))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 12, "</td><td>")
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
for _, cap := range p["capabilities"].([]string) {
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 13, "<span class=\"badge bg-info me-1\">")
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
var templ_7745c5c3_Var7 string
|
||||
templ_7745c5c3_Var7, templ_7745c5c3_Err = templ.JoinStringErrs(cap)
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `weed/admin/view/app/plugins.templ`, Line: 77, Col: 38}
|
||||
}
|
||||
_, templ_7745c5c3_Err = templ_7745c5c3_Buffer.WriteString(templ.EscapeString(templ_7745c5c3_Var7))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 14, "</span>")
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
}
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 15, "</td><td><a href=\"")
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
var templ_7745c5c3_Var8 templ.SafeURL
|
||||
templ_7745c5c3_Var8, templ_7745c5c3_Err = templ.JoinURLErrs(templ.SafeURL("/plugins/jobs/" + p["id"].(string)))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `weed/admin/view/app/plugins.templ`, Line: 81, Col: 60}
|
||||
}
|
||||
_, templ_7745c5c3_Err = templ_7745c5c3_Buffer.WriteString(templ.EscapeString(templ_7745c5c3_Var8))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 16, "\" class=\"btn btn-sm btn-primary\">Jobs</a> <a href=\"")
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
var templ_7745c5c3_Var9 templ.SafeURL
|
||||
templ_7745c5c3_Var9, templ_7745c5c3_Err = templ.JoinURLErrs(templ.SafeURL("/plugins/config/" + p["id"].(string)))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ.Error{Err: templ_7745c5c3_Err, FileName: `weed/admin/view/app/plugins.templ`, Line: 82, Col: 62}
|
||||
}
|
||||
_, templ_7745c5c3_Err = templ_7745c5c3_Buffer.WriteString(templ.EscapeString(templ_7745c5c3_Var9))
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 17, "\" class=\"btn btn-sm btn-secondary\">Config</a></td></tr>")
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
}
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 5, "</tbody></table></div>")
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 18, "</tbody></table></div>")
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
}
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 6, "</div></div>")
|
||||
templ_7745c5c3_Err = templruntime.WriteString(templ_7745c5c3_Buffer, 19, "</div></div>")
|
||||
if templ_7745c5c3_Err != nil {
|
||||
return templ_7745c5c3_Err
|
||||
}
|
||||
|
||||
@@ -38,6 +38,7 @@ var Commands = []*Command{
|
||||
cmdMqBroker,
|
||||
cmdMqKafkaGateway,
|
||||
cmdDB,
|
||||
cmdPluginWorker,
|
||||
cmdS3,
|
||||
cmdScaffold,
|
||||
cmdServer,
|
||||
|
||||
@@ -0,0 +1,313 @@
|
||||
package command
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"os/signal"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"sync"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/plugin_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/security"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util/grace"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util/version"
|
||||
"google.golang.org/grpc"
|
||||
)
|
||||
|
||||
var cmdPluginWorker = &Command{
|
||||
UsageLine: "plugin_worker -admin=<grpc_address> [-plugins=<plugin_types>] [-workingDir=<path>]",
|
||||
Short: "start a plugin-based worker using the new plugin system",
|
||||
Long: `Start a worker using the new plugin system. This worker connects to the admin server
|
||||
via gRPC and registers capabilities for handling maintenance tasks.
|
||||
|
||||
Supported plugins: erasure_coding, vacuum, balance
|
||||
|
||||
Examples:
|
||||
weed plugin_worker -admin=localhost:23646
|
||||
weed plugin_worker -admin=admin.example.com:23646
|
||||
weed plugin_worker -admin=localhost:23646 -plugins=erasure_coding,vacuum
|
||||
weed plugin_worker -admin=localhost:23646 -workingDir=/tmp/worker
|
||||
weed plugin_worker -admin=localhost:23646 -debug
|
||||
`,
|
||||
}
|
||||
|
||||
var (
|
||||
pluginWorkerAdminServer = cmdPluginWorker.Flag.String("admin", "localhost:23646",
|
||||
"admin server gRPC address (usually admin HTTP port + 10000)")
|
||||
pluginWorkerPlugins = cmdPluginWorker.Flag.String("plugins", "erasure_coding,vacuum,balance",
|
||||
"comma-separated list of plugin types to enable")
|
||||
pluginWorkerWorkingDir = cmdPluginWorker.Flag.String("workingDir", "",
|
||||
"working directory for the plugin worker")
|
||||
pluginWorkerMaxConcurrent = cmdPluginWorker.Flag.Int("maxConcurrent", 2,
|
||||
"maximum number of concurrent jobs")
|
||||
pluginWorkerDebug = cmdPluginWorker.Flag.Bool("debug", false,
|
||||
"enable debug logging")
|
||||
pluginWorkerDebugPort = cmdPluginWorker.Flag.Int("debug.port", 6061,
|
||||
"http port for debugging")
|
||||
pluginWorkerTimeout = cmdPluginWorker.Flag.Duration("timeout", 30*time.Second,
|
||||
"gRPC connection timeout")
|
||||
)
|
||||
|
||||
func init() {
|
||||
cmdPluginWorker.Run = runPluginWorker
|
||||
}
|
||||
|
||||
// GenericPluginWorker is a multi-plugin worker that connects to admin server
|
||||
type GenericPluginWorker struct {
|
||||
ID string
|
||||
AdminServer pb.ServerAddress
|
||||
Plugins []string
|
||||
MaxConcurrentJobs int
|
||||
WorkingDir string
|
||||
PluginServiceClient plugin_pb.PluginServiceClient
|
||||
Conn *grpc.ClientConn
|
||||
Context context.Context
|
||||
Cancel context.CancelFunc
|
||||
mu sync.RWMutex
|
||||
isRunning bool
|
||||
}
|
||||
|
||||
// NewGenericPluginWorker creates a new generic plugin worker
|
||||
func NewGenericPluginWorker(adminServer pb.ServerAddress, plugins []string, workingDir string, maxConcurrent int) *GenericPluginWorker {
|
||||
workerID := fmt.Sprintf("worker-%s-%d", hostname(), time.Now().UnixNano())
|
||||
return &GenericPluginWorker{
|
||||
ID: string(workerID),
|
||||
AdminServer: adminServer,
|
||||
Plugins: plugins,
|
||||
MaxConcurrentJobs: maxConcurrent,
|
||||
WorkingDir: workingDir,
|
||||
}
|
||||
}
|
||||
|
||||
// Start connects to admin server and registers plugins
|
||||
func (w *GenericPluginWorker) Start(ctx context.Context) error {
|
||||
w.mu.Lock()
|
||||
defer w.mu.Unlock()
|
||||
|
||||
if w.isRunning {
|
||||
return fmt.Errorf("worker is already running")
|
||||
}
|
||||
|
||||
w.Context, w.Cancel = context.WithCancel(ctx)
|
||||
|
||||
// Create gRPC connection
|
||||
dialCtx, cancel := context.WithTimeout(w.Context, *pluginWorkerTimeout)
|
||||
defer cancel()
|
||||
|
||||
grpcDialOption := security.LoadClientTLS(util.GetViper(), "grpc.worker")
|
||||
|
||||
conn, err := grpc.DialContext(dialCtx, w.AdminServer.ToGrpcAddress(), grpcDialOption)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to connect to admin grpc server at %s: %v", w.AdminServer.ToGrpcAddress(), err)
|
||||
}
|
||||
|
||||
w.Conn = conn
|
||||
w.PluginServiceClient = plugin_pb.NewPluginServiceClient(conn)
|
||||
|
||||
glog.Infof("Successfully connected to admin grpc server at %s", w.AdminServer)
|
||||
|
||||
// Register plugin with admin server
|
||||
capabilities := []string{}
|
||||
for _, plugin := range w.Plugins {
|
||||
plugin = strings.TrimSpace(plugin)
|
||||
if plugin != "" {
|
||||
capabilities = append(capabilities, plugin)
|
||||
}
|
||||
}
|
||||
|
||||
connectReq := &plugin_pb.PluginConnectRequest{
|
||||
PluginId: w.ID,
|
||||
PluginName: fmt.Sprintf("GenericWorker-%s", strings.Join(w.Plugins, ",")),
|
||||
Version: version.VERSION,
|
||||
Capabilities: capabilities,
|
||||
MaxConcurrentJobs: int32(w.MaxConcurrentJobs),
|
||||
SupportsStreaming: false,
|
||||
Port: 0,
|
||||
Metadata: map[string]string{
|
||||
"working_dir": w.WorkingDir,
|
||||
},
|
||||
}
|
||||
|
||||
connectCtx, cancel := context.WithTimeout(w.Context, 10*time.Second)
|
||||
defer cancel()
|
||||
|
||||
connectResp, err := w.PluginServiceClient.Connect(connectCtx, connectReq)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to register with admin server: %v", err)
|
||||
}
|
||||
|
||||
if !connectResp.Success {
|
||||
return fmt.Errorf("plugin registration failed: %s", connectResp.Message)
|
||||
}
|
||||
|
||||
glog.Infof("Plugin registered successfully. Assigned types: %v", connectResp.AssignedTypes)
|
||||
|
||||
w.isRunning = true
|
||||
|
||||
// Start background health reporting
|
||||
go w.healthReportLoop()
|
||||
|
||||
glog.Infof("Plugin worker started successfully with ID: %s", w.ID)
|
||||
return nil
|
||||
}
|
||||
|
||||
// Stop gracefully stops the worker
|
||||
func (w *GenericPluginWorker) Stop() error {
|
||||
w.mu.Lock()
|
||||
defer w.mu.Unlock()
|
||||
|
||||
if !w.isRunning {
|
||||
return nil
|
||||
}
|
||||
|
||||
w.isRunning = false
|
||||
if w.Cancel != nil {
|
||||
w.Cancel()
|
||||
}
|
||||
|
||||
if w.Conn != nil {
|
||||
return w.Conn.Close()
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// healthReportLoop sends periodic health reports to admin server
|
||||
func (w *GenericPluginWorker) healthReportLoop() {
|
||||
ticker := time.NewTicker(30 * time.Second)
|
||||
defer ticker.Stop()
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-w.Context.Done():
|
||||
return
|
||||
case <-ticker.C:
|
||||
w.sendHealthReport()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// sendHealthReport sends a health report to the admin server
|
||||
func (w *GenericPluginWorker) sendHealthReport() {
|
||||
healthReport := &plugin_pb.HealthReport{
|
||||
PluginId: w.ID,
|
||||
TimestampMs: time.Now().UnixMilli(),
|
||||
Status: plugin_pb.HealthStatus_HEALTH_STATUS_HEALTHY,
|
||||
ActiveJobs: 0,
|
||||
CpuPercent: 0,
|
||||
MemoryBytes: 0,
|
||||
JobProgress: []*plugin_pb.JobProgress{},
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||
defer cancel()
|
||||
|
||||
_, err := w.PluginServiceClient.ReportHealth(ctx, healthReport)
|
||||
if err != nil {
|
||||
glog.Warningf("Failed to send health report: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// IsRunning checks if the worker is currently running
|
||||
func (w *GenericPluginWorker) IsRunning() bool {
|
||||
w.mu.RLock()
|
||||
defer w.mu.RUnlock()
|
||||
return w.isRunning
|
||||
}
|
||||
|
||||
func hostname() string {
|
||||
hostname, err := os.Hostname()
|
||||
if err != nil {
|
||||
hostname = "unknown"
|
||||
}
|
||||
return hostname
|
||||
}
|
||||
|
||||
func runPluginWorker(cmd *Command, args []string) bool {
|
||||
if *pluginWorkerDebug {
|
||||
grace.StartDebugServer(*pluginWorkerDebugPort)
|
||||
}
|
||||
|
||||
util.LoadConfiguration("security", false)
|
||||
|
||||
glog.Infof("Starting plugin worker (v%s)", version.VERSION)
|
||||
glog.Infof("Admin server: %s", *pluginWorkerAdminServer)
|
||||
glog.Infof("Plugins: %s", *pluginWorkerPlugins)
|
||||
|
||||
// Parse plugins
|
||||
plugins := strings.Split(*pluginWorkerPlugins, ",")
|
||||
validPlugins := []string{}
|
||||
for _, p := range plugins {
|
||||
p = strings.TrimSpace(p)
|
||||
if p != "" {
|
||||
validPlugins = append(validPlugins, p)
|
||||
}
|
||||
}
|
||||
|
||||
if len(validPlugins) == 0 {
|
||||
glog.Fatalf("No valid plugins specified. Valid options: erasure_coding, vacuum, balance")
|
||||
return false
|
||||
}
|
||||
|
||||
// Set up working directory
|
||||
workingDir := *pluginWorkerWorkingDir
|
||||
if workingDir == "" {
|
||||
var err error
|
||||
workingDir, err = os.Getwd()
|
||||
if err != nil {
|
||||
glog.Fatalf("Failed to get current working directory: %v", err)
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
// Create and validate working directory
|
||||
if err := os.MkdirAll(workingDir, 0755); err != nil {
|
||||
glog.Fatalf("Failed to create working directory: %v", err)
|
||||
return false
|
||||
}
|
||||
|
||||
glog.Infof("Working directory: %s", workingDir)
|
||||
|
||||
// Create plugin-specific subdirectories
|
||||
for _, plugin := range validPlugins {
|
||||
pluginDir := filepath.Join(workingDir, plugin)
|
||||
if err := os.MkdirAll(pluginDir, 0755); err != nil {
|
||||
glog.Fatalf("Failed to create plugin directory %s: %v", pluginDir, err)
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
// Create and start the worker
|
||||
adminServerAddress := pb.ServerAddress(*pluginWorkerAdminServer)
|
||||
worker := NewGenericPluginWorker(adminServerAddress, validPlugins, workingDir, *pluginWorkerMaxConcurrent)
|
||||
|
||||
ctx := context.Background()
|
||||
if err := worker.Start(ctx); err != nil {
|
||||
glog.Fatalf("Failed to start plugin worker: %v", err)
|
||||
return false
|
||||
}
|
||||
|
||||
// Set up signal handling
|
||||
sigChan := make(chan os.Signal, 1)
|
||||
signal.Notify(sigChan, syscall.SIGINT, syscall.SIGTERM)
|
||||
|
||||
glog.Infof("Plugin worker is running. Press Ctrl+C to stop")
|
||||
|
||||
// Wait for shutdown signal
|
||||
<-sigChan
|
||||
glog.Infof("Shutdown signal received, stopping plugin worker...")
|
||||
|
||||
if err := worker.Stop(); err != nil {
|
||||
glog.Errorf("Error stopping plugin worker: %v", err)
|
||||
}
|
||||
|
||||
glog.Infof("Plugin worker stopped")
|
||||
return true
|
||||
}
|
||||
Reference in new issue
Block a user