From 05cd8e154b612bbf86be0dde72e4e805fb7c952c Mon Sep 17 00:00:00 2001 From: mrhid6 Date: Mon, 20 Jul 2026 14:34:50 +0100 Subject: [PATCH] feat(workflows): step input params model + cascade step delete --- server/internal/models/workflow.go | 25 ++++++++++++------ server/internal/services/workflows.go | 38 +++++++++++++++++++++++++-- 2 files changed, 53 insertions(+), 10 deletions(-) diff --git a/server/internal/models/workflow.go b/server/internal/models/workflow.go index 8789a0c..7b9afd5 100644 --- a/server/internal/models/workflow.go +++ b/server/internal/models/workflow.go @@ -6,6 +6,12 @@ import ( "go.mongodb.org/mongo-driver/v2/bson" ) +type InputParam struct { + Name string `bson:"name" json:"name"` + Default string `bson:"default" json:"default"` + Description string `bson:"description" json:"description"` +} + type WorkflowStep struct { ID bson.ObjectID `bson:"_id,omitempty" json:"-"` StepID string `bson:"step_id" json:"step_id"` @@ -14,17 +20,19 @@ type WorkflowStep struct { Interpreter string `bson:"interpreter" json:"interpreter"` // "bash" | "powershell" Script string `bson:"script" json:"script"` DeclaredOutputs []string `bson:"declared_outputs" json:"declared_outputs"` + DeclaredInputs []InputParam `bson:"declared_inputs" json:"declared_inputs"` SecretRefs []string `bson:"secret_refs" json:"secret_refs"` CreatedAt time.Time `bson:"created_at" json:"created_at"` UpdatedAt time.Time `bson:"updated_at" json:"updated_at"` } type WorkflowStepRef struct { - StepID string `bson:"step_id" json:"step_id"` - Order int `bson:"order" json:"order"` - OnFailure string `bson:"on_failure" json:"on_failure"` // "stop" | "continue" | "retry" - MaxRetries int `bson:"max_retries" json:"max_retries"` - Overrides *StepOverride `bson:"overrides,omitempty" json:"overrides,omitempty"` + StepID string `bson:"step_id" json:"step_id"` + Order int `bson:"order" json:"order"` + OnFailure string `bson:"on_failure" json:"on_failure"` // "stop" | "continue" | "retry" + MaxRetries int `bson:"max_retries" json:"max_retries"` + Overrides *StepOverride `bson:"overrides,omitempty" json:"overrides,omitempty"` + Inputs map[string]string `bson:"inputs,omitempty" json:"inputs,omitempty"` } type StepOverride struct { @@ -48,9 +56,10 @@ type ResolvedStep struct { Name string `bson:"name" json:"name"` Interpreter string `bson:"interpreter" json:"interpreter"` Script string `bson:"script" json:"script"` - SecretRefs []string `bson:"secret_refs" json:"secret_refs"` - OnFailure string `bson:"on_failure" json:"on_failure"` - MaxRetries int `bson:"max_retries" json:"max_retries"` + SecretRefs []string `bson:"secret_refs" json:"secret_refs"` + OnFailure string `bson:"on_failure" json:"on_failure"` + MaxRetries int `bson:"max_retries" json:"max_retries"` + Inputs map[string]string `bson:"inputs" json:"inputs"` } type StepRun struct { diff --git a/server/internal/services/workflows.go b/server/internal/services/workflows.go index b7c2b7b..36ea584 100644 --- a/server/internal/services/workflows.go +++ b/server/internal/services/workflows.go @@ -66,6 +66,9 @@ func CreateStep(s models.WorkflowStep) (*models.WorkflowStep, error) { if s.SecretRefs == nil { s.SecretRefs = []string{} } + if s.DeclaredInputs == nil { + s.DeclaredInputs = []models.InputParam{} + } if _, err := db.Col("workflow_steps").InsertOne(ctx, s); err != nil { return nil, err } @@ -81,6 +84,7 @@ func UpdateStep(stepID string, s models.WorkflowStep) error { "interpreter": s.Interpreter, "script": s.Script, "declared_outputs": s.DeclaredOutputs, + "declared_inputs": s.DeclaredInputs, "secret_refs": s.SecretRefs, "updated_at": time.Now(), }}) @@ -90,8 +94,38 @@ func UpdateStep(stepID string, s models.WorkflowStep) error { func DeleteStep(stepID string) error { ctx, cancel := wfCtx() defer cancel() - _, err := db.Col("workflow_steps").DeleteOne(ctx, bson.M{"step_id": stepID}) - return err + if _, err := db.Col("workflow_steps").DeleteOne(ctx, bson.M{"step_id": stepID}); err != nil { + return err + } + // Cascade: remove this step from every workflow that references it, re-sequencing orders. + cur, err := db.Col("workflows").Find(ctx, bson.M{"steps.step_id": stepID}) + if err != nil { + return err + } + defer cur.Close(ctx) + var wfs []models.Workflow + if err := cur.All(ctx, &wfs); err != nil { + return err + } + for _, w := range wfs { + kept := make([]models.WorkflowStepRef, 0, len(w.Steps)) + for _, ref := range w.Steps { + if ref.StepID == stepID { + continue + } + kept = append(kept, ref) + } + for i := range kept { + kept[i].Order = i + } + if _, err := db.Col("workflows").UpdateOne(ctx, + bson.M{"workflow_id": w.WorkflowID}, + bson.M{"$set": bson.M{"steps": kept, "updated_at": time.Now()}}, + ); err != nil { + return err + } + } + return nil } func getStep(ctx context.Context, stepID string) (*models.WorkflowStep, error) {