diff --git a/automation/service/service.go b/automation/service/service.go index e139a6540..01beed031 100644 --- a/automation/service/service.go +++ b/automation/service/service.go @@ -123,7 +123,6 @@ func Initialize(ctx context.Context, log *zap.Logger, s store.Storer, ws websock automation.LoopHandler(Registry(), DefaultWorkflow.parser) automation.CorredorHandler(Registry(), corredor.Service()) automation.EmailHandler(Registry()) - return } diff --git a/automation/types/workflow.go b/automation/types/workflow.go index 4bd0f6ddb..e9eb557f5 100644 --- a/automation/types/workflow.go +++ b/automation/types/workflow.go @@ -4,9 +4,10 @@ import ( "database/sql/driver" "encoding/json" "fmt" + "time" + "github.com/cortezaproject/corteza-server/pkg/expr" "github.com/cortezaproject/corteza-server/pkg/filter" - "time" ) type ( diff --git a/tests/workflows/0001_basics_test.go b/tests/workflows/0001_basics_test.go new file mode 100644 index 000000000..24828b2d0 --- /dev/null +++ b/tests/workflows/0001_basics_test.go @@ -0,0 +1,53 @@ +package workflows + +import ( + "context" + "fmt" + "sync" + "testing" + "time" + + "github.com/cortezaproject/corteza-server/automation/types" + "github.com/stretchr/testify/require" +) + +func Test0001_basics(t *testing.T) { + ctx := bypassRBAC(context.Background()) + ctx, fn := context.WithTimeout(ctx, time.Second) + defer fn() + loadScenario(ctx, t) + + var ( + req = require.New(t) + aux = struct{ Foo int64 }{} + vars, _ = mustExecWorkflow(ctx, t, "basic", types.WorkflowExecParams{}) + ) + + req.NoError(vars.Decode(&aux)) + req.Equal(int64(42), aux.Foo) +} + +func Test0001_basics_detect_datarace_issues(t *testing.T) { + ctx := bypassRBAC(context.Background()) + loadScenarioWithName(ctx, t, "S0001_basics") + + var ( + wg = &sync.WaitGroup{} + ) + + for i := 0; i < 10; i++ { + wg.Add(1) + go t.Run(fmt.Sprintf("%d", i), func(t *testing.T) { + var ( + aux = struct{ Foo int64 }{} + req = require.New(t) + vars, _ = mustExecWorkflow(ctx, t, "basic", types.WorkflowExecParams{}) + ) + req.NoError(vars.Decode(&aux)) + req.Equal(int64(42), aux.Foo) + wg.Done() + }) + } + + wg.Wait() +} diff --git a/tests/workflows/README.md b/tests/workflows/README.md new file mode 100644 index 000000000..ebc7f3bc2 --- /dev/null +++ b/tests/workflows/README.md @@ -0,0 +1 @@ +Cross component workflow tests diff --git a/tests/workflows/main_test.go b/tests/workflows/main_test.go new file mode 100644 index 000000000..cea534a65 --- /dev/null +++ b/tests/workflows/main_test.go @@ -0,0 +1,139 @@ +package workflows + +import ( + "context" + "fmt" + "path" + "testing" + + "github.com/cortezaproject/corteza-server/app" + "github.com/cortezaproject/corteza-server/automation/service" + autTypes "github.com/cortezaproject/corteza-server/automation/types" + "github.com/cortezaproject/corteza-server/pkg/auth" + "github.com/cortezaproject/corteza-server/pkg/envoy" + "github.com/cortezaproject/corteza-server/pkg/envoy/csv" + "github.com/cortezaproject/corteza-server/pkg/envoy/directory" + "github.com/cortezaproject/corteza-server/pkg/envoy/json" + envoyStore "github.com/cortezaproject/corteza-server/pkg/envoy/store" + "github.com/cortezaproject/corteza-server/pkg/envoy/yaml" + "github.com/cortezaproject/corteza-server/pkg/errors" + "github.com/cortezaproject/corteza-server/pkg/eventbus" + "github.com/cortezaproject/corteza-server/pkg/expr" + "github.com/cortezaproject/corteza-server/pkg/id" + "github.com/cortezaproject/corteza-server/store" + sysTypes "github.com/cortezaproject/corteza-server/system/types" + "github.com/cortezaproject/corteza-server/tests/helpers" +) + +var ( + defApp *app.CortezaApp + defStore store.Storer + eventBus = eventbus.New() +) + +func init() { + helpers.RecursiveDotEnvLoad() +} + +func TestMain(m *testing.M) { + ctx := context.Background() + + defApp = helpers.NewIntegrationTestApp(ctx, func(app *app.CortezaApp) (err error) { + defStore = app.Store + eventbus.Set(eventBus) + return nil + }) + + if err := defApp.Activate(ctx); err != nil { + panic(fmt.Errorf("could not activate corteza: %v", err)) + } + + m.Run() +} + +func cleanup(t *testing.T) { + var ( + ctx = context.Background() + ) + + if err := defStore.TruncateAutomationWorkflows(ctx); err != nil { + t.Fatalf("failed to decode scenario data: %v", err) + } +} + +func loadScenario(ctx context.Context, t *testing.T) { + loadScenarioWithName(ctx, t, "S"+t.Name()[4:]) +} + +func loadScenarioWithName(ctx context.Context, t *testing.T, scenario string) { + var ( + err error + ) + + cleanup(t) + + decoded, err := directory.Decode( + ctx, + path.Join("testdata", scenario), + yaml.Decoder(), + csv.Decoder(), + json.Decoder(), + ) + if err != nil { + t.Fatalf("failed to decode scenario data: %v", err) + } + + storeEnc := envoyStore.NewStoreEncoder(defStore, &envoyStore.EncoderConfig{}) + + b := envoy.NewBuilder(storeEnc) + g, err := b.Build(ctx, decoded...) + if err != nil { + t.Fatalf("failed to build structure graph: %v", err) + } + + if err = envoy.Encode(ctx, g, storeEnc); err != nil { + t.Fatalf("failed to build structure graph: %v", err) + } + + // Reload and register workflows + if err = service.DefaultWorkflow.Load(ctx); err != nil { + t.Fatalf("failed to reload workflows: %v", err) + } +} + +func bypassRBAC(ctx context.Context) context.Context { + u := &sysTypes.User{ + ID: id.Next(), + } + + u.SetRoles(auth.BypassRoles().IDs()...) + + return auth.SetIdentityToContext(ctx, u) +} + +func execWorkflow(ctx context.Context, name string, p autTypes.WorkflowExecParams) (*expr.Vars, autTypes.Stacktrace, error) { + wf, err := defStore.LookupAutomationWorkflowByHandle(ctx, name) + if err != nil { + return nil, nil, err + } + + return service.DefaultWorkflow.Exec(ctx, wf.ID, p) +} + +func mustExecWorkflow(ctx context.Context, t *testing.T, name string, p autTypes.WorkflowExecParams) (vars *expr.Vars, strace autTypes.Stacktrace) { + var err error + vars, strace, err = execWorkflow(ctx, name, p) + if err != nil { + if issues, is := err.(autTypes.WorkflowIssueSet); is { + for _, i := range issues { + t.Logf("issue: %s", i.Description) + t.Logf(" %v", i.Culprit) + } + } + + t.Fatalf("could not exec %q: %v", name, errors.Unwrap(err)) + + } + + return +} diff --git a/tests/workflows/testdata/S0001_basics/workflow.yaml b/tests/workflows/testdata/S0001_basics/workflow.yaml new file mode 100644 index 000000000..57dc7b858 --- /dev/null +++ b/tests/workflows/testdata/S0001_basics/workflow.yaml @@ -0,0 +1,18 @@ +workflows: + basic: + enabled: true + trace: true + triggers: + - enabled: true + stepID: 1 + + steps: + - stepID: 1 + kind: expressions + arguments: [ { target: foo, type: Integer, expr: "40" } ] + - stepID: 2 + kind: expressions + arguments: [ { target: foo, type: Integer, expr: "foo + 2" } ] + + paths: + - { parentID: 1, childID: 2 }