From 67d0d9627605b732cbda3e16cd8f4f4111476a26 Mon Sep 17 00:00:00 2001 From: Seongmin Lee Date: Fri, 17 Jul 2026 01:07:56 +0900 Subject: [PATCH] spindle/{config,xrpc}: add pipeline cancel timeout Signed-off-by: Seongmin Lee --- spindle/config/config.go | 2 + spindle/xrpc/pipeline_cancel_pipeline.go | 51 +++++++++++++----------- 2 files changed, 29 insertions(+), 24 deletions(-) diff --git a/spindle/config/config.go b/spindle/config/config.go index 3bae65d1..64585956 100644 --- a/spindle/config/config.go +++ b/spindle/config/config.go @@ -25,6 +25,8 @@ type Server struct { QueueSize int `env:"QUEUE_SIZE, default=100"` MaxJobCount int `env:"MAX_JOB_COUNT, default=2"` // max number of pipelines that run at a time DockerSocket string `env:"DOCKER_SOCKET"` // path to a docker socket to expose to workflow containers + + PipelineCancelTimeout time.Duration `env:"PIPELINE_CANCEL_TIMEOUT, default=10min"` } type Tap struct { diff --git a/spindle/xrpc/pipeline_cancel_pipeline.go b/spindle/xrpc/pipeline_cancel_pipeline.go index ac1cdaab..10261062 100644 --- a/spindle/xrpc/pipeline_cancel_pipeline.go +++ b/spindle/xrpc/pipeline_cancel_pipeline.go @@ -1,6 +1,7 @@ package xrpc import ( + "context" "encoding/json" "fmt" "net/http" @@ -74,34 +75,36 @@ func (x *Xrpc) CancelPipeline(w http.ResponseWriter, r *http.Request) { } } - canceled := false - defer func() { - l.Debug("canceled pipeline", "canceled", canceled) - }() - - for _, wName := range workflows { - wid := models.WorkflowId{ - PipelineId: pipelineId, - Name: wName, - } - l.Debug("cancel pipeline", "wid", wid) + w.WriteHeader(http.StatusAccepted) - for _, engine := range x.Engines { - l.Debug("destroying workflow", "wid", wid) - err := engine.DestroyWorkflow(r.Context(), wid) - if err != nil { - fail(xrpcerr.GenericError(fmt.Errorf("failed to destroy workflow: %w", err))) - return + go func() { + ctx, cancel := context.WithTimeout(context.Background(), x.Config.Server.PipelineCancelTimeout) + defer cancel() + for _, wName := range workflows { + wid := models.WorkflowId{ + PipelineId: pipelineId, + Name: wName, } - err = x.Db.StatusCancelled(wid, "User canceled the workflow", -1, x.Notifier) - if err != nil { - fail(xrpcerr.GenericError(fmt.Errorf("failed to emit status failed: %w", err))) - return + l := l.With("wid", wid) + l.Debug("cancel workflow") + + for name, engine := range x.Engines { + l := l.With("engine", name) + l.Debug("destroying workflow") + err := engine.DestroyWorkflow(ctx, wid) + if err != nil { + l.Error("failed to destroy workflow", "err", err) + return + } + err = x.Db.StatusCancelled(wid, "User canceled the workflow", -1, x.Notifier) + if err != nil { + l.Error("failed to update pipeline status in db", "err", err) + return + } + l.Debug("canceled workflow") } } - } - - canceled = true + }() w.WriteHeader(http.StatusOK) } -- 2.51.2