Operator actions in the report reply (D1, R-314/R-279/R-177, decision 185)
The hub may ask the running controller for a CLOSED list of actions, carried in the report ACK (operator_actions) and answered on the next report (operator_action_results): offsite_backup_now, abandon_stop, abandon_extend (1-30 days), run_job (fill-watch, offsite-integrity, offsite-proof, disk-health-check). Unknown action/job/argument -> refused, nothing called. Once per id (in memory; every action is safe to repeat). - internal/report/opactions.go: the executor; results re-sent until the hub stops listing the id. - scheduler.RunNow: refuses unknown / already-running jobs; OnDemand(ctx) makes an operator's offsite-integrity run even when not due. - ExtendAbandon never shortens the countdown and refuses in the hub phase; StopAbandon reports failure when the hub cancel failed (it said success). - Report ACK read cap 4 KiB -> 64 KiB (an ACK over the cap dropped every field in it). Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_0159rPz1ZhFKsS53msqPYxtS
This commit is contained in:
@@ -2,6 +2,7 @@ package scheduler
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log"
|
||||
"sync"
|
||||
@@ -359,7 +360,74 @@ func (s *Scheduler) executeJob(job *Job, quiet bool) {
|
||||
}
|
||||
job.Running = true
|
||||
s.mu.Unlock()
|
||||
s.runClaimed(s.ctx, job, quiet)
|
||||
}
|
||||
|
||||
// ErrUnknownJob / ErrJobRunning / ErrNotStarted are RunNow's refusals — values, so a caller can
|
||||
// branch on them (R-224's lesson: a string is not something a caller can branch on).
|
||||
var (
|
||||
ErrUnknownJob = errors.New("no job with this name is registered on this box")
|
||||
ErrJobRunning = errors.New("the job is already running")
|
||||
ErrNotStarted = errors.New("the scheduler is not running")
|
||||
)
|
||||
|
||||
// onDemandKey marks the context of a RunNow execution (OnDemand reads it).
|
||||
type onDemandKey struct{}
|
||||
|
||||
// OnDemand reports whether the job is running because someone asked for it now (RunNow), not on
|
||||
// its schedule. A job whose scheduled run checks due-ness first (the off-site integrity check) uses
|
||||
// it to run anyway — the operator asked for a check, not for "is a check due".
|
||||
func OnDemand(ctx context.Context) bool {
|
||||
v, _ := ctx.Value(onDemandKey{}).(bool)
|
||||
return v
|
||||
}
|
||||
|
||||
// RunNow runs a registered job once, now, outside its schedule (R-314/R-177, `09` §3 decision 185:
|
||||
// the operator's run_job action). It REFUSES an unknown name and a job that is already running —
|
||||
// never queues, never runs it twice at once (the same Running flag the schedule uses, claimed under
|
||||
// the mutex, so a scheduled tick and RunNow cannot both start it). On success the job runs in its own
|
||||
// goroutine and its error (nil = completed) arrives on the returned channel, which is then closed.
|
||||
func (s *Scheduler) RunNow(name string) (<-chan error, error) {
|
||||
s.mu.Lock()
|
||||
if !s.started || s.ctx == nil {
|
||||
s.mu.Unlock()
|
||||
return nil, ErrNotStarted
|
||||
}
|
||||
var job *Job
|
||||
for _, j := range s.jobs {
|
||||
if j.Name == name {
|
||||
job = j
|
||||
break
|
||||
}
|
||||
}
|
||||
if job == nil {
|
||||
s.mu.Unlock()
|
||||
s.dbg("RunNow %q refused: unknown job", name)
|
||||
return nil, ErrUnknownJob
|
||||
}
|
||||
if job.Running {
|
||||
s.mu.Unlock()
|
||||
s.dbg("RunNow %q refused: already running", name)
|
||||
return nil, ErrJobRunning
|
||||
}
|
||||
job.Running = true
|
||||
ctx := context.WithValue(s.ctx, onDemandKey{}, true)
|
||||
s.wg.Add(1)
|
||||
s.mu.Unlock()
|
||||
|
||||
s.logger.Printf("[INFO] [scheduler] Job %s started on demand", name)
|
||||
done := make(chan error, 1)
|
||||
go func() {
|
||||
defer s.wg.Done()
|
||||
defer close(done)
|
||||
done <- s.runClaimed(ctx, job, false)
|
||||
}()
|
||||
return done, nil
|
||||
}
|
||||
|
||||
// runClaimed executes a job whose Running flag the caller has already set. It clears the flag,
|
||||
// records LastRun/LastErr, and returns the job's error.
|
||||
func (s *Scheduler) runClaimed(ctx context.Context, job *Job, quiet bool) (err error) {
|
||||
defer func() {
|
||||
s.mu.Lock()
|
||||
job.Running = false
|
||||
@@ -369,8 +437,9 @@ func (s *Scheduler) executeJob(job *Job, quiet bool) {
|
||||
// Panic recovery
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
err = fmt.Errorf("panic: %v", r)
|
||||
s.mu.Lock()
|
||||
job.LastErr = fmt.Errorf("panic: %v", r)
|
||||
job.LastErr = err
|
||||
job.LastRun = time.Now()
|
||||
s.mu.Unlock()
|
||||
s.logger.Printf("[ERROR] [scheduler] Job %s panicked: %v", job.Name, r)
|
||||
@@ -383,7 +452,7 @@ func (s *Scheduler) executeJob(job *Job, quiet bool) {
|
||||
s.dbg("job %s: execution starting", job.Name)
|
||||
|
||||
start := time.Now()
|
||||
err := job.Fn(s.ctx)
|
||||
err = job.Fn(ctx)
|
||||
elapsed := time.Since(start)
|
||||
|
||||
s.mu.Lock()
|
||||
@@ -400,6 +469,7 @@ func (s *Scheduler) executeJob(job *Job, quiet bool) {
|
||||
// Routine per-cycle timing line → TRACE (dropped from the ring; the completed/failed lines above
|
||||
// carry the outcome). This was the biggest ring filler under load (fix-6).
|
||||
s.trace("job %s: finished in %s (err=%v)", job.Name, elapsed.Round(time.Millisecond), err)
|
||||
return err
|
||||
}
|
||||
|
||||
// parseDailyTime parses "HH:MM" and returns hour and minute.
|
||||
|
||||
Reference in New Issue
Block a user