Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -281,6 +281,8 @@ atmos workflow [workflow_name] [flags]
|------|-------|-------------|
| `--file` | `-f` | Workflow file (relative to `workflows.base_path`) |
| `--stack` | `-s` | Override stack for all Atmos-type steps |
| `--tags` | | Component tag selector forwarded to all Atmos-type steps (matches any provided tag) |
| `--labels` | | Component label selector forwarded to all Atmos-type steps (matches all provided key-value pairs) |
| `--from-step` | | Start execution from the named step |
| `--dry-run` | | Preview steps without executing |
| `--identity` | | Default identity for steps without explicit identity |
Expand Down
2 changes: 2 additions & 0 deletions cmd/workflow/workflow.go
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,8 @@ func init() {
flags.WithStringFlag("file", "f", "", "Specify workflow file to run (optional if workflow name is unique)"),
flags.WithBoolFlag("dry-run", "", false, "Simulate the workflow without making any changes"),
flags.WithStringFlag("stack", "s", "", "Stack name"),
flags.WithStringSliceFlag("tags", "", nil, "Filter by tags (comma-separated, matches any): --tags=production,tier-1"),
flags.WithStringFlag("labels", "", "", "Filter by labels (comma-separated key=value or key:value pairs, matches all): --labels=cost-center=platform,compliance=sox"),
flags.WithStringFlag("from-step", "", "", "Resume the workflow from the specified step"),
flags.WithStringFlag("identity", "", "", "Identity to use for workflow steps that don't specify their own identity"),
flags.WithEnvVars("file", "ATMOS_WORKFLOW_FILE"),
Expand Down
18 changes: 18 additions & 0 deletions cmd/workflow/workflow_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
package workflow

import (
"testing"

"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)

func TestWorkflowSelectorFlags(t *testing.T) {
tags := workflowCmd.Flags().Lookup("tags")
require.NotNil(t, tags)
assert.Equal(t, "stringSlice", tags.Value.Type())

labels := workflowCmd.Flags().Lookup("labels")
require.NotNil(t, labels)
assert.Equal(t, "string", labels.Value.Type())
}
11 changes: 10 additions & 1 deletion internal/exec/workflow.go
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,14 @@ func ExecuteWorkflowCmd(cmd *cobra.Command, args []string) error {
if err != nil {
return err
}
commandLineTags, err := flags.GetStringSlice("tags")
if err != nil {
return err
}
commandLineLabels, err := flags.GetString("labels")
if err != nil {
return err
}
processStacks := commandLineStack != ""

// InitCliConfig finds and merges CLI configurations in the following order:
Expand Down Expand Up @@ -192,7 +200,8 @@ func ExecuteWorkflowCmd(cmd *cobra.Command, args []string) error {
workflowDefinition = i
}

err = ExecuteWorkflow(atmosConfig, workflowName, workflowPath, &workflowDefinition, dryRun, commandLineStack, fromStep, commandLineIdentity)
err = ExecuteWorkflow(atmosConfig, workflowName, workflowPath, &workflowDefinition, dryRun, commandLineStack, fromStep, commandLineIdentity,
workflowCommandFilters{tags: commandLineTags, labels: commandLineLabels})
if err != nil {
return err
}
Expand Down
4 changes: 4 additions & 0 deletions internal/exec/workflow_control_adapter.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,8 @@ type workflowControlContext struct {
workflowDefinition *schema.WorkflowDefinition
dryRun bool
commandLineStack string
commandLineTags []string
commandLineLabels string
commandLineIdentity string
baseEnv []string
persistentEnv map[string]string
Expand All @@ -28,6 +30,8 @@ func executeWorkflowControlStep(ctx context.Context, control *workflowControlCon
BasePath: control.atmosConfig.BasePath,
BaseEnv: control.baseEnv,
CommandLineStack: control.commandLineStack,
CommandLineTags: control.commandLineTags,
CommandLineLabels: control.commandLineLabels,
CommandLineIdentity: control.commandLineIdentity,
PrepareEnv: func(baseEnv []string, identity string, stepName string, workflowEnv map[string]string, stepEnv map[string]string) ([]string, error) {
return prepareStepEnvironment(baseEnv, identity, stepName, control.authManager, workflowEnv, control.persistentEnv, stepEnv)
Expand Down
2 changes: 2 additions & 0 deletions internal/exec/workflow_no_stacks_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -166,6 +166,8 @@ func createWorkflowCmdForTest() *cobra.Command {
cmd.PersistentFlags().StringP("file", "f", "", "Workflow file")
cmd.PersistentFlags().Bool("dry-run", false, "Dry run")
cmd.PersistentFlags().StringP("stack", "s", "", "Stack")
cmd.PersistentFlags().StringSlice("tags", nil, "Tags")
cmd.PersistentFlags().String("labels", "", "Labels")
cmd.PersistentFlags().String("from-step", "", "From step")
cmd.PersistentFlags().String("identity", "", "Identity")

Expand Down
2 changes: 2 additions & 0 deletions internal/exec/workflow_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -428,6 +428,8 @@ func TestExecuteWorkflowCmd(t *testing.T) {
cmd.PersistentFlags().StringP("file", "f", "", "Workflow file")
cmd.PersistentFlags().Bool("dry-run", false, "Dry run")
cmd.PersistentFlags().StringP("stack", "s", "", "Stack")
cmd.PersistentFlags().StringSlice("tags", nil, "Tags")
cmd.PersistentFlags().String("labels", "", "Labels")
cmd.PersistentFlags().String("from-step", "", "From step")
cmd.PersistentFlags().String("identity", "", "Identity")

Expand Down
34 changes: 19 additions & 15 deletions internal/exec/workflow_utils.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,6 @@ import (
"io"
"os"
"path/filepath"
"slices"
"sort"
"strings"

Expand Down Expand Up @@ -319,6 +318,11 @@ func checkAndMergeDefaultIdentity(atmosConfig *schema.AtmosConfiguration) bool {
return false
}

type workflowCommandFilters struct {
tags []string
labels string
}

// ExecuteWorkflow executes an Atmos workflow.
func ExecuteWorkflow(
atmosConfig schema.AtmosConfiguration,
Expand All @@ -329,8 +333,13 @@ func ExecuteWorkflow(
commandLineStack string,
fromStep string,
commandLineIdentity string,
commandLineFilters ...workflowCommandFilters,
) (retErr error) {
defer perf.Track(&atmosConfig, "exec.ExecuteWorkflow")()
commandFilters := workflowCommandFilters{}
if len(commandLineFilters) > 0 {
commandFilters = commandLineFilters[0]
}
var activeContainer *workflowPkg.ContainerSession
defer func() {
if activeContainer == nil {
Expand Down Expand Up @@ -707,6 +716,8 @@ func ExecuteWorkflow(
workflowDefinition: workflowDefinition,
dryRun: dryRun,
commandLineStack: commandLineStack,
commandLineTags: commandFilters.tags,
commandLineLabels: commandFilters.labels,
commandLineIdentity: stepIdentity,
baseEnv: baseEnv,
persistentEnv: persistentEnv,
Expand Down Expand Up @@ -811,24 +822,17 @@ func ExecuteWorkflow(
args = strings.Fields(command)
}

args = workflowPkg.AppendAtmosStepFlags(args, workflowPkg.AtmosStepFlags{
Stack: finalStack,
Tags: commandFilters.tags,
Labels: commandFilters.labels,
})
Comment thread
coderabbitai[bot] marked this conversation as resolved.
if finalStack != "" {
if idx := slices.Index(args, "--"); idx != -1 {
// Insert before the "--"
// Take everything up to idx, then add "-s", finalStack, then tack on the rest
args = append(args[:idx], append([]string{"-s", finalStack}, args[idx:]...)...)
} else {
// just append at the end
args = append(args, []string{"-s", finalStack}...)
}

log.Debug("Using stack", "stack", finalStack)
}

// Build display command for RenderCommand.
displayCmd := "atmos " + command
if finalStack != "" {
displayCmd = fmt.Sprintf("atmos %s -s %s", command, finalStack)
}
// Build display command from the final arguments so it matches execution.
displayCmd := "atmos " + strings.Join(args, " ")
// Render command before execution if show.command is enabled.
stepPkg.RenderCommand(&step, workflowDefinition, displayCmd)

Expand Down
93 changes: 93 additions & 0 deletions internal/exec/workflow_utils_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package exec
import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"os"
Expand All @@ -25,6 +26,7 @@ import (
githubprovider "github.com/cloudposse/atmos/pkg/ci/providers/github"
cfg "github.com/cloudposse/atmos/pkg/config"
"github.com/cloudposse/atmos/pkg/dependencies"
"github.com/cloudposse/atmos/pkg/diagnostics"
stepPkg "github.com/cloudposse/atmos/pkg/runner/step"
"github.com/cloudposse/atmos/pkg/schema"
)
Expand Down Expand Up @@ -1716,6 +1718,97 @@ func TestExecuteWorkflow_DryRunAtmosStepStackBeforeSeparator(t *testing.T) {
require.NoError(t, err)
}

func TestExecuteWorkflow_ForwardsCommandLineFilters(t *testing.T) {
stacksPath := "../../tests/fixtures/scenarios/workflows"
t.Setenv("ATMOS_CLI_CONFIG_PATH", stacksPath)
t.Setenv("ATMOS_BASE_PATH", stacksPath)

tests := []struct {
name string
step schema.WorkflowStep
expectedStarts int
}{
{
name: "atmos step",
step: schema.WorkflowStep{
Name: "apply",
Type: schema.TaskTypeAtmos,
Command: "terraform apply -- -auto-approve",
},
expectedStarts: 1,
},
{
name: "parallel atmos step",
step: schema.WorkflowStep{
Name: "parallel",
Type: schema.TaskTypeParallel,
Steps: []schema.WorkflowStep{{
Name: "apply",
Type: schema.TaskTypeAtmos,
Command: "terraform apply -- -auto-approve",
}},
},
expectedStarts: 1,
},
{
name: "matrix atmos step",
step: schema.WorkflowStep{
Name: "matrix",
Type: schema.TaskTypeMatrix,
Matrix: map[string][]string{"region": {"us-east-1", "us-west-2"}},
Steps: []schema.WorkflowStep{{
Name: "apply",
Type: schema.TaskTypeAtmos,
Command: "terraform apply -- -auto-approve",
}},
},
expectedStarts: 2,
},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
atmosConfig, err := cfg.InitCliConfig(schema.ConfigAndStacksInfo{}, false)
require.NoError(t, err)
diagnosticsPath := filepath.Join(t.TempDir(), "events.jsonl")
atmosConfig.Diagnostics = schema.Diagnostics{Enabled: true, File: diagnosticsPath}

err = ExecuteWorkflow(
atmosConfig,
"selector-forwarding",
"/path/to/workflow.yaml",
&schema.WorkflowDefinition{Steps: []schema.WorkflowStep{tt.step}},
true,
"tenant1-ue2-dev",
"",
"",
workflowCommandFilters{tags: []string{"networking"}, labels: "deployment:dev"},
)
require.NoError(t, err)

data, err := os.ReadFile(diagnosticsPath)
require.NoError(t, err)
var starts []diagnostics.Event
for _, line := range strings.Split(strings.TrimSpace(string(data)), "\n") {
var event diagnostics.Event
require.NoError(t, json.Unmarshal([]byte(line), &event))
if event.Type == "process.start" && event.Command == "atmos" {
starts = append(starts, event)
}
}

require.Len(t, starts, tt.expectedStarts)
for _, start := range starts {
assert.Equal(t, []string{
"terraform", "apply", "-s", "tenant1-ue2-dev",
"--tags=networking", "--labels=deployment:dev",
"--", "-auto-approve",
}, start.Args)
}
})
}
}

type workflowRecordingGitHubProvider struct {
*githubprovider.Provider
output *bytes.Buffer
Expand Down
37 changes: 32 additions & 5 deletions pkg/workflow/control_executor.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import (
"mvdan.cc/sh/v3/shell"

iolib "github.com/cloudposse/atmos/pkg/io"
"github.com/cloudposse/atmos/pkg/perf"
"github.com/cloudposse/atmos/pkg/process"
"github.com/cloudposse/atmos/pkg/retry"
"github.com/cloudposse/atmos/pkg/schema"
Expand Down Expand Up @@ -56,6 +57,8 @@ type ControlCommandExecutor struct {
BasePath string
BaseEnv []string
CommandLineStack string
CommandLineTags []string
CommandLineLabels string
CommandLineIdentity string
PrepareEnv ControlEnvironmentFunc
RunCommand ControlCommandRunner
Expand Down Expand Up @@ -191,7 +194,11 @@ func (executor *ControlCommandExecutor) executeAtmos(ctx context.Context, step *
if parseErr != nil {
args = strings.Fields(step.Command)
}
args = appendControlStack(args, executor.finalStack(step))
args = AppendAtmosStepFlags(args, AtmosStepFlags{
Stack: executor.finalStack(step),
Tags: executor.CommandLineTags,
Labels: executor.CommandLineLabels,
})
dir := executor.workingDirectory(step)

ioSpec := executor.commandStreams(output)
Expand Down Expand Up @@ -291,14 +298,34 @@ func executeControlSleep(ctx context.Context, step *schema.WorkflowStep) (*Contr
}
}

func appendControlStack(args []string, stack string) []string {
if stack == "" {
// AtmosStepFlags are command-line flags applied to an Atmos workflow step.
type AtmosStepFlags struct {
Stack string
Tags []string
Labels string
}

// AppendAtmosStepFlags inserts workflow command flags before a pass-through separator.
func AppendAtmosStepFlags(args []string, flags AtmosStepFlags) []string {
defer perf.Track(nil, "workflow.AppendAtmosStepFlags")()

injected := make([]string, 0, 4)
if flags.Stack != "" {
injected = append(injected, "-s", flags.Stack)
}
if len(flags.Tags) > 0 {
injected = append(injected, "--tags="+strings.Join(flags.Tags, ","))
}
if flags.Labels != "" {
injected = append(injected, "--labels="+flags.Labels)
}
if len(injected) == 0 {
return args
}
if idx := indexOfControlArg(args, "--"); idx != -1 {
return append(args[:idx], append([]string{"-s", stack}, args[idx:]...)...)
return append(args[:idx], append(injected, args[idx:]...)...)
}
return append(args, "-s", stack)
return append(args, injected...)
}

func indexOfControlArg(values []string, needle string) int {
Expand Down
Loading
Loading