Compare commits

...
Author SHA1 Message Date
Chris LuandClaude Opus 4.6 b1c862e7e3 fix: prevent unhandled promise rejection for scheduler status fetch
Attach .catch at promise creation to mark it as handled by the runtime,
then check the result type after await. This prevents an
unhandledrejection event if the fetch rejects before the await.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-03 19:27:17 -08:00
Chris LuandClaude Opus 4.6 05902c7101 fix: use lifecycle context for job-type budget and retry backoff
- Use r.ctx (cancelled on shutdown) as parent for per-job-type budget
  context instead of context.Background(), so in-flight work is
  cancelled promptly on plugin shutdown.
- Replace waitForShutdownOrTimer with waitForShutdownOrCtx in retry
  backoff so retries respect the per-job-type budget context.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-03 19:27:06 -08:00
Chris LuandClaude Opus 4.6 e0cccae794 fix: add -scheduler.idleSleep flag to prevent EC integration test hang
The new 17-minute default idle sleep caused TestEcEndToEnd to hang
because the scheduler would not re-check for work frequently enough
after the first iteration found nothing.

Add -scheduler.idleSleep CLI flag (in seconds) to configure the idle
sleep duration. The EC integration test now passes -scheduler.idleSleep=2
so detection runs every 2 seconds when idle.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-03 19:05:56 -08:00
Chris LuandClaude Opus 4.6 6289beb8f5 fix: add ClusterContextProvider to shutdown test and handle status fetch errors
Address PR review nitpicks:
- Add ClusterContextProvider to TestGracefulShutdownDuringIteration so
  the scheduler loop actually starts (New() requires it).
- Wrap schedulerStatusPromise await in try/catch in plugin.templ so a
  failed status fetch does not break refreshJobsAndActivities rendering.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-03 18:49:37 -08:00
Chris LuandClaude Opus 4.6 35786d72e0 fix: unconditional lock release and context-aware executor reservation
Address PR review findings:
- Make defer releaseLock() unconditional in runSchedulerIteration since
  acquireAdminLock guarantees a non-nil release func on success.
- Add ctx parameter to reserveScheduledExecutor so it honors the
  per-job-type budget context, preventing waits from exceeding the
  30-minute budget.
- Add waitForShutdownOrCtx helper that also selects on ctx.Done().
- Update all reserveScheduledExecutor call sites in tests.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-03 18:49:28 -08:00
Chris LuandClaude Opus 4.6 51d715b4e6 plugin: update and add tests for sequential iteration model
Update TestListSchedulerStatesIncludesPolicyAndState to use
currentJobType for DetectionInFlight and assert MaxJobTypeDurationSeconds.

Add new tests:
- TestIterationAcquiresLockOnce: verify single lock acquire per iteration
- TestIterationReturnsfalseWhenLockFails: lock failure returns no-work
- TestSchedulerPhaseTracking: phase transitions and finishIteration
- TestGetSchedulerStatusIncludesIterationFields: status field population
- TestGracefulShutdownDuringIteration: clean shutdown mid-iteration
- TestIdleSleepDurationDefault: default 17m idle sleep
- TestIdleSleepDurationCustom: custom idle sleep from Options

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-03 16:42:48 -08:00
Chris LuandClaude Opus 4.6 4f97f7d413 plugin: add scheduler iteration status card to UI
Add a Scheduler Iteration card showing the current phase, active job
type, idle sleep duration, and last iteration result.

Update the Scheduler State table to replace In Flight / Next Detection
/ Interval columns with Status (Processing/Idle) and Max Duration
columns to match the new iteration model.

Add renderSchedulerIterationStatus JS function and update
refreshJobsAndActivities to fetch and render iteration status.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-03 16:42:26 -08:00
Chris LuandClaude Opus 4.6 c8ad7b528d plugin: update scheduler status to expose iteration state
Add Phase, CurrentJobType, IdleSleepSeconds, IterationStartedAt,
LastIterationEndedAt, and LastIterationWorkDetected to SchedulerStatus.

Set SchedulerTickSeconds to match IdleSleepSeconds for backward
compatibility.

Add MaxJobTypeDurationSeconds to SchedulerJobTypeStatus and remove
NextDetectionAt and DetectionIntervalSeconds.

Update GetSchedulerStatus to read iteration tracking fields and
adapt waiting status logic for the new model.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-03 16:42:10 -08:00
Chris LuandClaude Opus 4.6 41f22f5c92 plugin: rewrite scheduler to sequential lock-based iteration model
Replace the ticker-based parallel goroutine scheduler with a
single-lock sequential-processing iteration model:

- schedulerLoop now calls runSchedulerIteration in a loop, sleeping
  idleSleepDuration (default 17m) when no work is found, or looping
  immediately when work is detected.

- runSchedulerIteration acquires the admin lock once, loads cluster
  context once, then processes all enabled job types sequentially.

- runJobTypeIteration runs detection then execution for one job type
  within a per-job-type budget (default 30m) enforced via context.

- Remove markDetectionDue, finishDetection, pruneSchedulerState,
  runSchedulerTick, and runScheduledDetection which are no longer
  needed.

- dispatchScheduledProposals and executeScheduledJobWithExecutor now
  accept a parent context for budget enforcement.

- Remove DetectionInterval from schedulerPolicy, add MaxJobTypeDuration.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-03 16:41:59 -08:00
Chris LuandClaude Opus 4.6 05138c6031 plugin: add iteration tracking fields to Plugin struct
Replace per-job-type detection maps (nextDetectionAt, detectionInFlight)
with iteration-level tracking fields: schedulerPhase, currentJobType,
iterationStartedAt, lastIterationEndedAt, lastIterationWorkDetected.

Add IdleSleepDuration to Options and idleSleepDuration to Plugin.
Add MaxJobTypeDurationSeconds to SchedulerJobTypeState and remove
NextDetectionAt and DetectionIntervalSeconds which are no longer
needed in the sequential iteration model.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-03 16:40:50 -08:00
9 changed files with 769 additions and 285 deletions
@@ -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")
+5 -3
View File
@@ -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
},
+14 -5
View File
@@ -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),
+250 -200
View File
@@ -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.
+345 -18
View File
@@ -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)
}
}
+52 -29
View File
@@ -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
View File
@@ -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 {
+80 -12
View File
@@ -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();
}
+7 -1
View File
@@ -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()