* admin: add plugin runtime UI page and route wiring * pb: add plugin gRPC contract and generated bindings * admin/plugin: implement worker registry, runtime, monitoring, and config store * admin/dash: wire plugin runtime and expose plugin workflow APIs * command: add flags to enable plugin runtime * admin: rename remaining plugin v2 wording to plugin * admin/plugin: add detectable job type registry helper * admin/plugin: add scheduled detection and dispatch orchestration * admin/plugin: prefetch job type descriptors when workers connect * admin/plugin: add known job type discovery API and UI * admin/plugin: refresh design doc to match current implementation * admin/plugin: enforce per-worker scheduler concurrency limits * admin/plugin: use descriptor runtime defaults for scheduler policy * admin/ui: auto-load first known plugin job type on page open * admin/plugin: bootstrap persisted config from descriptor defaults * admin/plugin: dedupe scheduled proposals by dedupe key * admin/ui: add job type and state filters for plugin monitoring * admin/ui: add per-job-type plugin activity summary * admin/plugin: split descriptor read API from schema refresh * admin/ui: keep plugin summary metrics global while tables are filtered * admin/plugin: retry executor reservation before timing out * admin/plugin: expose scheduler states for monitoring * admin/ui: show per-job-type scheduler states in plugin monitor * pb/plugin: rename protobuf package to plugin * admin/plugin: rename pluginRuntime wiring to plugin * admin/plugin: remove runtime naming from plugin APIs and UI * admin/plugin: rename runtime files to plugin naming * admin/plugin: persist jobs and activities for monitor recovery * admin/plugin: lease one detector worker per job type * admin/ui: show worker load from plugin heartbeats * admin/plugin: skip stale workers for detector and executor picks * plugin/worker: add plugin worker command and stream runtime scaffold * plugin/worker: implement vacuum detect and execute handlers * admin/plugin: document external vacuum plugin worker starter * command: update plugin.worker help to reflect implemented flow * command/admin: drop legacy Plugin V2 label * plugin/worker: validate vacuum job type and respect min interval * plugin/worker: test no-op detect when min interval not elapsed * command/admin: document plugin.worker external process * plugin/worker: advertise configured concurrency in hello * command/plugin.worker: add jobType handler selection * command/plugin.worker: test handler selection by job type * command/plugin.worker: persist worker id in workingDir * admin/plugin: document plugin.worker jobType and workingDir flags * plugin/worker: support cancel request for in-flight work * plugin/worker: test cancel request acknowledgements * command/plugin.worker: document workingDir and jobType behavior * plugin/worker: emit executor activity events for monitor * plugin/worker: test executor activity builder * admin/plugin: send last successful run in detection request * admin/plugin: send cancel request when detect or execute context ends * admin/plugin: document worker cancel request responsibility * admin/handlers: expose plugin scheduler states API in no-auth mode * admin/handlers: test plugin scheduler states route registration * admin/plugin: keep worker id on worker-generated activity records * admin/plugin: test worker id propagation in monitor activities * admin/dash: always initialize plugin service * command/admin: remove plugin enable flags and default to enabled * admin/dash: drop pluginEnabled constructor parameter * admin/plugin UI: stop checking plugin enabled state * admin/plugin: remove docs for plugin enable flags * admin/dash: remove unused plugin enabled check method * admin/dash: fallback to in-memory plugin init when dataDir fails * admin/plugin API: expose worker gRPC port in status * command/plugin.worker: resolve admin gRPC port via plugin status * split plugin UI into overview/configuration/monitoring pages * Update layout_templ.go * add volume_balance plugin worker handler * wire plugin.worker CLI for volume_balance job type * add erasure_coding plugin worker handler * wire plugin.worker CLI for erasure_coding job type * support multi-job handlers in plugin worker runtime * allow plugin.worker jobType as comma-separated list * admin/plugin UI: rename to Workers and simplify config view * plugin worker: queue detection requests instead of capacity reject * Update plugin_worker.go * plugin volume_balance: remove force_move/timeout from worker config UI * plugin erasure_coding: enforce local working dir and cleanup * admin/plugin UI: rename admin settings to job scheduling * admin/plugin UI: persist and robustly render detection results * admin/plugin: record and return detection trace metadata * admin/plugin UI: show detection process and decision trace * plugin: surface detector decision trace as activities * mini: start a plugin worker by default * admin/plugin UI: split monitoring into detection and execution tabs * plugin worker: emit detection decision trace for EC and balance * admin workers UI: split monitoring into detection and execution pages * plugin scheduler: skip proposals for active assigned/running jobs * admin workers UI: add job queue tab * plugin worker: add dummy stress detector and executor job type * admin workers UI: reorder tabs to detection queue execution * admin workers UI: regenerate plugin template * plugin defaults: include dummy stress and add stress tests * plugin dummy stress: rotate detection selections across runs * plugin scheduler: remove cross-run proposal dedupe * plugin queue: track pending scheduled jobs * plugin scheduler: wait for executor capacity before dispatch * plugin scheduler: skip detection when waiting backlog is high * plugin: add disk-backed job detail API and persistence * admin ui: show plugin job detail modal from job id links * plugin: generate unique job ids instead of reusing proposal ids * plugin worker: emit heartbeats on work state changes * plugin registry: round-robin tied executor and detector picks * add temporary EC overnight stress runner * plugin job details: persist and render EC execution plans * ec volume details: color data and parity shard badges * shard labels: keep parity ids numeric and color-only distinction * admin: remove legacy maintenance UI routes and templates * admin: remove dead maintenance endpoint helpers * Update layout_templ.go * remove dummy_stress worker and command support * refactor plugin UI to job-type top tabs and sub-tabs * migrate weed worker command to plugin runtime * remove plugin.worker command and keep worker runtime with metrics * update helm worker args for jobType and execution flags * set plugin scheduling defaults to global 16 and per-worker 4 * stress: fix RPC context reuse and remove redundant variables in ec_stress_runner * admin/plugin: fix lifecycle races, safe channel operations, and terminal state constants * admin/dash: randomize job IDs and fix priority zero-value overwrite in plugin API * admin/handlers: implement buffered rendering to prevent response corruption * admin/plugin: implement debounced persistence flusher and optimize BuildJobDetail memory lookups * admin/plugin: fix priority overwrite and implement bounded wait in scheduler reserve * admin/plugin: implement atomic file writes and fix run record side effects * admin/plugin: use P prefix for parity shard labels in execution plans * admin/plugin: enable parallel execution for cancellation tests * admin: refactor time.Time fields to pointers for better JSON omitempty support * admin/plugin: implement pointer-safe time assignments and comparisons in plugin core * admin/plugin: fix time assignment and sorting logic in plugin monitor after pointer refactor * admin/plugin: update scheduler activity tracking to use time pointers * admin/plugin: fix time-based run history trimming after pointer refactor * admin/dash: fix JobSpec struct literal in plugin API after pointer refactor * admin/view: add D/P prefixes to EC shard badges for UI consistency * admin/plugin: use lifecycle-aware context for schema prefetching * Update ec_volume_details_templ.go * admin/stress: fix proposal sorting and log volume cleanup errors * stress: refine ec stress runner with math/rand and collection name - Added Collection field to VolumeEcShardsDeleteRequest for correct filename construction. - Replaced crypto/rand with seeded math/rand PRNG for bulk payloads. - Added documentation for EcMinAge zero-value behavior. - Added logging for ignored errors in volume/shard deletion. * admin: return internal server error for plugin store failures Changed error status code from 400 Bad Request to 500 Internal Server Error for failures in GetPluginJobDetail to correctly reflect server-side errors. * admin: implement safe channel sends and graceful shutdown sync - Added sync.WaitGroup to Plugin struct to manage background goroutines. - Implemented safeSendCh helper using recover() to prevent panics on closed channels. - Ensured Shutdown() waits for all background operations to complete. * admin: robustify plugin monitor with nil-safe time and record init - Standardized nil-safe assignment for *time.Time pointers (CreatedAt, UpdatedAt, CompletedAt). - Ensured persistJobDetailSnapshot initializes new records correctly if they don't exist on disk. - Fixed debounced persistence to trigger immediate write on job completion. * admin: improve scheduler shutdown behavior and logic guards - Replaced brittle error string matching with explicit r.shutdownCh selection for shutdown detection. - Removed redundant nil guard in buildScheduledJobSpec. - Standardized WaitGroup usage for schedulerLoop. * admin: implement deep copy for job parameters and atomic write fixes - Implemented deepCopyGenericValue and used it in cloneTrackedJob to prevent shared state. - Ensured atomicWriteFile creates parent directories before writing. * admin: remove unreachable branch in shard classification Removed an unreachable 'totalShards <= 0' check in classifyShardID as dataShards and parityShards are already guarded. * admin: secure UI links and use canonical shard constants - Added rel="noopener noreferrer" to external links for security. - Replaced magic number 14 with erasure_coding.TotalShardsCount. - Used renderEcShardBadge for missing shard list consistency. * admin: stabilize plugin tests and fix regressions - Composed a robust plugin_monitor_test.go to handle asynchronous persistence. - Updated all time.Time literals to use timeToPtr helper. - Added explicit Shutdown() calls in tests to synchronize with debounced writes. - Fixed syntax errors and orphaned struct literals in tests. * Potential fix for code scanning alert no. 278: Slice memory allocation with excessive size value Co-authored-by: Copilot Autofix powered by AI <62310815+github-advanced-security[bot]@users.noreply.github.com> * Potential fix for code scanning alert no. 283: Uncontrolled data used in path expression Co-authored-by: Copilot Autofix powered by AI <62310815+github-advanced-security[bot]@users.noreply.github.com> * admin: finalize refinements for error handling, scheduler, and race fixes - Standardized HTTP 500 status codes for store failures in plugin_api.go. - Tracked scheduled detection goroutines with sync.WaitGroup for safe shutdown. - Fixed race condition in safeSendDetectionComplete by extracting channel under lock. - Implemented deep copy for JobActivity details. - Used defaultDirPerm constant in atomicWriteFile. * test(ec): migrate admin dockertest to plugin APIs * admin/plugin_api: fix RunPluginJobTypeAPI to return 500 for server-side detection/filter errors * admin/plugin_api: fix ExecutePluginJobAPI to return 500 for job execution failures * admin/plugin_api: limit parseProtoJSONBody request body to 1MB to prevent unbounded memory usage * admin/plugin: consolidate regex to package-level validJobTypePattern; add char validation to sanitizeJobID * admin/plugin: fix racy Shutdown channel close with sync.Once * admin/plugin: track sendLoop and recv goroutines in WorkerStream with r.wg * admin/plugin: document writeProtoFiles atomicity — .pb is source of truth, .json is human-readable only * admin/plugin: extract activityLess helper to deduplicate nil-safe OccurredAt sort comparators * test/ec: check http.NewRequest errors to prevent nil req panics * test/ec: replace deprecated ioutil/math/rand, fix stale step comment 5.1→3.1 * plugin(ec): raise default detection and scheduling throughput limits * topology: include empty disks in volume list and EC capacity fallback * topology: remove hard 10-task cap for detection planning * Update ec_volume_details_templ.go * adjust default * fix tests --------- Co-authored-by: Copilot Autofix powered by AI <62310815+github-advanced-security[bot]@users.noreply.github.com>
349 lines
9.9 KiB
Go
349 lines
9.9 KiB
Go
package command
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"net/http"
|
|
"os"
|
|
"os/signal"
|
|
"path/filepath"
|
|
"strconv"
|
|
"strings"
|
|
"syscall"
|
|
"time"
|
|
|
|
"github.com/prometheus/client_golang/prometheus/promhttp"
|
|
"github.com/seaweedfs/seaweedfs/weed/glog"
|
|
pluginworker "github.com/seaweedfs/seaweedfs/weed/plugin/worker"
|
|
"github.com/seaweedfs/seaweedfs/weed/security"
|
|
statsCollect "github.com/seaweedfs/seaweedfs/weed/stats"
|
|
"github.com/seaweedfs/seaweedfs/weed/util"
|
|
"github.com/seaweedfs/seaweedfs/weed/util/version"
|
|
"google.golang.org/grpc"
|
|
)
|
|
|
|
const defaultPluginWorkerJobTypes = "vacuum,volume_balance,erasure_coding"
|
|
|
|
type pluginWorkerRunOptions struct {
|
|
AdminServer string
|
|
WorkerID string
|
|
WorkingDir string
|
|
JobTypes string
|
|
Heartbeat time.Duration
|
|
Reconnect time.Duration
|
|
MaxDetect int
|
|
MaxExecute int
|
|
Address string
|
|
MetricsPort int
|
|
MetricsIP string
|
|
}
|
|
|
|
func runPluginWorkerWithOptions(options pluginWorkerRunOptions) bool {
|
|
util.LoadConfiguration("security", false)
|
|
|
|
options.AdminServer = strings.TrimSpace(options.AdminServer)
|
|
if options.AdminServer == "" {
|
|
options.AdminServer = "localhost:23646"
|
|
}
|
|
|
|
options.JobTypes = strings.TrimSpace(options.JobTypes)
|
|
if options.JobTypes == "" {
|
|
options.JobTypes = defaultPluginWorkerJobTypes
|
|
}
|
|
|
|
if options.Heartbeat <= 0 {
|
|
options.Heartbeat = 15 * time.Second
|
|
}
|
|
if options.Reconnect <= 0 {
|
|
options.Reconnect = 5 * time.Second
|
|
}
|
|
if options.MaxDetect <= 0 {
|
|
options.MaxDetect = 1
|
|
}
|
|
if options.MaxExecute <= 0 {
|
|
options.MaxExecute = 4
|
|
}
|
|
options.MetricsIP = strings.TrimSpace(options.MetricsIP)
|
|
if options.MetricsIP == "" {
|
|
options.MetricsIP = "0.0.0.0"
|
|
}
|
|
|
|
resolvedAdminServer := resolvePluginWorkerAdminServer(options.AdminServer)
|
|
if resolvedAdminServer != options.AdminServer {
|
|
fmt.Printf("Resolved admin worker gRPC endpoint: %s -> %s\n", options.AdminServer, resolvedAdminServer)
|
|
}
|
|
|
|
dialOption := security.LoadClientTLS(util.GetViper(), "grpc.worker")
|
|
workerID, err := resolvePluginWorkerID(options.WorkerID, options.WorkingDir)
|
|
if err != nil {
|
|
glog.Errorf("Failed to resolve plugin worker ID: %v", err)
|
|
return false
|
|
}
|
|
|
|
handlers, err := buildPluginWorkerHandlers(options.JobTypes, dialOption)
|
|
if err != nil {
|
|
glog.Errorf("Failed to build plugin worker handlers: %v", err)
|
|
return false
|
|
}
|
|
worker, err := pluginworker.NewWorker(pluginworker.WorkerOptions{
|
|
AdminServer: resolvedAdminServer,
|
|
WorkerID: workerID,
|
|
WorkerVersion: version.Version(),
|
|
WorkerAddress: options.Address,
|
|
HeartbeatInterval: options.Heartbeat,
|
|
ReconnectDelay: options.Reconnect,
|
|
MaxDetectionConcurrency: options.MaxDetect,
|
|
MaxExecutionConcurrency: options.MaxExecute,
|
|
GrpcDialOption: dialOption,
|
|
Handlers: handlers,
|
|
})
|
|
if err != nil {
|
|
glog.Errorf("Failed to create plugin worker: %v", err)
|
|
return false
|
|
}
|
|
|
|
if options.MetricsPort > 0 {
|
|
go startPluginWorkerMetricsServer(options.MetricsIP, options.MetricsPort, worker)
|
|
}
|
|
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
defer cancel()
|
|
|
|
sigCh := make(chan os.Signal, 1)
|
|
signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM)
|
|
defer signal.Stop(sigCh)
|
|
|
|
go func() {
|
|
sig := <-sigCh
|
|
fmt.Printf("\nReceived signal %v, stopping plugin worker...\n", sig)
|
|
cancel()
|
|
}()
|
|
|
|
fmt.Printf("Starting plugin worker (admin=%s)\n", resolvedAdminServer)
|
|
if err := worker.Run(ctx); err != nil {
|
|
glog.Errorf("Plugin worker stopped with error: %v", err)
|
|
return false
|
|
}
|
|
fmt.Println("Plugin worker stopped")
|
|
return true
|
|
}
|
|
|
|
func resolvePluginWorkerID(explicitID string, workingDir string) (string, error) {
|
|
id := strings.TrimSpace(explicitID)
|
|
if id != "" {
|
|
return id, nil
|
|
}
|
|
|
|
workingDir = strings.TrimSpace(workingDir)
|
|
if workingDir == "" {
|
|
return "", nil
|
|
}
|
|
if err := os.MkdirAll(workingDir, 0755); err != nil {
|
|
return "", err
|
|
}
|
|
|
|
workerIDPath := filepath.Join(workingDir, "plugin.worker.id")
|
|
if data, err := os.ReadFile(workerIDPath); err == nil {
|
|
if persisted := strings.TrimSpace(string(data)); persisted != "" {
|
|
return persisted, nil
|
|
}
|
|
}
|
|
|
|
generated := fmt.Sprintf("plugin-%d", time.Now().UnixNano())
|
|
if err := os.WriteFile(workerIDPath, []byte(generated+"\n"), 0644); err != nil {
|
|
return "", err
|
|
}
|
|
return generated, nil
|
|
}
|
|
|
|
func buildPluginWorkerHandler(jobType string, dialOption grpc.DialOption) (pluginworker.JobHandler, error) {
|
|
canonicalJobType, err := canonicalPluginWorkerJobType(jobType)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
switch canonicalJobType {
|
|
case "vacuum":
|
|
return pluginworker.NewVacuumHandler(dialOption), nil
|
|
case "volume_balance":
|
|
return pluginworker.NewVolumeBalanceHandler(dialOption), nil
|
|
case "erasure_coding":
|
|
return pluginworker.NewErasureCodingHandler(dialOption), nil
|
|
default:
|
|
return nil, fmt.Errorf("unsupported plugin job type %q", canonicalJobType)
|
|
}
|
|
}
|
|
|
|
func buildPluginWorkerHandlers(jobTypes string, dialOption grpc.DialOption) ([]pluginworker.JobHandler, error) {
|
|
parsedJobTypes, err := parsePluginWorkerJobTypes(jobTypes)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
handlers := make([]pluginworker.JobHandler, 0, len(parsedJobTypes))
|
|
for _, jobType := range parsedJobTypes {
|
|
handler, buildErr := buildPluginWorkerHandler(jobType, dialOption)
|
|
if buildErr != nil {
|
|
return nil, buildErr
|
|
}
|
|
handlers = append(handlers, handler)
|
|
}
|
|
return handlers, nil
|
|
}
|
|
|
|
func parsePluginWorkerJobTypes(jobTypes string) ([]string, error) {
|
|
jobTypes = strings.TrimSpace(jobTypes)
|
|
if jobTypes == "" {
|
|
return []string{"vacuum"}, nil
|
|
}
|
|
|
|
parts := strings.Split(jobTypes, ",")
|
|
parsed := make([]string, 0, len(parts))
|
|
seen := make(map[string]struct{}, len(parts))
|
|
|
|
for _, part := range parts {
|
|
part = strings.TrimSpace(part)
|
|
if part == "" {
|
|
continue
|
|
}
|
|
canonical, err := canonicalPluginWorkerJobType(part)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if _, found := seen[canonical]; found {
|
|
continue
|
|
}
|
|
seen[canonical] = struct{}{}
|
|
parsed = append(parsed, canonical)
|
|
}
|
|
|
|
if len(parsed) == 0 {
|
|
return []string{"vacuum"}, nil
|
|
}
|
|
return parsed, nil
|
|
}
|
|
|
|
func canonicalPluginWorkerJobType(jobType string) (string, error) {
|
|
switch strings.ToLower(strings.TrimSpace(jobType)) {
|
|
case "", "vacuum":
|
|
return "vacuum", nil
|
|
case "volume_balance", "balance", "volume.balance", "volume-balance":
|
|
return "volume_balance", nil
|
|
case "erasure_coding", "erasure-coding", "erasure.coding", "ec":
|
|
return "erasure_coding", nil
|
|
default:
|
|
return "", fmt.Errorf("unsupported plugin job type %q", jobType)
|
|
}
|
|
}
|
|
|
|
func resolvePluginWorkerAdminServer(adminServer string) string {
|
|
adminServer = strings.TrimSpace(adminServer)
|
|
host, httpPort, hasExplicitGrpcPort, err := parsePluginWorkerAdminAddress(adminServer)
|
|
if err != nil || hasExplicitGrpcPort {
|
|
return adminServer
|
|
}
|
|
|
|
workerGrpcPort, err := fetchPluginWorkerGrpcPort(host, httpPort)
|
|
if err != nil || workerGrpcPort <= 0 {
|
|
return adminServer
|
|
}
|
|
|
|
// Keep canonical host:http form when admin gRPC follows the default +10000 rule.
|
|
if workerGrpcPort == httpPort+10000 {
|
|
return adminServer
|
|
}
|
|
|
|
return fmt.Sprintf("%s:%d.%d", host, httpPort, workerGrpcPort)
|
|
}
|
|
|
|
func parsePluginWorkerAdminAddress(adminServer string) (host string, httpPort int, hasExplicitGrpcPort bool, err error) {
|
|
adminServer = strings.TrimSpace(adminServer)
|
|
colonIndex := strings.LastIndex(adminServer, ":")
|
|
if colonIndex <= 0 || colonIndex >= len(adminServer)-1 {
|
|
return "", 0, false, fmt.Errorf("invalid admin address %q", adminServer)
|
|
}
|
|
|
|
host = adminServer[:colonIndex]
|
|
portPart := adminServer[colonIndex+1:]
|
|
if dotIndex := strings.LastIndex(portPart, "."); dotIndex > 0 && dotIndex < len(portPart)-1 {
|
|
if _, parseErr := strconv.Atoi(portPart[dotIndex+1:]); parseErr == nil {
|
|
hasExplicitGrpcPort = true
|
|
portPart = portPart[:dotIndex]
|
|
}
|
|
}
|
|
|
|
httpPort, err = strconv.Atoi(portPart)
|
|
if err != nil || httpPort <= 0 {
|
|
return "", 0, false, fmt.Errorf("invalid admin http port in %q", adminServer)
|
|
}
|
|
return host, httpPort, hasExplicitGrpcPort, nil
|
|
}
|
|
|
|
func fetchPluginWorkerGrpcPort(host string, httpPort int) (int, error) {
|
|
client := &http.Client{Timeout: 2 * time.Second}
|
|
address := util.JoinHostPort(host, httpPort)
|
|
var lastErr error
|
|
|
|
for _, scheme := range []string{"http", "https"} {
|
|
statusURL := fmt.Sprintf("%s://%s/api/plugin/status", scheme, address)
|
|
resp, err := client.Get(statusURL)
|
|
if err != nil {
|
|
lastErr = err
|
|
continue
|
|
}
|
|
|
|
var payload struct {
|
|
WorkerGrpcPort int `json:"worker_grpc_port"`
|
|
}
|
|
decodeErr := json.NewDecoder(resp.Body).Decode(&payload)
|
|
resp.Body.Close()
|
|
|
|
if resp.StatusCode != http.StatusOK {
|
|
lastErr = fmt.Errorf("status code %d from %s", resp.StatusCode, statusURL)
|
|
continue
|
|
}
|
|
if decodeErr != nil {
|
|
lastErr = fmt.Errorf("decode plugin status from %s: %w", statusURL, decodeErr)
|
|
continue
|
|
}
|
|
if payload.WorkerGrpcPort <= 0 {
|
|
lastErr = fmt.Errorf("plugin status from %s returned empty worker_grpc_port", statusURL)
|
|
continue
|
|
}
|
|
|
|
return payload.WorkerGrpcPort, nil
|
|
}
|
|
|
|
if lastErr == nil {
|
|
lastErr = fmt.Errorf("plugin status endpoint unavailable")
|
|
}
|
|
return 0, lastErr
|
|
}
|
|
|
|
func pluginWorkerHealthHandler(w http.ResponseWriter, _ *http.Request) {
|
|
w.WriteHeader(http.StatusOK)
|
|
}
|
|
|
|
func pluginWorkerReadyHandler(pluginRuntime *pluginworker.Worker) http.HandlerFunc {
|
|
return func(w http.ResponseWriter, _ *http.Request) {
|
|
if pluginRuntime == nil || !pluginRuntime.IsConnected() {
|
|
w.WriteHeader(http.StatusServiceUnavailable)
|
|
return
|
|
}
|
|
w.WriteHeader(http.StatusOK)
|
|
}
|
|
}
|
|
|
|
func startPluginWorkerMetricsServer(ip string, port int, pluginRuntime *pluginworker.Worker) {
|
|
mux := http.NewServeMux()
|
|
mux.HandleFunc("/health", pluginWorkerHealthHandler)
|
|
mux.HandleFunc("/ready", pluginWorkerReadyHandler(pluginRuntime))
|
|
mux.Handle("/metrics", promhttp.HandlerFor(statsCollect.Gather, promhttp.HandlerOpts{}))
|
|
|
|
glog.V(0).Infof("Starting plugin worker metrics server at %s", statsCollect.JoinHostPort(ip, port))
|
|
if err := http.ListenAndServe(statsCollect.JoinHostPort(ip, port), mux); err != nil {
|
|
glog.Errorf("Plugin worker metrics server failed to start: %v", err)
|
|
}
|
|
}
|