diff --git a/Makefile b/Makefile index f0b19c9..3c9ea76 100644 --- a/Makefile +++ b/Makefile @@ -25,7 +25,6 @@ apps: @temporal workflow result --workflow-id apps-manual test: - cd controller && go test ./... cd test && go test fmt: @@ -36,9 +35,7 @@ fmt: . terragrunt hcl format cd infra/_modules && tofu fmt -recursive - cd controller && go fmt ./... cd infra/_modules/tfstate && go fmt ./... - cd infra/staging && go fmt ./... cd test && go fmt ./... tidy: fmt diff --git a/controller/Dockerfile b/controller/Dockerfile deleted file mode 100644 index d03aa9a..0000000 --- a/controller/Dockerfile +++ /dev/null @@ -1,26 +0,0 @@ -FROM docker.io/golang:1.24.3-alpine AS builder - -WORKDIR /src - -COPY go.mod go.sum ./ - -RUN go mod download - -COPY . . - -RUN go build -o /bin/worker ./worker - -FROM docker.io/nixos/nix - -RUN echo "experimental-features = flakes nix-command" >> /etc/nix/nix.conf - -# TODO use native nix develop, currently it's a bit slow -RUN nix-env --install --quiet --attr \ - nixpkgs.kubernetes-helm \ - nixpkgs.opentofu \ - nixpkgs.oras \ - nixpkgs.terragrunt - -COPY --from=builder /bin/worker /bin/worker - -CMD ["/bin/worker"] diff --git a/controller/activities/activity_test.go b/controller/activities/activity_test.go deleted file mode 100644 index 99f058b..0000000 --- a/controller/activities/activity_test.go +++ /dev/null @@ -1,264 +0,0 @@ -package activities - -import ( - "context" - "strings" - "testing" - "time" - - "github.com/stretchr/testify/assert" - "github.com/stretchr/testify/suite" - "go.temporal.io/sdk/testsuite" -) - -type ActivityTestSuite struct { - suite.Suite - testsuite.WorkflowTestSuite - - env *testsuite.TestActivityEnvironment -} - -func (s *ActivityTestSuite) SetupTest() { - s.env = s.NewTestActivityEnvironment() - s.env.SetTestTimeout(30 * time.Second) -} - -// Test graph pruning logic with direct calls -func (s *ActivityTestSuite) TestPruneGraph_Success() { - ctx := context.Background() - originalGraph := &Graph{ - Nodes: map[string]bool{ - "vpc": true, - "database": true, - "app": true, - }, - Edges: map[string][]string{ - "database": {"vpc"}, - "app": {"database"}, - }, - } - changedFiles := []string{"database"} - - prunedGraph, err := PruneGraph(ctx, originalGraph, changedFiles) - - s.NoError(err) - s.True(prunedGraph.Nodes["database"]) // changed module should be included - s.True(prunedGraph.Nodes["app"]) // dependent should be included - s.False(prunedGraph.Nodes["vpc"]) // non-dependent should be pruned -} - -func (s *ActivityTestSuite) TestPruneGraph_EmptyChanges() { - ctx := context.Background() - originalGraph := &Graph{ - Nodes: map[string]bool{ - "vpc": true, - "database": true, - }, - Edges: map[string][]string{ - "database": {"vpc"}, - }, - } - changedFiles := []string{} // No changes - - prunedGraph, err := PruneGraph(ctx, originalGraph, changedFiles) - - s.NoError(err) - s.Empty(prunedGraph.Nodes) // no changes means empty graph -} - -func (s *ActivityTestSuite) TestPruneGraph_ComplexDependencies() { - ctx := context.Background() - // Complex graph: monitoring -> app -> [database, cache] -> vpc - originalGraph := &Graph{ - Nodes: map[string]bool{ - "vpc": true, - "database": true, - "cache": true, - "app": true, - "monitoring": true, - }, - Edges: map[string][]string{ - "database": {"vpc"}, - "cache": {"vpc"}, - "app": {"database", "cache"}, - "monitoring": {"app"}, - }, - } - changedFiles := []string{"database"} // Only database changed - - prunedGraph, err := PruneGraph(ctx, originalGraph, changedFiles) - - s.NoError(err) - s.True(prunedGraph.Nodes["database"]) // changed module - s.True(prunedGraph.Nodes["app"]) // direct dependent - s.True(prunedGraph.Nodes["monitoring"]) // transitive dependent - s.False(prunedGraph.Nodes["vpc"]) // not a dependent - s.False(prunedGraph.Nodes["cache"]) // not a dependent -} - -// Test using TestActivityEnvironment for activities that need proper context -func (s *ActivityTestSuite) TestTerragruntPrune_WithActivityEnvironment() { - // Test data - graph := &Graph{ - Nodes: map[string]bool{ - "vpc": true, - "database": true, - "app": true, - "monitoring": true, - }, - Edges: map[string][]string{ - "app": {"database", "vpc"}, - "database": {"vpc"}, - "monitoring": {"app"}, - }, - } - - changedModules := []string{"database"} - - s.env.RegisterActivity(PruneGraph) - - val, err := s.env.ExecuteActivity(PruneGraph, graph, changedModules) - s.NoError(err) - - var result *Graph - s.NoError(val.Get(&result)) - - // Only database (changed) and its dependents (app, monitoring) should be included - // vpc is not included because nothing depends on it - expectedNodes := []string{"database", "app", "monitoring"} - actualNodes := result.GetNodes() - s.ElementsMatch(expectedNodes, actualNodes) - - s.Contains(result.Nodes, "database") - s.Contains(result.Nodes, "app") - s.Contains(result.Nodes, "monitoring") - s.NotContains(result.Nodes, "vpc") -} - -func TestActivityTestSuite(t *testing.T) { - suite.Run(t, new(ActivityTestSuite)) -} - -// Additional comprehensive unit tests for graph functions -func TestNewGraphFromDot_EmptyGraph(t *testing.T) { - dotString := `digraph { - }` - - graph, err := NewGraphFromDot(dotString) - - assert.NoError(t, err) - assert.Empty(t, graph.Nodes) -} - -func TestNewGraphFromDot_InvalidFormat(t *testing.T) { - dotString := `not a valid dot format` - - graph, err := NewGraphFromDot(dotString) - - assert.NoError(t, err) // Should not error, just ignore invalid lines - assert.Empty(t, graph.Nodes) -} - -func TestGraph_TopologicalSort_CyclicGraph(t *testing.T) { - // Create a graph with a cycle: A -> B -> C -> A - graph := &Graph{ - Nodes: map[string]bool{ - "a": true, - "b": true, - "c": true, - }, - Edges: map[string][]string{ - "a": {"b"}, - "b": {"c"}, - "c": {"a"}, // Creates cycle - }, - } - - levels := graph.TopologicalSort() - - // Should handle cycles gracefully by putting remaining nodes in final level - assert.Greater(t, len(levels), 0) - - // All nodes should be present somewhere in the levels - allNodes := make(map[string]bool) - for _, level := range levels { - for _, node := range level { - allNodes[node] = true - } - } - assert.True(t, allNodes["a"]) - assert.True(t, allNodes["b"]) - assert.True(t, allNodes["c"]) -} - -func TestExtractQuoted(t *testing.T) { - // Test the extractQuoted function indirectly through NewGraphFromDot - dotString := `digraph { - "hello" -> "world"; - "test"; - }` - - graph, err := NewGraphFromDot(dotString) - - assert.NoError(t, err) - assert.True(t, graph.Nodes["hello"]) - assert.True(t, graph.Nodes["world"]) - assert.True(t, graph.Nodes["test"]) -} - -func TestGraph_AddEdge_CreatesNodes(t *testing.T) { - graph := NewGraph() - - graph.AddEdge("a", "b") - - assert.True(t, graph.Nodes["a"]) - assert.True(t, graph.Nodes["b"]) - assert.Contains(t, graph.Edges["a"], "b") -} - -func TestGraph_GetNodes(t *testing.T) { - graph := &Graph{ - Nodes: map[string]bool{ - "vpc": true, - "database": true, - "app": true, - }, - Edges: map[string][]string{}, - } - - nodes := graph.GetNodes() - - assert.Len(t, nodes, 3) - assert.Contains(t, nodes, "vpc") - assert.Contains(t, nodes, "database") - assert.Contains(t, nodes, "app") -} - -func TestClone_PathGeneration(t *testing.T) { - // Test that generateRepoPath creates deterministic paths - url1 := "https://github.com/example/repo.git" - revision1 := "main" - - path1 := generateRepoPath(url1, revision1) - path2 := generateRepoPath(url1, revision1) - - // Same inputs should generate same path - assert.Equal(t, path1, path2) - - // Different inputs should generate different paths - path3 := generateRepoPath(url1, "develop") - assert.NotEqual(t, path1, path3) - - path4 := generateRepoPath("https://github.com/other/repo.git", revision1) - assert.NotEqual(t, path1, path4) - - // Paths should be under /tmp/cloudlab-repos/ - assert.True(t, strings.HasPrefix(path1, "/tmp/cloudlab-repos/")) -} - -func TestClone_CheckRepoStatus(t *testing.T) { - // Test hasCorrectRevision with non-existent directory - nonExistentPath := "/tmp/non-existent-repo-12345" - hasCorrect := hasCorrectRevision(context.Background(), nonExistentPath, "main") - assert.False(t, hasCorrect) -} diff --git a/controller/activities/app.go b/controller/activities/app.go deleted file mode 100644 index 186144c..0000000 --- a/controller/activities/app.go +++ /dev/null @@ -1,170 +0,0 @@ -package activities - -import ( - "bytes" - "context" - "fmt" - "os" - "os/exec" - "path" - "path/filepath" - "strings" - - "go.temporal.io/sdk/activity" - "gopkg.in/yaml.v3" -) - -func PushRenderedApp(ctx context.Context, appsPath, namespace, app, cluster, registry string) (*PushResult, error) { - logger := activity.GetLogger(ctx) - - tmpDir, err := os.MkdirTemp("", fmt.Sprintf("%s-%s-", app, cluster)) - if err != nil { - logger.Error("failed to create temp dir", "error", err) - return nil, err - } - defer os.RemoveAll(tmpDir) - - cmd := exec.CommandContext( - ctx, - "helm", "template", - "--namespace", namespace, - app, - "oci://ghcr.io/bjw-s-labs/helm/app-template:4.1.1", - "--values", path.Join(namespace, app, cluster+".yaml"), - ) - cmd.Dir = appsPath - - var stdout, stderr bytes.Buffer - cmd.Stdout = &stdout - cmd.Stderr = &stderr - - logger.Info("running helm template", "cmd", cmd.String()) - - if err := cmd.Run(); err != nil { - logger.Error("helm template failed", "error", err, "stderr", stderr.String()) - return nil, err - } - - if err := os.WriteFile(filepath.Join(tmpDir, "rendered.yaml"), stdout.Bytes(), 0644); err != nil { - logger.Error("failed to write rendered output to file", "error", err) - return nil, err - } - - outputPath, err := filepath.Abs(tmpDir) - if err != nil { - logger.Error("failed to get absolute path to rendered manifests", "error", err) - return nil, err - } - - imageRef := fmt.Sprintf("%s/%s/%s:%s", registry, namespace, app, cluster) - result, err := PushManifests(ctx, outputPath, imageRef) - if err != nil { - logger.Error("failed to push manifests", "error", err) - return nil, err - } - - return result, nil -} - -func DiscoverApps(ctx context.Context, appsDir string, cluster string) ([]string, error) { - // TODO logs - _ = activity.GetLogger(ctx) - var matched []string - err := filepath.Walk(appsDir, func(path string, info os.FileInfo, err error) error { - if err != nil || info.IsDir() { - return nil - } - if strings.HasSuffix(info.Name(), cluster+".yaml") { - matched = append(matched, path) - } - return nil - }) - if err != nil { - return nil, err - } - return matched, nil -} - -type Image struct { - Repository string - Tag string -} - -func updateImageTags(node *yaml.Node, newImages []Image) (bool, error) { - changed := false - var walk func(n *yaml.Node) - walk = func(n *yaml.Node) { - if n.Kind != yaml.MappingNode { - for _, child := range n.Content { - walk(child) - } - return - } - for i := 0; i < len(n.Content)-1; i += 2 { - key := n.Content[i] - val := n.Content[i+1] - if key.Value == "image" && val.Kind == yaml.MappingNode { - var repoNode, tagNode *yaml.Node - for j := 0; j < len(val.Content)-1; j += 2 { - k := val.Content[j] - v := val.Content[j+1] - switch k.Value { - case "repository": - repoNode = v - case "tag": - tagNode = v - } - } - if repoNode != nil && tagNode != nil { - for _, img := range newImages { - if repoNode.Value == img.Repository && tagNode.Value != img.Tag { - tagNode.Value = img.Tag - changed = true - } - } - } - } else { - walk(val) - } - } - } - walk(node) - return changed, nil -} - -func UpdateAppVersion(ctx context.Context, appsDir, namespace, app, cluster string, newImages []Image) (bool, error) { - path := filepath.Join(appsDir, namespace, app, fmt.Sprintf("%s.yaml", cluster)) - - data, err := os.ReadFile(path) - if err != nil { - return false, fmt.Errorf("failed to read file: %w", err) - } - - var node yaml.Node - if err := yaml.Unmarshal(data, &node); err != nil { - return false, fmt.Errorf("failed to unmarshal YAML: %w", err) - } - - changed, err := updateImageTags(&node, newImages) - if err != nil { - return false, fmt.Errorf("failed to update image tags: %w", err) - } - - if changed { - var buf bytes.Buffer - encoder := yaml.NewEncoder(&buf) - encoder.SetIndent(2) - - if err := encoder.Encode(&node); err != nil { - return false, fmt.Errorf("failed to encode YAML: %w", err) - } - encoder.Close() - - newData := buf.Bytes() - if err := os.WriteFile(path, newData, 0644); err != nil { - return false, fmt.Errorf("failed to write YAML file: %w", err) - } - } - - return changed, nil -} diff --git a/controller/activities/app_test.go b/controller/activities/app_test.go deleted file mode 100644 index d86588f..0000000 --- a/controller/activities/app_test.go +++ /dev/null @@ -1,473 +0,0 @@ -package activities - -import ( - "context" - "fmt" - "os" - "path/filepath" - "testing" - - "github.com/stretchr/testify/assert" - "github.com/stretchr/testify/require" - "gopkg.in/yaml.v3" -) - -func TestUpdateAppVersion(t *testing.T) { - tests := []struct { - name string - yamlContent string - newImages []Image - expectedUpdate bool - expectError bool - }{ - { - name: "blog app production update", - yamlContent: `defaultPodOptions: - labels: - "istio.io/dataplane-mode": "ambient" -controllers: - main: - replicas: 2 - strategy: RollingUpdate - containers: - main: - image: - repository: docker.io/khuedoan/blog - tag: 6fbd90b77a81e0bcb330fddaa230feff744a7010 -service: - main: - controller: main - ports: - http: - port: 3000 - protocol: HTTP`, - newImages: []Image{ - {Repository: "docker.io/khuedoan/blog", Tag: "abc123def456789"}, - }, - expectedUpdate: true, - }, - { - name: "actualbudget app version update", - yamlContent: `defaultPodOptions: - labels: - "istio.io/dataplane-mode": "ambient" -controllers: - main: - containers: - main: - image: - repository: docker.io/actualbudget/actual-server - tag: 25.6.1-alpine -service: - main: - controller: main - ports: - http: - port: 5006 - protocol: HTTP`, - newImages: []Image{ - {Repository: "docker.io/actualbudget/actual-server", Tag: "25.7.0-alpine"}, - }, - expectedUpdate: true, - }, - { - name: "notes app with ghcr registry", - yamlContent: `defaultPodOptions: - labels: - istio.io/dataplane-mode: ambient -controllers: - main: - type: statefulset - containers: - main: - image: - repository: ghcr.io/silverbulletmd/silverbullet - tag: v2 - envFrom: - - secret: silverbullet -service: - main: - controller: main - ports: - http: - port: 3000 - protocol: HTTP`, - newImages: []Image{ - {Repository: "ghcr.io/silverbulletmd/silverbullet", Tag: "v3"}, - }, - expectedUpdate: true, - }, - { - name: "example service with local registry", - yamlContent: `defaultPodOptions: - labels: - istio.io/dataplane-mode: ambient -controllers: - main: - replicas: 2 - strategy: RollingUpdate - containers: - main: - image: - repository: registry.registry.svc.cluster.local/example-service - tag: 828c31f942e8913ab2af53a2841c180586c5b7e1 -service: - main: - controller: main - ports: - http: - port: 8080 - protocol: HTTP`, - newImages: []Image{ - {Repository: "registry.registry.svc.cluster.local/example-service", Tag: "abc123def456789012345678901234567890abcd"}, - }, - expectedUpdate: true, - }, - { - name: "no matching repository", - yamlContent: `defaultPodOptions: - labels: - istio.io/dataplane-mode: ambient -controllers: - main: - containers: - main: - image: - repository: docker.io/khuedoan/blog - tag: 6fbd90b77a81e0bcb330fddaa230feff744a7010`, - newImages: []Image{ - {Repository: "docker.io/different/app", Tag: "newversion"}, - }, - expectedUpdate: false, - }, - { - name: "multiple images same yaml - partial update", - yamlContent: `defaultPodOptions: - labels: - istio.io/dataplane-mode: ambient -controllers: - frontend: - containers: - main: - image: - repository: docker.io/khuedoan/blog - tag: 6fbd90b77a81e0bcb330fddaa230feff744a7010 - backend: - containers: - main: - image: - repository: ghcr.io/silverbulletmd/silverbullet - tag: v2`, - newImages: []Image{ - {Repository: "docker.io/khuedoan/blog", Tag: "newcommithash123"}, - }, - expectedUpdate: true, - }, - { - name: "malformed yaml structure", - yamlContent: `controllers: - main: - containers: - main: - image: - repository: docker.io/test/app - # missing tag field -service: - main: invalid yaml structure`, - newImages: []Image{ - {Repository: "docker.io/test/app", Tag: "v2.0.0"}, - }, - expectError: false, // Should handle gracefully - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - // Create temporary directory structure - tempDir, err := os.MkdirTemp("", "test-update-app-") - require.NoError(t, err) - defer os.RemoveAll(tempDir) - - namespace := "test-ns" - app := "test-app" - cluster := "test-cluster" - - // Create directory structure - appDir := filepath.Join(tempDir, namespace, app) - err = os.MkdirAll(appDir, 0755) - require.NoError(t, err) - - // Write test YAML file - yamlPath := filepath.Join(appDir, fmt.Sprintf("%s.yaml", cluster)) - err = os.WriteFile(yamlPath, []byte(tt.yamlContent), 0644) - require.NoError(t, err) - - // Execute UpdateAppVersion - ctx := context.Background() - changed, err := UpdateAppVersion(ctx, tempDir, namespace, app, cluster, tt.newImages) - - if tt.expectError { - assert.Error(t, err) - return - } - - require.NoError(t, err) - assert.Equal(t, tt.expectedUpdate, changed, "Expected change result doesn't match") - - // Read the updated file - updatedContent, err := os.ReadFile(yamlPath) - require.NoError(t, err) - - // Parse the updated YAML - var updatedData map[string]interface{} - err = yaml.Unmarshal(updatedContent, &updatedData) - require.NoError(t, err) - - // Verify updates were applied correctly - if tt.expectedUpdate { - verifyImageUpdates(t, updatedData, tt.newImages) - } - }) - } -} - -func TestUpdateAppVersion_FileErrors(t *testing.T) { - ctx := context.Background() - tempDir, err := os.MkdirTemp("", "test-update-app-errors-") - require.NoError(t, err) - defer os.RemoveAll(tempDir) - - t.Run("non-existent file", func(t *testing.T) { - _, err := UpdateAppVersion(ctx, tempDir, "ns", "app", "cluster", []Image{}) - assert.Error(t, err) - assert.Contains(t, err.Error(), "failed to read file") - }) - - t.Run("invalid yaml", func(t *testing.T) { - namespace := "test-ns" - app := "test-app" - cluster := "test-cluster" - - appDir := filepath.Join(tempDir, namespace, app) - err = os.MkdirAll(appDir, 0755) - require.NoError(t, err) - - yamlPath := filepath.Join(appDir, fmt.Sprintf("%s.yaml", cluster)) - err = os.WriteFile(yamlPath, []byte("invalid: yaml: content: ["), 0644) - require.NoError(t, err) - - _, err = UpdateAppVersion(ctx, tempDir, namespace, app, cluster, []Image{}) - assert.Error(t, err) - assert.Contains(t, err.Error(), "failed to unmarshal YAML") - }) -} - -func TestUpdateImageTags(t *testing.T) { - tests := []struct { - name string - yamlContent string - newImages []Image - expectedTags map[string]string // repository -> expected tag - }{ - { - name: "blog app git hash update", - yamlContent: `controllers: - main: - containers: - main: - image: - repository: docker.io/khuedoan/blog - tag: 6fbd90b77a81e0bcb330fddaa230feff744a7010`, - newImages: []Image{ - {Repository: "docker.io/khuedoan/blog", Tag: "abc123def456789"}, - }, - expectedTags: map[string]string{ - "docker.io/khuedoan/blog": "abc123def456789", - }, - }, - { - name: "actualbudget version update", - yamlContent: `controllers: - main: - containers: - main: - image: - repository: docker.io/actualbudget/actual-server - tag: 25.6.1-alpine`, - newImages: []Image{ - {Repository: "docker.io/actualbudget/actual-server", Tag: "25.7.0-alpine"}, - }, - expectedTags: map[string]string{ - "docker.io/actualbudget/actual-server": "25.7.0-alpine", - }, - }, - { - name: "mixed registries partial update", - yamlContent: `controllers: - main: - containers: - main: - image: - repository: ghcr.io/silverbulletmd/silverbullet - tag: v2 - worker: - containers: - worker: - image: - repository: docker.io/actualbudget/actual-server - tag: 25.6.1-alpine`, - newImages: []Image{ - {Repository: "ghcr.io/silverbulletmd/silverbullet", Tag: "v3"}, - }, - expectedTags: map[string]string{ - "ghcr.io/silverbulletmd/silverbullet": "v3", - "docker.io/actualbudget/actual-server": "25.6.1-alpine", // unchanged - }, - }, - { - name: "local registry with full real structure", - yamlContent: `defaultPodOptions: - labels: - istio.io/dataplane-mode: ambient -controllers: - main: - replicas: 2 - strategy: RollingUpdate - containers: - main: - image: - repository: registry.registry.svc.cluster.local/example-service - tag: 828c31f942e8913ab2af53a2841c180586c5b7e1 -service: - main: - controller: main - ports: - http: - port: 8080 - protocol: HTTP`, - newImages: []Image{ - {Repository: "registry.registry.svc.cluster.local/example-service", Tag: "newgithash12345678901234567890"}, - }, - expectedTags: map[string]string{ - "registry.registry.svc.cluster.local/example-service": "newgithash12345678901234567890", - }, - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - var node yaml.Node - err := yaml.Unmarshal([]byte(tt.yamlContent), &node) - require.NoError(t, err) - - _, err = updateImageTags(&node, tt.newImages) - require.NoError(t, err) - - // Marshall back to verify changes - updatedYAML, err := yaml.Marshal(&node) - require.NoError(t, err) - - var updatedData map[string]interface{} - err = yaml.Unmarshal(updatedYAML, &updatedData) - require.NoError(t, err) - - // Verify the expected tag updates - for expectedRepo, expectedTag := range tt.expectedTags { - found := false - findImageTag(updatedData, expectedRepo, expectedTag, &found) - assert.True(t, found, "Expected to find repository %s with tag %s", expectedRepo, expectedTag) - } - }) - } -} - -// Helper function to verify image updates in parsed YAML data -func verifyImageUpdates(t *testing.T, data map[string]interface{}, expectedImages []Image) { - for _, img := range expectedImages { - found := false - findImageTag(data, img.Repository, img.Tag, &found) - assert.True(t, found, "Expected to find repository %s with tag %s", img.Repository, img.Tag) - } -} - -// Recursive helper to find image tags in nested YAML structure -func findImageTag(data interface{}, targetRepo, expectedTag string, found *bool) { - switch v := data.(type) { - case map[string]interface{}: - if imageMap, ok := v["image"].(map[string]interface{}); ok { - if repo, repoOk := imageMap["repository"].(string); repoOk && repo == targetRepo { - if tag, tagOk := imageMap["tag"].(string); tagOk && tag == expectedTag { - *found = true - return - } - } - } - for _, value := range v { - findImageTag(value, targetRepo, expectedTag, found) - } - case []interface{}: - for _, item := range v { - findImageTag(item, targetRepo, expectedTag, found) - } - } -} - -func TestUpdateAppVersion_YAMLIndentation(t *testing.T) { - // Test that YAML is written with 2-space indentation - tempDir, err := os.MkdirTemp("", "test-yaml-indent-") - require.NoError(t, err) - defer os.RemoveAll(tempDir) - - namespace := "test" - app := "indent-test" - cluster := "local" - - // Create directory structure - appDir := filepath.Join(tempDir, namespace, app) - err = os.MkdirAll(appDir, 0755) - require.NoError(t, err) - - // Create a test YAML file with nested structure - yamlContent := `controllers: - main: - containers: - main: - image: - repository: docker.io/test/app - tag: v1.0.0 -service: - main: - controller: main - ports: - http: - port: 8080` - - yamlPath := filepath.Join(appDir, fmt.Sprintf("%s.yaml", cluster)) - err = os.WriteFile(yamlPath, []byte(yamlContent), 0644) - require.NoError(t, err) - - // Update with new image - newImages := []Image{ - {Repository: "docker.io/test/app", Tag: "v2.0.0"}, - } - - ctx := context.Background() - _, err = UpdateAppVersion(ctx, tempDir, namespace, app, cluster, newImages) - require.NoError(t, err) - - // Read the updated file and check indentation - updatedContent, err := os.ReadFile(yamlPath) - require.NoError(t, err) - - contentStr := string(updatedContent) - - // Check that nested elements use 2-space indentation - assert.Contains(t, contentStr, "controllers:\n main:") - assert.Contains(t, contentStr, " main:\n containers:") - assert.Contains(t, contentStr, " containers:\n main:") - assert.Contains(t, contentStr, " main:\n image:") - assert.Contains(t, contentStr, " image:\n repository:") - - // Verify the tag was actually updated - assert.Contains(t, contentStr, "tag: v2.0.0") -} diff --git a/controller/activities/git.go b/controller/activities/git.go deleted file mode 100644 index deb1bf1..0000000 --- a/controller/activities/git.go +++ /dev/null @@ -1,143 +0,0 @@ -package activities - -import ( - "context" - "crypto/sha256" - "fmt" - "os" - "os/exec" - "path/filepath" - "strings" - - "go.temporal.io/sdk/activity" -) - -func generateRepoPath(url string, revision string) string { - hash := sha256.Sum256([]byte(url + ":" + revision)) - return filepath.Join("/tmp", "cloudlab-repos", fmt.Sprintf("%x", hash)[:16]) -} - -func hasCorrectRevision(ctx context.Context, path, revision string) bool { - if _, err := os.Stat(filepath.Join(path, ".git")); os.IsNotExist(err) { - return false - } - - cmd := exec.CommandContext(ctx, "git", "rev-parse", revision) - cmd.Dir = path - return cmd.Run() == nil -} - -func Clone(ctx context.Context, url string, revision string) (string, error) { - logger := activity.GetLogger(ctx) - path := generateRepoPath(url, revision) - - if hasCorrectRevision(ctx, path, revision) { - logger.Info("Repository already available", "path", path) - return path, nil - } - - if err := os.MkdirAll(filepath.Dir(path), 0755); err != nil { - return "", fmt.Errorf("failed to create parent directory: %w", err) - } - os.RemoveAll(path) - - logger.Info("Cloning repository", "url", url, "revision", revision) - - cmd := exec.CommandContext(ctx, "git", "clone", "--depth", "1", "--branch", revision, url, path) - if err := cmd.Run(); err != nil { - os.RemoveAll(path) - return "", fmt.Errorf("failed to clone repository: %w", err) - } - - return path, nil -} - -func ChangedModules(ctx context.Context, repoPath string, oldRevision string) ([]string, error) { - logger := activity.GetLogger(ctx) - - // Since we now clone with depth 1, we need to fetch the oldRevision before we can diff against it - logger.Info("Fetching old revision for comparison", "oldRevision", oldRevision) - fetchCmd := exec.CommandContext(ctx, "git", "fetch", "origin", oldRevision) - fetchCmd.Dir = repoPath - if err := fetchCmd.Run(); err != nil { - return nil, fmt.Errorf("failed to fetch old revision %s: %w", oldRevision, err) - } - - cmd := exec.CommandContext(ctx, "git", "diff", "--name-only", oldRevision, "HEAD") - cmd.Dir = repoPath - output, err := cmd.Output() - if err != nil { - return nil, fmt.Errorf("failed to run git diff: %w", err) - } - - seen := make(map[string]struct{}) - var modules []string - - for _, file := range strings.Fields(string(output)) { - if file == "" { - continue - } - - for dir := filepath.Dir(file); dir != "." && dir != "/"; dir = filepath.Dir(dir) { - if _, err := os.Stat(filepath.Join(repoPath, dir, "terragrunt.hcl")); err == nil { - // Remove infra/stack prefix to get module path - if parts := strings.Split(filepath.ToSlash(dir), "/"); len(parts) >= 3 && parts[0] == "infra" { - if module := strings.Join(parts[2:], "/"); module != "" { - if _, exists := seen[module]; !exists { - modules = append(modules, module) - seen[module] = struct{}{} - } - } - } - break - } - } - } - - return modules, nil -} - -func GitAdd(ctx context.Context, path string) error { - logger := activity.GetLogger(ctx) - - dir := filepath.Dir(path) - relPath := filepath.Base(path) - - cmd := exec.Command("git", "-C", dir, "add", relPath) - cmd.Stdout = os.Stdout - cmd.Stderr = os.Stderr - if err := cmd.Run(); err != nil { - logger.Error("git add failed", "error", err) - return fmt.Errorf("git add failed: %w", err) - } - - return nil -} - -func GitCommit(ctx context.Context, dir string, message string) error { - logger := activity.GetLogger(ctx) - - cmd := exec.Command("git", "-C", dir, "commit", "-m", message) - cmd.Stdout = os.Stdout - cmd.Stderr = os.Stderr - if err := cmd.Run(); err != nil { - logger.Error("git commit failed", "error", err) - return fmt.Errorf("git commit failed: %w", err) - } - - return nil -} - -func GitPush(ctx context.Context, dir string) error { - logger := activity.GetLogger(ctx) - - cmd := exec.Command("git", "-C", dir, "push") - cmd.Stdout = os.Stdout - cmd.Stderr = os.Stderr - if err := cmd.Run(); err != nil { - logger.Error("git push failed", "error", err) - return fmt.Errorf("git push failed: %w", err) - } - - return nil -} diff --git a/controller/activities/git_test.go b/controller/activities/git_test.go deleted file mode 100644 index f7f708d..0000000 --- a/controller/activities/git_test.go +++ /dev/null @@ -1,422 +0,0 @@ -package activities - -import ( - "os" - "path/filepath" - "reflect" - "sort" - "strings" - "testing" -) - -func TestChangedModules(t *testing.T) { - // Create a temporary directory structure for testing - tempDir, err := os.MkdirTemp("", "test-git-") - if err != nil { - t.Fatalf("Failed to create temp dir: %v", err) - } - defer os.RemoveAll(tempDir) - - // Create test directory structure - testDirs := []string{ - "infra/dev/core", - "infra/dev/networking", - "infra/dev/databases/postgres", - "infra/prod/core", - "infra/prod/monitoring", - "shared/modules/vpc", - "shared/modules/security", - "docs", - } - - for _, dir := range testDirs { - err := os.MkdirAll(filepath.Join(tempDir, dir), 0755) - if err != nil { - t.Fatalf("Failed to create dir %s: %v", dir, err) - } - } - - // Create terragrunt.hcl files in specific directories - terragruntDirs := []string{ - "infra/dev/core", - "infra/dev/networking", - "infra/dev/databases/postgres", - "infra/prod/core", - "infra/prod/monitoring", - "shared/modules/vpc", - } - - for _, dir := range terragruntDirs { - terragruntPath := filepath.Join(tempDir, dir, "terragrunt.hcl") - err := os.WriteFile(terragruntPath, []byte("# terragrunt config"), 0644) - if err != nil { - t.Fatalf("Failed to create terragrunt.hcl in %s: %v", dir, err) - } - } - - // Test cases with mock changed files - testCases := []struct { - name string - changedFiles []string - expected []string - }{ - { - name: "Single module change", - changedFiles: []string{ - "infra/dev/core/main.tf", - "infra/dev/core/variables.tf", - }, - expected: []string{"core"}, - }, - { - name: "Multiple modules in same environment", - changedFiles: []string{ - "infra/dev/core/main.tf", - "infra/dev/networking/vpc.tf", - "infra/dev/databases/postgres/db.tf", - }, - expected: []string{"core", "networking", "databases/postgres"}, - }, - { - name: "Modules across different environments", - changedFiles: []string{ - "infra/dev/core/main.tf", - "infra/prod/core/main.tf", - "infra/prod/monitoring/alerts.tf", - }, - expected: []string{"core", "monitoring"}, - }, - { - name: "Shared modules without infra prefix", - changedFiles: []string{ - "shared/modules/vpc/main.tf", - "shared/modules/vpc/outputs.tf", - }, - expected: []string{"shared/modules/vpc"}, - }, - { - name: "Files without terragrunt.hcl", - changedFiles: []string{ - "docs/README.md", - "shared/modules/security/policy.tf", // no terragrunt.hcl in security dir - }, - expected: []string{}, - }, - { - name: "Mixed files with and without terragrunt.hcl", - changedFiles: []string{ - "infra/dev/core/main.tf", - "docs/README.md", - "shared/modules/vpc/vpc.tf", - }, - expected: []string{"core", "shared/modules/vpc"}, - }, - } - - for _, tc := range testCases { - t.Run(tc.name, func(t *testing.T) { - // Create a mock ChangedModules function that uses our test data - // We'll create a custom function that simulates the file system checks - modules := getChangedModulesFromFiles(tempDir, tc.changedFiles) - - sort.Strings(modules) - sort.Strings(tc.expected) - - if !reflect.DeepEqual(modules, tc.expected) { - t.Errorf("Expected modules %v, but got %v", tc.expected, modules) - } - }) - } -} - -// Helper function to simulate ChangedModules logic without Git -func getChangedModulesFromFiles(repoPath string, changedFiles []string) []string { - seen := make(map[string]struct{}) - modules := make([]string, 0) // Initialize as empty slice instead of nil - - for _, file := range changedFiles { - // Get the directory of the changed file - dir := filepath.Dir(file) - - // Walk up the directory tree to find the closest directory containing terragrunt.hcl - currentDir := dir - for { - terragruntPath := filepath.Join(repoPath, currentDir, "terragrunt.hcl") - if _, err := os.Stat(terragruntPath); err == nil { - // Found terragrunt.hcl, this is a module directory - modulePath := currentDir - - // Remove infra// prefix if present - if len(modulePath) > 0 && filepath.HasPrefix(modulePath, "infra/") { - parts := strings.Split(filepath.ToSlash(modulePath), "/") - if len(parts) >= 3 && parts[0] == "infra" { - // Remove "infra" and environment (e.g., "dev", "prod") - modulePath = strings.Join(parts[2:], "/") - } - } - - // Skip empty paths - if modulePath != "" && modulePath != "." { - // Normalize path separators to forward slashes - modulePath = filepath.ToSlash(modulePath) - - if _, exists := seen[modulePath]; !exists { - modules = append(modules, modulePath) - seen[modulePath] = struct{}{} - } - } - break - } - - // Move up one directory level - parent := filepath.Dir(currentDir) - if parent == currentDir || parent == "." { - // Reached the root, no terragrunt.hcl found - break - } - currentDir = parent - } - } - - return modules -} - -func TestGitAdd_PathParsing(t *testing.T) { - // Test the path parsing logic in GitAdd without requiring actual git commands - tests := []struct { - name string - inputPath string - expectedDir string - expectedFile string - }{ - { - name: "simple file", - inputPath: "/tmp/test.yaml", - expectedDir: "/tmp", - expectedFile: "test.yaml", - }, - { - name: "nested file", - inputPath: "/apps/namespace/app/cluster.yaml", - expectedDir: "/apps/namespace/app", - expectedFile: "cluster.yaml", - }, - { - name: "relative path", - inputPath: "relative/file.yaml", - expectedDir: "relative", - expectedFile: "file.yaml", - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - // Test the path manipulation logic that GitAdd uses - actualDir := filepath.Dir(tt.inputPath) - actualFile := filepath.Base(tt.inputPath) - - if actualDir != tt.expectedDir { - t.Errorf("Expected directory '%s', got '%s'", tt.expectedDir, actualDir) - } - - if actualFile != tt.expectedFile { - t.Errorf("Expected filename '%s', got '%s'", tt.expectedFile, actualFile) - } - }) - } -} - -func TestGitCommit_PathParsing(t *testing.T) { - // Test the path parsing logic in GitCommit - tests := []struct { - name string - inputPath string - expectedDir string - message string - }{ - { - name: "simple file with default message", - inputPath: "/tmp/test.yaml", - expectedDir: "/tmp", - message: "chore(test/app): update local version", - }, - { - name: "nested file with custom message", - inputPath: "/apps/namespace/app/cluster.yaml", - expectedDir: "/apps/namespace/app", - message: "feat: update application configuration", - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - // Test the path manipulation logic that GitCommit uses - actualDir := filepath.Dir(tt.inputPath) - - if actualDir != tt.expectedDir { - t.Errorf("Expected directory '%s', got '%s'", tt.expectedDir, actualDir) - } - - // Verify message is not empty - if tt.message == "" { - t.Error("Commit message should not be empty") - } - }) - } -} - -func TestGitPush_PathParsing(t *testing.T) { - // Test the path parsing logic in GitPush - tests := []struct { - name string - inputPath string - expectedDir string - }{ - { - name: "simple file", - inputPath: "/tmp/test.yaml", - expectedDir: "/tmp", - }, - { - name: "nested file", - inputPath: "/apps/namespace/app/cluster.yaml", - expectedDir: "/apps/namespace/app", - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - // Test the path manipulation logic that GitPush uses - actualDir := filepath.Dir(tt.inputPath) - - if actualDir != tt.expectedDir { - t.Errorf("Expected directory '%s', got '%s'", tt.expectedDir, actualDir) - } - }) - } -} - -func TestGitActivities_CommandStructure(t *testing.T) { - // Test that the separate git activities construct the expected commands - testPath := "/tmp/test/app/cluster.yaml" - expectedDir := "/tmp/test/app" - expectedFile := "cluster.yaml" - commitMessage := "chore(khuedoan/blog): update production version" - - // Verify the path parsing logic - actualDir := filepath.Dir(testPath) - actualFile := filepath.Base(testPath) - - if actualDir != expectedDir { - t.Errorf("Expected directory '%s', got '%s'", expectedDir, actualDir) - } - - if actualFile != expectedFile { - t.Errorf("Expected filename '%s', got '%s'", expectedFile, actualFile) - } - - // Verify the expected command structures for each activity - tests := []struct { - name string - expectedCommand []string - description string - }{ - { - name: "GitAdd command", - expectedCommand: []string{"git", "-C", expectedDir, "add", expectedFile}, - description: "GitAdd should construct git add command", - }, - { - name: "GitCommit command", - expectedCommand: []string{"git", "-C", expectedDir, "commit", "-m", commitMessage}, - description: "GitCommit should construct git commit command with message", - }, - { - name: "GitPush command", - expectedCommand: []string{"git", "-C", expectedDir, "push"}, - description: "GitPush should construct git push command", - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - cmd := tt.expectedCommand - - if len(cmd) < 3 { - t.Errorf("%s should have at least 3 parts, got %d", tt.description, len(cmd)) - return - } - - if cmd[0] != "git" { - t.Errorf("%s should start with 'git', got '%s'", tt.description, cmd[0]) - } - - if cmd[1] != "-C" { - t.Errorf("%s should have '-C' as second argument, got '%s'", tt.description, cmd[1]) - } - - if cmd[2] != expectedDir { - t.Errorf("%s should use directory '%s', got '%s'", tt.description, expectedDir, cmd[2]) - } - }) - } -} - -func TestGenerateRepoPath(t *testing.T) { - tests := []struct { - name string - url string - revision string - wantPath bool // whether we expect a valid path - }{ - { - name: "simple repo", - url: "https://github.com/user/repo.git", - revision: "main", - wantPath: true, - }, - { - name: "same repo different revision", - url: "https://github.com/user/repo.git", - revision: "develop", - wantPath: true, - }, - { - name: "empty inputs", - url: "", - revision: "", - wantPath: true, // Should still generate a path - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - path := generateRepoPath(tt.url, tt.revision) - - if tt.wantPath { - if path == "" { - t.Error("Expected non-empty path") - } - if !strings.Contains(path, "/tmp/cloudlab-repos/") { - t.Errorf("Expected path to contain '/tmp/cloudlab-repos/', got: %s", path) - } - if len(filepath.Base(path)) != 16 { - t.Errorf("Expected base path to be 16 characters, got: %s", filepath.Base(path)) - } - } - }) - } - - // Test that same inputs generate same path - path1 := generateRepoPath("https://github.com/test/repo.git", "main") - path2 := generateRepoPath("https://github.com/test/repo.git", "main") - if path1 != path2 { - t.Errorf("Same inputs should generate same path: %s != %s", path1, path2) - } - - // Test that different inputs generate different paths - path3 := generateRepoPath("https://github.com/test/repo.git", "develop") - if path1 == path3 { - t.Error("Different revisions should generate different paths") - } -} diff --git a/controller/activities/graph.go b/controller/activities/graph.go deleted file mode 100644 index cf220e6..0000000 --- a/controller/activities/graph.go +++ /dev/null @@ -1,192 +0,0 @@ -package activities - -import ( - "context" - "strings" -) - -type Graph struct { - Nodes map[string]bool `json:"nodes"` - Edges map[string][]string `json:"edges"` -} - -func NewGraph() *Graph { - return &Graph{ - Nodes: make(map[string]bool), - Edges: make(map[string][]string), - } -} - -func (g *Graph) AddNode(name string) { - g.Nodes[name] = true -} - -func (g *Graph) AddEdge(src, dest string) { - g.AddNode(src) - g.AddNode(dest) - g.Edges[src] = append(g.Edges[src], dest) -} - -func (g *Graph) GetNodes() []string { - nodes := make([]string, 0, len(g.Nodes)) - for name := range g.Nodes { - nodes = append(nodes, name) - } - return nodes -} - -func PruneGraph(ctx context.Context, graph *Graph, changed []string) (*Graph, error) { - dependents := make(map[string][]string) - for src, dests := range graph.Edges { - for _, dest := range dests { - dependents[dest] = append(dependents[dest], src) - } - } - - keep := make(map[string]bool) - var visit func(string) - visit = func(node string) { - if keep[node] { - return - } - keep[node] = true - for _, dep := range dependents[node] { - visit(dep) - } - } - - for _, nodeName := range changed { - if graph.Nodes[nodeName] { - visit(nodeName) - } - } - - prunedGraph := NewGraph() - for node := range keep { - prunedGraph.AddNode(node) - } - for src, dests := range graph.Edges { - if keep[src] { - for _, dest := range dests { - if keep[dest] { - prunedGraph.AddEdge(src, dest) - } - } - } - } - - return prunedGraph, nil -} - -func NewGraphFromDot(dot string) (*Graph, error) { - graph := NewGraph() - - lines := strings.SplitSeq(dot, "\n") - for line := range lines { - line = strings.TrimSpace(line) - if line == "" || strings.HasPrefix(line, "//") || line == "digraph {" || line == "}" { - continue - } - - line = strings.TrimSuffix(line, ";") - line = strings.TrimSpace(line) - - if strings.Contains(line, "->") { - parts := strings.Split(line, "->") - if len(parts) == 2 { - src := extractQuotedString(strings.TrimSpace(parts[0])) - dest := extractQuotedString(strings.TrimSpace(parts[1])) - if src != "" && dest != "" { - graph.AddEdge(src, dest) - } - } - } else { - // Parse standalone nodes: "C" - nodeName := extractQuotedString(line) - if nodeName != "" { - graph.AddNode(nodeName) - } - } - } - - return graph, nil -} - -// extractQuotedString extracts the content between quotes from a string like "hello" -func extractQuotedString(s string) string { - s = strings.TrimSpace(s) - if len(s) >= 2 && s[0] == '"' && s[len(s)-1] == '"' { - return s[1 : len(s)-1] - } - return "" -} - -// TopologicalSort returns modules grouped by dependency levels for parallel execution. -// Edge from A to B means A depends on B, so B must run before A. -func (g *Graph) TopologicalSort() [][]string { - // Build adjacency list and in-degree count - adjList := make(map[string][]string) - inDegree := make(map[string]int) - - // Initialize all nodes with in-degree 0 - for node := range g.Nodes { - inDegree[node] = 0 - adjList[node] = []string{} - } - - // Build the graph and calculate in-degrees - // Edge from Src to Dest means Src depends on Dest - // So Dest should run before Src - for src, dests := range g.Edges { - for _, dest := range dests { - adjList[dest] = append(adjList[dest], src) - inDegree[src]++ - } - } - - var levels [][]string - remaining := make(map[string]bool) - for node := range g.Nodes { - remaining[node] = true - } - - // Process nodes level by level - for len(remaining) > 0 { - var currentLevel []string - - // Find all nodes with in-degree 0 (no dependencies) - for nodeName := range remaining { - if inDegree[nodeName] == 0 { - currentLevel = append(currentLevel, nodeName) - } - } - - // If no nodes found with in-degree 0, there's a cycle - if len(currentLevel) == 0 { - // Return remaining nodes as the final level to handle cycles gracefully - var cycleNodes []string - for nodeName := range remaining { - cycleNodes = append(cycleNodes, nodeName) - } - if len(cycleNodes) > 0 { - levels = append(levels, cycleNodes) - } - break - } - - // Add current level - levels = append(levels, currentLevel) - - // Remove processed nodes and update in-degrees - for _, nodeName := range currentLevel { - delete(remaining, nodeName) - for _, dependent := range adjList[nodeName] { - if remaining[dependent] { - inDegree[dependent]-- - } - } - } - } - - return levels -} diff --git a/controller/activities/graph_test.go b/controller/activities/graph_test.go deleted file mode 100644 index f9d2b0a..0000000 --- a/controller/activities/graph_test.go +++ /dev/null @@ -1,404 +0,0 @@ -package activities - -import ( - "context" - "reflect" - "sort" - "testing" -) - -func TestPruneGraphSimple(t *testing.T) { - dot := ` -digraph { - "A" -> "B"; - "B" -> "C"; - "D" -> "B"; - "E" -> "F"; - "C"; - "F"; -} -` - graph, err := NewGraphFromDot(dot) - if err != nil { - t.Fatalf("Failed to create graph from dot: %v", err) - } - - testCases := []struct { - name string - changed []string - expectedNodes []string - expectedEdges map[string][]string - }{ - { - name: "Prune to C and its dependencies", - changed: []string{"C"}, - expectedNodes: []string{"A", "B", "C", "D"}, - expectedEdges: map[string][]string{ - "A": {"B"}, - "B": {"C"}, - "D": {"B"}, - }, - }, - { - name: "Prune to F and its dependencies", - changed: []string{"F"}, - expectedNodes: []string{"E", "F"}, - expectedEdges: map[string][]string{ - "E": {"F"}, - }, - }, - { - name: "Prune to B and its dependencies", - changed: []string{"B"}, - expectedNodes: []string{"A", "B", "D"}, - expectedEdges: map[string][]string{ - "A": {"B"}, - "D": {"B"}, - }, - }, - { - name: "No nodes changed", - changed: []string{}, - expectedNodes: []string{}, - expectedEdges: map[string][]string{}, - }, - { - name: "Changed node not in graph", - changed: []string{"Z"}, - expectedNodes: []string{}, - expectedEdges: map[string][]string{}, - }, - { - name: "Multiple changed nodes", - changed: []string{"C", "F"}, - expectedNodes: []string{"A", "B", "C", "D", "E", "F"}, - expectedEdges: map[string][]string{ - "A": {"B"}, - "B": {"C"}, - "D": {"B"}, - "E": {"F"}, - }, - }, - } - - for _, tc := range testCases { - t.Run(tc.name, func(t *testing.T) { - prunedGraph, err := PruneGraph(context.Background(), graph, tc.changed) - if err != nil { - t.Fatalf("PruneGraph failed: %v", err) - } - - prunedNodes := prunedGraph.GetNodes() - sort.Strings(prunedNodes) - sort.Strings(tc.expectedNodes) - - if !reflect.DeepEqual(prunedNodes, tc.expectedNodes) { - t.Errorf("Expected nodes %v, but got %v", tc.expectedNodes, prunedNodes) - } - - // Compare edges - if !reflect.DeepEqual(prunedGraph.Edges, tc.expectedEdges) { - t.Errorf("Expected edges %v, but got %v", tc.expectedEdges, prunedGraph.Edges) - } - }) - } -} - -func TestPruneGraphRealWorld(t *testing.T) { - realWorldDot := `digraph { - "aks-windows-node-exporter" ; - "azuresql" ; - "azuresql" -> "core"; - "azuresqlusers" ; - "azuresqlusers" -> "azuresql"; - "bootstrap-va" ; - "bootstrap-va" -> "cluster-va"; - "bootstrap2-va" ; - "bootstrap2-va" -> "cluster2-va"; - "cluster-va" ; - "cluster-va" -> "core"; - "cluster2-va" ; - "core" ; - "db/auror-integration" ; - "db/auror-integration" -> "core"; - "db/auror-integration" -> "azuresql"; - "db/doc-chat" ; - "db/doc-chat" -> "core"; - "db/doc-chat" -> "azuresql"; - "dems-cluster-identity" ; - "dems-cluster-identity" -> "cluster-va"; - "dems-search-grpc/cosmosdb-cassandra" ; - "dems-search-grpc/cosmosdb-cassandra" -> "core"; - "doc-chat/openai" ; - "doc-chat/openai" -> "core"; - "doc-chat/openai-fallback" ; - "doc-chat/openai-fallback" -> "core"; - "doc-chat/search-service-va" ; - "doc-chat/search-service-va" -> "core"; - "ecom/arkham-hsm-als-endpoint" ; - "ecom/arkham-hsm-legacy-endpoint" ; - "ecom/redis" ; - "ecom/redis" -> "core"; - "ecom/redis-case" ; - "ecom/redis-case" -> "core"; - "ecom/redis-legacy-endpoint" ; - "ecom/redis-legacy-endpoint" -> "core"; - "ecom/redis-legacy-endpoint" -> "ecom/redis"; - "ecom/redis-sharon" ; - "ecom/redis-sharon" -> "core"; - "ecom/redis-webhooks-premium" ; - "ecom/redis-webhooks-premium" -> "core"; - "endpoints/azuresql-legacy-endpoint-tx" ; - "endpoints/azuresql-legacy-endpoint-tx" -> "core"; - "endpoints/azuresql-legacy-endpoint-tx" -> "azuresql"; - "endpoints/azuresql-legacy-endpoint-va" ; - "endpoints/azuresql-legacy-endpoint-va" -> "core"; - "endpoints/azuresql-legacy-endpoint-va" -> "azuresql"; - "endpoints/storage-accounts" ; - "endpoints/storage-accounts" -> "storage-accounts/ingestion"; - "enterprise/app-identity/auror" ; - "enterprise/app-identity/auror" -> "core"; - "enterprise/keyvault/auror" ; - "enterprise/keyvault/auror" -> "core"; - "enterprise/keyvault/auror" -> "enterprise/app-identity/auror"; - "enterprise/redis/auror" ; - "enterprise/redis/auror" -> "core"; - "espio/az-openai" ; - "espio/az-openai" -> "core"; - "espio/espio-redis" ; - "espio/espio-redis" -> "core"; - "espio/openai" ; - "espio/openai-b" ; - "espio/openai-b" -> "core"; - "eventgrid-subscription" ; - "eventgrid-subscription" -> "core"; - "evp/hyperscale" ; - "evp/hyperscale" -> "core"; - "evp/hyperscaleusers" ; - "evp/hyperscaleusers" -> "evp/hyperscale"; - "performance/lakehouse" ; - "performance/lakehouse" -> "core"; - "performance/redis-jarvis" ; - "performance/redis-jarvis" -> "core"; - "performance/redis-pipeline" ; - "performance/redis-pipeline" -> "core"; - "performance/redis-starhopper" ; - "performance/redis-starhopper" -> "core"; - "pes/keyvault" ; - "pes/keyvault" -> "core"; - "pes/keyvault" -> "dems-cluster-identity"; - "ratelimit/redis" ; - "ratelimit/redis" -> "core"; - "sage/datafactory" ; - "sage/datafactory" -> "core"; - "sage/datafactory/alerts" ; - "sage/datafactory/alerts" -> "sage/datafactory"; - "sage/datafactory/alerts" -> "core"; - "sage/datafactory/evidence-domain-migration-internal-pipeline" ; - "sage/datafactory/evidence-domain-migration-internal-pipeline" -> "sage/datafactory"; - "sage/datafactory/evidence-domain-migration-main-pipeline" ; - "sage/datafactory/evidence-domain-migration-main-pipeline" -> "sage/datafactory"; - "sage/endpoints/hyperscale-legacy-endpoint-tx" ; - "sage/endpoints/hyperscale-legacy-endpoint-tx" -> "core"; - "sage/endpoints/hyperscale-legacy-endpoint-tx" -> "sage/hyperscale"; - "sage/endpoints/hyperscale-legacy-endpoint-va" ; - "sage/endpoints/hyperscale-legacy-endpoint-va" -> "core"; - "sage/endpoints/hyperscale-legacy-endpoint-va" -> "sage/hyperscale"; - "sage/hyperscale" ; - "sage/hyperscale" -> "core"; - "sage/hyperscale/named-replica" ; - "sage/hyperscale/named-replica" -> "sage/hyperscale"; - "sage/hyperscale/named-replica" -> "core"; - "sage/hyperscaleusers" ; - "sage/hyperscaleusers" -> "sage/hyperscale"; - "sage/redis" ; - "sage/redis" -> "core"; - "servicebus-premium" ; - "servicebus-premium" -> "core"; - "sonic/rev-storage" ; - "sonic/rev-storage" -> "core"; - "sonic/sonic" ; - "sonic/sonic" -> "core"; - "sonic/sonic-redis" ; - "sonic/sonic-redis" -> "core"; - "sonic/translation" ; - "sonic/translation" -> "core"; - "storage-accounts/case-share" ; - "storage-accounts/ingestion" ; - "storage-accounts/rtiworker" ; - "storage-accounts/sage" ; - "system-status/cosmosdb-cassandra" ; - "system-status/cosmosdb-cassandra" -> "core"; - "user-settings/cosmosdb-cassandra" ; - "user-settings/cosmosdb-cassandra" -> "core"; - "visionsearchpoc/vision" ; - "visionsearchpoc/vision" -> "core"; - "visualization/cosmosdb-cassandra" ; - "visualization/cosmosdb-cassandra" -> "core"; - "visualization/redis-cluster" ; - "visualization/redis-cluster" -> "core"; - "visualization/redis-cluster-rtm" ; - "visualization/redis-cluster-rtm" -> "core"; - "vm-apps/lsln-500" ; - "vm-apps/lsln-500" -> "core"; - "vm-apps/lsln-500" -> "vm-apps/solr8-j11-lb"; - "vm-apps/solr8-j11-lb" ; - "vm-apps/solr8-j11-lb" -> "core"; - "webhooks/cosmosdb-cassandra-dispatch" ; - "webhooks/cosmosdb-cassandra-dispatch" -> "core"; - "xshare/azuresql" ; - "xshare/azuresql" -> "core"; - "xshare/azuresqlusers" ; - "xshare/azuresqlusers" -> "xshare/azuresql"; -} -` - graph, err := NewGraphFromDot(realWorldDot) - if err != nil { - t.Fatalf("Failed to create graph from real-world DOT: %v", err) - } - - // Test case: cluster-va changed - prunedGraph, err := PruneGraph(context.Background(), graph, []string{"cluster-va"}) - if err != nil { - t.Fatalf("PruneGraph failed: %v", err) - } - - // Expected result: - // digraph { - // "bootstrap-va" -> "cluster-va"; - // "dems-cluster-identity" -> "cluster-va"; - // "pes/keyvault" -> "dems-cluster-identity"; - // } - expectedNodes := []string{"bootstrap-va", "cluster-va", "dems-cluster-identity", "pes/keyvault"} - expectedEdges := map[string][]string{ - "bootstrap-va": {"cluster-va"}, - "dems-cluster-identity": {"cluster-va"}, - "pes/keyvault": {"dems-cluster-identity"}, - } - - prunedNodes := prunedGraph.GetNodes() - sort.Strings(prunedNodes) - sort.Strings(expectedNodes) - - if !reflect.DeepEqual(prunedNodes, expectedNodes) { - t.Errorf("Expected nodes %v, but got %v", expectedNodes, prunedNodes) - } - - if !reflect.DeepEqual(prunedGraph.Edges, expectedEdges) { - t.Errorf("Expected edges %v, but got %v", expectedEdges, prunedGraph.Edges) - } -} - -func TestTopologicalSort(t *testing.T) { - testCases := []struct { - name string - nodes []string - edges map[string][]string - expectedLevels [][]string - }{ - { - name: "Simple linear dependency", - nodes: []string{"A", "B", "C"}, - edges: map[string][]string{ - "B": {"A"}, - "C": {"B"}, - }, - expectedLevels: [][]string{ - {"A"}, - {"B"}, - {"C"}, - }, - }, - { - name: "Parallel dependencies", - nodes: []string{"A", "B", "C", "D"}, - edges: map[string][]string{ - "C": {"A", "B"}, - "D": {"C"}, - }, - expectedLevels: [][]string{ - {"A", "B"}, - {"C"}, - {"D"}, - }, - }, - { - name: "Complex dependency graph", - nodes: []string{"A", "B", "C", "D", "E", "F"}, - edges: map[string][]string{ - "C": {"A", "B"}, - "D": {"C"}, - "E": {"C"}, - "F": {"D", "E"}, - }, - expectedLevels: [][]string{ - {"A", "B"}, - {"C"}, - {"D", "E"}, - {"F"}, - }, - }, - { - name: "No dependencies", - nodes: []string{"A", "B", "C"}, - edges: map[string][]string{}, - expectedLevels: [][]string{{"A", "B", "C"}}, - }, - { - name: "Single node", - nodes: []string{"A"}, - edges: map[string][]string{}, - expectedLevels: [][]string{{"A"}}, - }, - { - name: "Real world example: bootstrap depends on cluster", - nodes: []string{"bootstrap", "cluster"}, - edges: map[string][]string{ - "bootstrap": {"cluster"}, - }, - expectedLevels: [][]string{ - {"cluster"}, - {"bootstrap"}, - }, - }, - } - - for _, tc := range testCases { - t.Run(tc.name, func(t *testing.T) { - // Create graph - graph := NewGraph() - - for _, nodeName := range tc.nodes { - graph.AddNode(nodeName) - } - - for src, dests := range tc.edges { - for _, dest := range dests { - graph.AddEdge(src, dest) - } - } - - // Get topological sort - levels := graph.TopologicalSort() - - // Verify number of levels - if len(levels) != len(tc.expectedLevels) { - t.Errorf("Expected %d levels, but got %d", len(tc.expectedLevels), len(levels)) - return - } - - // Verify each level - for levelIndex, expectedLevel := range tc.expectedLevels { - actualLevel := levels[levelIndex] - - // Sort both slices for comparison - sort.Strings(expectedLevel) - sort.Strings(actualLevel) - - if !reflect.DeepEqual(actualLevel, expectedLevel) { - t.Errorf("Level %d: expected %v, but got %v", levelIndex, expectedLevel, actualLevel) - } - } - }) - } -} diff --git a/controller/activities/oci.go b/controller/activities/oci.go deleted file mode 100644 index d9d51bb..0000000 --- a/controller/activities/oci.go +++ /dev/null @@ -1,43 +0,0 @@ -package activities - -import ( - "bytes" - "context" - "encoding/json" - "os/exec" - - "go.temporal.io/sdk/activity" -) - -type PushResult struct { - Reference string `json:"reference"` - MediaType string `json:"mediaType"` - Digest string `json:"digest"` - Size int `json:"size"` - Annotations map[string]string `json:"annotations"` - ArtifactType string `json:"artifactType"` - ReferenceAsTags []string `json:"referenceAsTags"` -} - -func PushManifests(ctx context.Context, path string, image string) (*PushResult, error) { - logger := activity.GetLogger(ctx) - cmd := exec.CommandContext(ctx, "oras", "push", "--format=json", "--plain-http", image, ".") - cmd.Dir = path - - var stdout, stderr bytes.Buffer - cmd.Stdout = &stdout - cmd.Stderr = &stderr - - if err := cmd.Run(); err != nil { - logger.Error("oras push failed", "error", err, "stderr", stderr.String()) - return nil, err - } - - var result PushResult - if err := json.Unmarshal(stdout.Bytes(), &result); err != nil { - logger.Error("failed to parse oras output", "error", err, "output", stdout.String()) - return nil, err - } - - return &result, nil -} diff --git a/controller/activities/terragrunt.go b/controller/activities/terragrunt.go deleted file mode 100644 index 4d1e427..0000000 --- a/controller/activities/terragrunt.go +++ /dev/null @@ -1,105 +0,0 @@ -package activities - -import ( - "bufio" - "context" - "fmt" - "os/exec" - "path/filepath" - "strings" - "time" - - "go.temporal.io/sdk/activity" -) - -func TerragruntGraph(ctx context.Context, path string) (*Graph, error) { - cmd := exec.CommandContext(ctx, "terragrunt", "dag", "graph") - cmd.Dir = path - output, err := cmd.Output() - if err != nil { - return nil, fmt.Errorf("failed to run terragrunt dag graph: %w", err) - } - - return NewGraphFromDot(string(output)) -} - -func TerragruntPrune(ctx context.Context, graph *Graph, changedFiles []string) (*Graph, error) { - return PruneGraph(ctx, graph, changedFiles) -} - -func TerragruntApply(ctx context.Context, repoUrl string, revision string, modulePath string, stack string) error { - logger := activity.GetLogger(ctx) - logger.Info("Running terragrunt apply", "module", modulePath, "stack", stack) - - repoPath, err := Clone(ctx, repoUrl, revision) - if err != nil { - return fmt.Errorf("failed to ensure repository is available: %w", err) - } - - fullPath := filepath.Join(repoPath, "infra", stack, modulePath) - - cmd := exec.CommandContext(ctx, "terragrunt", "apply", "--backend-bootstrap", "--auto-approve") - cmd.Dir = fullPath - - // Create pipes to capture output and send heartbeats - stdout, err := cmd.StdoutPipe() - if err != nil { - return fmt.Errorf("failed to create stdout pipe: %w", err) - } - stderr, err := cmd.StderrPipe() - if err != nil { - return fmt.Errorf("failed to create stderr pipe: %w", err) - } - - if err := cmd.Start(); err != nil { - return fmt.Errorf("failed to start terragrunt apply: %w", err) - } - - // Monitor output and send heartbeats - done := make(chan error, 1) - go func() { - done <- cmd.Wait() - }() - - // Send heartbeats while monitoring output - heartbeatTicker := time.NewTicker(25 * time.Second) // Send heartbeat every 25s (before 30s timeout) - defer heartbeatTicker.Stop() - - var lastOutput string - outputScanner := bufio.NewScanner(stdout) - errorScanner := bufio.NewScanner(stderr) - - for { - select { - case err := <-done: - if err != nil { - return fmt.Errorf("terragrunt apply failed for module %s: %w", modulePath, err) - } - safeHeartbeat(ctx, fmt.Sprintf("Terragrunt apply completed for %s", modulePath)) - return nil - - case <-heartbeatTicker.C: - safeHeartbeat(ctx, fmt.Sprintf("Terragrunt apply in progress for %s - %s", modulePath, lastOutput)) - - default: - // Check for new output - if outputScanner.Scan() { - line := strings.TrimSpace(outputScanner.Text()) - if line != "" { - lastOutput = line - logger.Info("Terragrunt output", "module", modulePath, "output", line) - } - } - if errorScanner.Scan() { - line := strings.TrimSpace(errorScanner.Text()) - if line != "" { - lastOutput = line - logger.Info("Terragrunt error output", "module", modulePath, "error", line) - } - } - - // Small sleep to prevent busy waiting - time.Sleep(100 * time.Millisecond) - } - } -} diff --git a/controller/activities/utils.go b/controller/activities/utils.go deleted file mode 100644 index 0b07c48..0000000 --- a/controller/activities/utils.go +++ /dev/null @@ -1,17 +0,0 @@ -package activities - -import ( - "context" - - "go.temporal.io/sdk/activity" -) - -// safeHeartbeat sends a heartbeat only if we're in an activity context -func safeHeartbeat(ctx context.Context, details string) { - defer func() { - if r := recover(); r != nil { - // Ignore panic - we're not in an activity context - } - }() - activity.RecordHeartbeat(ctx, details) -} diff --git a/controller/go.mod b/controller/go.mod deleted file mode 100644 index e8e8a3f..0000000 --- a/controller/go.mod +++ /dev/null @@ -1,35 +0,0 @@ -module cloudlab/controller - -go 1.24.3 - -require ( - github.com/stretchr/testify v1.10.0 - go.temporal.io/sdk v1.34.0 -) - -require ( - github.com/davecgh/go-spew v1.1.1 // indirect - github.com/facebookgo/clock v0.0.0-20150410010913-600d898af40a // indirect - github.com/gogo/protobuf v1.3.2 // indirect - github.com/golang/mock v1.6.0 // indirect - github.com/google/go-cmp v0.7.0 // indirect - github.com/google/uuid v1.6.0 // indirect - github.com/grpc-ecosystem/go-grpc-middleware v1.4.0 // indirect - github.com/grpc-ecosystem/grpc-gateway/v2 v2.22.0 // indirect - github.com/nexus-rpc/sdk-go v0.3.0 // indirect - github.com/pmezard/go-difflib v1.0.0 // indirect - github.com/robfig/cron v1.2.0 // indirect - github.com/rogpeppe/go-internal v1.13.1 // indirect - github.com/stretchr/objx v0.5.2 // indirect - go.temporal.io/api v1.46.0 // indirect - golang.org/x/net v0.39.0 // indirect - golang.org/x/sync v0.13.0 // indirect - golang.org/x/sys v0.32.0 // indirect - golang.org/x/text v0.24.0 // indirect - golang.org/x/time v0.9.0 // indirect - google.golang.org/genproto/googleapis/api v0.0.0-20240827150818-7e3bb234dfed // indirect - google.golang.org/genproto/googleapis/rpc v0.0.0-20240827150818-7e3bb234dfed // indirect - google.golang.org/grpc v1.66.0 // indirect - google.golang.org/protobuf v1.36.5 // indirect - gopkg.in/yaml.v3 v3.0.1 // indirect -) diff --git a/controller/go.sum b/controller/go.sum deleted file mode 100644 index 2a15a98..0000000 --- a/controller/go.sum +++ /dev/null @@ -1,174 +0,0 @@ -cloud.google.com/go v0.26.0/go.mod h1:aQUYkXzVsufM+DwF1aE+0xfcU+56JwCaLick0ClmMTw= -github.com/BurntSushi/toml v0.3.1/go.mod h1:xHWCNGjB5oqiDr8zfno3MHue2Ht5sIBksp03qcyfWMU= -github.com/benbjohnson/clock v1.1.0/go.mod h1:J11/hYXuz8f4ySSvYwY0FKfm+ezbsZBKZxNJlLklBHA= -github.com/census-instrumentation/opencensus-proto v0.2.1/go.mod h1:f6KPmirojxKA12rnyqOA5BBL4O983OfeGPqjHWSTneU= -github.com/client9/misspell v0.3.4/go.mod h1:qj6jICC3Q7zFZvVWo7KLAzC3yx5G7kyvSDkc90ppPyw= -github.com/cncf/udpa/go v0.0.0-20191209042840-269d4d468f6f/go.mod h1:M8M6+tZqaGXZJjfX53e64911xZQV5JYwmTeXPW+k8Sc= -github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= -github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= -github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= -github.com/envoyproxy/go-control-plane v0.9.0/go.mod h1:YTl/9mNaCwkRvm6d1a2C3ymFceY/DCBVvsKhRF0iEA4= -github.com/envoyproxy/go-control-plane v0.9.1-0.20191026205805-5f8ba28d4473/go.mod h1:YTl/9mNaCwkRvm6d1a2C3ymFceY/DCBVvsKhRF0iEA4= -github.com/envoyproxy/go-control-plane v0.9.4/go.mod h1:6rpuAdCZL397s3pYoYcLgu1mIlRU8Am5FuJP05cCM98= -github.com/envoyproxy/protoc-gen-validate v0.1.0/go.mod h1:iSmxcyjqTsJpI2R4NaDN7+kN2VEUnK/pcBlmesArF7c= -github.com/facebookgo/clock v0.0.0-20150410010913-600d898af40a h1:yDWHCSQ40h88yih2JAcL6Ls/kVkSE8GFACTGVnMPruw= -github.com/facebookgo/clock v0.0.0-20150410010913-600d898af40a/go.mod h1:7Ga40egUymuWXxAe151lTNnCv97MddSOVsjpPPkityA= -github.com/go-kit/log v0.1.0/go.mod h1:zbhenjAZHb184qTLMA9ZjW7ThYL0H2mk7Q6pNt4vbaY= -github.com/go-logfmt/logfmt v0.5.0/go.mod h1:wCYkCAKZfumFQihp8CzCvQ3paCTfi41vtzG1KdI/P7A= -github.com/go-stack/stack v1.8.0/go.mod h1:v0f6uXyyMGvRgIKkXu+yp6POWl0qKG85gN/melR3HDY= -github.com/gogo/protobuf v1.3.2 h1:Ov1cvc58UF3b5XjBnZv7+opcTcQFZebYjWzi34vdm4Q= -github.com/gogo/protobuf v1.3.2/go.mod h1:P1XiOD3dCwIKUDQYPy72D8LYyHL2YPYrpS2s69NZV8Q= -github.com/golang/glog v0.0.0-20160126235308-23def4e6c14b/go.mod h1:SBH7ygxi8pfUlaOkMMuAQtPIUF8ecWP5IEl/CR7VP2Q= -github.com/golang/mock v1.1.1/go.mod h1:oTYuIxOrZwtPieC+H1uAHpcLFnEyAGVDL/k47Jfbm0A= -github.com/golang/mock v1.6.0 h1:ErTB+efbowRARo13NNdxyJji2egdxLGQhRaY+DUumQc= -github.com/golang/mock v1.6.0/go.mod h1:p6yTPP+5HYm5mzsMV8JkE6ZKdX+/wYM6Hr+LicevLPs= -github.com/golang/protobuf v1.2.0/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U= -github.com/golang/protobuf v1.3.2/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U= -github.com/golang/protobuf v1.3.3/go.mod h1:vzj43D7+SQXF/4pzW/hwtAqwc6iTitCiVSaWz5lYuqw= -github.com/golang/protobuf v1.5.0 h1:LUVKkCeviFUMKqHa4tXIIij/lbhnMbP7Fn5wKdKkRh4= -github.com/golang/protobuf v1.5.0/go.mod h1:FsONVRAS9T7sI+LIUmWTfcYkHO4aIWwzhcaSAoJOfIk= -github.com/google/go-cmp v0.2.0/go.mod h1:oXzfMopK8JAjlY9xF4vHSVASa0yLyX7SntLO5aqRK0M= -github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= -github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= -github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= -github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= -github.com/grpc-ecosystem/go-grpc-middleware v1.4.0 h1:UH//fgunKIs4JdUbpDl1VZCDaL56wXCB/5+wF6uHfaI= -github.com/grpc-ecosystem/go-grpc-middleware v1.4.0/go.mod h1:g5qyo/la0ALbONm6Vbp88Yd8NsDy6rZz+RcrMPxvld8= -github.com/grpc-ecosystem/grpc-gateway/v2 v2.22.0 h1:asbCHRVmodnJTuQ3qamDwqVOIjwqUPTYmYuemVOx+Ys= -github.com/grpc-ecosystem/grpc-gateway/v2 v2.22.0/go.mod h1:ggCgvZ2r7uOoQjOyu2Y1NhHmEPPzzuhWgcza5M1Ji1I= -github.com/kisielk/errcheck v1.5.0/go.mod h1:pFxgyoBC7bSaBwPgfKdkLd5X25qrDl4LWUI2bnpBCr8= -github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+oQHNcck= -github.com/konsorten/go-windows-terminal-sequences v1.0.1/go.mod h1:T0+1ngSBFLxvqU3pZ+m/2kptfBszLMUkC4ZK/EgS/cQ= -github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo= -github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= -github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= -github.com/kr/pty v1.1.1/go.mod h1:pFQYn66WHrOpPYNljwOMqo10TkYh1fy3cYio2l3bCsQ= -github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI= -github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= -github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= -github.com/nexus-rpc/sdk-go v0.3.0 h1:Y3B0kLYbMhd4C2u00kcYajvmOrfozEtTV/nHSnV57jA= -github.com/nexus-rpc/sdk-go v0.3.0/go.mod h1:TpfkM2Cw0Rlk9drGkoiSMpFqflKTiQLWUNyKJjF8mKQ= -github.com/opentracing/opentracing-go v1.1.0/go.mod h1:UkNAQd3GIcIGf0SeVgPpRdFStlNbqXla1AfSYxPUl2o= -github.com/pkg/errors v0.8.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= -github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= -github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= -github.com/prometheus/client_model v0.0.0-20190812154241-14fe0d1b01d4/go.mod h1:xMI15A0UPsDsEKsMN9yxemIoYk6Tm2C1GtYGdfGttqA= -github.com/robfig/cron v1.2.0 h1:ZjScXvvxeQ63Dbyxy76Fj3AT3Ut0aKsyd2/tl3DTMuQ= -github.com/robfig/cron v1.2.0/go.mod h1:JGuDeoQd7Z6yL4zQhZ3OPEVHB7fL6Ka6skscFHfmt2k= -github.com/rogpeppe/go-internal v1.13.1 h1:KvO1DLK/DRN07sQ1LQKScxyZJuNnedQ5/wKSR38lUII= -github.com/rogpeppe/go-internal v1.13.1/go.mod h1:uMEvuHeurkdAXX61udpOXGD/AzZDWNMNyH2VO9fmH0o= -github.com/sirupsen/logrus v1.4.2/go.mod h1:tLMulIdttU9McNUspp0xgXVQah82FyeX6MwdIuYE2rE= -github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= -github.com/stretchr/objx v0.1.1/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= -github.com/stretchr/objx v0.5.2 h1:xuMeJ0Sdp5ZMRXx/aWO6RZxdr3beISkG5/G/aIRr3pY= -github.com/stretchr/objx v0.5.2/go.mod h1:FRsXN1f5AsAjCGJKqEizvkpNtU+EGNCLh3NxZ/8L+MA= -github.com/stretchr/testify v1.2.2/go.mod h1:a8OnRcib4nhh0OaRAV+Yts87kKdq0PP7pXfy6kDkUVs= -github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= -github.com/stretchr/testify v1.4.0/go.mod h1:j7eGeouHqKxXV5pUuKE4zz7dFj8WfuZ+81PSLYec5m4= -github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= -github.com/stretchr/testify v1.10.0 h1:Xv5erBjTwe/5IxqUQTdXv5kgmIvbHo3QQyRwhJsOfJA= -github.com/stretchr/testify v1.10.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY= -github.com/yuin/goldmark v1.1.27/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= -github.com/yuin/goldmark v1.2.1/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= -github.com/yuin/goldmark v1.3.5/go.mod h1:mwnBkeHKe2W/ZEtQ+71ViKU8L12m81fl3OWwC1Zlc8k= -go.temporal.io/api v1.46.0 h1:O1efPDB6O2B8uIeCDIa+3VZC7tZMvYsMZYQapSbHvCg= -go.temporal.io/api v1.46.0/go.mod h1:iaxoP/9OXMJcQkETTECfwYq4cw/bj4nwov8b3ZLVnXM= -go.temporal.io/sdk v1.34.0 h1:VLg/h6ny7GvLFVoQPqz2NcC93V9yXboQwblkRvZ1cZE= -go.temporal.io/sdk v1.34.0/go.mod h1:iE4U5vFrH3asOhqpBBphpj9zNtw8btp8+MSaf5A0D3w= -go.uber.org/atomic v1.7.0/go.mod h1:fEN4uk6kAWBTFdckzkM89CLk9XfWZrxpCo0nPH17wJc= -go.uber.org/goleak v1.1.10/go.mod h1:8a7PlsEVH3e/a/GLqe5IIrQx6GzcnRmZEufDUTk4A7A= -go.uber.org/multierr v1.6.0/go.mod h1:cdWPpRnG4AhwMwsgIHip0KRBQjJy5kYEpYjJxpXp9iU= -go.uber.org/zap v1.18.1/go.mod h1:xg/QME4nWcxGxrpdeYfq7UvYrLh66cuVKdrbD1XF/NI= -golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= -golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI= -golang.org/x/crypto v0.0.0-20200622213623-75b288015ac9/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= -golang.org/x/exp v0.0.0-20190121172915-509febef88a4/go.mod h1:CJ0aWSM057203Lf6IL+f9T1iT9GByDxfZKAQTCR3kQA= -golang.org/x/lint v0.0.0-20181026193005-c67002cb31c3/go.mod h1:UVdnD1Gm6xHRNCYTkRU2/jEulfH38KcIWyp/GAMgvoE= -golang.org/x/lint v0.0.0-20190227174305-5b3e6a55c961/go.mod h1:wehouNa3lNwaWXcvxsM5YxQ5yQlVC4a0KAMCusXpPoU= -golang.org/x/lint v0.0.0-20190313153728-d0100b6bd8b3/go.mod h1:6SW0HCj/g11FgYtHlgUYUwCkIfeOF89ocIRzGO/8vkc= -golang.org/x/lint v0.0.0-20190930215403-16217165b5de/go.mod h1:6SW0HCj/g11FgYtHlgUYUwCkIfeOF89ocIRzGO/8vkc= -golang.org/x/mod v0.2.0/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA= -golang.org/x/mod v0.3.0/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA= -golang.org/x/mod v0.4.2/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA= -golang.org/x/net v0.0.0-20180724234803-3673e40ba225/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= -golang.org/x/net v0.0.0-20180826012351-8a410e7b638d/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= -golang.org/x/net v0.0.0-20190213061140-3a22650c66bd/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= -golang.org/x/net v0.0.0-20190311183353-d8887717615a/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= -golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= -golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= -golang.org/x/net v0.0.0-20200226121028-0de0cce0169b/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= -golang.org/x/net v0.0.0-20201021035429-f5854403a974/go.mod h1:sp8m0HH+o8qH0wwXwYZr8TS3Oi6o0r6Gce1SSxlDquU= -golang.org/x/net v0.0.0-20210405180319-a5a99cb37ef4/go.mod h1:p54w0d4576C0XHj96bSt6lcn1PtDYWL6XObtHCRCNQM= -golang.org/x/net v0.39.0 h1:ZCu7HMWDxpXpaiKdhzIfaltL9Lp31x/3fCP11bc6/fY= -golang.org/x/net v0.39.0/go.mod h1:X7NRbYVEA+ewNkCNyJ513WmMdQ3BineSwVtN2zD/d+E= -golang.org/x/oauth2 v0.0.0-20180821212333-d2e6202438be/go.mod h1:N/0e6XlmueqKjAGxoOufVs8QHGRruUQn6yWY3a++T0U= -golang.org/x/sync v0.0.0-20180314180146-1d60e4601c6f/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= -golang.org/x/sync v0.0.0-20181108010431-42b317875d0f/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= -golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= -golang.org/x/sync v0.0.0-20190911185100-cd5d95a43a6e/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= -golang.org/x/sync v0.0.0-20201020160332-67f06af15bc9/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= -golang.org/x/sync v0.0.0-20210220032951-036812b2e83c/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= -golang.org/x/sync v0.13.0 h1:AauUjRAJ9OSnvULf/ARrrVywoJDy0YS2AwQ98I37610= -golang.org/x/sync v0.13.0/go.mod h1:1dzgHSNfp02xaA81J2MS99Qcpr2w7fw1gpm99rleRqA= -golang.org/x/sys v0.0.0-20180830151530-49385e6e1522/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= -golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= -golang.org/x/sys v0.0.0-20190412213103-97732733099d/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= -golang.org/x/sys v0.0.0-20190422165155-953cdadca894/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= -golang.org/x/sys v0.0.0-20200930185726-fdedc70b468f/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= -golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= -golang.org/x/sys v0.0.0-20210330210617-4fbd30eecc44/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= -golang.org/x/sys v0.0.0-20210510120138-977fb7262007/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.0.0-20211025201205-69cdffdb9359/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.32.0 h1:s77OFDvIQeibCmezSnk/q6iAfkdiQaJi4VzroCFrN20= -golang.org/x/sys v0.32.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k= -golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= -golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= -golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= -golang.org/x/text v0.24.0 h1:dd5Bzh4yt5KYA8f9CJHCP4FB4D51c2c6JvN37xJJkJ0= -golang.org/x/text v0.24.0/go.mod h1:L8rBsPeo2pSS+xqN0d5u2ikmjtmoJbDBT1b7nHvFCdU= -golang.org/x/time v0.9.0 h1:EsRrnYcQiGH+5FfbgvV4AP7qEZstoyrHB0DzarOQ4ZY= -golang.org/x/time v0.9.0/go.mod h1:3BpzKBy/shNhVucY/MWOyx10tF3SFh9QdLuxbVysPQM= -golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= -golang.org/x/tools v0.0.0-20190114222345-bf090417da8b/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= -golang.org/x/tools v0.0.0-20190226205152-f727befe758c/go.mod h1:9Yl7xja0Znq3iFh3HoIrodX9oNMXvdceNzlUR8zjMvY= -golang.org/x/tools v0.0.0-20190311212946-11955173bddd/go.mod h1:LCzVGOaR6xXOjkQ3onu1FJEFr0SW1gC7cKk1uF8kGRs= -golang.org/x/tools v0.0.0-20190524140312-2c0ae7006135/go.mod h1:RgjU9mgBXZiqYHBnxXauZ1Gv1EHHAz9KjViQ78xBX0Q= -golang.org/x/tools v0.0.0-20191108193012-7d206e10da11/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo= -golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo= -golang.org/x/tools v0.0.0-20200619180055-7c47624df98f/go.mod h1:EkVYQZoAsY45+roYkvgYkIh4xh/qjgUK9TdY2XT94GE= -golang.org/x/tools v0.0.0-20210106214847-113979e3529a/go.mod h1:emZCQorbCU4vsT4fOWvOPXz4eW1wZW4PmDk9uLelYpA= -golang.org/x/tools v0.1.1/go.mod h1:o0xws9oXOQQZyjljx8fwUC0k7L1pTE6eaCbjGeHmOkk= -golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= -golang.org/x/xerrors v0.0.0-20191011141410-1b5146add898/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= -golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= -golang.org/x/xerrors v0.0.0-20200804184101-5ec99f83aff1/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= -google.golang.org/appengine v1.1.0/go.mod h1:EbEs0AVv82hx2wNQdGPgUI5lhzA/G0D9YwlJXL52JkM= -google.golang.org/appengine v1.4.0/go.mod h1:xpcJRLb0r/rnEns0DIKYYv+WjYCduHsrkT7/EB5XEv4= -google.golang.org/genproto v0.0.0-20180817151627-c66870c02cf8/go.mod h1:JiN7NxoALGmiZfu7CAH4rXhgtRTLTxftemlI0sWmxmc= -google.golang.org/genproto v0.0.0-20190819201941-24fa4b261c55/go.mod h1:DMBHOl98Agz4BDEuKkezgsaosCRResVns1a3J2ZsMNc= -google.golang.org/genproto v0.0.0-20200423170343-7949de9c1215/go.mod h1:55QSHmfGQM9UVYDPBsyGGes0y52j32PQ3BqQfXhyH3c= -google.golang.org/genproto/googleapis/api v0.0.0-20240827150818-7e3bb234dfed h1:3RgNmBoI9MZhsj3QxC+AP/qQhNwpCLOvYDYYsFrhFt0= -google.golang.org/genproto/googleapis/api v0.0.0-20240827150818-7e3bb234dfed/go.mod h1:OCdP9MfskevB/rbYvHTsXTtKC+3bHWajPdoKgjcYkfo= -google.golang.org/genproto/googleapis/rpc v0.0.0-20240827150818-7e3bb234dfed h1:J6izYgfBXAI3xTKLgxzTmUltdYaLsuBxFCgDHWJ/eXg= -google.golang.org/genproto/googleapis/rpc v0.0.0-20240827150818-7e3bb234dfed/go.mod h1:UqMtugtsSgubUsoxbuAoiCXvqvErP7Gf0so0mK9tHxU= -google.golang.org/grpc v1.19.0/go.mod h1:mqu4LbDTu4XGKhr4mRzUsmM4RtVoemTSY81AxZiDr8c= -google.golang.org/grpc v1.23.0/go.mod h1:Y5yQAOtifL1yxbo5wqy6BxZv8vAUGQwXBOALyacEbxg= -google.golang.org/grpc v1.25.1/go.mod h1:c3i+UQWmh7LiEpx4sFZnkU36qjEYZ0imhYfXVyQciAY= -google.golang.org/grpc v1.27.0/go.mod h1:qbnxyOmOxrQa7FizSgH+ReBfzJrCY1pSN7KXBS8abTk= -google.golang.org/grpc v1.29.1/go.mod h1:itym6AZVZYACWQqET3MqgPpjcuV5QH3BxFS3IjizoKk= -google.golang.org/grpc v1.66.0 h1:DibZuoBznOxbDQxRINckZcUvnCEvrW9pcWIE2yF9r1c= -google.golang.org/grpc v1.66.0/go.mod h1:s3/l6xSSCURdVfAnL+TqCNMyTDAGN6+lZeVxnZR128Y= -google.golang.org/protobuf v1.36.5 h1:tPhr+woSbjfYvY6/GPufUoYizxw1cF/yFoxJ2fmpwlM= -google.golang.org/protobuf v1.36.5/go.mod h1:9fA7Ob0pmnwhb644+1+CVWFRbNajQ6iRojtC/QF5bRE= -gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= -gopkg.in/check.v1 v1.0.0-20180628173108-788fd7840127/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= -gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= -gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= -gopkg.in/yaml.v2 v2.2.2/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= -gopkg.in/yaml.v2 v2.2.8/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= -gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= -gopkg.in/yaml.v3 v3.0.0-20210107192922-496545a6307b/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= -gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= -gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= -honnef.co/go/tools v0.0.0-20190102054323-c2f93a96b099/go.mod h1:rf3lG4BRIbNafJWhAfAdb/ePZxsR/4RtNHQocxwk9r4= -honnef.co/go/tools v0.0.0-20190523083050-ea95bdfd59fc/go.mod h1:rf3lG4BRIbNafJWhAfAdb/ePZxsR/4RtNHQocxwk9r4= diff --git a/controller/worker/main.go b/controller/worker/main.go deleted file mode 100644 index cec6354..0000000 --- a/controller/worker/main.go +++ /dev/null @@ -1,48 +0,0 @@ -package main - -import ( - "log" - "os" - - "cloudlab/controller/activities" - "cloudlab/controller/workflows" - - "go.temporal.io/sdk/client" - "go.temporal.io/sdk/worker" -) - -func main() { - // The client and worker are heavyweight objects that should be created once per process. - temporalClient, err := client.Dial(client.Options{ - HostPort: os.Getenv("TEMPORAL_HOST"), - }) - if err != nil { - log.Fatalln("Unable to create client", err) - } - defer temporalClient.Close() - - w := worker.New(temporalClient, "cloudlab", worker.Options{}) - - w.RegisterActivity(activities.Clone) - w.RegisterActivity(activities.ChangedModules) - w.RegisterActivity(activities.TerragruntGraph) - w.RegisterActivity(activities.PruneGraph) - w.RegisterActivity(activities.TerragruntApply) - w.RegisterActivity(activities.PushManifests) - w.RegisterActivity(activities.PushRenderedApp) - w.RegisterActivity(activities.DiscoverApps) - w.RegisterActivity(activities.UpdateAppVersion) - w.RegisterActivity(activities.GitAdd) - w.RegisterActivity(activities.GitCommit) - w.RegisterActivity(activities.GitPush) - - w.RegisterWorkflow(workflows.Infra) - w.RegisterWorkflow(workflows.Platform) - w.RegisterWorkflow(workflows.Apps) - w.RegisterWorkflow(workflows.AppUpdate) - - err = w.Run(worker.InterruptCh()) - if err != nil { - log.Fatalln("Unable to start Worker", err) - } -} diff --git a/controller/workflows/app_update.go b/controller/workflows/app_update.go deleted file mode 100644 index 2dc6816..0000000 --- a/controller/workflows/app_update.go +++ /dev/null @@ -1,148 +0,0 @@ -package workflows - -import ( - "fmt" - "os" - "path/filepath" - "time" - - "cloudlab/controller/activities" - - "go.temporal.io/sdk/workflow" -) - -type AppUpdateInput struct { - Url string - Revision string - Namespace string - App string - Cluster string - Registry string - NewImages []activities.Image -} - -// AppUpdate workflow clones a repository, updates app versions, and syncs changes back to git -func AppUpdate(ctx workflow.Context, input AppUpdateInput) error { - logger := workflow.GetLogger(ctx) - logger.Info("AppUpdate workflow started", "input", input) - - var workspace string - if err := workflow.ExecuteActivity( - workflow.WithActivityOptions(ctx, workflow.ActivityOptions{ - StartToCloseTimeout: 2 * time.Minute, - }), - activities.Clone, - input.Url, - input.Revision, - ).Get(ctx, &workspace); err != nil { - logger.Error("Failed to clone repository", "error", err) - return fmt.Errorf("failed to clone repository: %w", err) - } - - logger.Info("Repository cloned successfully", "workspace", workspace) - - defer func() { - if err := os.RemoveAll(workspace); err != nil { - logger.Error("Failed to cleanup workspace", "workspace", workspace, "error", err) - } - }() - - appsDir := filepath.Join(workspace, "apps") - var changed bool - if err := workflow.ExecuteActivity( - workflow.WithActivityOptions(ctx, workflow.ActivityOptions{ - StartToCloseTimeout: 30 * time.Second, - }), - activities.UpdateAppVersion, - appsDir, - input.Namespace, - input.App, - input.Cluster, - input.NewImages, - ).Get(ctx, &changed); err != nil { - logger.Error("failed to update app version", "error", err) - return fmt.Errorf("failed to update app version: %w", err) - } - - logger.Info("App version updated successfully", - "namespace", input.Namespace, - "app", input.App, - "cluster", input.Cluster, - "changed", changed) - - // Skip remaining steps if no changes were made - if !changed { - logger.Info("No changes detected, skipping remaining steps") - return nil - } - - // Step 3: Git add changes - appFilePath := filepath.Join(appsDir, input.Namespace, input.App, fmt.Sprintf("%s.yaml", input.Cluster)) - if err := workflow.ExecuteActivity( - workflow.WithActivityOptions(ctx, workflow.ActivityOptions{ - StartToCloseTimeout: 30 * time.Second, - }), - activities.GitAdd, - appFilePath, - ).Get(ctx, nil); err != nil { - logger.Error("Failed to add changes to git", "error", err) - return fmt.Errorf("failed to add changes to git: %w", err) - } - - // Step 4: Git commit changes - commitMessage := fmt.Sprintf("chore(%s/%s): update %s version", input.Namespace, input.App, input.Cluster) - if err := workflow.ExecuteActivity( - workflow.WithActivityOptions(ctx, workflow.ActivityOptions{ - StartToCloseTimeout: 30 * time.Second, - }), - activities.GitCommit, - workspace, - commitMessage, - ).Get(ctx, nil); err != nil { - logger.Error("Failed to commit changes to git", "error", err) - return fmt.Errorf("failed to commit changes to git: %w", err) - } - - // Step 5 & 6: Execute GitPush and PushRenderedApp concurrently - gitPushFuture := workflow.ExecuteActivity( - workflow.WithActivityOptions(ctx, workflow.ActivityOptions{ - StartToCloseTimeout: 1 * time.Minute, - }), - activities.GitPush, - workspace, - ) - - pushRenderedAppFuture := workflow.ExecuteActivity( - workflow.WithActivityOptions(ctx, workflow.ActivityOptions{ - StartToCloseTimeout: 2 * time.Minute, - }), - activities.PushRenderedApp, - appsDir, - input.Namespace, - input.App, - input.Cluster, - input.Registry, - ) - - // Wait for GitPush to complete - if err := gitPushFuture.Get(ctx, nil); err != nil { - logger.Error("Failed to push changes to git", "error", err) - return fmt.Errorf("failed to push changes to git: %w", err) - } - - // Wait for PushRenderedApp to complete - var pushResult *activities.PushResult - if err := pushRenderedAppFuture.Get(ctx, &pushResult); err != nil { - logger.Error("Failed to push rendered app to registry", "error", err) - return fmt.Errorf("failed to push rendered app to registry: %w", err) - } - - logger.Info("AppUpdate workflow completed successfully", - "namespace", input.Namespace, - "app", input.App, - "cluster", input.Cluster, - "updated_images", len(input.NewImages), - "rendered_app_digest", pushResult.Digest) - - return nil -} diff --git a/controller/workflows/app_update_test.go b/controller/workflows/app_update_test.go deleted file mode 100644 index e63cd88..0000000 --- a/controller/workflows/app_update_test.go +++ /dev/null @@ -1,417 +0,0 @@ -package workflows - -import ( - "errors" - "testing" - "time" - - "cloudlab/controller/activities" - - "github.com/stretchr/testify/mock" - "github.com/stretchr/testify/suite" - "go.temporal.io/sdk/testsuite" -) - -type AppUpdateWorkflowTestSuite struct { - suite.Suite - testsuite.WorkflowTestSuite - - env *testsuite.TestWorkflowEnvironment -} - -func (s *AppUpdateWorkflowTestSuite) SetupTest() { - s.env = s.NewTestWorkflowEnvironment() - s.env.SetTestTimeout(30 * time.Second) -} - -func (s *AppUpdateWorkflowTestSuite) AfterTest(suiteName, testName string) { - s.env.AssertExpectations(s.T()) -} - -func (s *AppUpdateWorkflowTestSuite) TestAppUpdate_Success() { - input := AppUpdateInput{ - Url: "https://github.com/example/cloudlab.git", - Revision: "main", - Namespace: "khuedoan", - App: "blog", - Cluster: "production", - Registry: "registry.example.com", - NewImages: []activities.Image{ - {Repository: "docker.io/khuedoan/blog", Tag: "abc123def456789"}, - }, - } - workspace := "/tmp/cloudlab-repos/abc123" - appFilePath := workspace + "/apps/khuedoan/blog/production.yaml" - mockPushResult := &activities.PushResult{ - Reference: "registry.example.com/khuedoan/blog:production", - Digest: "sha256:abc123def456", - } - - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return(workspace, nil) - s.env.OnActivity(activities.UpdateAppVersion, mock.Anything, - workspace+"/apps", input.Namespace, input.App, input.Cluster, input.NewImages).Return(true, nil) - s.env.OnActivity(activities.GitAdd, mock.Anything, appFilePath).Return(nil) - s.env.OnActivity(activities.GitCommit, mock.Anything, workspace, "chore(khuedoan/blog): update production version").Return(nil) - s.env.OnActivity(activities.GitPush, mock.Anything, workspace).Return(nil) - s.env.OnActivity(activities.PushRenderedApp, mock.Anything, - workspace+"/apps", input.Namespace, input.App, input.Cluster, input.Registry).Return(mockPushResult, nil) - - s.env.ExecuteWorkflow(AppUpdate, input) - - s.True(s.env.IsWorkflowCompleted()) - s.NoError(s.env.GetWorkflowError()) -} - -func (s *AppUpdateWorkflowTestSuite) TestAppUpdate_CloneFailure() { - input := AppUpdateInput{ - Url: "https://github.com/example/invalid-repo.git", - Revision: "main", - Namespace: "test", - App: "app", - Cluster: "local", - Registry: "registry.127.0.0.1.sslip.io", - NewImages: []activities.Image{ - {Repository: "test/app", Tag: "v1.0.0"}, - }, - } - - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return("", errors.New("repository not found")) - - s.env.ExecuteWorkflow(AppUpdate, input) - - s.True(s.env.IsWorkflowCompleted()) - s.Error(s.env.GetWorkflowError()) - s.Contains(s.env.GetWorkflowError().Error(), "failed to clone repository") -} - -func (s *AppUpdateWorkflowTestSuite) TestAppUpdate_UpdateAppVersionFailure() { - input := AppUpdateInput{ - Url: "https://github.com/example/cloudlab.git", - Revision: "main", - Namespace: "finance", - App: "actualbudget", - Cluster: "local", - Registry: "registry.127.0.0.1.sslip.io", - NewImages: []activities.Image{ - {Repository: "docker.io/actualbudget/actual-server", Tag: "25.7.0-alpine"}, - }, - } - workspace := "/tmp/cloudlab-repos/def456" - - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return(workspace, nil) - s.env.OnActivity(activities.UpdateAppVersion, mock.Anything, - workspace+"/apps", input.Namespace, input.App, input.Cluster, input.NewImages).Return(false, - errors.New("failed to read file: no such file or directory")) - - s.env.ExecuteWorkflow(AppUpdate, input) - - s.True(s.env.IsWorkflowCompleted()) - s.Error(s.env.GetWorkflowError()) - s.Contains(s.env.GetWorkflowError().Error(), "failed to update app version") -} - -func (s *AppUpdateWorkflowTestSuite) TestAppUpdate_GitAddFailure() { - input := AppUpdateInput{ - Url: "https://github.com/example/cloudlab.git", - Revision: "main", - Namespace: "khuedoan", - App: "notes", - Cluster: "production", - Registry: "registry.example.com", - NewImages: []activities.Image{ - {Repository: "ghcr.io/silverbulletmd/silverbullet", Tag: "v3"}, - }, - } - workspace := "/tmp/cloudlab-repos/ghi789" - appFilePath := workspace + "/apps/khuedoan/notes/production.yaml" - - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return(workspace, nil) - s.env.OnActivity(activities.UpdateAppVersion, mock.Anything, - workspace+"/apps", input.Namespace, input.App, input.Cluster, input.NewImages).Return(true, nil) - s.env.OnActivity(activities.GitAdd, mock.Anything, appFilePath).Return( - errors.New("git add failed: file not found")) - - s.env.ExecuteWorkflow(AppUpdate, input) - - s.True(s.env.IsWorkflowCompleted()) - s.Error(s.env.GetWorkflowError()) - s.Contains(s.env.GetWorkflowError().Error(), "failed to add changes to git") -} - -func (s *AppUpdateWorkflowTestSuite) TestAppUpdate_GitCommitFailure() { - input := AppUpdateInput{ - Url: "https://github.com/example/cloudlab.git", - Revision: "main", - Namespace: "khuedoan", - App: "notes", - Cluster: "production", - Registry: "registry.example.com", - NewImages: []activities.Image{ - {Repository: "ghcr.io/silverbulletmd/silverbullet", Tag: "v3"}, - }, - } - workspace := "/tmp/cloudlab-repos/ghi789" - appFilePath := workspace + "/apps/khuedoan/notes/production.yaml" - - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return(workspace, nil) - s.env.OnActivity(activities.UpdateAppVersion, mock.Anything, - workspace+"/apps", input.Namespace, input.App, input.Cluster, input.NewImages).Return(true, nil) - s.env.OnActivity(activities.GitAdd, mock.Anything, appFilePath).Return(nil) - s.env.OnActivity(activities.GitCommit, mock.Anything, workspace, "chore(khuedoan/notes): update production version").Return( - errors.New("git commit failed: nothing to commit")) - - s.env.ExecuteWorkflow(AppUpdate, input) - - s.True(s.env.IsWorkflowCompleted()) - s.Error(s.env.GetWorkflowError()) - s.Contains(s.env.GetWorkflowError().Error(), "failed to commit changes to git") -} - -func (s *AppUpdateWorkflowTestSuite) TestAppUpdate_GitPushFailure() { - input := AppUpdateInput{ - Url: "https://github.com/example/cloudlab.git", - Revision: "main", - Namespace: "khuedoan", - App: "notes", - Cluster: "production", - Registry: "registry.example.com", - NewImages: []activities.Image{ - {Repository: "ghcr.io/silverbulletmd/silverbullet", Tag: "v3"}, - }, - } - workspace := "/tmp/cloudlab-repos/ghi789" - appFilePath := workspace + "/apps/khuedoan/notes/production.yaml" - - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return(workspace, nil) - s.env.OnActivity(activities.UpdateAppVersion, mock.Anything, - workspace+"/apps", input.Namespace, input.App, input.Cluster, input.NewImages).Return(true, nil) - s.env.OnActivity(activities.GitAdd, mock.Anything, appFilePath).Return(nil) - s.env.OnActivity(activities.GitCommit, mock.Anything, workspace, "chore(khuedoan/notes): update production version").Return(nil) - s.env.OnActivity(activities.GitPush, mock.Anything, workspace).Return( - errors.New("git push failed: authentication required")) - - s.env.ExecuteWorkflow(AppUpdate, input) - - s.True(s.env.IsWorkflowCompleted()) - s.Error(s.env.GetWorkflowError()) - s.Contains(s.env.GetWorkflowError().Error(), "failed to push changes to git") -} - -func (s *AppUpdateWorkflowTestSuite) TestAppUpdate_PushRenderedAppFailure() { - input := AppUpdateInput{ - Url: "https://github.com/example/cloudlab.git", - Revision: "main", - Namespace: "khuedoan", - App: "notes", - Cluster: "production", - Registry: "registry.example.com", - NewImages: []activities.Image{ - {Repository: "ghcr.io/silverbulletmd/silverbullet", Tag: "v3"}, - }, - } - workspace := "/tmp/cloudlab-repos/ghi789" - appFilePath := workspace + "/apps/khuedoan/notes/production.yaml" - - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return(workspace, nil) - s.env.OnActivity(activities.UpdateAppVersion, mock.Anything, - workspace+"/apps", input.Namespace, input.App, input.Cluster, input.NewImages).Return(true, nil) - s.env.OnActivity(activities.GitAdd, mock.Anything, appFilePath).Return(nil) - s.env.OnActivity(activities.GitCommit, mock.Anything, workspace, "chore(khuedoan/notes): update production version").Return(nil) - s.env.OnActivity(activities.GitPush, mock.Anything, workspace).Return(nil) - s.env.OnActivity(activities.PushRenderedApp, mock.Anything, - workspace+"/apps", input.Namespace, input.App, input.Cluster, input.Registry).Return( - nil, errors.New("helm template failed: chart not found")) - - s.env.ExecuteWorkflow(AppUpdate, input) - - s.True(s.env.IsWorkflowCompleted()) - s.Error(s.env.GetWorkflowError()) - s.Contains(s.env.GetWorkflowError().Error(), "failed to push rendered app to registry") -} - -func (s *AppUpdateWorkflowTestSuite) TestAppUpdate_MultipleImages() { - input := AppUpdateInput{ - Url: "https://github.com/example/cloudlab.git", - Revision: "develop", - Namespace: "test", - App: "example", - Cluster: "local", - Registry: "registry.registry.svc.cluster.local", - NewImages: []activities.Image{ - {Repository: "registry.registry.svc.cluster.local/example-service", Tag: "newcommithash123"}, - {Repository: "docker.io/redis", Tag: "7.0-alpine"}, - }, - } - workspace := "/tmp/cloudlab-repos/jkl012" - appFilePath := workspace + "/apps/test/example/local.yaml" - mockPushResult := &activities.PushResult{ - Reference: "registry.registry.svc.cluster.local/test/example:local", - Digest: "sha256:def789abc123", - } - - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return(workspace, nil) - s.env.OnActivity(activities.UpdateAppVersion, mock.Anything, - workspace+"/apps", input.Namespace, input.App, input.Cluster, input.NewImages).Return(true, nil) - s.env.OnActivity(activities.GitAdd, mock.Anything, appFilePath).Return(nil) - s.env.OnActivity(activities.GitCommit, mock.Anything, workspace, "chore(test/example): update local version").Return(nil) - s.env.OnActivity(activities.GitPush, mock.Anything, workspace).Return(nil) - s.env.OnActivity(activities.PushRenderedApp, mock.Anything, - workspace+"/apps", input.Namespace, input.App, input.Cluster, input.Registry).Return(mockPushResult, nil) - - s.env.ExecuteWorkflow(AppUpdate, input) - - s.True(s.env.IsWorkflowCompleted()) - s.NoError(s.env.GetWorkflowError()) -} - -func (s *AppUpdateWorkflowTestSuite) TestAppUpdate_RealWorldExample() { - // Test with realistic data from the actual apps directory - input := AppUpdateInput{ - Url: "https://github.com/khuedoan/cloudlab.git", - Revision: "main", - Namespace: "khuedoan", - App: "blog", - Cluster: "production", - Registry: "registry.cloudlab.khuedoan.com", - NewImages: []activities.Image{ - {Repository: "docker.io/khuedoan/blog", Tag: "1234567890abcdef1234567890abcdef12345678"}, - }, - } - workspace := "/tmp/cloudlab-repos/realworld123" - appFilePath := workspace + "/apps/khuedoan/blog/production.yaml" - mockPushResult := &activities.PushResult{ - Reference: "registry.cloudlab.khuedoan.com/khuedoan/blog:production", - Digest: "sha256:realworld789", - } - - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return(workspace, nil) - s.env.OnActivity(activities.UpdateAppVersion, mock.Anything, - workspace+"/apps", input.Namespace, input.App, input.Cluster, input.NewImages).Return(true, nil) - s.env.OnActivity(activities.GitAdd, mock.Anything, appFilePath).Return(nil) - s.env.OnActivity(activities.GitCommit, mock.Anything, workspace, "chore(khuedoan/blog): update production version").Return(nil) - s.env.OnActivity(activities.GitPush, mock.Anything, workspace).Return(nil) - s.env.OnActivity(activities.PushRenderedApp, mock.Anything, - workspace+"/apps", input.Namespace, input.App, input.Cluster, input.Registry).Return(mockPushResult, nil) - - s.env.ExecuteWorkflow(AppUpdate, input) - - s.True(s.env.IsWorkflowCompleted()) - s.NoError(s.env.GetWorkflowError()) -} - -func (s *AppUpdateWorkflowTestSuite) TestAppUpdate_ActivityTimeout() { - input := AppUpdateInput{ - Url: "https://github.com/example/cloudlab.git", - Revision: "main", - Namespace: "test", - App: "slow-app", - Cluster: "production", - Registry: "registry.example.com", - NewImages: []activities.Image{ - {Repository: "test/slow-app", Tag: "v1.0.0"}, - }, - } - - // Simulate a timeout - we'll just return an error since the test timeout catches this - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return("", errors.New("timeout")) - - s.env.ExecuteWorkflow(AppUpdate, input) - - s.True(s.env.IsWorkflowCompleted()) - s.Error(s.env.GetWorkflowError()) -} - -func (s *AppUpdateWorkflowTestSuite) TestAppUpdate_EmptyImages() { - input := AppUpdateInput{ - Url: "https://github.com/example/cloudlab.git", - Revision: "main", - Namespace: "test", - App: "app", - Cluster: "local", - Registry: "registry.127.0.0.1.sslip.io", - NewImages: []activities.Image{}, // Empty images array - } - workspace := "/tmp/cloudlab-repos/empty123" - appFilePath := workspace + "/apps/test/app/local.yaml" - mockPushResult := &activities.PushResult{ - Reference: "registry.127.0.0.1.sslip.io/test/app:local", - Digest: "sha256:empty456", - } - - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return(workspace, nil) - s.env.OnActivity(activities.UpdateAppVersion, mock.Anything, - workspace+"/apps", input.Namespace, input.App, input.Cluster, input.NewImages).Return(true, nil) - s.env.OnActivity(activities.GitAdd, mock.Anything, appFilePath).Return(nil) - s.env.OnActivity(activities.GitCommit, mock.Anything, workspace, "chore(test/app): update local version").Return(nil) - s.env.OnActivity(activities.GitPush, mock.Anything, workspace).Return(nil) - s.env.OnActivity(activities.PushRenderedApp, mock.Anything, - workspace+"/apps", input.Namespace, input.App, input.Cluster, input.Registry).Return(mockPushResult, nil) - - s.env.ExecuteWorkflow(AppUpdate, input) - - s.True(s.env.IsWorkflowCompleted()) - s.NoError(s.env.GetWorkflowError()) -} - -func (s *AppUpdateWorkflowTestSuite) TestAppUpdate_NoChanges() { - input := AppUpdateInput{ - Url: "https://github.com/example/cloudlab.git", - Revision: "main", - Namespace: "test", - App: "app", - Cluster: "local", - Registry: "registry.example.com", - NewImages: []activities.Image{ - {Repository: "docker.io/test/app", Tag: "existing-tag"}, // Same tag as already in file - }, - } - workspace := "/tmp/cloudlab-repos/nochange123" - - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return(workspace, nil) - s.env.OnActivity(activities.UpdateAppVersion, mock.Anything, - workspace+"/apps", input.Namespace, input.App, input.Cluster, input.NewImages).Return(false, nil) - // Note: No other activities should be called when there are no changes - - s.env.ExecuteWorkflow(AppUpdate, input) - - s.True(s.env.IsWorkflowCompleted()) - s.NoError(s.env.GetWorkflowError()) -} - -func (s *AppUpdateWorkflowTestSuite) TestAppUpdate_SpecialCharactersInPath() { - input := AppUpdateInput{ - Url: "https://github.com/example/cloudlab.git", - Revision: "feature/special-branch-name", - Namespace: "test-namespace", - App: "app-with-dashes", - Cluster: "staging-env", - Registry: "registry.example.com", - NewImages: []activities.Image{ - {Repository: "registry.example.com/test/app-with-dashes", Tag: "v1.2.3-rc1"}, - }, - } - workspace := "/tmp/cloudlab-repos/special456" - appFilePath := workspace + "/apps/test-namespace/app-with-dashes/staging-env.yaml" - mockPushResult := &activities.PushResult{ - Reference: "registry.example.com/test-namespace/app-with-dashes:staging-env", - Digest: "sha256:special123", - } - - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return(workspace, nil) - s.env.OnActivity(activities.UpdateAppVersion, mock.Anything, - workspace+"/apps", input.Namespace, input.App, input.Cluster, input.NewImages).Return(true, nil) - s.env.OnActivity(activities.GitAdd, mock.Anything, appFilePath).Return(nil) - s.env.OnActivity(activities.GitCommit, mock.Anything, workspace, "chore(test-namespace/app-with-dashes): update staging-env version").Return(nil) - s.env.OnActivity(activities.GitPush, mock.Anything, workspace).Return(nil) - s.env.OnActivity(activities.PushRenderedApp, mock.Anything, - workspace+"/apps", input.Namespace, input.App, input.Cluster, input.Registry).Return(mockPushResult, nil) - - s.env.ExecuteWorkflow(AppUpdate, input) - - s.True(s.env.IsWorkflowCompleted()) - s.NoError(s.env.GetWorkflowError()) -} - -func TestAppUpdateWorkflowTestSuite(t *testing.T) { - suite.Run(t, new(AppUpdateWorkflowTestSuite)) -} diff --git a/controller/workflows/apps.go b/controller/workflows/apps.go deleted file mode 100644 index 07f3c52..0000000 --- a/controller/workflows/apps.go +++ /dev/null @@ -1,94 +0,0 @@ -package workflows - -import ( - "os" - "path/filepath" - "strings" - "time" - - "cloudlab/controller/activities" - - "go.temporal.io/sdk/workflow" -) - -type AppsInput struct { - Url string - Revision string - Registry string - Cluster string -} - -func Apps(ctx workflow.Context, input PlatformInput) error { - logger := workflow.GetLogger(ctx) - logger.Info("Platform workflow started", "platform", input) - - var workspace string - if err := workflow.ExecuteActivity( - workflow.WithActivityOptions(ctx, workflow.ActivityOptions{ - StartToCloseTimeout: 1 * time.Minute, - }), - activities.Clone, - input.Url, - input.Revision, - ).Get(ctx, &workspace); err != nil { - return err - } - - defer os.RemoveAll(workspace) - - appsDir := workspace + "/apps" - - // TODO this should be a separate activity - var matchedPaths []string - if err := workflow.ExecuteActivity( - workflow.WithActivityOptions(ctx, workflow.ActivityOptions{ - StartToCloseTimeout: 10 * time.Second, - }), - activities.DiscoverApps, - appsDir, - input.Cluster, - ).Get(ctx, &matchedPaths); err != nil { - return err - } - ctx = workflow.WithActivityOptions(ctx, workflow.ActivityOptions{ - StartToCloseTimeout: 1 * time.Minute, - }) - - var futures []workflow.Future - var results []activities.PushResult - - for _, yamlPath := range matchedPaths { - parts := strings.Split(filepath.ToSlash(yamlPath), "/") - if len(parts) < 4 { - logger.Warn("Skipping invalid path", "path", yamlPath) - continue - } - - namespace := parts[len(parts)-3] - app := parts[len(parts)-2] - - logger.Info("Dispatching PushRenderedHelm", "path", yamlPath, "namespace", namespace, "app", app) - - fut := workflow.ExecuteActivity( - ctx, - activities.PushRenderedApp, - appsDir, - namespace, - app, - input.Cluster, - input.Registry, - ) - futures = append(futures, fut) - } - - for _, fut := range futures { - var result activities.PushResult - if err := fut.Get(ctx, &result); err != nil { - return err - } - results = append(results, result) - } - - logger.Info("Finished pushing all matching apps", "count", len(results)) - return nil -} diff --git a/controller/workflows/infra.go b/controller/workflows/infra.go deleted file mode 100644 index 5d04ea0..0000000 --- a/controller/workflows/infra.go +++ /dev/null @@ -1,101 +0,0 @@ -package workflows - -import ( - "fmt" - "os" - "time" - - "cloudlab/controller/activities" - - "go.temporal.io/sdk/temporal" - "go.temporal.io/sdk/workflow" -) - -type InfraInputs struct { - Url string - Revision string - OldRevision string - Stack string -} - -func Infra(ctx workflow.Context, input InfraInputs) (*activities.Graph, error) { - logger := workflow.GetLogger(ctx) - logger.Info("Infra workflow started", "infra", input) - - cloneCtx := workflow.WithActivityOptions(ctx, workflow.ActivityOptions{ - StartToCloseTimeout: 1 * time.Minute, - }) - - var workspace string - if err := workflow.ExecuteActivity(cloneCtx, activities.Clone, input.Url, input.Revision).Get(ctx, &workspace); err != nil { - return nil, err - } - - defer os.RemoveAll(workspace) - - // Graph and analysis activities: moderate timeout - analysisCtx := workflow.WithActivityOptions(ctx, workflow.ActivityOptions{ - StartToCloseTimeout: 5 * time.Second, - RetryPolicy: &temporal.RetryPolicy{ - MaximumAttempts: 1, - }, - }) - - var graph *activities.Graph - var prunedGraph *activities.Graph - - // Get the terragrunt graph - if err := workflow.ExecuteActivity(analysisCtx, activities.TerragruntGraph, workspace+"/infra/"+input.Stack).Get(ctx, &graph); err != nil { - return nil, err - } - - // If oldRevision is not provided, use the full graph (no pruning) - if input.OldRevision == "" { - logger.Info("No oldRevision provided, using full graph", "nodes", len(graph.Nodes)) - prunedGraph = graph - } else { - // Determine changed modules and prune graph - var changedModules []string - if err := workflow.ExecuteActivity(analysisCtx, activities.ChangedModules, workspace, input.OldRevision).Get(ctx, &changedModules); err != nil { - return nil, err - } - - if err := workflow.ExecuteActivity(analysisCtx, activities.PruneGraph, graph, changedModules).Get(ctx, &prunedGraph); err != nil { - return nil, err - } - - logger.Info("Graph pruning completed", "nodes", len(prunedGraph.Nodes)) - } - - for levelIndex, level := range prunedGraph.TopologicalSort() { - logger.Info("Starting terragrunt apply", "level", levelIndex, "modules", level) - - var futures []workflow.Future - for _, module := range level { - moduleCtx := workflow.WithActivityOptions(ctx, workflow.ActivityOptions{ - StartToCloseTimeout: 30 * time.Minute, - HeartbeatTimeout: 2 * time.Minute, - Summary: fmt.Sprintf("%s/%s", input.Stack, module), - RetryPolicy: &temporal.RetryPolicy{ - MaximumAttempts: 2, - NonRetryableErrorTypes: []string{ - "TerraformValidationError", - "TerraformPlanError", - }, - }, - }) - futures = append(futures, workflow.ExecuteActivity(moduleCtx, activities.TerragruntApply, input.Url, input.Revision, module, input.Stack)) - } - - for i, future := range futures { - if err := future.Get(ctx, nil); err != nil { - logger.Error("TerragruntApply failed", "module", level[i], "level", levelIndex, "error", err) - return nil, err - } - logger.Info("Module apply completed", "module", level[i], "level", levelIndex) - } - } - - logger.Info("Infra workflow completed", "levels", len(prunedGraph.TopologicalSort()), "modules", len(prunedGraph.Nodes)) - return prunedGraph, nil -} diff --git a/controller/workflows/infra_test.go b/controller/workflows/infra_test.go deleted file mode 100644 index ca2a0cb..0000000 --- a/controller/workflows/infra_test.go +++ /dev/null @@ -1,465 +0,0 @@ -package workflows - -import ( - "context" - "errors" - "testing" - "time" - - "cloudlab/controller/activities" - - "github.com/stretchr/testify/mock" - "github.com/stretchr/testify/suite" - "go.temporal.io/sdk/testsuite" -) - -type InfraWorkflowTestSuite struct { - suite.Suite - testsuite.WorkflowTestSuite - - env *testsuite.TestWorkflowEnvironment -} - -func (s *InfraWorkflowTestSuite) SetupTest() { - s.env = s.NewTestWorkflowEnvironment() - // Set a reasonable timeout for tests - s.env.SetTestTimeout(30 * time.Second) -} - -func (s *InfraWorkflowTestSuite) AfterTest(suiteName, testName string) { - s.env.AssertExpectations(s.T()) -} - -func (s *InfraWorkflowTestSuite) TestInfraWorkflow_Success() { - input := InfraInputs{ - Url: "https://github.com/example/repo.git", - Revision: "main", - OldRevision: "HEAD~1", - Stack: "dev", - } - repoPath := "/tmp/infra-12345" - changedModules := []string{"module1", "module2"} - - graph := &activities.Graph{ - Nodes: map[string]bool{ - "module1": true, - "module2": true, - }, - Edges: map[string][]string{ - "module1": {"module2"}, // module1 depends on module2 - }, - } - - prunedGraph := graph // Both modules changed - - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return(repoPath, nil) - s.env.OnActivity(activities.TerragruntGraph, mock.Anything, repoPath+"/infra/"+input.Stack).Return(graph, nil) - s.env.OnActivity(activities.ChangedModules, mock.Anything, repoPath, input.OldRevision).Return(changedModules, nil) - s.env.OnActivity(activities.PruneGraph, mock.Anything, graph, changedModules).Return(prunedGraph, nil) - s.env.OnActivity(activities.TerragruntApply, mock.Anything, input.Url, input.Revision, "module2", input.Stack).Return(nil) - s.env.OnActivity(activities.TerragruntApply, mock.Anything, input.Url, input.Revision, "module1", input.Stack).Return(nil) - - s.env.ExecuteWorkflow(Infra, input) - - s.True(s.env.IsWorkflowCompleted()) - s.NoError(s.env.GetWorkflowError()) -} - -func (s *InfraWorkflowTestSuite) TestInfraWorkflow_CloneFailure() { - input := InfraInputs{ - Url: "https://github.com/example/invalid-repo.git", - Revision: "main", - OldRevision: "HEAD~1", - Stack: "dev", - } - - // Mock Clone to return error - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return("", errors.New("repository not found")) - - s.env.ExecuteWorkflow(Infra, input) - - s.True(s.env.IsWorkflowCompleted()) - s.Error(s.env.GetWorkflowError()) - s.Contains(s.env.GetWorkflowError().Error(), "repository not found") -} - -func (s *InfraWorkflowTestSuite) TestInfraWorkflow_TerragruntGraphFailure() { - input := InfraInputs{ - Url: "https://github.com/example/repo.git", - Revision: "main", - OldRevision: "HEAD~1", - Stack: "dev", - } - repoPath := "/tmp/infra-12345" - - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return(repoPath, nil) - s.env.OnActivity(activities.TerragruntGraph, mock.Anything, repoPath+"/infra/"+input.Stack).Return( - (*activities.Graph)(nil), errors.New("terragrunt dag graph failed")) - - s.env.ExecuteWorkflow(Infra, input) - - s.True(s.env.IsWorkflowCompleted()) - s.Error(s.env.GetWorkflowError()) - s.Contains(s.env.GetWorkflowError().Error(), "terragrunt dag graph failed") -} - -func (s *InfraWorkflowTestSuite) TestInfraWorkflow_ChangedModulesFailure() { - input := InfraInputs{ - Url: "https://github.com/example/repo.git", - Revision: "main", - OldRevision: "HEAD~1", - Stack: "dev", - } - repoPath := "/tmp/infra-12345" - graph := &activities.Graph{ - Nodes: map[string]bool{"module1": true}, - Edges: map[string][]string{}, - } - - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return(repoPath, nil) - s.env.OnActivity(activities.TerragruntGraph, mock.Anything, repoPath+"/infra/"+input.Stack).Return(graph, nil) - s.env.OnActivity(activities.ChangedModules, mock.Anything, repoPath, input.OldRevision).Return( - []string{}, errors.New("git diff failed")) - - s.env.ExecuteWorkflow(Infra, input) - - s.True(s.env.IsWorkflowCompleted()) - s.Error(s.env.GetWorkflowError()) - s.Contains(s.env.GetWorkflowError().Error(), "git diff failed") -} - -func (s *InfraWorkflowTestSuite) TestInfraWorkflow_TerragruntApplyFailure() { - input := InfraInputs{ - Url: "https://github.com/example/repo.git", - Revision: "main", - OldRevision: "HEAD~1", - Stack: "dev", - } - repoPath := "/tmp/infra-12345" - changedModules := []string{"module1"} - - graph := &activities.Graph{ - Nodes: map[string]bool{"module1": true}, - Edges: map[string][]string{}, - } - - prunedGraph := &activities.Graph{ - Nodes: map[string]bool{"module1": true}, - Edges: map[string][]string{}, - } - - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return(repoPath, nil) - s.env.OnActivity(activities.TerragruntGraph, mock.Anything, repoPath+"/infra/"+input.Stack).Return(graph, nil) - s.env.OnActivity(activities.ChangedModules, mock.Anything, repoPath, input.OldRevision).Return(changedModules, nil) - s.env.OnActivity(activities.PruneGraph, mock.Anything, graph, changedModules).Return(prunedGraph, nil) - s.env.OnActivity(activities.TerragruntApply, mock.Anything, input.Url, input.Revision, "module1", input.Stack).Return( - errors.New("terragrunt apply failed: resource conflict")) - - s.env.ExecuteWorkflow(Infra, input) - - s.True(s.env.IsWorkflowCompleted()) - s.Error(s.env.GetWorkflowError()) - s.Contains(s.env.GetWorkflowError().Error(), "terragrunt apply failed") -} - -func (s *InfraWorkflowTestSuite) TestInfraWorkflow_ComplexDependencyGraph() { - input := InfraInputs{ - Url: "https://github.com/example/repo.git", - Revision: "main", - OldRevision: "HEAD~1", - Stack: "prod", - } - repoPath := "/tmp/infra-67890" - changedModules := []string{"vpc", "database"} - - // Complex dependency graph: - // app -> [database, loadbalancer] - // database -> vpc - // loadbalancer -> vpc - // monitoring -> app - graph := &activities.Graph{ - Nodes: map[string]bool{ - "vpc": true, - "database": true, - "loadbalancer": true, - "app": true, - "monitoring": true, - }, - Edges: map[string][]string{ - "app": {"database", "loadbalancer"}, - "database": {"vpc"}, - "loadbalancer": {"vpc"}, - "monitoring": {"app"}, - }, - } - - // Pruned graph should contain changed modules and their dependents - prunedGraph := &activities.Graph{ - Nodes: map[string]bool{ - "vpc": true, - "database": true, - "app": true, - "monitoring": true, - }, - Edges: map[string][]string{ - "app": {"database"}, - "database": {"vpc"}, - "monitoring": {"app"}, - }, - } - - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return(repoPath, nil) - s.env.OnActivity(activities.TerragruntGraph, mock.Anything, repoPath+"/infra/"+input.Stack).Return(graph, nil) - s.env.OnActivity(activities.ChangedModules, mock.Anything, repoPath, input.OldRevision).Return(changedModules, nil) - s.env.OnActivity(activities.PruneGraph, mock.Anything, graph, changedModules).Return(prunedGraph, nil) - - // Mock TerragruntApply calls in dependency order - // Level 0: vpc - s.env.OnActivity(activities.TerragruntApply, mock.Anything, input.Url, input.Revision, "vpc", input.Stack).Return(nil) - // Level 1: database - s.env.OnActivity(activities.TerragruntApply, mock.Anything, input.Url, input.Revision, "database", input.Stack).Return(nil) - // Level 2: app - s.env.OnActivity(activities.TerragruntApply, mock.Anything, input.Url, input.Revision, "app", input.Stack).Return(nil) - // Level 3: monitoring - s.env.OnActivity(activities.TerragruntApply, mock.Anything, input.Url, input.Revision, "monitoring", input.Stack).Return(nil) - - s.env.ExecuteWorkflow(Infra, input) - - s.True(s.env.IsWorkflowCompleted()) - s.NoError(s.env.GetWorkflowError()) - - var result *activities.Graph - s.NoError(s.env.GetWorkflowResult(&result)) - s.True(result.Nodes["vpc"]) - s.True(result.Nodes["database"]) - s.True(result.Nodes["app"]) - s.True(result.Nodes["monitoring"]) -} - -func (s *InfraWorkflowTestSuite) TestInfraWorkflow_NoChangedModules() { - input := InfraInputs{ - Url: "https://github.com/example/repo.git", - Revision: "main", - OldRevision: "HEAD~1", - Stack: "dev", - } - repoPath := "/tmp/infra-12345" - changedModules := []string{} // No changes - - graph := &activities.Graph{ - Nodes: map[string]bool{ - "module1": true, - "module2": true, - }, - Edges: map[string][]string{ - "module1": {"module2"}, - }, - } - - // Pruned graph should be empty - prunedGraph := &activities.Graph{ - Nodes: map[string]bool{}, - Edges: map[string][]string{}, - } - - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return(repoPath, nil) - s.env.OnActivity(activities.TerragruntGraph, mock.Anything, repoPath+"/infra/"+input.Stack).Return(graph, nil) - s.env.OnActivity(activities.ChangedModules, mock.Anything, repoPath, input.OldRevision).Return(changedModules, nil) - s.env.OnActivity(activities.PruneGraph, mock.Anything, graph, changedModules).Return(prunedGraph, nil) - - // No TerragruntApply calls should be made since no modules to deploy - - s.env.ExecuteWorkflow(Infra, input) - - s.True(s.env.IsWorkflowCompleted()) - s.NoError(s.env.GetWorkflowError()) - - var result *activities.Graph - s.NoError(s.env.GetWorkflowResult(&result)) - s.Empty(result.Nodes) -} - -func (s *InfraWorkflowTestSuite) TestInfraWorkflow_ActivityTimeout() { - input := InfraInputs{ - Url: "https://github.com/example/repo.git", - Revision: "main", - OldRevision: "HEAD~1", - Stack: "dev", - } - - // Mock Clone to simulate a timeout scenario - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return( - "", errors.New("activity timeout")) - - s.env.ExecuteWorkflow(Infra, input) - - s.True(s.env.IsWorkflowCompleted()) - s.Error(s.env.GetWorkflowError()) -} - -func (s *InfraWorkflowTestSuite) TestInfraWorkflow_ParallelExecution() { - // Test that modules at the same dependency level are executed in parallel - input := InfraInputs{ - Url: "https://github.com/example/repo.git", - Revision: "main", - OldRevision: "HEAD~1", - Stack: "dev", - } - repoPath := "/tmp/infra-12345" - changedModules := []string{"module-a", "module-b", "module-c"} - - // Graph with parallel modules: - // module-a and module-b can run in parallel (both depend on module-c) - graph := &activities.Graph{ - Nodes: map[string]bool{ - "module-a": true, - "module-b": true, - "module-c": true, - }, - Edges: map[string][]string{ - "module-a": {"module-c"}, - "module-b": {"module-c"}, - }, - } - - prunedGraph := graph // All modules changed - - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return(repoPath, nil) - s.env.OnActivity(activities.TerragruntGraph, mock.Anything, repoPath+"/infra/"+input.Stack).Return(graph, nil) - s.env.OnActivity(activities.ChangedModules, mock.Anything, repoPath, input.OldRevision).Return(changedModules, nil) - s.env.OnActivity(activities.PruneGraph, mock.Anything, graph, changedModules).Return(prunedGraph, nil) - - // Level 0: module-c - s.env.OnActivity(activities.TerragruntApply, mock.Anything, input.Url, input.Revision, "module-c", input.Stack).Return(nil) - // Level 1: module-a and module-b (should execute in parallel) - s.env.OnActivity(activities.TerragruntApply, mock.Anything, input.Url, input.Revision, "module-a", input.Stack).Return(nil) - s.env.OnActivity(activities.TerragruntApply, mock.Anything, input.Url, input.Revision, "module-b", input.Stack).Return(nil) - - s.env.ExecuteWorkflow(Infra, input) - - s.True(s.env.IsWorkflowCompleted()) - s.NoError(s.env.GetWorkflowError()) -} - -func (s *InfraWorkflowTestSuite) TestInfraWorkflow_WorkerFailureRetry() { - // Test that TerragruntApply can handle worker failure and retry on a different worker - // by ensuring the activity is self-contained (clones repo internally) - input := InfraInputs{ - Url: "https://github.com/example/repo.git", - Revision: "main", - OldRevision: "HEAD~1", - Stack: "dev", - } - repoPath := "/tmp/infra-12345" - newWorkerRepoPath := "/tmp/infra-67890" - changedModules := []string{"module1"} - - graph := &activities.Graph{ - Nodes: map[string]bool{ - "module1": true, - }, - Edges: map[string][]string{}, - } - - prunedGraph := graph // Only module1 changed - - // Initial workflow activities (successful) - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return(repoPath, nil) - s.env.OnActivity(activities.TerragruntGraph, mock.Anything, repoPath+"/infra/"+input.Stack).Return(graph, nil) - s.env.OnActivity(activities.ChangedModules, mock.Anything, repoPath, input.OldRevision).Return(changedModules, nil) - s.env.OnActivity(activities.PruneGraph, mock.Anything, graph, changedModules).Return(prunedGraph, nil) - - // Simulate worker failure and retry on different worker - applyCallCount := 0 - s.env.OnActivity(activities.TerragruntApply, mock.Anything, input.Url, input.Revision, "module1", input.Stack).Return( - func(ctx context.Context, repoUrl, revision, modulePath, stack string) error { - applyCallCount++ - if applyCallCount == 1 { - // First attempt fails (simulating worker failure) - return errors.New("worker failed: connection lost") - } - // Second attempt succeeds (activity is self-contained and clones repo again) - return nil - }) - - // Mock additional Clone calls for TerragruntApply retries - // The activity will call Clone internally to ensure repo availability - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return(newWorkerRepoPath, nil).Maybe() - - s.env.ExecuteWorkflow(Infra, input) - - s.True(s.env.IsWorkflowCompleted()) - s.NoError(s.env.GetWorkflowError()) - - var result *activities.Graph - s.NoError(s.env.GetWorkflowResult(&result)) - s.True(result.Nodes["module1"]) -} - -func (s *InfraWorkflowTestSuite) TestInfraWorkflow_NoOldRevisionProvided() { - // Test that when oldRevision is not provided, all modules are deployed - input := InfraInputs{ - Url: "https://github.com/example/repo.git", - Revision: "main", - OldRevision: "", // Empty string - no old revision provided - Stack: "dev", - } - repoPath := "/tmp/infra-12345" - - // Full graph with all modules - graph := &activities.Graph{ - Nodes: map[string]bool{ - "vpc": true, - "database": true, - "loadbalancer": true, - "app": true, - "monitoring": true, - }, - Edges: map[string][]string{ - "app": {"database", "loadbalancer"}, - "database": {"vpc"}, - "loadbalancer": {"vpc"}, - "monitoring": {"app"}, - }, - } - - // When oldRevision is empty, the workflow should use the full graph - // without calling ChangedModules or PruneGraph activities - - s.env.OnActivity(activities.Clone, mock.Anything, input.Url, input.Revision).Return(repoPath, nil) - s.env.OnActivity(activities.TerragruntGraph, mock.Anything, repoPath+"/infra/"+input.Stack).Return(graph, nil) - - // Mock TerragruntApply calls in dependency order for all modules - // Level 0: vpc - s.env.OnActivity(activities.TerragruntApply, mock.Anything, input.Url, input.Revision, "vpc", input.Stack).Return(nil) - // Level 1: database, loadbalancer - s.env.OnActivity(activities.TerragruntApply, mock.Anything, input.Url, input.Revision, "database", input.Stack).Return(nil) - s.env.OnActivity(activities.TerragruntApply, mock.Anything, input.Url, input.Revision, "loadbalancer", input.Stack).Return(nil) - // Level 2: app - s.env.OnActivity(activities.TerragruntApply, mock.Anything, input.Url, input.Revision, "app", input.Stack).Return(nil) - // Level 3: monitoring - s.env.OnActivity(activities.TerragruntApply, mock.Anything, input.Url, input.Revision, "monitoring", input.Stack).Return(nil) - - s.env.ExecuteWorkflow(Infra, input) - - s.True(s.env.IsWorkflowCompleted()) - s.NoError(s.env.GetWorkflowError()) - - var result *activities.Graph - s.NoError(s.env.GetWorkflowResult(&result)) - - // Verify that all modules are in the result graph - s.True(result.Nodes["vpc"]) - s.True(result.Nodes["database"]) - s.True(result.Nodes["loadbalancer"]) - s.True(result.Nodes["app"]) - s.True(result.Nodes["monitoring"]) - - // Verify that the result graph has the same structure as the original - s.Equal(len(graph.Nodes), len(result.Nodes)) - s.Equal(len(graph.Edges), len(result.Edges)) -} - -func TestInfraWorkflowTestSuite(t *testing.T) { - suite.Run(t, new(InfraWorkflowTestSuite)) -} diff --git a/controller/workflows/platform.go b/controller/workflows/platform.go deleted file mode 100644 index 1087e81..0000000 --- a/controller/workflows/platform.go +++ /dev/null @@ -1,50 +0,0 @@ -package workflows - -import ( - "os" - "time" - - "cloudlab/controller/activities" - - "go.temporal.io/sdk/workflow" -) - -type PlatformInput struct { - Url string - Revision string - Registry string - Cluster string -} - -func Platform(ctx workflow.Context, input PlatformInput) error { - logger := workflow.GetLogger(ctx) - logger.Info("Platform workflow started", "platform", input) - - var workspace string - if err := workflow.ExecuteActivity( - workflow.WithActivityOptions(ctx, workflow.ActivityOptions{ - StartToCloseTimeout: 1 * time.Minute, - }), - activities.Clone, - input.Url, - input.Revision, - ).Get(ctx, &workspace); err != nil { - return err - } - - defer os.RemoveAll(workspace) - - var pushResult *activities.PushResult - if err := workflow.ExecuteActivity( - workflow.WithActivityOptions(ctx, workflow.ActivityOptions{ - StartToCloseTimeout: 1 * time.Minute, - }), - activities.PushManifests, - workspace+"/platform/"+input.Cluster, - input.Registry+"/platform:"+input.Cluster, - ).Get(ctx, &pushResult); err != nil { - return err - } - - return nil -} diff --git a/platform/local/app-engine.yaml b/platform/local/app-engine.yaml deleted file mode 100644 index 5c52520..0000000 --- a/platform/local/app-engine.yaml +++ /dev/null @@ -1,60 +0,0 @@ -apiVersion: argoproj.io/v1alpha1 -kind: Application -metadata: - finalizers: - - resources-finalizer.argocd.argoproj.io - name: app-engine -spec: - destination: - name: in-cluster - namespace: app-engine - project: default - syncPolicy: - automated: - prune: true - selfHeal: true - syncOptions: - - CreateNamespace=true - - ApplyOutOfSyncOnly=true - - ServerSideApply=true - source: - repoURL: https://bjw-s-labs.github.io/helm-charts - chart: app-template - targetRevision: 3.7.3 - helm: - valuesObject: - defaultPodOptions: - restartPolicy: Always - labels: - istio.io/dataplane-mode: ambient - hostNetwork: true - controllers: - worker: - strategy: RollingUpdate - containers: - app: - image: - # TODO bootstrap and build itself - # repository: registry.registry.svc.cluster.local/khuedoan/app-engine - repository: docker.io/khuedoan/app-engine - tag: 4118f906ab07a17f3dac608f1a690b2215e4d2a5 - pullPolicy: Always - env: - TEMPORAL_URL: http://temporal-frontend.temporal:7233 - REGISTRY: registry.registry.svc.cluster.local - docker: - image: - repository: docker.io/library/docker - tag: 27-dind - command: - - dockerd - - --host=unix:///var/run/docker.sock - - --insecure-registry=registry.registry.svc.cluster.local - securityContext: - privileged: true - persistence: - socket: - type: emptyDir - globalMounts: - - path: /var/run - subPath: docker.sock