Add pause and cancel to a running migration
A live run could only be stopped one account at a time, and stopping it at all meant losing the queue: the accounts that had not started yet stayed idle with no record that they were meant to run. A run now carries a handle holding the context that stops every account under it plus the reason it was stopped. Pause and cancel take the same path and differ only in the status left behind — paused accounts are what Resume re-runs, and the migration journal makes each one continue where it stopped instead of re-copying. Accounts still queued when the stop lands get the same status as the interrupted ones, so the whole remainder is resumable after a pause and cancelled after a cancel. Database writes keep using the uncancellable context, so statuses and counters survive the stop. The scheduler skips paused tasks: auto-starting a full run would defeat the pause. An operator stopping a run no longer trips the schedule breaker either — that is for failures, not for intent. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -31,6 +31,9 @@ func (s *Server) Router() http.Handler {
|
||||
api.HandleFunc("POST /api/tasks/{id}/import", s.handleImportCSV)
|
||||
api.HandleFunc("POST /api/tasks/{id}/test", s.handleTestAccounts)
|
||||
api.HandleFunc("POST /api/tasks/{id}/run", s.handleRun)
|
||||
api.HandleFunc("POST /api/tasks/{id}/pause", s.handlePauseRun)
|
||||
api.HandleFunc("POST /api/tasks/{id}/cancel", s.handleCancelRun)
|
||||
api.HandleFunc("POST /api/tasks/{id}/resume", s.handleResumeRun)
|
||||
api.HandleFunc("POST /api/tasks/{id}/accounts/{accountId}/cancel", s.handleCancelAccount)
|
||||
api.HandleFunc("POST /api/tasks/{id}/accounts/{accountId}/probe", s.handleProbeAccountFolders)
|
||||
api.HandleFunc("PUT /api/tasks/{id}/accounts/{accountId}/folder-mapping", s.handleSetAccountFolderMapping)
|
||||
|
||||
@@ -127,6 +127,56 @@ func (s *Server) handleRun(w http.ResponseWriter, r *http.Request) {
|
||||
writeJSON(w, http.StatusAccepted, map[string]int64{"run_id": runID})
|
||||
}
|
||||
|
||||
// handlePauseRun stops the live run but keeps the unfinished accounts
|
||||
// resumable; handleCancelRun stops it for good. Both are no-ops (409) when the
|
||||
// task has no run in flight.
|
||||
func (s *Server) handlePauseRun(w http.ResponseWriter, r *http.Request) {
|
||||
s.stopRun(w, r, s.orch.PauseTask)
|
||||
}
|
||||
|
||||
func (s *Server) handleCancelRun(w http.ResponseWriter, r *http.Request) {
|
||||
s.stopRun(w, r, s.orch.CancelTask)
|
||||
}
|
||||
|
||||
func (s *Server) stopRun(w http.ResponseWriter, r *http.Request, stop func(int64) bool) {
|
||||
taskID, err := pathID(r, "id")
|
||||
if err != nil {
|
||||
http.Error(w, "bad id", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
if !stop(taskID) {
|
||||
http.Error(w, "task is not running", http.StatusConflict)
|
||||
return
|
||||
}
|
||||
w.WriteHeader(http.StatusAccepted)
|
||||
}
|
||||
|
||||
// handleResumeRun restarts a paused task with the accounts its pause left
|
||||
// unfinished; already-copied messages are skipped by the migration journal.
|
||||
func (s *Server) handleResumeRun(w http.ResponseWriter, r *http.Request) {
|
||||
taskID, err := pathID(r, "id")
|
||||
if err != nil {
|
||||
http.Error(w, "bad id", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
runID, err := s.orch.ResumeTask(r.Context(), taskID)
|
||||
switch {
|
||||
case errors.Is(err, orchestrator.ErrNothingToResume):
|
||||
http.Error(w, "no paused accounts to resume", http.StatusConflict)
|
||||
return
|
||||
case errors.Is(err, orchestrator.ErrNotTested):
|
||||
http.Error(w, "accounts must pass connection tests first", http.StatusConflict)
|
||||
return
|
||||
case errors.Is(err, orchestrator.ErrAlreadyRunning):
|
||||
http.Error(w, "task is already running", http.StatusConflict)
|
||||
return
|
||||
case err != nil:
|
||||
http.Error(w, err.Error(), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
writeJSON(w, http.StatusAccepted, map[string]int64{"run_id": runID})
|
||||
}
|
||||
|
||||
func (s *Server) handleCancelAccount(w http.ResponseWriter, r *http.Request) {
|
||||
taskID, err := pathID(r, "id")
|
||||
if err != nil {
|
||||
|
||||
@@ -17,6 +17,7 @@ import (
|
||||
var ErrNotTested = errors.New("accounts not fully tested")
|
||||
var ErrAlreadyRunning = errors.New("task already running")
|
||||
var ErrNoAccountsSelected = errors.New("no matching accounts selected")
|
||||
var ErrNothingToResume = errors.New("no paused accounts to resume")
|
||||
|
||||
// maxAccountErrors caps how many individual error rows one account records per
|
||||
// run, so a corrupt mailbox producing thousands of failures can't bloat the
|
||||
@@ -62,6 +63,38 @@ func planFolders(folders []string, mapping map[string]string, excluded []string)
|
||||
return plan
|
||||
}
|
||||
|
||||
// A run stops either on its own or because the operator intervened. Pausing and
|
||||
// cancelling take the same path — stop the in-flight work — and differ only in
|
||||
// the status left behind: paused accounts are what Resume picks up again.
|
||||
type stopReason int32
|
||||
|
||||
const (
|
||||
stopNone stopReason = iota
|
||||
stopPaused
|
||||
stopCancelled
|
||||
)
|
||||
|
||||
func (r stopReason) accountStatus() string {
|
||||
if r == stopPaused {
|
||||
return "paused"
|
||||
}
|
||||
return "cancelled"
|
||||
}
|
||||
|
||||
// runHandle is the live state of one task's run: the cancel that stops every
|
||||
// account under it, plus why it was stopped.
|
||||
type runHandle struct {
|
||||
cancel context.CancelFunc
|
||||
reason atomic.Int32
|
||||
}
|
||||
|
||||
func (h *runHandle) stopWith(r stopReason) {
|
||||
h.reason.CompareAndSwap(int32(stopNone), int32(r))
|
||||
h.cancel()
|
||||
}
|
||||
|
||||
func (h *runHandle) stopReason() stopReason { return stopReason(h.reason.Load()) }
|
||||
|
||||
type Orchestrator struct {
|
||||
store *store.Store
|
||||
hub *wshub.Hub
|
||||
@@ -70,10 +103,66 @@ type Orchestrator struct {
|
||||
|
||||
mu sync.Mutex
|
||||
cancels map[int64]context.CancelFunc // account_id -> cancel of its in-flight copy
|
||||
runs map[int64]*runHandle // task_id -> live run
|
||||
}
|
||||
|
||||
func New(s *store.Store, hub *wshub.Hub, encKey []byte, concurrency int) *Orchestrator {
|
||||
return &Orchestrator{store: s, hub: hub, encKey: encKey, concurrency: concurrency, cancels: map[int64]context.CancelFunc{}}
|
||||
return &Orchestrator{
|
||||
store: s, hub: hub, encKey: encKey, concurrency: concurrency,
|
||||
cancels: map[int64]context.CancelFunc{},
|
||||
runs: map[int64]*runHandle{},
|
||||
}
|
||||
}
|
||||
|
||||
// PauseTask stops the task's live run, leaving every unfinished account
|
||||
// "paused" so ResumeTask can pick them up. Returns false if nothing is running.
|
||||
func (o *Orchestrator) PauseTask(taskID int64) bool { return o.stopRun(taskID, stopPaused) }
|
||||
|
||||
// CancelTask stops the task's live run and marks every unfinished account
|
||||
// "cancelled". Returns false if nothing is running.
|
||||
func (o *Orchestrator) CancelTask(taskID int64) bool { return o.stopRun(taskID, stopCancelled) }
|
||||
|
||||
func (o *Orchestrator) stopRun(taskID int64, reason stopReason) bool {
|
||||
o.mu.Lock()
|
||||
h, ok := o.runs[taskID]
|
||||
o.mu.Unlock()
|
||||
if !ok {
|
||||
return false
|
||||
}
|
||||
h.stopWith(reason)
|
||||
return true
|
||||
}
|
||||
|
||||
func (o *Orchestrator) registerRun(taskID int64, h *runHandle) {
|
||||
o.mu.Lock()
|
||||
o.runs[taskID] = h
|
||||
o.mu.Unlock()
|
||||
}
|
||||
|
||||
func (o *Orchestrator) unregisterRun(taskID int64) {
|
||||
o.mu.Lock()
|
||||
delete(o.runs, taskID)
|
||||
o.mu.Unlock()
|
||||
}
|
||||
|
||||
// ResumeTask restarts a paused task with exactly the accounts the pause left
|
||||
// unfinished. Everything already copied is skipped by the migration journal, so
|
||||
// each account continues where it stopped.
|
||||
func (o *Orchestrator) ResumeTask(ctx context.Context, taskID int64) (int64, error) {
|
||||
accs, err := o.store.ListAccountsByTask(ctx, taskID)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
ids := make([]int64, 0, len(accs))
|
||||
for _, a := range accs {
|
||||
if a.Status == "paused" {
|
||||
ids = append(ids, a.ID)
|
||||
}
|
||||
}
|
||||
if len(ids) == 0 {
|
||||
return 0, ErrNothingToResume
|
||||
}
|
||||
return o.Run(ctx, taskID, "manual", ids)
|
||||
}
|
||||
|
||||
// CancelAccount aborts the in-flight copy for one account, if it is running.
|
||||
@@ -231,11 +320,21 @@ func (o *Orchestrator) Run(ctx context.Context, taskID int64, trigger string, ac
|
||||
}
|
||||
o.hub.Publish(wshub.Event{Type: "run_started", TaskID: taskID, Data: map[string]any{"run_id": runID}})
|
||||
|
||||
go o.runAll(context.WithoutCancel(ctx), task, runID, accs, srcEP, dstEP, trigger)
|
||||
// dbCtx outlives the request so status/counter writes still land after a
|
||||
// pause or cancel; runCtx is what Pause/Cancel actually stop, and every
|
||||
// account's IMAP work hangs off it.
|
||||
dbCtx := context.WithoutCancel(ctx)
|
||||
runCtx, runCancel := context.WithCancel(dbCtx)
|
||||
h := &runHandle{cancel: runCancel}
|
||||
o.registerRun(taskID, h)
|
||||
|
||||
go o.runAll(dbCtx, runCtx, h, task, runID, accs, srcEP, dstEP, trigger)
|
||||
return runID, nil
|
||||
}
|
||||
|
||||
func (o *Orchestrator) runAll(ctx context.Context, task store.Task, runID int64, accs []store.Account, srcEP, dstEP imapx.Endpoint, trigger string) {
|
||||
func (o *Orchestrator) runAll(ctx, runCtx context.Context, h *runHandle, task store.Task, runID int64, accs []store.Account, srcEP, dstEP imapx.Endpoint, trigger string) {
|
||||
defer o.unregisterRun(task.ID)
|
||||
defer h.cancel()
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
slog.Error("run coordinator panicked", "task", task.ID, "run", runID, "panic", r)
|
||||
@@ -256,7 +355,17 @@ func (o *Orchestrator) runAll(ctx context.Context, task store.Task, runID int64,
|
||||
sem := make(chan struct{}, o.concurrency)
|
||||
var wg sync.WaitGroup
|
||||
|
||||
for _, a := range accs {
|
||||
for i, a := range accs {
|
||||
// Stopped mid-queue: the accounts that never started are marked with the
|
||||
// same status as the ones that were interrupted, so a pause leaves the
|
||||
// whole remainder resumable and a cancel leaves it cancelled.
|
||||
if runCtx.Err() != nil {
|
||||
st := h.stopReason().accountStatus()
|
||||
for _, rest := range accs[i:] {
|
||||
_ = o.store.SetAccountStatus(ctx, rest.ID, st)
|
||||
}
|
||||
break
|
||||
}
|
||||
wg.Add(1)
|
||||
sem <- struct{}{}
|
||||
go func(a store.Account) {
|
||||
@@ -273,7 +382,7 @@ func (o *Orchestrator) runAll(ctx context.Context, task store.Task, runID int64,
|
||||
mu.Unlock()
|
||||
}
|
||||
}()
|
||||
c, s, e := o.runAccount(ctx, task, runID, a, srcEP, dstEP)
|
||||
c, s, e := o.runAccount(ctx, runCtx, h, task, runID, a, srcEP, dstEP)
|
||||
mu.Lock()
|
||||
totCopied += c
|
||||
totSkipped += s
|
||||
@@ -283,23 +392,32 @@ func (o *Orchestrator) runAll(ctx context.Context, task store.Task, runID int64,
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
reason := h.stopReason()
|
||||
status := "done"
|
||||
if totErr > 0 {
|
||||
switch {
|
||||
case reason == stopPaused:
|
||||
status = "paused"
|
||||
case reason == stopCancelled:
|
||||
status = "cancelled"
|
||||
case totErr > 0:
|
||||
status = "done_with_errors"
|
||||
}
|
||||
_ = o.store.FinishRun(ctx, runID, status, totCopied, totSkipped, totErr)
|
||||
_ = o.store.SetTaskStatus(ctx, task.ID, status)
|
||||
o.hub.Publish(wshub.Event{Type: "run_done", TaskID: task.ID,
|
||||
Data: map[string]any{"run_id": runID, "copied": totCopied, "skipped": totSkipped, "errors": totErr}})
|
||||
Data: map[string]any{"run_id": runID, "status": status,
|
||||
"copied": totCopied, "skipped": totSkipped, "errors": totErr}})
|
||||
|
||||
if shouldBreak(trigger, totErr) {
|
||||
// An operator stopping the run is not a schedule failure, so leave the
|
||||
// breaker alone even when the accounts that did run reported errors.
|
||||
if reason == stopNone && shouldBreak(trigger, totErr) {
|
||||
_ = o.store.SetTaskBroken(ctx, task.ID)
|
||||
o.hub.Publish(wshub.Event{Type: "task_broken", TaskID: task.ID,
|
||||
Data: map[string]any{"task_id": task.ID, "errors": totErr}})
|
||||
}
|
||||
}
|
||||
|
||||
func (o *Orchestrator) runAccount(ctx context.Context, task store.Task, runID int64, a store.Account, srcEP, dstEP imapx.Endpoint) (int64, int64, int64) {
|
||||
func (o *Orchestrator) runAccount(ctx, runCtx context.Context, h *runHandle, task store.Task, runID int64, a store.Account, srcEP, dstEP imapx.Endpoint) (int64, int64, int64) {
|
||||
o.hub.Publish(wshub.Event{Type: "account_started", TaskID: task.ID, Data: map[string]any{
|
||||
"account_id": a.ID,
|
||||
"src_login": a.SrcLogin, "src_host": srcEP.Host, "src_port": srcEP.Port,
|
||||
@@ -310,10 +428,11 @@ func (o *Orchestrator) runAccount(ctx context.Context, task store.Task, runID in
|
||||
_ = o.store.ResetAccountCounters(ctx, a.ID) // start from zero; IncAccountCounters is additive
|
||||
_ = o.store.ClearAccountErrors(ctx, a.ID) // drop last run's per-error rows
|
||||
|
||||
// Per-account cancellable context: IMAP work uses actx (so CancelAccount
|
||||
// stops it); DB writes keep the parent ctx so status/counters persist even
|
||||
// after cancellation. ctx is context.WithoutCancel from runAll.
|
||||
actx, cancel := context.WithCancel(ctx)
|
||||
// Per-account cancellable context: IMAP work uses actx, so both CancelAccount
|
||||
// and a task-wide pause/cancel (which cancels runCtx) stop it. DB writes keep
|
||||
// ctx — the uncancellable one from runAll — so status/counters persist even
|
||||
// after cancellation.
|
||||
actx, cancel := context.WithCancel(runCtx)
|
||||
o.registerCancel(a.ID, cancel)
|
||||
defer func() {
|
||||
o.unregisterCancel(a.ID)
|
||||
@@ -322,20 +441,20 @@ func (o *Orchestrator) runAccount(ctx context.Context, task store.Task, runID in
|
||||
|
||||
srcPass, err := crypto.Decrypt(o.encKey, a.SrcPassEnc)
|
||||
if err != nil {
|
||||
return o.accountFailed(ctx, task.ID, runID, a, srcEP, dstEP, "src", err)
|
||||
return o.accountFailed(ctx, runCtx, h, task.ID, runID, a, srcEP, dstEP, "src", err)
|
||||
}
|
||||
dstPass, err := crypto.Decrypt(o.encKey, a.DstPassEnc)
|
||||
if err != nil {
|
||||
return o.accountFailed(ctx, task.ID, runID, a, srcEP, dstEP, "dst", err)
|
||||
return o.accountFailed(ctx, runCtx, h, task.ID, runID, a, srcEP, dstEP, "dst", err)
|
||||
}
|
||||
|
||||
src, err := imapx.Connect(actx, srcEP)
|
||||
if err != nil {
|
||||
return o.accountFailed(ctx, task.ID, runID, a, srcEP, dstEP, "src", err)
|
||||
return o.accountFailed(ctx, runCtx, h, task.ID, runID, a, srcEP, dstEP, "src", err)
|
||||
}
|
||||
if err := src.Login(a.SrcLogin, string(srcPass)).Wait(); err != nil {
|
||||
_ = src.Logout().Wait()
|
||||
return o.accountFailed(ctx, task.ID, runID, a, srcEP, dstEP, "src", err)
|
||||
return o.accountFailed(ctx, runCtx, h, task.ID, runID, a, srcEP, dstEP, "src", err)
|
||||
}
|
||||
// srcClient holds the LIVE source connection. A body-read timeout (server
|
||||
// under-delivering a literal) forces a mid-run reconnect via reconnectSrc,
|
||||
@@ -348,11 +467,11 @@ func (o *Orchestrator) runAccount(ctx context.Context, task store.Task, runID in
|
||||
|
||||
dst, err := imapx.Connect(actx, dstEP)
|
||||
if err != nil {
|
||||
return o.accountFailed(ctx, task.ID, runID, a, srcEP, dstEP, "dst", err)
|
||||
return o.accountFailed(ctx, runCtx, h, task.ID, runID, a, srcEP, dstEP, "dst", err)
|
||||
}
|
||||
defer func() { _ = dst.Logout().Wait() }()
|
||||
if err := dst.Login(a.DstLogin, string(dstPass)).Wait(); err != nil {
|
||||
return o.accountFailed(ctx, task.ID, runID, a, srcEP, dstEP, "dst", err)
|
||||
return o.accountFailed(ctx, runCtx, h, task.ID, runID, a, srcEP, dstEP, "dst", err)
|
||||
}
|
||||
|
||||
// reconnectSrc dials and logs in a fresh source client, swaps it in as the
|
||||
@@ -422,7 +541,7 @@ func (o *Orchestrator) runAccount(ctx context.Context, task store.Task, runID in
|
||||
folders, err := imapx.ListFolders(src)
|
||||
touch()
|
||||
if err != nil {
|
||||
return o.accountFailed(ctx, task.ID, runID, a, srcEP, dstEP, "src", err)
|
||||
return o.accountFailed(ctx, runCtx, h, task.ID, runID, a, srcEP, dstEP, "src", err)
|
||||
}
|
||||
|
||||
// Planning pass: decide folders from the account's own config, then EXAMINE
|
||||
@@ -547,11 +666,18 @@ func (o *Orchestrator) runAccount(ctx context.Context, task store.Task, runID in
|
||||
}
|
||||
|
||||
if actx.Err() != nil {
|
||||
_ = o.store.SetAccountStatus(ctx, a.ID, "cancelled")
|
||||
o.hub.Publish(wshub.Event{Type: "cancelled", TaskID: task.ID,
|
||||
// A task-wide pause leaves the account resumable; anything else (per-account
|
||||
// cancel, stall watchdog, task-wide cancel) leaves it cancelled.
|
||||
st := "cancelled"
|
||||
if runCtx.Err() != nil {
|
||||
st = h.stopReason().accountStatus()
|
||||
}
|
||||
_ = o.store.SetAccountStatus(ctx, a.ID, st)
|
||||
o.hub.Publish(wshub.Event{Type: st, TaskID: task.ID,
|
||||
Data: map[string]any{"account_id": a.ID, "src_login": a.SrcLogin,
|
||||
"copied": copied, "skipped": skipped, "errors": errs}})
|
||||
slog.Info("account cancelled", "account", a.ID, "src_login", a.SrcLogin, "copied", copied, "skipped", skipped)
|
||||
slog.Info("account stopped", "account", a.ID, "src_login", a.SrcLogin,
|
||||
"status", st, "copied", copied, "skipped", skipped)
|
||||
return copied, skipped, errs
|
||||
}
|
||||
|
||||
@@ -567,11 +693,16 @@ func (o *Orchestrator) runAccount(ctx context.Context, task store.Task, runID in
|
||||
return copied, skipped, errs
|
||||
}
|
||||
|
||||
func (o *Orchestrator) accountFailed(ctx context.Context, taskID, runID int64, a store.Account, srcEP, dstEP imapx.Endpoint, side string, err error) (int64, int64, int64) {
|
||||
// A cancellation surfacing as an error is a cancel, not a failure.
|
||||
func (o *Orchestrator) accountFailed(ctx, runCtx context.Context, h *runHandle, taskID, runID int64, a store.Account, srcEP, dstEP imapx.Endpoint, side string, err error) (int64, int64, int64) {
|
||||
// A cancellation surfacing as an error is a stop, not a failure — and a
|
||||
// task-wide pause must still leave the account resumable.
|
||||
if errors.Is(err, context.Canceled) {
|
||||
_ = o.store.SetAccountStatus(ctx, a.ID, "cancelled")
|
||||
o.hub.Publish(wshub.Event{Type: "cancelled", TaskID: taskID,
|
||||
st := "cancelled"
|
||||
if runCtx.Err() != nil {
|
||||
st = h.stopReason().accountStatus()
|
||||
}
|
||||
_ = o.store.SetAccountStatus(ctx, a.ID, st)
|
||||
o.hub.Publish(wshub.Event{Type: st, TaskID: taskID,
|
||||
Data: map[string]any{"account_id": a.ID, "src_login": a.SrcLogin}})
|
||||
return 0, 0, 0
|
||||
}
|
||||
|
||||
@@ -0,0 +1,83 @@
|
||||
package orchestrator
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestStopReasonAccountStatus(t *testing.T) {
|
||||
if got := stopPaused.accountStatus(); got != "paused" {
|
||||
t.Fatalf("stopPaused = %q want paused", got)
|
||||
}
|
||||
if got := stopCancelled.accountStatus(); got != "cancelled" {
|
||||
t.Fatalf("stopCancelled = %q want cancelled", got)
|
||||
}
|
||||
// A run stopped without an operator reason (per-account cancel, stall
|
||||
// watchdog) must not look like a pause, or Resume would pick it up.
|
||||
if got := stopNone.accountStatus(); got != "cancelled" {
|
||||
t.Fatalf("stopNone = %q want cancelled", got)
|
||||
}
|
||||
}
|
||||
|
||||
// The first stop wins: a cancel arriving after a pause must not downgrade the
|
||||
// accounts a pause already promised to keep resumable, and vice versa.
|
||||
func TestRunHandleFirstStopWins(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
first, later stopReason
|
||||
}{
|
||||
{"pause then cancel", stopPaused, stopCancelled},
|
||||
{"cancel then pause", stopCancelled, stopPaused},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
_, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
h := &runHandle{cancel: cancel}
|
||||
h.stopWith(tc.first)
|
||||
h.stopWith(tc.later)
|
||||
if got := h.stopReason(); got != tc.first {
|
||||
t.Fatalf("reason = %v want %v", got, tc.first)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestRunHandleStopCancelsContext(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
h := &runHandle{cancel: cancel}
|
||||
if ctx.Err() != nil {
|
||||
t.Fatal("context cancelled before stop")
|
||||
}
|
||||
h.stopWith(stopPaused)
|
||||
if ctx.Err() == nil {
|
||||
t.Fatal("stop must cancel the run context")
|
||||
}
|
||||
}
|
||||
|
||||
// Pause/Cancel report false for a task with no live run, which the HTTP layer
|
||||
// turns into 409 instead of pretending it stopped something.
|
||||
func TestStopRunWithoutLiveRun(t *testing.T) {
|
||||
o := &Orchestrator{runs: map[int64]*runHandle{}}
|
||||
if o.PauseTask(1) {
|
||||
t.Fatal("PauseTask must report false with no live run")
|
||||
}
|
||||
if o.CancelTask(1) {
|
||||
t.Fatal("CancelTask must report false with no live run")
|
||||
}
|
||||
|
||||
_, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
h := &runHandle{cancel: cancel}
|
||||
o.registerRun(1, h)
|
||||
if !o.PauseTask(1) {
|
||||
t.Fatal("PauseTask must report true for a live run")
|
||||
}
|
||||
if got := h.stopReason(); got != stopPaused {
|
||||
t.Fatalf("reason = %v want stopPaused", got)
|
||||
}
|
||||
o.unregisterRun(1)
|
||||
if o.CancelTask(1) {
|
||||
t.Fatal("unregistered run must not be stoppable")
|
||||
}
|
||||
}
|
||||
@@ -120,13 +120,16 @@ type SchedulableTask struct {
|
||||
}
|
||||
|
||||
// ListSchedulableTasks returns tasks eligible to auto-run: schedule on, not
|
||||
// broken, not currently running — each joined with its last finished run time.
|
||||
// broken, neither running nor paused — each joined with its last finished run
|
||||
// time. A paused task waits for the operator to resume it; auto-starting a full
|
||||
// run behind their back would defeat the pause.
|
||||
func (s *Store) ListSchedulableTasks(ctx context.Context) ([]SchedulableTask, error) {
|
||||
rows, err := s.Pool.Query(ctx,
|
||||
`SELECT t.id, t.schedule_interval_seconds, t.schedule_anchor,
|
||||
(SELECT max(finished_at) FROM runs r WHERE r.task_id=t.id AND r.finished_at IS NOT NULL)
|
||||
FROM tasks t
|
||||
WHERE t.schedule_interval_seconds > 0 AND NOT t.broken AND t.status <> 'running'`)
|
||||
WHERE t.schedule_interval_seconds > 0 AND NOT t.broken
|
||||
AND t.status <> 'running' AND t.status <> 'paused'`)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user