Monorepo for Tangled
0

Configure Feed

Select the types of activity you want to include in your feed.

spindle/engine: properly interrupt running workflows on cancel, dont overwrite cancelled status

Signed-off-by: dawn <dawn@tangled.org>

author
dawn
date (Jul 23, 2026, 8:44 PM +0300) commit 50574018 parent bba3c2f9 change-id srytztpu
+206 -53
+81 -45
spindle/engine/engine.go
··· 16 16 ) 17 17 18 18 var ( 19 - ErrTimedOut = errors.New("timed out") 20 - ErrWorkflowFailed = errors.New("workflow failed") 19 + ErrTimedOut = errors.New("timed out") 20 + ErrWorkflowFailed = errors.New("workflow failed") 21 + ErrWorkflowCanceled = errors.New("workflow canceled") 21 22 ) 22 23 24 + var ( 25 + activeMu sync.Mutex 26 + activeCancels = make(map[models.WorkflowId]context.CancelCauseFunc) 27 + ) 28 + 29 + func CancelWorkflow(wid models.WorkflowId) { 30 + activeMu.Lock() 31 + cancel, ok := activeCancels[wid] 32 + activeMu.Unlock() 33 + if ok { 34 + cancel(ErrWorkflowCanceled) 35 + } 36 + } 37 + 38 + // user cancel, timeout is DeadlineExceeded 39 + func isCanceled(wfCtx context.Context) bool { 40 + return errors.Is(context.Cause(wfCtx), ErrWorkflowCanceled) 41 + } 42 + 43 + // for when recording early wf cancellations 44 + func writeWfError(db *db.DB, n *notifier.Notifier, l *slog.Logger, wfCtx context.Context, wid models.WorkflowId, phase string, err error) { 45 + l = l.With("wid", wid, "phase", phase) 46 + switch { 47 + case isCanceled(wfCtx): 48 + l.Info("workflow canceled") 49 + if dbErr := db.StatusCancelled(wid, "User canceled the workflow", -1, n); dbErr != nil { 50 + l.Error("failed to set workflow status to cancelled", "err", dbErr) 51 + } 52 + case errors.Is(err, ErrTimedOut) || errors.Is(wfCtx.Err(), context.DeadlineExceeded): 53 + l.Info("workflow timed out") 54 + if dbErr := db.StatusTimeout(wid, n); dbErr != nil { 55 + l.Error("failed to set workflow status to timeout", "err", dbErr) 56 + } 57 + default: 58 + l.Error("workflow failed", "err", err) 59 + if dbErr := db.StatusFailed(wid, err.Error(), -1, n); dbErr != nil { 60 + l.Error("failed to set workflow status to failed", "err", dbErr) 61 + } 62 + } 63 + } 64 + 23 65 type workflowFinalizer interface { 24 66 FinalizeWorkflow(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, wfLogger models.WorkflowLogger) error 25 67 } ··· 82 124 } 83 125 84 126 wg.Go(func() { 85 - 127 + if st, err := db.GetStatus(wid); err == nil && models.StatusKind(st.Status).IsFinish() { 128 + l.Info("skipping finished workflow", "wid", wid, "status", st.Status) 129 + return 130 + } 86 131 defer func() { 87 132 if s3 != nil { 88 133 logFile := filepath.Join(cfg.Server.LogDir, fmt.Sprintf("%s.log", wid.String())) ··· 101 146 defer wfLogger.Close() 102 147 } 103 148 149 + timeoutCtx, timeoutCancel := context.WithTimeout(ctx, workflowTimeout) 150 + defer timeoutCancel() 151 + 152 + wfCtx, userCancel := context.WithCancelCause(timeoutCtx) 153 + defer userCancel(nil) 154 + 155 + // allow wf context to be cancelled properly by manual cancel 156 + activeMu.Lock() 157 + activeCancels[wid] = userCancel 158 + activeMu.Unlock() 159 + defer func() { 160 + activeMu.Lock() 161 + delete(activeCancels, wid) 162 + activeMu.Unlock() 163 + }() 164 + 104 165 l.Info("waiting for slot", "wid", wid) 105 166 slot := WorkflowSlot(NoopSlot{}) 106 167 if s, ok := eng.(WorkflowSlotter); ok { 107 - var err error 108 - slot, err = s.AcquireWorkflowSlot(ctx, wid, &w) 168 + slot, err = s.AcquireWorkflowSlot(wfCtx, wid, &w) 109 169 if err != nil { 110 - l.Error("failed to acquire slot", "wid", wid, "err", err) 111 - dbErr := db.StatusFailed(wid, err.Error(), -1, n) 112 - if dbErr != nil { 113 - l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr) 114 - } 170 + writeWfError(db, n, l, wfCtx, wid, "waiting for slot", err) 115 171 return 116 172 } 117 173 } ··· 123 179 return 124 180 } 125 181 126 - err = eng.SetupWorkflow(ctx, wid, &w, wfLogger) 182 + err = eng.SetupWorkflow(wfCtx, wid, &w, wfLogger) 127 183 if err != nil { 128 - // TODO(winter): Should this always set StatusFailed? 129 - // In the original, we only do in a subset of cases. 130 - l.Error("setting up workflow", "wid", wid, "err", err) 131 - 132 - destroyErr := eng.DestroyWorkflow(ctx, wid) 133 - if destroyErr != nil { 134 - l.Error("failed to destroy workflow after setup failure", "error", destroyErr) 135 - } 136 - 137 - dbErr := db.StatusFailed(wid, err.Error(), -1, n) 138 - if dbErr != nil { 139 - l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr) 184 + if !isCanceled(wfCtx) { 185 + if destroyErr := eng.DestroyWorkflow(ctx, wid); destroyErr != nil { 186 + l.Error("failed to destroy workflow after setup failure", "error", destroyErr) 187 + } 140 188 } 189 + writeWfError(db, n, l, wfCtx, wid, "setting up workflow", err) 141 190 return 142 191 } 143 192 defer eng.DestroyWorkflow(ctx, wid) 144 193 145 - ctx, cancel := context.WithTimeout(ctx, workflowTimeout) 146 - defer cancel() 147 - 148 194 for stepIdx, step := range w.Steps { 149 - // log start of step 150 195 if wfLogger != nil { 151 196 wfLogger. 152 197 ControlWriter(stepIdx, step, models.StepStatusStart). 153 198 Write([]byte{0}) 154 199 } 155 200 156 - err = eng.RunStep(ctx, wid, &w, stepIdx, allSecrets, wfLogger) 201 + err = eng.RunStep(wfCtx, wid, &w, stepIdx, allSecrets, wfLogger) 157 202 158 - // log end of step 159 203 if wfLogger != nil { 160 204 wfLogger. 161 205 ControlWriter(stepIdx, step, models.StepStatusEnd). ··· 163 207 } 164 208 165 209 if err != nil { 166 - if errors.Is(err, ErrTimedOut) { 167 - dbErr := db.StatusTimeout(wid, n) 168 - if dbErr != nil { 169 - l.Error("failed to set workflow status to timeout", "wid", wid, "err", dbErr) 170 - } 171 - } else { 172 - dbErr := db.StatusFailed(wid, err.Error(), -1, n) 173 - if dbErr != nil { 174 - l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr) 175 - } 176 - } 210 + writeWfError(db, n, l, wfCtx, wid, "running step", err) 177 211 return 178 212 } 179 213 } 180 214 181 215 if finalizer, ok := eng.(workflowFinalizer); ok { 182 - if err := finalizer.FinalizeWorkflow(ctx, wid, &w, wfLogger); err != nil { 183 - dbErr := db.StatusFailed(wid, err.Error(), -1, n) 184 - if dbErr != nil { 185 - l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr) 186 - } 216 + if err := finalizer.FinalizeWorkflow(wfCtx, wid, &w, wfLogger); err != nil { 217 + writeWfError(db, n, l, wfCtx, wid, "finalizing", err) 187 218 return 188 219 } 220 + } 221 + 222 + if isCanceled(wfCtx) { 223 + writeWfError(db, n, l, wfCtx, wid, "before success", nil) 224 + return 189 225 } 190 226 191 227 err = db.StatusSuccess(wid, n)
+109
spindle/engine/engine_test.go
··· 155 155 t.Fatalf("expected unique status to be success, got status=%v err=%v", statusUnique, err) 156 156 } 157 157 } 158 + 159 + func TestCancelWorkflow_NotOverwritten(t *testing.T) { 160 + t.Parallel() 161 + 162 + testDB := newTestDB(t) 163 + logger := slog.New(slog.NewTextHandler(os.Stderr, nil)) 164 + 165 + stepStarted := make(chan struct{}) 166 + eng := &mockEngine{ 167 + runStepFunc: func(ctx context.Context, wid models.WorkflowId, idx int) error { 168 + close(stepStarted) 169 + <-ctx.Done() 170 + return ctx.Err() 171 + }, 172 + } 173 + 174 + pipelineId := models.PipelineId{ 175 + Knot: "test-knot", 176 + Rkey: "test-rkey", 177 + } 178 + 179 + wid := models.WorkflowId{ 180 + PipelineId: pipelineId, 181 + Name: "cancel_test_job", 182 + } 183 + 184 + pipeline := &models.Pipeline{ 185 + Workflows: map[models.Engine][]models.Workflow{ 186 + eng: { 187 + { 188 + Name: "cancel_test_job", 189 + Steps: []models.Step{mockStep{name: "step1"}}, 190 + }, 191 + }, 192 + }, 193 + } 194 + 195 + cfg := &config.Config{Server: config.Server{LogDir: t.TempDir()}} 196 + doneChan := make(chan struct{}) 197 + go func() { 198 + StartWorkflows(logger, nil, cfg, testDB, nil, context.Background(), pipeline, pipelineId) 199 + close(doneChan) 200 + }() 201 + 202 + select { 203 + case <-stepStarted: 204 + case <-time.After(5 * time.Second): 205 + t.Fatal("timed out waiting for step to start") 206 + } 207 + 208 + _ = testDB.StatusCancelled(wid, "User canceled the workflow", -1, nil) 209 + CancelWorkflow(wid) 210 + 211 + select { 212 + case <-doneChan: 213 + case <-time.After(5 * time.Second): 214 + t.Fatal("timed out waiting for StartWorkflows to complete") 215 + } 216 + 217 + // the runner writes StatusCancelled itself when it sees the canceled ctx 218 + // the handler writes nothing for a live wf, so nothing lands after to overwrite it 219 + st, err := testDB.GetStatus(wid) 220 + if err != nil { 221 + t.Fatalf("GetStatus error = %v", err) 222 + } 223 + if st.Status != string(models.StatusKindCancelled) { 224 + t.Fatalf("expected status to be cancelled, got %s", st.Status) 225 + } 226 + } 227 + 228 + func TestSetupTimeout_ReportsTimeout(t *testing.T) { 229 + t.Parallel() 230 + 231 + testDB := newTestDB(t) 232 + logger := slog.New(slog.NewTextHandler(os.Stderr, nil)) 233 + 234 + // setup blocks past the workflow timeout, so it should land as timeout not failed 235 + eng := &mockEngine{ 236 + timeout: 100 * time.Millisecond, 237 + setupFunc: func(ctx context.Context, wid models.WorkflowId) error { 238 + <-ctx.Done() 239 + return ctx.Err() 240 + }, 241 + } 242 + 243 + pipelineId := models.PipelineId{Knot: "test-knot", Rkey: "test-rkey"} 244 + wid := models.WorkflowId{PipelineId: pipelineId, Name: "timeout_job"} 245 + 246 + pipeline := &models.Pipeline{ 247 + Workflows: map[models.Engine][]models.Workflow{ 248 + eng: {{Name: "timeout_job", Steps: []models.Step{mockStep{name: "step1"}}}}, 249 + }, 250 + } 251 + 252 + cfg := &config.Config{Server: config.Server{LogDir: t.TempDir()}} 253 + StartWorkflows(logger, nil, cfg, testDB, nil, context.Background(), pipeline, pipelineId) 254 + 255 + st, err := testDB.GetStatus(wid) 256 + if err != nil { 257 + t.Fatalf("GetStatus error = %v", err) 258 + } 259 + if st.Status != string(models.StatusKindTimeout) { 260 + t.Fatalf("expected status to be timeout, got %s", st.Status) 261 + } 262 + 263 + if len(eng.runStepCalls) != 0 { 264 + t.Fatalf("expected no steps to run after setup timeout, got %d", len(eng.runStepCalls)) 265 + } 266 + }
+16 -8
spindle/xrpc/pipeline_cancel_pipeline.go
··· 7 7 8 8 "github.com/bluesky-social/indigo/atproto/syntax" 9 9 "tangled.org/core/api/tangled" 10 + "tangled.org/core/spindle/engine" 10 11 "tangled.org/core/spindle/models" 11 12 xrpcerr "tangled.org/core/xrpc/errors" 12 13 ) ··· 86 87 } 87 88 l.Debug("cancel pipeline", "wid", wid) 88 89 89 - for _, engine := range x.Engines { 90 + // dont cancel a workflow that already finished 91 + st, err := x.Db.GetStatus(wid) 92 + if err == nil && models.StatusKind(st.Status).IsFinish() { 93 + continue 94 + } 95 + 96 + if err := x.Db.StatusCancelled(wid, "User canceled the workflow", -1, x.Notifier); err != nil { 97 + fail(xrpcerr.GenericError(fmt.Errorf("failed to emit status cancelled: %w", err))) 98 + return 99 + } 100 + 101 + engine.CancelWorkflow(wid) 102 + 103 + for _, eng := range x.Engines { 90 104 l.Debug("destroying workflow", "wid", wid) 91 - err := engine.DestroyWorkflow(r.Context(), wid) 92 - if err != nil { 105 + if err := eng.DestroyWorkflow(r.Context(), wid); err != nil { 93 106 fail(xrpcerr.GenericError(fmt.Errorf("failed to destroy workflow: %w", err))) 94 - return 95 - } 96 - err = x.Db.StatusCancelled(wid, "User canceled the workflow", -1, x.Notifier) 97 - if err != nil { 98 - fail(xrpcerr.GenericError(fmt.Errorf("failed to emit status failed: %w", err))) 99 107 return 100 108 } 101 109 }