From 15c9da1b01fa095163d8a546b0c6ba73ca90ab5f Mon Sep 17 00:00:00 2001 From: mrhid6 Date: Tue, 21 Jul 2026 10:25:41 +0100 Subject: [PATCH] feat(server): default steps seed-on-boot + admin re-sync --- server/cmd/main.go | 6 ++ server/internal/api/workflows.go | 11 +++ server/internal/services/defaults.go | 94 +++++++++++++++++++++++ server/internal/services/defaults_test.go | 38 +++++++++ server/internal/services/workflows.go | 7 ++ 5 files changed, 156 insertions(+) create mode 100644 server/internal/services/defaults.go create mode 100644 server/internal/services/defaults_test.go diff --git a/server/cmd/main.go b/server/cmd/main.go index f5b12c2..4aa877f 100644 --- a/server/cmd/main.go +++ b/server/cmd/main.go @@ -31,6 +31,12 @@ func main() { log.Printf("warning: failed to ensure workflow indexes: %v", err) } + if created, updated, err := services.SeedDefaultSteps(); err != nil { + log.Printf("warning: failed to seed default steps: %v", err) + } else { + log.Printf("default steps seeded: %d created, %d updated", created, updated) + } + services.StartLogSweeper() redisAddr := getEnv("REDIS_ADDR", "localhost:6379") diff --git a/server/internal/api/workflows.go b/server/internal/api/workflows.go index 1e061e0..76f183e 100644 --- a/server/internal/api/workflows.go +++ b/server/internal/api/workflows.go @@ -22,6 +22,7 @@ func registerWorkflowRoutes(g *gin.RouterGroup) { g.DELETE("/steps/:id", deleteStep) g.GET("/steps/:id/export", exportStep) g.POST("/steps/import", importStep) + g.POST("/steps/seed-defaults", seedDefaults) g.POST("/steps/parse", parseStep) g.GET("/workflows", listWorkflows) @@ -201,6 +202,16 @@ func exportStep(c *gin.Context) { c.Data(http.StatusOK, "application/json", b) } +func seedDefaults(c *gin.Context) { + created, updated, err := services.SeedDefaultSteps() + if err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + services.LogEvent("workflow.defaults_synced", actorFromCtx(c), "", "", fmt.Sprintf("default steps synced: %d created, %d updated", created, updated)) + c.JSON(http.StatusOK, gin.H{"created": created, "updated": updated}) +} + func importStep(c *gin.Context) { body, err := io.ReadAll(c.Request.Body) if err != nil { diff --git a/server/internal/services/defaults.go b/server/internal/services/defaults.go new file mode 100644 index 0000000..068c335 --- /dev/null +++ b/server/internal/services/defaults.go @@ -0,0 +1,94 @@ +package services + +import ( + "os" + "path/filepath" + "time" + + "github.com/google/uuid" + "github.com/mrhid6/vantage/server/internal/db" + "github.com/mrhid6/vantage/server/internal/models" + "go.mongodb.org/mongo-driver/v2/bson" + "go.mongodb.org/mongo-driver/v2/mongo/options" +) + +// DefaultStepsDir returns the directory holding default step JSON files. +func DefaultStepsDir() string { + dir := os.Getenv("VANTAGE_DEFAULT_STEPS_DIR") + if dir == "" { + dir = filepath.Join("data", "default-steps") + } + _ = os.MkdirAll(dir, 0700) + return dir +} + +// readDefaultStepFiles parses every *.json in the defaults dir into +// source=default library steps (with slug set). Non-json and invalid files are +// skipped silently; a slug is derived from the step name. +func readDefaultStepFiles() ([]models.WorkflowStep, error) { + matches, err := filepath.Glob(filepath.Join(DefaultStepsDir(), "*.json")) + if err != nil { + return nil, err + } + out := []models.WorkflowStep{} + for _, path := range matches { + b, err := os.ReadFile(path) + if err != nil { + continue + } + s, err := ParseStepDoc(b) + if err != nil { + continue + } + s.Source = "default" + s.Slug = Slugify(s.Name) + if s.Slug == "" { + continue + } + out = append(out, s) + } + return out, nil +} + +// SeedDefaultSteps upserts default steps from disk keyed on {slug, source}. +// Re-sync overwrites default-step content; user steps are never touched. +func SeedDefaultSteps() (created, updated int, err error) { + steps, err := readDefaultStepFiles() + if err != nil { + return 0, 0, err + } + ctx, cancel := wfCtx() + defer cancel() + col := db.Col("workflow_steps") + for _, s := range steps { + filter := bson.M{"slug": s.Slug, "source": "default"} + set := bson.M{ + "name": s.Name, + "description": s.Description, + "interpreter": s.Interpreter, + "script": s.Script, + "declared_outputs": s.DeclaredOutputs, + "declared_inputs": s.DeclaredInputs, + "secret_refs": s.SecretRefs, + "updated_at": time.Now(), + } + res, uerr := col.UpdateOne(ctx, filter, bson.M{ + "$set": set, + "$setOnInsert": bson.M{ + "step_id": uuid.New().String(), + "slug": s.Slug, + "source": "default", + "created_at": time.Now(), + }, + }, options.UpdateOne().SetUpsert(true)) + if uerr != nil { + return created, updated, uerr + } + if res.UpsertedCount > 0 { + created++ + } else if res.ModifiedCount > 0 { + updated++ + } + } + return created, updated, nil +} diff --git a/server/internal/services/defaults_test.go b/server/internal/services/defaults_test.go new file mode 100644 index 0000000..8f6a463 --- /dev/null +++ b/server/internal/services/defaults_test.go @@ -0,0 +1,38 @@ +package services + +import ( + "os" + "path/filepath" + "testing" +) + +func TestDefaultStepsDirEnv(t *testing.T) { + dir := filepath.Join(t.TempDir(), "ds") + t.Setenv("VANTAGE_DEFAULT_STEPS_DIR", dir) + got := DefaultStepsDir() + if got != dir { + t.Fatalf("got %q want %q", got, dir) + } + if _, err := os.Stat(dir); err != nil { + t.Fatalf("dir not created: %v", err) + } +} + +func TestReadDefaultStepFiles(t *testing.T) { + dir := t.TempDir() + t.Setenv("VANTAGE_DEFAULT_STEPS_DIR", dir) + good := `{"kind":"vantage.step/v1","name":"Ping Host","interpreter":"bash","script":"ping -c1 x=1 >> $WORKFLOW_ENV"}` + os.WriteFile(filepath.Join(dir, "ping.json"), []byte(good), 0600) + os.WriteFile(filepath.Join(dir, "notes.txt"), []byte("ignore me"), 0600) + + steps, err := readDefaultStepFiles() + if err != nil { + t.Fatal(err) + } + if len(steps) != 1 { + t.Fatalf("want 1 step, got %d", len(steps)) + } + if steps[0].Slug != "ping-host" || steps[0].Source != "default" { + t.Fatalf("bad seed step: %+v", steps[0]) + } +} diff --git a/server/internal/services/workflows.go b/server/internal/services/workflows.go index eb0d519..124d1e1 100644 --- a/server/internal/services/workflows.go +++ b/server/internal/services/workflows.go @@ -25,6 +25,13 @@ func EnsureWorkflowIndexes() error { }); err != nil { return err } + if _, err := db.Col("workflow_steps").Indexes().CreateOne(ctx, mongo.IndexModel{ + Keys: bson.D{{Key: "slug", Value: 1}}, + Options: options.Index().SetUnique(true). + SetPartialFilterExpression(bson.M{"source": "default"}), + }); err != nil { + return err + } if _, err := db.Col("workflows").Indexes().CreateOne(ctx, mongo.IndexModel{ Keys: bson.D{{Key: "workflow_id", Value: 1}}, Options: options.Index().SetUnique(true), }); err != nil {