From d1b3cd2f74a4ab29ca19051ab6d766d549de9aeb Mon Sep 17 00:00:00 2001 From: mrhid6 Date: Tue, 4 Aug 2026 13:36:43 +0100 Subject: [PATCH] feat: target workflow runs by tag selector --- server/internal/api/workflows.go | 4 ++++ server/internal/models/workflow.go | 1 + server/internal/services/workflow_runner.go | 19 ++++++------------- server/internal/services/workflows.go | 7 +++++++ 4 files changed, 18 insertions(+), 13 deletions(-) diff --git a/server/internal/api/workflows.go b/server/internal/api/workflows.go index 4ecef14..2f2d6f9 100644 --- a/server/internal/api/workflows.go +++ b/server/internal/api/workflows.go @@ -343,6 +343,10 @@ func deleteWorkflow(c *gin.Context) { func runWorkflow(c *gin.Context) { runID, err := services.TriggerWorkflow(auth.InstanceID(c), c.Param("id"), actorFromCtx(c)) if err != nil { + if errors.Is(err, services.ErrNoTargets) { + c.JSON(http.StatusBadRequest, gin.H{"error": "this workflow matches no servers"}) + return + } c.JSON(http.StatusServiceUnavailable, gin.H{"error": err.Error()}) return } diff --git a/server/internal/models/workflow.go b/server/internal/models/workflow.go index a7ad86b..74e1906 100644 --- a/server/internal/models/workflow.go +++ b/server/internal/models/workflow.go @@ -50,6 +50,7 @@ type Workflow struct { WorkflowID string `bson:"workflow_id" json:"workflow_id"` Name string `bson:"name" json:"name"` TargetServerIDs []string `bson:"target_server_ids" json:"target_server_ids"` + TargetTags map[string]string `bson:"target_tags,omitempty" json:"target_tags,omitempty"` Steps []WorkflowStepRef `bson:"steps" json:"steps"` CreatedAt time.Time `bson:"created_at" json:"created_at"` UpdatedAt time.Time `bson:"updated_at" json:"updated_at"` diff --git a/server/internal/services/workflow_runner.go b/server/internal/services/workflow_runner.go index 34cb35e..4a31a70 100644 --- a/server/internal/services/workflow_runner.go +++ b/server/internal/services/workflow_runner.go @@ -22,17 +22,14 @@ func TriggerWorkflow(instanceID, workflowID, actor string) (string, error) { if err != nil { return "", err } - if len(wf.TargetServerIDs) == 0 { - return "", fmt.Errorf("workflow has no target servers") + targets, err := ResolveTargets(instanceID, wf.TargetServerIDs, wf.TargetTags) + if err != nil { + return "", err } if len(wf.Steps) == 0 { return "", fmt.Errorf("workflow has no steps") } - if err := validateTargetServers(instanceID, wf.TargetServerIDs); err != nil { - return "", err - } - ctx, cancel := wfCtx() running := db.Col("workflow_runs").FindOne(ctx, bson.M{"instance_id": instanceID, "workflow_id": workflowID, "status": "running"}) cancel() @@ -54,14 +51,10 @@ func TriggerWorkflow(instanceID, workflowID, actor string) (string, error) { Status: "running", TriggeredBy: actor, StartedAt: time.Now(), - ServerRuns: make([]models.ServerRun, 0, len(wf.TargetServerIDs)), + ServerRuns: make([]models.ServerRun, 0, len(targets)), } - for _, sid := range wf.TargetServerIDs { - hostname := sid - if s, e := getServerByID(sid); e == nil { - hostname = s.Hostname - } - sr := models.ServerRun{ServerID: sid, Hostname: hostname, Status: "queued", RunEnv: map[string]string{}} + for _, srv := range targets { + sr := models.ServerRun{ServerID: srv.ServerID, Hostname: srv.Hostname, Status: "queued", RunEnv: map[string]string{}} for _, rs := range resolved { sr.Steps = append(sr.Steps, models.StepRun{Order: rs.Order, Name: rs.Name, Status: "queued", OutputEnv: map[string]string{}}) } diff --git a/server/internal/services/workflows.go b/server/internal/services/workflows.go index 4a3fabb..6161483 100644 --- a/server/internal/services/workflows.go +++ b/server/internal/services/workflows.go @@ -249,6 +249,9 @@ func CreateWorkflow(instanceID string, w models.Workflow) (*models.Workflow, err if err := ValidateWorkflow(w); err != nil { return nil, err } + if err := ValidateTags(w.TargetTags); err != nil { + return nil, err + } if err := validateTargetServers(instanceID, w.TargetServerIDs); err != nil { return nil, err } @@ -265,6 +268,9 @@ func UpdateWorkflow(instanceID, id string, w models.Workflow) error { if err := ValidateWorkflow(w); err != nil { return err } + if err := ValidateTags(w.TargetTags); err != nil { + return err + } if err := validateTargetServers(instanceID, w.TargetServerIDs); err != nil { return err } @@ -272,6 +278,7 @@ func UpdateWorkflow(instanceID, id string, w models.Workflow) error { _, err := db.Col("workflows").UpdateOne(ctx, bson.M{"workflow_id": id, "instance_id": instanceID}, bson.M{"$set": bson.M{ "name": w.Name, "target_server_ids": w.TargetServerIDs, + "target_tags": w.TargetTags, "steps": w.Steps, "updated_at": time.Now(), }})