From 039ac2f1dad9d466291238673a93ae5cf57de3e5 Mon Sep 17 00:00:00 2001 From: Vassiliy Yegorov Date: Wed, 29 Jul 2026 11:20:52 +0700 Subject: [PATCH] Add pause and cancel to a running migration MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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) --- internal/httpapi/router.go | 3 + internal/httpapi/run.go | 50 +++++++ internal/orchestrator/orchestrator.go | 185 ++++++++++++++++++++++---- internal/orchestrator/stop_test.go | 83 ++++++++++++ internal/store/tasks.go | 7 +- web/src/api.ts | 8 ++ web/src/pages/TaskDetail.tsx | 96 +++++++++++-- 7 files changed, 391 insertions(+), 41 deletions(-) create mode 100644 internal/orchestrator/stop_test.go diff --git a/internal/httpapi/router.go b/internal/httpapi/router.go index 0e842ae..e4b8d9e 100644 --- a/internal/httpapi/router.go +++ b/internal/httpapi/router.go @@ -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) diff --git a/internal/httpapi/run.go b/internal/httpapi/run.go index dbcecaa..a658d74 100644 --- a/internal/httpapi/run.go +++ b/internal/httpapi/run.go @@ -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 { diff --git a/internal/orchestrator/orchestrator.go b/internal/orchestrator/orchestrator.go index 21e6feb..b27325e 100644 --- a/internal/orchestrator/orchestrator.go +++ b/internal/orchestrator/orchestrator.go @@ -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 } diff --git a/internal/orchestrator/stop_test.go b/internal/orchestrator/stop_test.go new file mode 100644 index 0000000..a196e51 --- /dev/null +++ b/internal/orchestrator/stop_test.go @@ -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") + } +} diff --git a/internal/store/tasks.go b/internal/store/tasks.go index 4646aa6..dfd1186 100644 --- a/internal/store/tasks.go +++ b/internal/store/tasks.go @@ -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 } diff --git a/web/src/api.ts b/web/src/api.ts index 2891f51..9bcf2d6 100644 --- a/web/src/api.ts +++ b/web/src/api.ts @@ -151,6 +151,14 @@ export const testAccounts = (id: number) => api(`/api/tasks/${id}/test`, { metho export const runTask = (id: number, accountIds?: number[]) => api(`/api/tasks/${id}/run`, accountIds?.length ? jsonBody({ account_ids: accountIds }) : { method: 'POST' }) +// Pause stops the run but leaves its unfinished accounts resumable; cancel ends +// it and marks them cancelled. Resume re-runs exactly the paused accounts. +export const pauseTask = (id: number) => api(`/api/tasks/${id}/pause`, { method: 'POST' }) + +export const cancelTask = (id: number) => api(`/api/tasks/${id}/cancel`, { method: 'POST' }) + +export const resumeTask = (id: number) => api<{ run_id: number }>(`/api/tasks/${id}/resume`, { method: 'POST' }) + export interface Run { id: number task_id: number diff --git a/web/src/pages/TaskDetail.tsx b/web/src/pages/TaskDetail.tsx index 8ab5faf..968bc3e 100644 --- a/web/src/pages/TaskDetail.tsx +++ b/web/src/pages/TaskDetail.tsx @@ -1,5 +1,5 @@ import { useEffect, useRef, useState, type ChangeEvent, type FormEvent } from 'react' -import { cancelAccount, createAccount, deleteAccount, getTask, importCSV, importKerioCSV, probeAccountFolders, probeFolders, runTask, setAccountFolderMapping, setTaskSchedule, testAccounts, updateAccountCredentials, type Account, type TaskDetail as TaskDetailData } from '../api' +import { cancelAccount, cancelTask, createAccount, deleteAccount, getTask, importCSV, importKerioCSV, pauseTask, probeAccountFolders, probeFolders, resumeTask, runTask, setAccountFolderMapping, setTaskSchedule, testAccounts, updateAccountCredentials, type Account, type TaskDetail as TaskDetailData } from '../api' import { connectTaskWS, type TaskEvent } from '../ws' import { StatusBadge } from '../components/StatusBadge' import { useConfirm } from '../components/ConfirmProvider' @@ -59,6 +59,8 @@ function describeEvent(ev: TaskEvent): string { } case 'cancelled': return `CANCELLED #${d.account_id} (${d.src_login}): copied ${d.copied ?? 0}, skipped ${d.skipped ?? 0}` + case 'paused': + return `PAUSED #${d.account_id} (${d.src_login}): copied ${d.copied ?? 0}, skipped ${d.skipped ?? 0} — resumable` case 'error': { const where = d.folder ? ` folder "${d.folder}"` : d.side ? ` (${d.side} ${at})` : '' return `ERROR #${d.account_id}${where}: ${d.error}` @@ -66,7 +68,7 @@ function describeEvent(ev: TaskEvent): string { case 'run_started': return `RUN started (run #${d.run_id})` case 'run_done': - return `RUN finished: copied ${d.copied}, skipped ${d.skipped}, errors ${d.errors}` + return `RUN ${String(d.status ?? 'finished')}: copied ${d.copied}, skipped ${d.skipped}, errors ${d.errors}` default: return JSON.stringify(ev.data) } @@ -163,7 +165,7 @@ export function TaskDetail({ id }: { id: number }) { }, } }) - } else if (accId != null && (ev.type === 'account_started' || ev.type === 'account_done' || ev.type === 'cancelled' || (ev.type === 'error' && d.folder == null))) { + } else if (accId != null && (ev.type === 'account_started' || ev.type === 'account_done' || ev.type === 'cancelled' || ev.type === 'paused' || (ev.type === 'error' && d.folder == null))) { // terminal/reset for this account — drop live overlay, fall back to DB setLive((prev) => { if (!(accId in prev)) return prev @@ -174,7 +176,7 @@ export function TaskDetail({ id }: { id: number }) { } // Structural events refresh the persisted view; `progress` is covered by live state. - if (['account_started', 'account_test', 'account_done', 'run_started', 'run_done', 'error', 'folder', 'cancelled', 'plan', 'task_broken'].includes(ev.type)) { + if (['account_started', 'account_test', 'account_done', 'run_started', 'run_done', 'error', 'folder', 'cancelled', 'paused', 'plan', 'task_broken'].includes(ev.type)) { reload() } }), @@ -395,6 +397,51 @@ export function TaskDetail({ id }: { id: number }) { } } + async function onPause() { + setBusy('run') + setError(null) + try { + await pauseTask(id) + } catch (err) { + setError(err instanceof Error ? err.message : 'Failed to pause the run') + } finally { + setBusy(null) + } + } + + async function onCancelRun() { + const ok = await confirm({ + title: 'Cancel migration', + message: 'Stop the run and mark every unfinished account as cancelled? Copied messages are kept.', + confirmLabel: 'Cancel migration', + cancelLabel: 'Keep running', + danger: true, + }) + if (!ok) return + setBusy('run') + setError(null) + try { + await cancelTask(id) + } catch (err) { + setError(err instanceof Error ? err.message : 'Failed to cancel the run') + } finally { + setBusy(null) + } + } + + async function onResume() { + setBusy('run') + setError(null) + try { + await resumeTask(id) + reload() + } catch (err) { + setError(err instanceof Error ? err.message : 'Failed to resume the run') + } finally { + setBusy(null) + } + } + async function onSchedule(intervalSeconds: number) { setError(null) try { @@ -422,6 +469,8 @@ export function TaskDetail({ id }: { id: number }) { const { task, accounts } = data const isRunning = task.status === 'running' + // Accounts a pause left unfinished — what Resume picks up. + const pausedCount = accounts.filter((a) => a.status === 'paused').length // A failed connection test is the entry point for fixing the credentials that // caused it — an imported account is otherwise only deletable. @@ -515,14 +564,37 @@ export function TaskDetail({ id }: { id: number }) { - - {!runReady && accounts.length > 0 && ( + {isRunning ? ( + <> + + + pause keeps the unfinished accounts resumable + + ) : ( + <> + {pausedCount > 0 && ( + + )} + + + )} + {!isRunning && !runReady && accounts.length > 0 && ( {effectiveSelected.length > 0 ? 'selected accounts must pass both connection tests'