mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-06 14:45:51 +00:00
Compare commits
10
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b1c862e7e3 | ||
|
|
05902c7101 | ||
|
|
e0cccae794 | ||
|
|
6289beb8f5 | ||
|
|
35786d72e0 | ||
|
|
51d715b4e6 | ||
|
|
4f97f7d413 | ||
|
|
c8ad7b528d | ||
|
|
41f22f5c92 | ||
|
|
05138c6031 |
@@ -135,7 +135,7 @@ func ensureEnvironment(t *testing.T) {
|
||||
// 6. Start Admin
|
||||
os.RemoveAll(filepath.Join("tmp", "admin"))
|
||||
os.MkdirAll(filepath.Join("tmp", "admin"), 0755)
|
||||
startWeed(t, "admin", "admin", "-master=localhost:9333", "-port=23646", "-dataDir=./tmp/admin")
|
||||
startWeed(t, "admin", "admin", "-master=localhost:9333", "-port=23646", "-dataDir=./tmp/admin", "-scheduler.idleSleep=2")
|
||||
waitForUrl(t, AdminUrl+"/health", 60)
|
||||
|
||||
t.Log("Environment started successfully")
|
||||
|
||||
@@ -118,7 +118,7 @@ type AdminServer struct {
|
||||
|
||||
// Type definitions moved to types.go
|
||||
|
||||
func NewAdminServer(masters string, templateFS http.FileSystem, dataDir string, icebergPort int) *AdminServer {
|
||||
func NewAdminServer(masters string, templateFS http.FileSystem, dataDir string, icebergPort int, idleSleepDuration time.Duration) *AdminServer {
|
||||
grpcDialOption := security.LoadClientTLS(util.GetViper(), "grpc.admin")
|
||||
|
||||
// Create master client with multiple master support
|
||||
@@ -229,7 +229,8 @@ func NewAdminServer(masters string, templateFS http.FileSystem, dataDir string,
|
||||
}
|
||||
|
||||
plugin, err := adminplugin.New(adminplugin.Options{
|
||||
DataDir: dataDir,
|
||||
DataDir: dataDir,
|
||||
IdleSleepDuration: idleSleepDuration,
|
||||
ClusterContextProvider: func(_ context.Context) (*plugin_pb.ClusterContext, error) {
|
||||
return server.buildDefaultPluginClusterContext(), nil
|
||||
},
|
||||
@@ -238,7 +239,8 @@ func NewAdminServer(masters string, templateFS http.FileSystem, dataDir string,
|
||||
if err != nil && dataDir != "" {
|
||||
glog.Warningf("Failed to initialize plugin with dataDir=%q: %v. Falling back to in-memory plugin state.", dataDir, err)
|
||||
plugin, err = adminplugin.New(adminplugin.Options{
|
||||
DataDir: "",
|
||||
DataDir: "",
|
||||
IdleSleepDuration: idleSleepDuration,
|
||||
ClusterContextProvider: func(_ context.Context) (*plugin_pb.ClusterContext, error) {
|
||||
return server.buildDefaultPluginClusterContext(), nil
|
||||
},
|
||||
|
||||
@@ -32,6 +32,7 @@ type Options struct {
|
||||
OutgoingBufferSize int
|
||||
SendTimeout time.Duration
|
||||
SchedulerTick time.Duration
|
||||
IdleSleepDuration time.Duration
|
||||
ClusterContextProvider func(context.Context) (*plugin_pb.ClusterContext, error)
|
||||
LockManager LockManager
|
||||
}
|
||||
@@ -53,12 +54,16 @@ type Plugin struct {
|
||||
sendTimeout time.Duration
|
||||
|
||||
schedulerTick time.Duration
|
||||
idleSleepDuration time.Duration
|
||||
clusterContextProvider func(context.Context) (*plugin_pb.ClusterContext, error)
|
||||
lockManager LockManager
|
||||
|
||||
schedulerMu sync.Mutex
|
||||
nextDetectionAt map[string]time.Time
|
||||
detectionInFlight map[string]bool
|
||||
schedulerMu sync.Mutex
|
||||
schedulerPhase string
|
||||
currentJobType string
|
||||
iterationStartedAt time.Time
|
||||
lastIterationEndedAt time.Time
|
||||
lastIterationWorkDetected bool
|
||||
|
||||
detectorLeaseMu sync.Mutex
|
||||
detectorLeases map[string]string
|
||||
@@ -146,6 +151,10 @@ func New(options Options) (*Plugin, error) {
|
||||
if schedulerTick <= 0 {
|
||||
schedulerTick = defaultSchedulerTick
|
||||
}
|
||||
idleSleepDuration := options.IdleSleepDuration
|
||||
if idleSleepDuration <= 0 {
|
||||
idleSleepDuration = defaultIdleSleepDuration
|
||||
}
|
||||
|
||||
plugin := &Plugin{
|
||||
store: store,
|
||||
@@ -153,14 +162,14 @@ func New(options Options) (*Plugin, error) {
|
||||
outgoingBuffer: bufferSize,
|
||||
sendTimeout: sendTimeout,
|
||||
schedulerTick: schedulerTick,
|
||||
idleSleepDuration: idleSleepDuration,
|
||||
clusterContextProvider: options.ClusterContextProvider,
|
||||
lockManager: options.LockManager,
|
||||
schedulerPhase: "idle",
|
||||
sessions: make(map[string]*streamSession),
|
||||
pendingSchema: make(map[string]chan *plugin_pb.ConfigSchemaResponse),
|
||||
pendingDetection: make(map[string]*pendingDetectionState),
|
||||
pendingExecution: make(map[string]chan *plugin_pb.JobCompleted),
|
||||
nextDetectionAt: make(map[string]time.Time),
|
||||
detectionInFlight: make(map[string]bool),
|
||||
detectorLeases: make(map[string]string),
|
||||
schedulerExecReservations: make(map[string]int),
|
||||
schedulerDetection: make(map[string]*schedulerDetectionInfo),
|
||||
|
||||
@@ -17,7 +17,8 @@ var errExecutorAtCapacity = errors.New("executor is at capacity")
|
||||
|
||||
const (
|
||||
defaultSchedulerTick = 5 * time.Second
|
||||
defaultScheduledDetectionInterval = 300 * time.Second
|
||||
defaultIdleSleepDuration = 17 * time.Minute
|
||||
defaultMaxJobTypeDuration = 30 * time.Minute
|
||||
defaultScheduledDetectionTimeout = 45 * time.Second
|
||||
defaultScheduledExecutionTimeout = 90 * time.Second
|
||||
defaultScheduledMaxResults int32 = 1000
|
||||
@@ -31,9 +32,9 @@ const (
|
||||
)
|
||||
|
||||
type schedulerPolicy struct {
|
||||
DetectionInterval time.Duration
|
||||
DetectionTimeout time.Duration
|
||||
ExecutionTimeout time.Duration
|
||||
MaxJobTypeDuration time.Duration
|
||||
RetryBackoff time.Duration
|
||||
MaxResults int32
|
||||
ExecutionConcurrency int
|
||||
@@ -42,33 +43,63 @@ type schedulerPolicy struct {
|
||||
ExecutorReserveBackoff time.Duration
|
||||
}
|
||||
|
||||
func (r *Plugin) setSchedulerPhase(phase, jobType string) {
|
||||
r.schedulerMu.Lock()
|
||||
r.schedulerPhase = phase
|
||||
r.currentJobType = jobType
|
||||
r.schedulerMu.Unlock()
|
||||
}
|
||||
|
||||
func (r *Plugin) schedulerLoop() {
|
||||
defer r.wg.Done()
|
||||
ticker := time.NewTicker(r.schedulerTick)
|
||||
defer ticker.Stop()
|
||||
|
||||
// Try once immediately on startup.
|
||||
r.runSchedulerTick()
|
||||
workDetected := r.runSchedulerIteration()
|
||||
|
||||
for {
|
||||
if workDetected {
|
||||
// Immediate re-iteration when work was found.
|
||||
select {
|
||||
case <-r.shutdownCh:
|
||||
return
|
||||
default:
|
||||
}
|
||||
} else {
|
||||
r.setSchedulerPhase("idle", "")
|
||||
if !waitForShutdownOrTimer(r.shutdownCh, r.idleSleepDuration) {
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
select {
|
||||
case <-r.shutdownCh:
|
||||
return
|
||||
case <-ticker.C:
|
||||
r.runSchedulerTick()
|
||||
default:
|
||||
}
|
||||
|
||||
workDetected = r.runSchedulerIteration()
|
||||
}
|
||||
}
|
||||
|
||||
func (r *Plugin) runSchedulerTick() {
|
||||
r.expireStaleJobs(time.Now().UTC())
|
||||
func (r *Plugin) runSchedulerIteration() bool {
|
||||
now := time.Now().UTC()
|
||||
|
||||
r.schedulerMu.Lock()
|
||||
r.iterationStartedAt = now
|
||||
r.schedulerMu.Unlock()
|
||||
|
||||
r.expireStaleJobs(now)
|
||||
|
||||
jobTypes := r.registry.DetectableJobTypes()
|
||||
if len(jobTypes) == 0 {
|
||||
return
|
||||
r.finishIteration(false)
|
||||
return false
|
||||
}
|
||||
|
||||
active := make(map[string]struct{}, len(jobTypes))
|
||||
enabledJobTypes := make([]string, 0, len(jobTypes))
|
||||
policies := make(map[string]schedulerPolicy, len(jobTypes))
|
||||
|
||||
for _, jobType := range jobTypes {
|
||||
active[jobType] = struct{}{}
|
||||
|
||||
@@ -82,19 +113,171 @@ func (r *Plugin) runSchedulerTick() {
|
||||
continue
|
||||
}
|
||||
|
||||
if !r.markDetectionDue(jobType, policy.DetectionInterval) {
|
||||
continue
|
||||
}
|
||||
|
||||
r.wg.Add(1)
|
||||
go func(jt string, p schedulerPolicy) {
|
||||
defer r.wg.Done()
|
||||
r.runScheduledDetection(jt, p)
|
||||
}(jobType, policy)
|
||||
enabledJobTypes = append(enabledJobTypes, jobType)
|
||||
policies[jobType] = policy
|
||||
}
|
||||
|
||||
r.pruneSchedulerState(active)
|
||||
r.pruneDetectorLeases(active)
|
||||
|
||||
if len(enabledJobTypes) == 0 {
|
||||
r.finishIteration(false)
|
||||
return false
|
||||
}
|
||||
|
||||
// Acquire the lock ONCE for the entire iteration.
|
||||
r.setSchedulerPhase("acquiring_lock", "")
|
||||
releaseLock, lockErr := r.acquireAdminLock("plugin scheduler iteration")
|
||||
if lockErr != nil {
|
||||
glog.Warningf("Plugin scheduler failed to acquire lock: %v", lockErr)
|
||||
r.appendActivity(JobActivity{
|
||||
Source: "admin_scheduler",
|
||||
Message: fmt.Sprintf("scheduler iteration aborted: failed to acquire lock: %v", lockErr),
|
||||
Stage: "failed",
|
||||
OccurredAt: timeToPtr(time.Now().UTC()),
|
||||
})
|
||||
r.finishIteration(false)
|
||||
return false
|
||||
}
|
||||
defer releaseLock()
|
||||
|
||||
// Load cluster context ONCE for all job types.
|
||||
clusterContext, err := r.loadSchedulerClusterContext()
|
||||
if err != nil {
|
||||
r.appendActivity(JobActivity{
|
||||
Source: "admin_scheduler",
|
||||
Message: fmt.Sprintf("scheduler iteration aborted: %v", err),
|
||||
Stage: "failed",
|
||||
OccurredAt: timeToPtr(time.Now().UTC()),
|
||||
})
|
||||
r.finishIteration(false)
|
||||
return false
|
||||
}
|
||||
|
||||
// Process each job type sequentially.
|
||||
anyWorkDetected := false
|
||||
for _, jobType := range enabledJobTypes {
|
||||
select {
|
||||
case <-r.shutdownCh:
|
||||
r.finishIteration(anyWorkDetected)
|
||||
return anyWorkDetected
|
||||
default:
|
||||
}
|
||||
|
||||
policy := policies[jobType]
|
||||
r.setSchedulerPhase("processing", jobType)
|
||||
|
||||
if r.runJobTypeIteration(jobType, policy, clusterContext) {
|
||||
anyWorkDetected = true
|
||||
}
|
||||
}
|
||||
|
||||
r.finishIteration(anyWorkDetected)
|
||||
return anyWorkDetected
|
||||
}
|
||||
|
||||
func (r *Plugin) finishIteration(workDetected bool) {
|
||||
now := time.Now().UTC()
|
||||
r.schedulerMu.Lock()
|
||||
r.lastIterationEndedAt = now
|
||||
r.lastIterationWorkDetected = workDetected
|
||||
r.currentJobType = ""
|
||||
r.schedulerPhase = "idle"
|
||||
r.schedulerMu.Unlock()
|
||||
}
|
||||
|
||||
func (r *Plugin) runJobTypeIteration(
|
||||
jobType string,
|
||||
policy schedulerPolicy,
|
||||
clusterContext *plugin_pb.ClusterContext,
|
||||
) bool {
|
||||
budget := policy.MaxJobTypeDuration
|
||||
if budget <= 0 {
|
||||
budget = defaultMaxJobTypeDuration
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(r.ctx, budget)
|
||||
defer cancel()
|
||||
|
||||
start := time.Now().UTC()
|
||||
r.appendActivity(JobActivity{
|
||||
JobType: jobType,
|
||||
Source: "admin_scheduler",
|
||||
Message: "scheduled detection started",
|
||||
Stage: "detecting",
|
||||
OccurredAt: timeToPtr(start),
|
||||
})
|
||||
|
||||
if skip, waitingCount, waitingThreshold := r.shouldSkipDetectionForWaitingJobs(jobType, policy); skip {
|
||||
r.recordSchedulerDetectionSkip(jobType, fmt.Sprintf("waiting backlog %d reached threshold %d", waitingCount, waitingThreshold))
|
||||
r.appendActivity(JobActivity{
|
||||
JobType: jobType,
|
||||
Source: "admin_scheduler",
|
||||
Message: fmt.Sprintf("scheduled detection skipped: waiting backlog %d reached threshold %d", waitingCount, waitingThreshold),
|
||||
Stage: "skipped_waiting_backlog",
|
||||
OccurredAt: timeToPtr(time.Now().UTC()),
|
||||
})
|
||||
return false
|
||||
}
|
||||
|
||||
detectionTimeout := policy.DetectionTimeout
|
||||
if detectionTimeout <= 0 {
|
||||
detectionTimeout = defaultScheduledDetectionTimeout
|
||||
}
|
||||
detCtx, detCancel := context.WithTimeout(ctx, detectionTimeout)
|
||||
proposals, err := r.RunDetection(detCtx, jobType, clusterContext, policy.MaxResults)
|
||||
detCancel()
|
||||
if err != nil {
|
||||
r.recordSchedulerDetectionError(jobType, err)
|
||||
r.appendActivity(JobActivity{
|
||||
JobType: jobType,
|
||||
Source: "admin_scheduler",
|
||||
Message: fmt.Sprintf("scheduled detection failed: %v", err),
|
||||
Stage: "failed",
|
||||
OccurredAt: timeToPtr(time.Now().UTC()),
|
||||
})
|
||||
return false
|
||||
}
|
||||
|
||||
r.appendActivity(JobActivity{
|
||||
JobType: jobType,
|
||||
Source: "admin_scheduler",
|
||||
Message: fmt.Sprintf("scheduled detection completed: %d proposal(s)", len(proposals)),
|
||||
Stage: "detected",
|
||||
OccurredAt: timeToPtr(time.Now().UTC()),
|
||||
})
|
||||
r.recordSchedulerDetectionSuccess(jobType, len(proposals))
|
||||
|
||||
filteredByActive, skippedActive := r.filterProposalsWithActiveJobs(jobType, proposals)
|
||||
if skippedActive > 0 {
|
||||
r.appendActivity(JobActivity{
|
||||
JobType: jobType,
|
||||
Source: "admin_scheduler",
|
||||
Message: fmt.Sprintf("scheduled detection skipped %d proposal(s) due to active assigned/running jobs", skippedActive),
|
||||
Stage: "deduped_active_jobs",
|
||||
OccurredAt: timeToPtr(time.Now().UTC()),
|
||||
})
|
||||
}
|
||||
|
||||
if len(filteredByActive) == 0 {
|
||||
return false
|
||||
}
|
||||
|
||||
filtered := r.filterScheduledProposals(filteredByActive)
|
||||
if len(filtered) != len(filteredByActive) {
|
||||
r.appendActivity(JobActivity{
|
||||
JobType: jobType,
|
||||
Source: "admin_scheduler",
|
||||
Message: fmt.Sprintf("scheduled detection deduped %d proposal(s) within this run", len(filteredByActive)-len(filtered)),
|
||||
Stage: "deduped",
|
||||
OccurredAt: timeToPtr(time.Now().UTC()),
|
||||
})
|
||||
}
|
||||
|
||||
if len(filtered) == 0 {
|
||||
return false
|
||||
}
|
||||
|
||||
r.dispatchScheduledProposals(ctx, jobType, filtered, clusterContext, policy)
|
||||
return true
|
||||
}
|
||||
|
||||
func (r *Plugin) loadSchedulerPolicy(jobType string) (schedulerPolicy, bool, error) {
|
||||
@@ -116,9 +299,9 @@ func (r *Plugin) loadSchedulerPolicy(jobType string) (schedulerPolicy, bool, err
|
||||
}
|
||||
|
||||
policy := schedulerPolicy{
|
||||
DetectionInterval: durationFromSeconds(adminRuntime.DetectionIntervalSeconds, defaultScheduledDetectionInterval),
|
||||
DetectionTimeout: durationFromSeconds(adminRuntime.DetectionTimeoutSeconds, defaultScheduledDetectionTimeout),
|
||||
ExecutionTimeout: defaultScheduledExecutionTimeout,
|
||||
MaxJobTypeDuration: defaultMaxJobTypeDuration,
|
||||
RetryBackoff: durationFromSeconds(adminRuntime.RetryBackoffSeconds, defaultScheduledRetryBackoff),
|
||||
MaxResults: adminRuntime.MaxJobsPerDetection,
|
||||
ExecutionConcurrency: int(adminRuntime.GlobalExecutionConcurrency),
|
||||
@@ -127,9 +310,6 @@ func (r *Plugin) loadSchedulerPolicy(jobType string) (schedulerPolicy, bool, err
|
||||
ExecutorReserveBackoff: 200 * time.Millisecond,
|
||||
}
|
||||
|
||||
if policy.DetectionInterval < r.schedulerTick {
|
||||
policy.DetectionInterval = r.schedulerTick
|
||||
}
|
||||
if policy.MaxResults <= 0 {
|
||||
policy.MaxResults = defaultScheduledMaxResults
|
||||
}
|
||||
@@ -166,28 +346,18 @@ func (r *Plugin) ListSchedulerStates() ([]SchedulerJobTypeState, error) {
|
||||
}
|
||||
|
||||
r.schedulerMu.Lock()
|
||||
nextDetectionAt := make(map[string]time.Time, len(r.nextDetectionAt))
|
||||
for jobType, nextRun := range r.nextDetectionAt {
|
||||
nextDetectionAt[jobType] = nextRun
|
||||
}
|
||||
detectionInFlight := make(map[string]bool, len(r.detectionInFlight))
|
||||
for jobType, inFlight := range r.detectionInFlight {
|
||||
detectionInFlight[jobType] = inFlight
|
||||
}
|
||||
currentJobType := r.currentJobType
|
||||
r.schedulerMu.Unlock()
|
||||
|
||||
states := make([]SchedulerJobTypeState, 0, len(jobTypes))
|
||||
for _, jobTypeInfo := range jobTypes {
|
||||
jobType := jobTypeInfo.JobType
|
||||
state := SchedulerJobTypeState{
|
||||
JobType: jobType,
|
||||
DetectionInFlight: detectionInFlight[jobType],
|
||||
JobType: jobType,
|
||||
}
|
||||
|
||||
if nextRun, ok := nextDetectionAt[jobType]; ok && !nextRun.IsZero() {
|
||||
nextRunUTC := nextRun.UTC()
|
||||
state.NextDetectionAt = &nextRunUTC
|
||||
}
|
||||
// Mark as in-flight if this job type is being processed right now.
|
||||
state.DetectionInFlight = jobType == currentJobType && currentJobType != ""
|
||||
|
||||
policy, enabled, loadErr := r.loadSchedulerPolicy(jobType)
|
||||
|
||||
@@ -196,9 +366,9 @@ func (r *Plugin) ListSchedulerStates() ([]SchedulerJobTypeState, error) {
|
||||
} else {
|
||||
state.Enabled = enabled
|
||||
if enabled {
|
||||
state.DetectionIntervalSeconds = secondsFromDuration(policy.DetectionInterval)
|
||||
state.DetectionTimeoutSeconds = secondsFromDuration(policy.DetectionTimeout)
|
||||
state.ExecutionTimeoutSeconds = secondsFromDuration(policy.ExecutionTimeout)
|
||||
state.MaxJobTypeDurationSeconds = secondsFromDuration(policy.MaxJobTypeDuration)
|
||||
state.MaxJobsPerDetection = policy.MaxResults
|
||||
state.GlobalExecutionConcurrency = policy.ExecutionConcurrency
|
||||
state.PerWorkerExecutionConcurrency = policy.PerWorkerConcurrency
|
||||
@@ -261,49 +431,7 @@ func deriveSchedulerAdminRuntime(
|
||||
}
|
||||
}
|
||||
|
||||
func (r *Plugin) markDetectionDue(jobType string, interval time.Duration) bool {
|
||||
now := time.Now().UTC()
|
||||
|
||||
r.schedulerMu.Lock()
|
||||
defer r.schedulerMu.Unlock()
|
||||
|
||||
if r.detectionInFlight[jobType] {
|
||||
return false
|
||||
}
|
||||
|
||||
nextRun, exists := r.nextDetectionAt[jobType]
|
||||
if exists && now.Before(nextRun) {
|
||||
return false
|
||||
}
|
||||
|
||||
r.nextDetectionAt[jobType] = now.Add(interval)
|
||||
r.detectionInFlight[jobType] = true
|
||||
return true
|
||||
}
|
||||
|
||||
func (r *Plugin) finishDetection(jobType string) {
|
||||
r.schedulerMu.Lock()
|
||||
delete(r.detectionInFlight, jobType)
|
||||
r.schedulerMu.Unlock()
|
||||
}
|
||||
|
||||
func (r *Plugin) pruneSchedulerState(activeJobTypes map[string]struct{}) {
|
||||
r.schedulerMu.Lock()
|
||||
defer r.schedulerMu.Unlock()
|
||||
|
||||
for jobType := range r.nextDetectionAt {
|
||||
if _, ok := activeJobTypes[jobType]; !ok {
|
||||
delete(r.nextDetectionAt, jobType)
|
||||
delete(r.detectionInFlight, jobType)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (r *Plugin) clearSchedulerJobType(jobType string) {
|
||||
r.schedulerMu.Lock()
|
||||
delete(r.nextDetectionAt, jobType)
|
||||
delete(r.detectionInFlight, jobType)
|
||||
r.schedulerMu.Unlock()
|
||||
r.clearDetectorLease(jobType, "")
|
||||
}
|
||||
|
||||
@@ -318,116 +446,6 @@ func (r *Plugin) pruneDetectorLeases(activeJobTypes map[string]struct{}) {
|
||||
}
|
||||
}
|
||||
|
||||
func (r *Plugin) runScheduledDetection(jobType string, policy schedulerPolicy) {
|
||||
defer r.finishDetection(jobType)
|
||||
|
||||
releaseLock, lockErr := r.acquireAdminLock(fmt.Sprintf("plugin scheduled detection %s", jobType))
|
||||
if lockErr != nil {
|
||||
r.recordSchedulerDetectionError(jobType, lockErr)
|
||||
r.appendActivity(JobActivity{
|
||||
JobType: jobType,
|
||||
Source: "admin_scheduler",
|
||||
Message: fmt.Sprintf("scheduled detection aborted: failed to acquire lock: %v", lockErr),
|
||||
Stage: "failed",
|
||||
OccurredAt: timeToPtr(time.Now().UTC()),
|
||||
})
|
||||
return
|
||||
}
|
||||
if releaseLock != nil {
|
||||
defer releaseLock()
|
||||
}
|
||||
|
||||
start := time.Now().UTC()
|
||||
r.appendActivity(JobActivity{
|
||||
JobType: jobType,
|
||||
Source: "admin_scheduler",
|
||||
Message: "scheduled detection started",
|
||||
Stage: "detecting",
|
||||
OccurredAt: timeToPtr(start),
|
||||
})
|
||||
|
||||
if skip, waitingCount, waitingThreshold := r.shouldSkipDetectionForWaitingJobs(jobType, policy); skip {
|
||||
r.recordSchedulerDetectionSkip(jobType, fmt.Sprintf("waiting backlog %d reached threshold %d", waitingCount, waitingThreshold))
|
||||
r.appendActivity(JobActivity{
|
||||
JobType: jobType,
|
||||
Source: "admin_scheduler",
|
||||
Message: fmt.Sprintf("scheduled detection skipped: waiting backlog %d reached threshold %d", waitingCount, waitingThreshold),
|
||||
Stage: "skipped_waiting_backlog",
|
||||
OccurredAt: timeToPtr(time.Now().UTC()),
|
||||
})
|
||||
return
|
||||
}
|
||||
|
||||
clusterContext, err := r.loadSchedulerClusterContext()
|
||||
if err != nil {
|
||||
r.recordSchedulerDetectionError(jobType, err)
|
||||
r.appendActivity(JobActivity{
|
||||
JobType: jobType,
|
||||
Source: "admin_scheduler",
|
||||
Message: fmt.Sprintf("scheduled detection aborted: %v", err),
|
||||
Stage: "failed",
|
||||
OccurredAt: timeToPtr(time.Now().UTC()),
|
||||
})
|
||||
return
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), policy.DetectionTimeout)
|
||||
proposals, err := r.RunDetection(ctx, jobType, clusterContext, policy.MaxResults)
|
||||
cancel()
|
||||
if err != nil {
|
||||
r.recordSchedulerDetectionError(jobType, err)
|
||||
r.appendActivity(JobActivity{
|
||||
JobType: jobType,
|
||||
Source: "admin_scheduler",
|
||||
Message: fmt.Sprintf("scheduled detection failed: %v", err),
|
||||
Stage: "failed",
|
||||
OccurredAt: timeToPtr(time.Now().UTC()),
|
||||
})
|
||||
return
|
||||
}
|
||||
|
||||
r.appendActivity(JobActivity{
|
||||
JobType: jobType,
|
||||
Source: "admin_scheduler",
|
||||
Message: fmt.Sprintf("scheduled detection completed: %d proposal(s)", len(proposals)),
|
||||
Stage: "detected",
|
||||
OccurredAt: timeToPtr(time.Now().UTC()),
|
||||
})
|
||||
r.recordSchedulerDetectionSuccess(jobType, len(proposals))
|
||||
|
||||
filteredByActive, skippedActive := r.filterProposalsWithActiveJobs(jobType, proposals)
|
||||
if skippedActive > 0 {
|
||||
r.appendActivity(JobActivity{
|
||||
JobType: jobType,
|
||||
Source: "admin_scheduler",
|
||||
Message: fmt.Sprintf("scheduled detection skipped %d proposal(s) due to active assigned/running jobs", skippedActive),
|
||||
Stage: "deduped_active_jobs",
|
||||
OccurredAt: timeToPtr(time.Now().UTC()),
|
||||
})
|
||||
}
|
||||
|
||||
if len(filteredByActive) == 0 {
|
||||
return
|
||||
}
|
||||
|
||||
filtered := r.filterScheduledProposals(filteredByActive)
|
||||
if len(filtered) != len(filteredByActive) {
|
||||
r.appendActivity(JobActivity{
|
||||
JobType: jobType,
|
||||
Source: "admin_scheduler",
|
||||
Message: fmt.Sprintf("scheduled detection deduped %d proposal(s) within this run", len(filteredByActive)-len(filtered)),
|
||||
Stage: "deduped",
|
||||
OccurredAt: timeToPtr(time.Now().UTC()),
|
||||
})
|
||||
}
|
||||
|
||||
if len(filtered) == 0 {
|
||||
return
|
||||
}
|
||||
|
||||
r.dispatchScheduledProposals(jobType, filtered, clusterContext, policy)
|
||||
}
|
||||
|
||||
func (r *Plugin) loadSchedulerClusterContext() (*plugin_pb.ClusterContext, error) {
|
||||
if r.clusterContextProvider == nil {
|
||||
return nil, fmt.Errorf("cluster context provider is not configured")
|
||||
@@ -447,6 +465,7 @@ func (r *Plugin) loadSchedulerClusterContext() (*plugin_pb.ClusterContext, error
|
||||
}
|
||||
|
||||
func (r *Plugin) dispatchScheduledProposals(
|
||||
parentCtx context.Context,
|
||||
jobType string,
|
||||
proposals []*plugin_pb.JobProposal,
|
||||
clusterContext *plugin_pb.ClusterContext,
|
||||
@@ -460,6 +479,9 @@ func (r *Plugin) dispatchScheduledProposals(
|
||||
case <-r.shutdownCh:
|
||||
close(jobQueue)
|
||||
return
|
||||
case <-parentCtx.Done():
|
||||
close(jobQueue)
|
||||
return
|
||||
default:
|
||||
jobQueue <- job
|
||||
}
|
||||
@@ -485,6 +507,8 @@ func (r *Plugin) dispatchScheduledProposals(
|
||||
select {
|
||||
case <-r.shutdownCh:
|
||||
return
|
||||
case <-parentCtx.Done():
|
||||
return
|
||||
default:
|
||||
}
|
||||
|
||||
@@ -492,10 +516,12 @@ func (r *Plugin) dispatchScheduledProposals(
|
||||
select {
|
||||
case <-r.shutdownCh:
|
||||
return
|
||||
case <-parentCtx.Done():
|
||||
return
|
||||
default:
|
||||
}
|
||||
|
||||
executor, release, reserveErr := r.reserveScheduledExecutor(jobType, policy)
|
||||
executor, release, reserveErr := r.reserveScheduledExecutor(parentCtx, jobType, policy)
|
||||
if reserveErr != nil {
|
||||
select {
|
||||
case <-r.shutdownCh:
|
||||
@@ -515,7 +541,7 @@ func (r *Plugin) dispatchScheduledProposals(
|
||||
break
|
||||
}
|
||||
|
||||
err := r.executeScheduledJobWithExecutor(executor, job, clusterContext, policy)
|
||||
err := r.executeScheduledJobWithExecutor(parentCtx, executor, job, clusterContext, policy)
|
||||
release()
|
||||
if errors.Is(err, errExecutorAtCapacity) {
|
||||
r.trackExecutionQueued(job)
|
||||
@@ -560,6 +586,7 @@ func (r *Plugin) dispatchScheduledProposals(
|
||||
}
|
||||
|
||||
func (r *Plugin) reserveScheduledExecutor(
|
||||
ctx context.Context,
|
||||
jobType string,
|
||||
policy schedulerPolicy,
|
||||
) (*WorkerSession, func(), error) {
|
||||
@@ -572,6 +599,8 @@ func (r *Plugin) reserveScheduledExecutor(
|
||||
select {
|
||||
case <-r.shutdownCh:
|
||||
return nil, nil, fmt.Errorf("plugin is shutting down")
|
||||
case <-ctx.Done():
|
||||
return nil, nil, ctx.Err()
|
||||
default:
|
||||
}
|
||||
|
||||
@@ -581,8 +610,8 @@ func (r *Plugin) reserveScheduledExecutor(
|
||||
|
||||
executors, err := r.registry.ListExecutors(jobType)
|
||||
if err != nil {
|
||||
if !waitForShutdownOrTimer(r.shutdownCh, policy.ExecutorReserveBackoff) {
|
||||
return nil, nil, fmt.Errorf("plugin is shutting down")
|
||||
if !waitForShutdownOrCtx(r.shutdownCh, ctx, policy.ExecutorReserveBackoff) {
|
||||
return nil, nil, fmt.Errorf("plugin is shutting down or context canceled")
|
||||
}
|
||||
continue
|
||||
}
|
||||
@@ -595,8 +624,8 @@ func (r *Plugin) reserveScheduledExecutor(
|
||||
return executor, release, nil
|
||||
}
|
||||
|
||||
if !waitForShutdownOrTimer(r.shutdownCh, policy.ExecutorReserveBackoff) {
|
||||
return nil, nil, fmt.Errorf("plugin is shutting down")
|
||||
if !waitForShutdownOrCtx(r.shutdownCh, ctx, policy.ExecutorReserveBackoff) {
|
||||
return nil, nil, fmt.Errorf("plugin is shutting down or context canceled")
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -680,6 +709,7 @@ func schedulerWorkerExecutionLimit(executor *WorkerSession, jobType string, poli
|
||||
}
|
||||
|
||||
func (r *Plugin) executeScheduledJobWithExecutor(
|
||||
parentCtx context.Context,
|
||||
executor *WorkerSession,
|
||||
job *plugin_pb.JobSpec,
|
||||
clusterContext *plugin_pb.ClusterContext,
|
||||
@@ -695,10 +725,12 @@ func (r *Plugin) executeScheduledJobWithExecutor(
|
||||
select {
|
||||
case <-r.shutdownCh:
|
||||
return fmt.Errorf("plugin is shutting down")
|
||||
case <-parentCtx.Done():
|
||||
return parentCtx.Err()
|
||||
default:
|
||||
}
|
||||
|
||||
execCtx, cancel := context.WithTimeout(context.Background(), policy.ExecutionTimeout)
|
||||
execCtx, cancel := context.WithTimeout(parentCtx, policy.ExecutionTimeout)
|
||||
_, err := r.executeJobWithExecutor(execCtx, executor, job, clusterContext, int32(attempt))
|
||||
cancel()
|
||||
if err == nil {
|
||||
@@ -718,8 +750,8 @@ func (r *Plugin) executeScheduledJobWithExecutor(
|
||||
Stage: "retry",
|
||||
OccurredAt: timeToPtr(time.Now().UTC()),
|
||||
})
|
||||
if !waitForShutdownOrTimer(r.shutdownCh, policy.RetryBackoff) {
|
||||
return fmt.Errorf("plugin is shutting down")
|
||||
if !waitForShutdownOrCtx(r.shutdownCh, parentCtx, policy.RetryBackoff) {
|
||||
return fmt.Errorf("plugin is shutting down or context canceled")
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -861,6 +893,24 @@ func waitForShutdownOrTimer(shutdown <-chan struct{}, duration time.Duration) bo
|
||||
}
|
||||
}
|
||||
|
||||
func waitForShutdownOrCtx(shutdown <-chan struct{}, ctx context.Context, duration time.Duration) bool {
|
||||
if duration <= 0 {
|
||||
return true
|
||||
}
|
||||
|
||||
timer := time.NewTimer(duration)
|
||||
defer timer.Stop()
|
||||
|
||||
select {
|
||||
case <-shutdown:
|
||||
return false
|
||||
case <-ctx.Done():
|
||||
return false
|
||||
case <-timer.C:
|
||||
return true
|
||||
}
|
||||
}
|
||||
|
||||
// filterProposalsWithActiveJobs removes proposals whose dedupe keys already have active jobs.
|
||||
// It first expires stale tracked jobs via expireStaleJobs, which can mutate scheduler state,
|
||||
// so callers should treat this method as a stateful operation.
|
||||
|
||||
@@ -1,7 +1,9 @@
|
||||
package plugin
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
@@ -53,6 +55,9 @@ func TestLoadSchedulerPolicyUsesAdminConfig(t *testing.T) {
|
||||
if policy.RetryLimit != 4 {
|
||||
t.Fatalf("unexpected retry limit: got=%d", policy.RetryLimit)
|
||||
}
|
||||
if policy.MaxJobTypeDuration != defaultMaxJobTypeDuration {
|
||||
t.Fatalf("unexpected max job type duration: got=%v", policy.MaxJobTypeDuration)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadSchedulerPolicyUsesDescriptorDefaultsWhenConfigMissing(t *testing.T) {
|
||||
@@ -126,13 +131,13 @@ func TestReserveScheduledExecutorRespectsPerWorkerLimit(t *testing.T) {
|
||||
ExecutorReserveBackoff: time.Millisecond,
|
||||
}
|
||||
|
||||
executor1, release1, err := pluginSvc.reserveScheduledExecutor("balance", policy)
|
||||
executor1, release1, err := pluginSvc.reserveScheduledExecutor(context.Background(), "balance", policy)
|
||||
if err != nil {
|
||||
t.Fatalf("reserve executor 1: %v", err)
|
||||
}
|
||||
defer release1()
|
||||
|
||||
executor2, release2, err := pluginSvc.reserveScheduledExecutor("balance", policy)
|
||||
executor2, release2, err := pluginSvc.reserveScheduledExecutor(context.Background(), "balance", policy)
|
||||
if err != nil {
|
||||
t.Fatalf("reserve executor 2: %v", err)
|
||||
}
|
||||
@@ -254,7 +259,7 @@ func TestReserveScheduledExecutorTimesOutWhenNoExecutor(t *testing.T) {
|
||||
|
||||
start := time.Now()
|
||||
pluginSvc.Shutdown()
|
||||
_, _, err = pluginSvc.reserveScheduledExecutor("missing-job-type", policy)
|
||||
_, _, err = pluginSvc.reserveScheduledExecutor(context.Background(), "missing-job-type", policy)
|
||||
if err == nil {
|
||||
t.Fatalf("expected reservation shutdown error")
|
||||
}
|
||||
@@ -285,7 +290,7 @@ func TestReserveScheduledExecutorWaitsForWorkerCapacity(t *testing.T) {
|
||||
ExecutorReserveBackoff: 5 * time.Millisecond,
|
||||
}
|
||||
|
||||
_, release1, err := pluginSvc.reserveScheduledExecutor("balance", policy)
|
||||
_, release1, err := pluginSvc.reserveScheduledExecutor(context.Background(), "balance", policy)
|
||||
if err != nil {
|
||||
t.Fatalf("reserve executor 1: %v", err)
|
||||
}
|
||||
@@ -296,7 +301,7 @@ func TestReserveScheduledExecutorWaitsForWorkerCapacity(t *testing.T) {
|
||||
}
|
||||
secondReserveCh := make(chan reserveResult, 1)
|
||||
go func() {
|
||||
_, release2, reserveErr := pluginSvc.reserveScheduledExecutor("balance", policy)
|
||||
_, release2, reserveErr := pluginSvc.reserveScheduledExecutor(context.Background(), "balance", policy)
|
||||
if release2 != nil {
|
||||
release2()
|
||||
}
|
||||
@@ -407,10 +412,9 @@ func TestListSchedulerStatesIncludesPolicyAndState(t *testing.T) {
|
||||
},
|
||||
})
|
||||
|
||||
nextDetectionAt := time.Now().UTC().Add(2 * time.Minute).Round(time.Second)
|
||||
// Mark this job type as currently processing to test DetectionInFlight.
|
||||
pluginSvc.schedulerMu.Lock()
|
||||
pluginSvc.nextDetectionAt[jobType] = nextDetectionAt
|
||||
pluginSvc.detectionInFlight[jobType] = true
|
||||
pluginSvc.currentJobType = jobType
|
||||
pluginSvc.schedulerMu.Unlock()
|
||||
|
||||
states, err := pluginSvc.ListSchedulerStates()
|
||||
@@ -429,16 +433,7 @@ func TestListSchedulerStatesIncludesPolicyAndState(t *testing.T) {
|
||||
t.Fatalf("unexpected policy error: %s", state.PolicyError)
|
||||
}
|
||||
if !state.DetectionInFlight {
|
||||
t.Fatalf("expected detection in flight")
|
||||
}
|
||||
if state.NextDetectionAt == nil {
|
||||
t.Fatalf("expected next detection time")
|
||||
}
|
||||
if state.NextDetectionAt.Unix() != nextDetectionAt.Unix() {
|
||||
t.Fatalf("unexpected next detection time: got=%v want=%v", state.NextDetectionAt, nextDetectionAt)
|
||||
}
|
||||
if state.DetectionIntervalSeconds != 45 {
|
||||
t.Fatalf("unexpected detection interval: got=%d", state.DetectionIntervalSeconds)
|
||||
t.Fatalf("expected detection in flight when current job type matches")
|
||||
}
|
||||
if state.DetectionTimeoutSeconds != 30 {
|
||||
t.Fatalf("unexpected detection timeout: got=%d", state.DetectionTimeoutSeconds)
|
||||
@@ -446,6 +441,9 @@ func TestListSchedulerStatesIncludesPolicyAndState(t *testing.T) {
|
||||
if state.ExecutionTimeoutSeconds != 90 {
|
||||
t.Fatalf("unexpected execution timeout: got=%d", state.ExecutionTimeoutSeconds)
|
||||
}
|
||||
if state.MaxJobTypeDurationSeconds != int32(defaultMaxJobTypeDuration/time.Second) {
|
||||
t.Fatalf("unexpected max job type duration: got=%d", state.MaxJobTypeDurationSeconds)
|
||||
}
|
||||
if state.MaxJobsPerDetection != 80 {
|
||||
t.Fatalf("unexpected max jobs per detection: got=%d", state.MaxJobsPerDetection)
|
||||
}
|
||||
@@ -467,6 +465,23 @@ func TestListSchedulerStatesIncludesPolicyAndState(t *testing.T) {
|
||||
if state.ExecutorWorkerCount != 1 {
|
||||
t.Fatalf("unexpected executor worker count: got=%d", state.ExecutorWorkerCount)
|
||||
}
|
||||
|
||||
// Clear the current job type and verify DetectionInFlight is false.
|
||||
pluginSvc.schedulerMu.Lock()
|
||||
pluginSvc.currentJobType = ""
|
||||
pluginSvc.schedulerMu.Unlock()
|
||||
|
||||
states2, err := pluginSvc.ListSchedulerStates()
|
||||
if err != nil {
|
||||
t.Fatalf("ListSchedulerStates (2): %v", err)
|
||||
}
|
||||
state2 := findSchedulerState(states2, jobType)
|
||||
if state2 == nil {
|
||||
t.Fatalf("missing scheduler state for %s (2)", jobType)
|
||||
}
|
||||
if state2.DetectionInFlight {
|
||||
t.Fatalf("expected detection not in flight when current job type is empty")
|
||||
}
|
||||
}
|
||||
|
||||
func TestListSchedulerStatesShowsDisabledWhenNoPolicy(t *testing.T) {
|
||||
@@ -581,3 +596,315 @@ func TestPickDetectorReassignsWhenLeaseIsStale(t *testing.T) {
|
||||
t.Fatalf("expected detector lease to be updated to worker-a, got=%s", lease)
|
||||
}
|
||||
}
|
||||
|
||||
// mockLockManager records lock/release calls for testing.
|
||||
type mockLockManager struct {
|
||||
mu sync.Mutex
|
||||
acquireCount int
|
||||
releaseCount int
|
||||
failAcquire bool
|
||||
}
|
||||
|
||||
func (m *mockLockManager) Acquire(reason string) (func(), error) {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
if m.failAcquire {
|
||||
return nil, fmt.Errorf("mock lock acquisition failed")
|
||||
}
|
||||
m.acquireCount++
|
||||
return func() {
|
||||
m.mu.Lock()
|
||||
m.releaseCount++
|
||||
m.mu.Unlock()
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (m *mockLockManager) Status() interface{} {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (m *mockLockManager) getAcquireCount() int {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
return m.acquireCount
|
||||
}
|
||||
|
||||
func (m *mockLockManager) getReleaseCount() int {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
return m.releaseCount
|
||||
}
|
||||
|
||||
func TestIterationAcquiresLockOnce(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
lock := &mockLockManager{}
|
||||
pluginSvc, err := New(Options{
|
||||
LockManager: lock,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("New: %v", err)
|
||||
}
|
||||
defer pluginSvc.Shutdown()
|
||||
|
||||
// Register two enabled job types.
|
||||
for _, jt := range []string{"vacuum", "balance"} {
|
||||
err = pluginSvc.SaveJobTypeConfig(&plugin_pb.PersistedJobTypeConfig{
|
||||
JobType: jt,
|
||||
AdminRuntime: &plugin_pb.AdminRuntimeConfig{
|
||||
Enabled: true,
|
||||
DetectionTimeoutSeconds: 5,
|
||||
MaxJobsPerDetection: 10,
|
||||
GlobalExecutionConcurrency: 1,
|
||||
PerWorkerExecutionConcurrency: 1,
|
||||
},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("SaveJobTypeConfig(%s): %v", jt, err)
|
||||
}
|
||||
pluginSvc.registry.UpsertFromHello(&plugin_pb.WorkerHello{
|
||||
WorkerId: "worker-" + jt,
|
||||
Capabilities: []*plugin_pb.JobTypeCapability{
|
||||
{JobType: jt, CanDetect: true, CanExecute: true},
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
// runSchedulerIteration requires a cluster context provider.
|
||||
pluginSvc.clusterContextProvider = func(_ context.Context) (*plugin_pb.ClusterContext, error) {
|
||||
return &plugin_pb.ClusterContext{}, nil
|
||||
}
|
||||
|
||||
pluginSvc.runSchedulerIteration()
|
||||
|
||||
// Lock should have been acquired exactly once (not per-job-type).
|
||||
if lock.getAcquireCount() != 1 {
|
||||
t.Fatalf("expected 1 lock acquisition, got %d", lock.getAcquireCount())
|
||||
}
|
||||
if lock.getReleaseCount() != 1 {
|
||||
t.Fatalf("expected 1 lock release, got %d", lock.getReleaseCount())
|
||||
}
|
||||
}
|
||||
|
||||
func TestIterationReturnsfalseWhenLockFails(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
lock := &mockLockManager{failAcquire: true}
|
||||
pluginSvc, err := New(Options{
|
||||
LockManager: lock,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("New: %v", err)
|
||||
}
|
||||
defer pluginSvc.Shutdown()
|
||||
|
||||
err = pluginSvc.SaveJobTypeConfig(&plugin_pb.PersistedJobTypeConfig{
|
||||
JobType: "vacuum",
|
||||
AdminRuntime: &plugin_pb.AdminRuntimeConfig{
|
||||
Enabled: true,
|
||||
DetectionTimeoutSeconds: 5,
|
||||
MaxJobsPerDetection: 10,
|
||||
GlobalExecutionConcurrency: 1,
|
||||
PerWorkerExecutionConcurrency: 1,
|
||||
},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("SaveJobTypeConfig: %v", err)
|
||||
}
|
||||
pluginSvc.registry.UpsertFromHello(&plugin_pb.WorkerHello{
|
||||
WorkerId: "worker-a",
|
||||
Capabilities: []*plugin_pb.JobTypeCapability{
|
||||
{JobType: "vacuum", CanDetect: true, CanExecute: true},
|
||||
},
|
||||
})
|
||||
|
||||
pluginSvc.clusterContextProvider = func(_ context.Context) (*plugin_pb.ClusterContext, error) {
|
||||
return &plugin_pb.ClusterContext{}, nil
|
||||
}
|
||||
|
||||
result := pluginSvc.runSchedulerIteration()
|
||||
if result {
|
||||
t.Fatalf("expected false when lock acquisition fails")
|
||||
}
|
||||
}
|
||||
|
||||
func TestSchedulerPhaseTracking(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
pluginSvc, err := New(Options{})
|
||||
if err != nil {
|
||||
t.Fatalf("New: %v", err)
|
||||
}
|
||||
defer pluginSvc.Shutdown()
|
||||
|
||||
// Initial phase should be idle.
|
||||
pluginSvc.schedulerMu.Lock()
|
||||
phase := pluginSvc.schedulerPhase
|
||||
pluginSvc.schedulerMu.Unlock()
|
||||
|
||||
if phase != "idle" {
|
||||
t.Fatalf("expected initial phase to be idle, got=%s", phase)
|
||||
}
|
||||
|
||||
pluginSvc.setSchedulerPhase("processing", "vacuum")
|
||||
|
||||
pluginSvc.schedulerMu.Lock()
|
||||
phase = pluginSvc.schedulerPhase
|
||||
jobType := pluginSvc.currentJobType
|
||||
pluginSvc.schedulerMu.Unlock()
|
||||
|
||||
if phase != "processing" {
|
||||
t.Fatalf("expected phase processing, got=%s", phase)
|
||||
}
|
||||
if jobType != "vacuum" {
|
||||
t.Fatalf("expected current job type vacuum, got=%s", jobType)
|
||||
}
|
||||
|
||||
pluginSvc.finishIteration(true)
|
||||
|
||||
pluginSvc.schedulerMu.Lock()
|
||||
phase = pluginSvc.schedulerPhase
|
||||
jobType = pluginSvc.currentJobType
|
||||
workDetected := pluginSvc.lastIterationWorkDetected
|
||||
lastEnded := pluginSvc.lastIterationEndedAt
|
||||
pluginSvc.schedulerMu.Unlock()
|
||||
|
||||
if phase != "idle" {
|
||||
t.Fatalf("expected phase idle after finish, got=%s", phase)
|
||||
}
|
||||
if jobType != "" {
|
||||
t.Fatalf("expected empty job type after finish, got=%s", jobType)
|
||||
}
|
||||
if !workDetected {
|
||||
t.Fatalf("expected last iteration work detected to be true")
|
||||
}
|
||||
if lastEnded.IsZero() {
|
||||
t.Fatalf("expected last iteration ended at to be set")
|
||||
}
|
||||
}
|
||||
|
||||
func TestGetSchedulerStatusIncludesIterationFields(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
pluginSvc, err := New(Options{
|
||||
IdleSleepDuration: 10 * time.Minute,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("New: %v", err)
|
||||
}
|
||||
defer pluginSvc.Shutdown()
|
||||
|
||||
now := time.Now().UTC()
|
||||
pluginSvc.schedulerMu.Lock()
|
||||
pluginSvc.schedulerPhase = "processing"
|
||||
pluginSvc.currentJobType = "vacuum"
|
||||
pluginSvc.iterationStartedAt = now.Add(-5 * time.Second)
|
||||
pluginSvc.lastIterationEndedAt = now.Add(-20 * time.Second)
|
||||
pluginSvc.lastIterationWorkDetected = true
|
||||
pluginSvc.schedulerMu.Unlock()
|
||||
|
||||
status := pluginSvc.GetSchedulerStatus()
|
||||
|
||||
if status.Phase != "processing" {
|
||||
t.Fatalf("expected phase processing, got=%s", status.Phase)
|
||||
}
|
||||
if status.CurrentJobType != "vacuum" {
|
||||
t.Fatalf("expected current job type vacuum, got=%s", status.CurrentJobType)
|
||||
}
|
||||
if status.IdleSleepSeconds != 600 {
|
||||
t.Fatalf("expected idle sleep 600s, got=%d", status.IdleSleepSeconds)
|
||||
}
|
||||
if status.IterationStartedAt == nil {
|
||||
t.Fatalf("expected iteration started at to be set")
|
||||
}
|
||||
if status.LastIterationEndedAt == nil {
|
||||
t.Fatalf("expected last iteration ended at to be set")
|
||||
}
|
||||
if !status.LastIterationWorkDetected {
|
||||
t.Fatalf("expected last iteration work detected to be true")
|
||||
}
|
||||
// SchedulerTickSeconds should match IdleSleepSeconds for backward compat.
|
||||
if status.SchedulerTickSeconds != 600 {
|
||||
t.Fatalf("expected scheduler tick seconds to match idle sleep, got=%d", status.SchedulerTickSeconds)
|
||||
}
|
||||
}
|
||||
|
||||
func TestGracefulShutdownDuringIteration(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
pluginSvc, err := New(Options{
|
||||
IdleSleepDuration: time.Millisecond,
|
||||
ClusterContextProvider: func(_ context.Context) (*plugin_pb.ClusterContext, error) {
|
||||
return &plugin_pb.ClusterContext{}, nil
|
||||
},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("New: %v", err)
|
||||
}
|
||||
|
||||
// Register an enabled job type so the scheduler loop has work to consider.
|
||||
err = pluginSvc.SaveJobTypeConfig(&plugin_pb.PersistedJobTypeConfig{
|
||||
JobType: "vacuum",
|
||||
AdminRuntime: &plugin_pb.AdminRuntimeConfig{
|
||||
Enabled: true,
|
||||
DetectionTimeoutSeconds: 5,
|
||||
MaxJobsPerDetection: 10,
|
||||
GlobalExecutionConcurrency: 1,
|
||||
PerWorkerExecutionConcurrency: 1,
|
||||
},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("SaveJobTypeConfig: %v", err)
|
||||
}
|
||||
pluginSvc.registry.UpsertFromHello(&plugin_pb.WorkerHello{
|
||||
WorkerId: "worker-a",
|
||||
Capabilities: []*plugin_pb.JobTypeCapability{
|
||||
{JobType: "vacuum", CanDetect: true, CanExecute: true},
|
||||
},
|
||||
})
|
||||
|
||||
// Shutdown while the scheduler loop is actively iterating.
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
pluginSvc.Shutdown()
|
||||
close(done)
|
||||
}()
|
||||
|
||||
select {
|
||||
case <-done:
|
||||
// Good — clean shutdown.
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatalf("shutdown did not complete in time")
|
||||
}
|
||||
}
|
||||
|
||||
func TestIdleSleepDurationDefault(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
pluginSvc, err := New(Options{})
|
||||
if err != nil {
|
||||
t.Fatalf("New: %v", err)
|
||||
}
|
||||
defer pluginSvc.Shutdown()
|
||||
|
||||
if pluginSvc.idleSleepDuration != defaultIdleSleepDuration {
|
||||
t.Fatalf("expected default idle sleep %v, got=%v", defaultIdleSleepDuration, pluginSvc.idleSleepDuration)
|
||||
}
|
||||
}
|
||||
|
||||
func TestIdleSleepDurationCustom(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
customDuration := 5 * time.Minute
|
||||
pluginSvc, err := New(Options{
|
||||
IdleSleepDuration: customDuration,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("New: %v", err)
|
||||
}
|
||||
defer pluginSvc.Shutdown()
|
||||
|
||||
if pluginSvc.idleSleepDuration != customDuration {
|
||||
t.Fatalf("expected custom idle sleep %v, got=%v", customDuration, pluginSvc.idleSleepDuration)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -7,11 +7,17 @@ import (
|
||||
)
|
||||
|
||||
type SchedulerStatus struct {
|
||||
Now time.Time `json:"now"`
|
||||
SchedulerTickSeconds int `json:"scheduler_tick_seconds"`
|
||||
Waiting []SchedulerWaitingStatus `json:"waiting,omitempty"`
|
||||
InProcessJobs []SchedulerJobStatus `json:"in_process_jobs,omitempty"`
|
||||
JobTypes []SchedulerJobTypeStatus `json:"job_types,omitempty"`
|
||||
Now time.Time `json:"now"`
|
||||
SchedulerTickSeconds int `json:"scheduler_tick_seconds"`
|
||||
IdleSleepSeconds int `json:"idle_sleep_seconds"`
|
||||
Phase string `json:"phase"`
|
||||
CurrentJobType string `json:"current_job_type,omitempty"`
|
||||
IterationStartedAt *time.Time `json:"iteration_started_at,omitempty"`
|
||||
LastIterationEndedAt *time.Time `json:"last_iteration_ended_at,omitempty"`
|
||||
LastIterationWorkDetected bool `json:"last_iteration_work_detected"`
|
||||
Waiting []SchedulerWaitingStatus `json:"waiting,omitempty"`
|
||||
InProcessJobs []SchedulerJobStatus `json:"in_process_jobs,omitempty"`
|
||||
JobTypes []SchedulerJobTypeStatus `json:"job_types,omitempty"`
|
||||
}
|
||||
|
||||
type SchedulerWaitingStatus struct {
|
||||
@@ -36,15 +42,14 @@ type SchedulerJobStatus struct {
|
||||
}
|
||||
|
||||
type SchedulerJobTypeStatus struct {
|
||||
JobType string `json:"job_type"`
|
||||
Enabled bool `json:"enabled"`
|
||||
DetectionInFlight bool `json:"detection_in_flight"`
|
||||
NextDetectionAt *time.Time `json:"next_detection_at,omitempty"`
|
||||
DetectionIntervalSeconds int32 `json:"detection_interval_seconds,omitempty"`
|
||||
LastDetectedAt *time.Time `json:"last_detected_at,omitempty"`
|
||||
LastDetectedCount int `json:"last_detected_count,omitempty"`
|
||||
LastDetectionError string `json:"last_detection_error,omitempty"`
|
||||
LastDetectionSkipped string `json:"last_detection_skipped,omitempty"`
|
||||
JobType string `json:"job_type"`
|
||||
Enabled bool `json:"enabled"`
|
||||
DetectionInFlight bool `json:"detection_in_flight"`
|
||||
MaxJobTypeDurationSeconds int32 `json:"max_job_type_duration_seconds,omitempty"`
|
||||
LastDetectedAt *time.Time `json:"last_detected_at,omitempty"`
|
||||
LastDetectedCount int `json:"last_detected_count,omitempty"`
|
||||
LastDetectionError string `json:"last_detection_error,omitempty"`
|
||||
LastDetectionSkipped string `json:"last_detection_skipped,omitempty"`
|
||||
}
|
||||
|
||||
type schedulerDetectionInfo struct {
|
||||
@@ -124,10 +129,29 @@ func (r *Plugin) snapshotSchedulerDetection(jobType string) schedulerDetectionIn
|
||||
|
||||
func (r *Plugin) GetSchedulerStatus() SchedulerStatus {
|
||||
now := time.Now().UTC()
|
||||
|
||||
r.schedulerMu.Lock()
|
||||
phase := r.schedulerPhase
|
||||
currentJobType := r.currentJobType
|
||||
iterationStartedAt := r.iterationStartedAt
|
||||
lastIterationEndedAt := r.lastIterationEndedAt
|
||||
lastIterationWorkDetected := r.lastIterationWorkDetected
|
||||
r.schedulerMu.Unlock()
|
||||
|
||||
status := SchedulerStatus{
|
||||
Now: now,
|
||||
SchedulerTickSeconds: int(secondsFromDuration(r.schedulerTick)),
|
||||
InProcessJobs: r.listInProcessJobs(now),
|
||||
Now: now,
|
||||
SchedulerTickSeconds: int(secondsFromDuration(r.idleSleepDuration)),
|
||||
IdleSleepSeconds: int(secondsFromDuration(r.idleSleepDuration)),
|
||||
Phase: phase,
|
||||
CurrentJobType: currentJobType,
|
||||
LastIterationWorkDetected: lastIterationWorkDetected,
|
||||
InProcessJobs: r.listInProcessJobs(now),
|
||||
}
|
||||
if !iterationStartedAt.IsZero() {
|
||||
status.IterationStartedAt = timeToPtr(iterationStartedAt)
|
||||
}
|
||||
if !lastIterationEndedAt.IsZero() {
|
||||
status.LastIterationEndedAt = timeToPtr(lastIterationEndedAt)
|
||||
}
|
||||
|
||||
states, err := r.ListSchedulerStates()
|
||||
@@ -143,11 +167,10 @@ func (r *Plugin) GetSchedulerStatus() SchedulerStatus {
|
||||
info := r.snapshotSchedulerDetection(jobType)
|
||||
|
||||
jobStatus := SchedulerJobTypeStatus{
|
||||
JobType: jobType,
|
||||
Enabled: state.Enabled,
|
||||
DetectionInFlight: state.DetectionInFlight,
|
||||
NextDetectionAt: state.NextDetectionAt,
|
||||
DetectionIntervalSeconds: state.DetectionIntervalSeconds,
|
||||
JobType: jobType,
|
||||
Enabled: state.Enabled,
|
||||
DetectionInFlight: state.DetectionInFlight,
|
||||
MaxJobTypeDurationSeconds: state.MaxJobTypeDurationSeconds,
|
||||
}
|
||||
if !info.lastDetectedAt.IsZero() {
|
||||
jobStatus.LastDetectedAt = timeToPtr(info.lastDetectedAt)
|
||||
@@ -163,18 +186,18 @@ func (r *Plugin) GetSchedulerStatus() SchedulerStatus {
|
||||
|
||||
if state.DetectionInFlight {
|
||||
waiting = append(waiting, SchedulerWaitingStatus{
|
||||
Reason: "detection_in_flight",
|
||||
Reason: "processing",
|
||||
JobType: jobType,
|
||||
})
|
||||
} else if state.Enabled && state.NextDetectionAt != nil && now.Before(*state.NextDetectionAt) {
|
||||
waiting = append(waiting, SchedulerWaitingStatus{
|
||||
Reason: "next_detection_at",
|
||||
JobType: jobType,
|
||||
Until: state.NextDetectionAt,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
if phase == "idle" && !lastIterationEndedAt.IsZero() {
|
||||
waiting = append(waiting, SchedulerWaitingStatus{
|
||||
Reason: "idle_sleep",
|
||||
})
|
||||
}
|
||||
|
||||
sort.Slice(jobTypes, func(i, j int) bool {
|
||||
return jobTypes[i].JobType < jobTypes[j].JobType
|
||||
})
|
||||
|
||||
+15
-16
@@ -82,22 +82,21 @@ type JobDetail struct {
|
||||
}
|
||||
|
||||
type SchedulerJobTypeState struct {
|
||||
JobType string `json:"job_type"`
|
||||
Enabled bool `json:"enabled"`
|
||||
PolicyError string `json:"policy_error,omitempty"`
|
||||
DetectionInFlight bool `json:"detection_in_flight"`
|
||||
NextDetectionAt *time.Time `json:"next_detection_at,omitempty"`
|
||||
DetectionIntervalSeconds int32 `json:"detection_interval_seconds,omitempty"`
|
||||
DetectionTimeoutSeconds int32 `json:"detection_timeout_seconds,omitempty"`
|
||||
ExecutionTimeoutSeconds int32 `json:"execution_timeout_seconds,omitempty"`
|
||||
MaxJobsPerDetection int32 `json:"max_jobs_per_detection,omitempty"`
|
||||
GlobalExecutionConcurrency int `json:"global_execution_concurrency,omitempty"`
|
||||
PerWorkerExecutionConcurrency int `json:"per_worker_execution_concurrency,omitempty"`
|
||||
RetryLimit int `json:"retry_limit,omitempty"`
|
||||
RetryBackoffSeconds int32 `json:"retry_backoff_seconds,omitempty"`
|
||||
DetectorAvailable bool `json:"detector_available"`
|
||||
DetectorWorkerID string `json:"detector_worker_id,omitempty"`
|
||||
ExecutorWorkerCount int `json:"executor_worker_count"`
|
||||
JobType string `json:"job_type"`
|
||||
Enabled bool `json:"enabled"`
|
||||
PolicyError string `json:"policy_error,omitempty"`
|
||||
DetectionInFlight bool `json:"detection_in_flight"`
|
||||
DetectionTimeoutSeconds int32 `json:"detection_timeout_seconds,omitempty"`
|
||||
ExecutionTimeoutSeconds int32 `json:"execution_timeout_seconds,omitempty"`
|
||||
MaxJobTypeDurationSeconds int32 `json:"max_job_type_duration_seconds,omitempty"`
|
||||
MaxJobsPerDetection int32 `json:"max_jobs_per_detection,omitempty"`
|
||||
GlobalExecutionConcurrency int `json:"global_execution_concurrency,omitempty"`
|
||||
PerWorkerExecutionConcurrency int `json:"per_worker_execution_concurrency,omitempty"`
|
||||
RetryLimit int `json:"retry_limit,omitempty"`
|
||||
RetryBackoffSeconds int32 `json:"retry_backoff_seconds,omitempty"`
|
||||
DetectorAvailable bool `json:"detector_available"`
|
||||
DetectorWorkerID string `json:"detector_worker_id,omitempty"`
|
||||
ExecutorWorkerCount int `json:"executor_worker_count"`
|
||||
}
|
||||
|
||||
func timeToPtr(t time.Time) *time.Time {
|
||||
|
||||
@@ -115,12 +115,42 @@ templ Plugin(page string) {
|
||||
</div>
|
||||
</div>
|
||||
</div>
|
||||
<div class="row mb-4">
|
||||
<div class="col-12">
|
||||
<div class="card shadow-sm">
|
||||
<div class="card-header d-flex justify-content-between align-items-center flex-wrap gap-2">
|
||||
<h5 class="mb-0"><i class="fas fa-sync-alt me-2"></i>Scheduler Iteration</h5>
|
||||
<small class="text-muted">Current iteration phase and timing</small>
|
||||
</div>
|
||||
<div class="card-body" id="plugin-scheduler-iteration-card">
|
||||
<div class="row g-3">
|
||||
<div class="col-md-3 col-sm-6">
|
||||
<div class="small text-muted">Phase</div>
|
||||
<div class="fw-semibold" id="plugin-scheduler-phase">-</div>
|
||||
</div>
|
||||
<div class="col-md-3 col-sm-6">
|
||||
<div class="small text-muted">Current Job Type</div>
|
||||
<div class="fw-semibold" id="plugin-scheduler-current-jt">-</div>
|
||||
</div>
|
||||
<div class="col-md-3 col-sm-6">
|
||||
<div class="small text-muted">Idle Sleep</div>
|
||||
<div class="fw-semibold" id="plugin-scheduler-idle-sleep">-</div>
|
||||
</div>
|
||||
<div class="col-md-3 col-sm-6">
|
||||
<div class="small text-muted">Last Iteration</div>
|
||||
<div class="fw-semibold" id="plugin-scheduler-last-iteration">-</div>
|
||||
</div>
|
||||
</div>
|
||||
</div>
|
||||
</div>
|
||||
</div>
|
||||
</div>
|
||||
<div class="row mb-4">
|
||||
<div class="col-12">
|
||||
<div class="card shadow-sm">
|
||||
<div class="card-header d-flex justify-content-between align-items-center flex-wrap gap-2">
|
||||
<h5 class="mb-0"><i class="fas fa-clock me-2"></i>Scheduler State</h5>
|
||||
<small class="text-muted">Per job type detection schedule and execution limits</small>
|
||||
<small class="text-muted">Per job type execution limits and status</small>
|
||||
</div>
|
||||
<div class="card-body p-0">
|
||||
<div class="table-responsive">
|
||||
@@ -130,9 +160,8 @@ templ Plugin(page string) {
|
||||
<th>Job Type</th>
|
||||
<th>Enabled</th>
|
||||
<th>Detector</th>
|
||||
<th>In Flight</th>
|
||||
<th>Next Detection</th>
|
||||
<th>Interval</th>
|
||||
<th>Status</th>
|
||||
<th>Max Duration</th>
|
||||
<th>Exec Global</th>
|
||||
<th>Exec/Worker</th>
|
||||
<th>Executor Workers</th>
|
||||
@@ -140,7 +169,7 @@ templ Plugin(page string) {
|
||||
</tr>
|
||||
</thead>
|
||||
<tbody id="plugin-scheduler-table-body">
|
||||
<tr><td colspan="10" class="text-muted text-center py-3">Loading...</td></tr>
|
||||
<tr><td colspan="9" class="text-muted text-center py-3">Loading...</td></tr>
|
||||
</tbody>
|
||||
</table>
|
||||
</div>
|
||||
@@ -1432,7 +1461,7 @@ templ Plugin(page string) {
|
||||
|
||||
var states = Array.isArray(state.schedulerStates) ? state.schedulerStates : [];
|
||||
if (!states.length) {
|
||||
tbody.innerHTML = '<tr><td colspan="10" class="text-muted text-center py-3">No scheduler state available</td></tr>';
|
||||
tbody.innerHTML = '<tr><td colspan="9" class="text-muted text-center py-3">No scheduler state available</td></tr>';
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -1442,8 +1471,8 @@ templ Plugin(page string) {
|
||||
var enabled = !!item.enabled;
|
||||
var inFlight = !!item.detection_in_flight;
|
||||
var detector = item.detector_available ? textOrDash(item.detector_worker_id) : 'No detector';
|
||||
var intervalSeconds = Number(item.detection_interval_seconds || 0);
|
||||
var intervalText = intervalSeconds > 0 ? (String(intervalSeconds) + 's') : '-';
|
||||
var maxDurationSeconds = Number(item.max_job_type_duration_seconds || 0);
|
||||
var maxDurationText = maxDurationSeconds > 0 ? (String(Math.round(maxDurationSeconds / 60)) + 'm') : '-';
|
||||
var globalExec = Number(item.global_execution_concurrency || 0);
|
||||
var perWorkerExec = Number(item.per_worker_execution_concurrency || 0);
|
||||
var executorWorkers = Number(item.executor_worker_count || 0);
|
||||
@@ -1454,7 +1483,7 @@ templ Plugin(page string) {
|
||||
var effectiveExecText = enabled ? String(effectiveExec) : '-';
|
||||
|
||||
var enabledBadge = enabled ? '<span class="badge bg-success">Enabled</span>' : '<span class="badge bg-secondary">Disabled</span>';
|
||||
var inFlightBadge = inFlight ? '<span class="badge bg-warning text-dark">Yes</span>' : '<span class="badge bg-light text-dark">No</span>';
|
||||
var statusBadge = inFlight ? '<span class="badge bg-warning text-dark">Processing</span>' : '<span class="badge bg-light text-dark">Idle</span>';
|
||||
var policyErrorHtml = '';
|
||||
if (item.policy_error) {
|
||||
policyErrorHtml = '<div><small class="text-danger">' + escapeHtml(String(item.policy_error)) + '</small></div>';
|
||||
@@ -1464,9 +1493,8 @@ templ Plugin(page string) {
|
||||
'<td>' + escapeHtml(textOrDash(item.job_type)) + policyErrorHtml + '</td>' +
|
||||
'<td>' + enabledBadge + '</td>' +
|
||||
'<td><small>' + escapeHtml(detector) + '</small></td>' +
|
||||
'<td>' + inFlightBadge + '</td>' +
|
||||
'<td><small>' + escapeHtml(parseTime(item.next_detection_at) || '-') + '</small></td>' +
|
||||
'<td><small>' + escapeHtml(intervalText) + '</small></td>' +
|
||||
'<td>' + statusBadge + '</td>' +
|
||||
'<td><small>' + escapeHtml(maxDurationText) + '</small></td>' +
|
||||
'<td><small>' + escapeHtml(globalExecText) + '</small></td>' +
|
||||
'<td><small>' + escapeHtml(perWorkerExecText) + '</small></td>' +
|
||||
'<td><small>' + escapeHtml(executorWorkersText) + '</small></td>' +
|
||||
@@ -2780,6 +2808,37 @@ templ Plugin(page string) {
|
||||
}
|
||||
}
|
||||
|
||||
function renderSchedulerIterationStatus(schedulerStatus) {
|
||||
var phaseEl = document.getElementById('plugin-scheduler-phase');
|
||||
var currentJtEl = document.getElementById('plugin-scheduler-current-jt');
|
||||
var idleSleepEl = document.getElementById('plugin-scheduler-idle-sleep');
|
||||
var lastIterEl = document.getElementById('plugin-scheduler-last-iteration');
|
||||
if (!phaseEl) {
|
||||
return;
|
||||
}
|
||||
|
||||
var sched = (schedulerStatus && schedulerStatus.scheduler) ? schedulerStatus.scheduler : {};
|
||||
var phase = String(sched.phase || 'unknown');
|
||||
var phaseColors = { idle: 'bg-secondary', acquiring_lock: 'bg-info', processing: 'bg-primary' };
|
||||
var phaseColor = phaseColors[phase] || 'bg-dark';
|
||||
phaseEl.innerHTML = '<span class="badge ' + phaseColor + '">' + escapeHtml(phase) + '</span>';
|
||||
currentJtEl.textContent = sched.current_job_type || '-';
|
||||
|
||||
var idleSleep = Number(sched.idle_sleep_seconds || 0);
|
||||
idleSleepEl.textContent = idleSleep > 0 ? (String(Math.round(idleSleep / 60)) + ' min') : '-';
|
||||
|
||||
var lastEnded = parseTime(sched.last_iteration_ended_at);
|
||||
var workDetected = !!sched.last_iteration_work_detected;
|
||||
if (lastEnded) {
|
||||
var workBadge = workDetected
|
||||
? '<span class="badge bg-success ms-1">work found</span>'
|
||||
: '<span class="badge bg-light text-dark ms-1">no work</span>';
|
||||
lastIterEl.innerHTML = '<small>' + escapeHtml(lastEnded) + '</small>' + workBadge;
|
||||
} else {
|
||||
lastIterEl.textContent = '-';
|
||||
}
|
||||
}
|
||||
|
||||
async function refreshJobsAndActivities() {
|
||||
var executionStateFilter = document.getElementById('plugin-monitor-job-state-filter');
|
||||
|
||||
@@ -2788,10 +2847,18 @@ templ Plugin(page string) {
|
||||
var allJobsPromise = pluginRequest('GET', '/api/plugin/jobs?limit=500');
|
||||
var allActivitiesPromise = pluginRequest('GET', '/api/plugin/activities?limit=500');
|
||||
var schedulerPromise = pluginRequest('GET', '/api/plugin/scheduler-states');
|
||||
var schedulerStatusPromise = pluginRequest('GET', '/api/plugin/scheduler-status').catch(function(e) { return e; });
|
||||
|
||||
var allJobs = await allJobsPromise;
|
||||
var allActivities = await allActivitiesPromise;
|
||||
var schedulerStates = await schedulerPromise;
|
||||
var schedulerStatusResult = await schedulerStatusPromise;
|
||||
var schedulerStatus = null;
|
||||
if (schedulerStatusResult instanceof Error) {
|
||||
console.error('Failed to fetch scheduler status:', schedulerStatusResult);
|
||||
} else {
|
||||
schedulerStatus = schedulerStatusResult;
|
||||
}
|
||||
|
||||
state.jobs = Array.isArray(allJobs) ? allJobs : [];
|
||||
state.activities = Array.isArray(allActivities) ? allActivities : [];
|
||||
@@ -2803,6 +2870,7 @@ templ Plugin(page string) {
|
||||
renderExecutionJobs();
|
||||
renderExecutionActivities();
|
||||
renderSchedulerStates();
|
||||
renderSchedulerIterationStatus(schedulerStatus);
|
||||
renderStatus();
|
||||
renderJobTypeSummary();
|
||||
}
|
||||
|
||||
@@ -45,6 +45,7 @@ type AdminOptions struct {
|
||||
readOnlyPassword *string
|
||||
dataDir *string
|
||||
icebergPort *int
|
||||
idleSleepSeconds *int
|
||||
}
|
||||
|
||||
func init() {
|
||||
@@ -60,6 +61,7 @@ func init() {
|
||||
a.readOnlyUser = cmdAdmin.Flag.String("readOnlyUser", "", "read-only user username (optional, for view-only access)")
|
||||
a.readOnlyPassword = cmdAdmin.Flag.String("readOnlyPassword", "", "read-only user password (optional, for view-only access; requires adminPassword to be set)")
|
||||
a.icebergPort = cmdAdmin.Flag.Int("iceberg.port", 8181, "Iceberg REST Catalog port (0 to hide in UI)")
|
||||
a.idleSleepSeconds = cmdAdmin.Flag.Int("scheduler.idleSleep", 0, "scheduler idle sleep in seconds between iterations when no work is found (0 = default 17 minutes)")
|
||||
}
|
||||
|
||||
var cmdAdmin = &Command{
|
||||
@@ -290,7 +292,11 @@ func startAdminServer(ctx context.Context, options AdminOptions, enableUI bool,
|
||||
}
|
||||
|
||||
// Create admin server (plugin is always enabled)
|
||||
adminServer := dash.NewAdminServer(*options.master, nil, dataDir, icebergPort)
|
||||
var idleSleep time.Duration
|
||||
if options.idleSleepSeconds != nil && *options.idleSleepSeconds > 0 {
|
||||
idleSleep = time.Duration(*options.idleSleepSeconds) * time.Second
|
||||
}
|
||||
adminServer := dash.NewAdminServer(*options.master, nil, dataDir, icebergPort, idleSleep)
|
||||
|
||||
// Show discovered filers
|
||||
filers := adminServer.GetAllFilers()
|
||||
|
||||
Reference in New Issue
Block a user