diff --git a/cmd/ci/deploy.go b/cmd/ci/deploy.go index 4ebe8844..e1bed59a 100644 --- a/cmd/ci/deploy.go +++ b/cmd/ci/deploy.go @@ -65,7 +65,6 @@ func runDeployService(ctx context.Context, workspace *resources.Workspace, modul func initDeployService(ctx context.Context, workspace *resources.Workspace, module *resources.Module, service *resources.Service) (*orchestration.Flow, *platform.DeploymentManager, error) { w := wool.Get(ctx).In("deployService", wool.ThisField(resources.WithUnique(service))) - orchestration.SetDryRun(dryRun) env, err := orchestration.SelectEnvironment(workspace, envInput) if err != nil { return nil, nil, w.Wrap(err) @@ -87,9 +86,11 @@ func initDeployService(ctx context.Context, workspace *resources.Workspace, modu if err != nil { return nil, nil, stopFlowAfterError(flow, w.Wrap(err)) } - deploymentManager := platform.NewDeploymentManager(ctx, workspace, env) - - flow.WithDeploymentManager(deploymentManager) + var deploymentManager *platform.DeploymentManager + if !dryRun { + deploymentManager = platform.NewDeploymentManager(ctx, workspace, env) + flow.WithDeploymentManager(deploymentManager) + } return flow, deploymentManager, nil } @@ -99,9 +100,11 @@ func deployService(ctx context.Context, flow *orchestration.Flow, deploymentMana if err != nil { return w.Wrapf(err, "cannot start service") } - err = deploymentManager.Deploy(ctx, workspace) - if err != nil { - return w.Wrapf(err, "cannot deploy service") + if deploymentManager != nil { + err = deploymentManager.Deploy(ctx, workspace) + if err != nil { + return w.Wrapf(err, "cannot deploy service") + } } return nil diff --git a/cmd/deploy/environment_test.go b/cmd/deploy/environment_test.go index da686ae9..41c0734c 100644 --- a/cmd/deploy/environment_test.go +++ b/cmd/deploy/environment_test.go @@ -62,3 +62,70 @@ endpoints: require.Contains(t, err.Error(), `"production"`) require.Contains(t, err.Error(), "deploy-env") } + +func TestDirectApplyRequestedTreatsDryRunAndRenderOnlyAsNoMutation(t *testing.T) { + previousDryRun := dryRun + previousRenderOnly := renderOnly + t.Cleanup(func() { + dryRun = previousDryRun + renderOnly = previousRenderOnly + }) + + for _, test := range []struct { + name string + dryRun bool + renderOnly bool + want bool + }{ + {name: "default", want: true}, + {name: "dry run", dryRun: true}, + {name: "render only", renderOnly: true}, + {name: "both", dryRun: true, renderOnly: true}, + } { + t.Run(test.name, func(t *testing.T) { + dryRun = test.dryRun + renderOnly = test.renderOnly + require.Equal(t, test.want, directApplyRequested()) + }) + } +} + +func TestInitDeployServiceRejectsRemoteDirectApplyBeforeStartingFlow(t *testing.T) { + previousEnv := envInput + previousDryRun := dryRun + previousRenderOnly := renderOnly + t.Cleanup(func() { + envInput = previousEnv + dryRun = previousDryRun + renderOnly = previousRenderOnly + }) + envInput = "production" + dryRun = false + renderOnly = false + workspace := &resources.Workspace{ + Name: "deploy-env", + Environments: []*resources.Environment{{ + Name: "production", + Cluster: &resources.EnvironmentCluster{ + Kind: "eks", + Kubeconfig: "/does/not/exist", + Context: "k3d-production", + }, + }}, + } + service := &resources.Service{Name: "gateway"} + service.WithModule("web") + + flow, err := initDeployService( + context.Background(), + workspace, + &resources.Module{Name: "web"}, + service, + true, + ) + + require.Nil(t, flow) + require.Error(t, err) + require.Contains(t, err.Error(), "exact local k3d target") + require.Contains(t, err.Error(), "--render-only") +} diff --git a/cmd/deploy/module.go b/cmd/deploy/module.go index 8ea12c81..c603a2f3 100644 --- a/cmd/deploy/module.go +++ b/cmd/deploy/module.go @@ -1,12 +1,9 @@ package deploy import ( - "bytes" "context" "fmt" - "os/exec" "path" - "strings" "github.com/codefly-dev/cli/cmd/common" "github.com/codefly-dev/cli/pkg/cli" @@ -23,13 +20,14 @@ import ( // module-level kustomize layer (namespace, AppProject, Applications, // shared Ingress, etc.) at module/deployment/kustomize/overlays/. // -// Two modes via --render-only: +// Apply mode is the default. Both --render-only and --dry-run are +// no-mutation modes: // // (default) apply mode — agents render + LocalApplyManager kubectl-applies // each service, then the module-level kustomize is applied via // kubectl. Used for local-k3d direct-deploy. // -// --render-only — agents render but skip apply. Module-level +// --render-only/--dry-run — agents render but skip apply. Module-level // kustomize is NOT applied either (it's static YAML committed // to git). Used for the gitops/ArgoCD flow where ArgoCD syncs // the rendered tree from git. @@ -57,9 +55,21 @@ var ModuleCmd = &cobra.Command{ return err } + var deploymentManager deployments.Manager + var localApplyManager *deployments.LocalApplyManager + if directApplyRequested() { + localApplyManager, err = deployments.NewLocalApplyManager(ctx, workspace, env) + if err != nil { + return err + } + deploymentManager = localApplyManager + } else { + deploymentManager = deployments.NewRenderManager(workspace, env) + } + cli.Header(1, "Deploying module %s to env %s", module.Name, env.Name) - if renderOnly { - cli.Header(2, "render-only mode — manifests written to disk, no kubectl apply") + if !directApplyRequested() { + cli.Header(2, "no-mutation mode — manifests written to disk, no kubectl apply") } // Phase 1: per-service flows. Each service is its own Flow with @@ -69,17 +79,17 @@ var ModuleCmd = &cobra.Command{ // than a clean failure to investigate. for _, ref := range module.ServiceReferences { cli.Header(2, "Deploying service %s", ref.Name) - if err := deployOneService(ctx, workspace, module, ref.Name, env); err != nil { + if err := deployOneService(ctx, workspace, module, ref.Name, env, deploymentManager); err != nil { return fmt.Errorf("cannot deploy service %s: %w", ref.Name, err) } } - // Phase 2: module-level kustomize. Skipped in render-only mode + // Phase 2: module-level kustomize. Skipped in no-mutation modes // (the module-level YAML is static; nothing to render — it's // committed already). Skipped silently when the module hasn't // scaffolded a deployment/ folder yet. - if !renderOnly { - if err := applyModuleKustomize(ctx, module, env); err != nil { + if directApplyRequested() { + if err := applyModuleKustomize(ctx, module, env, localApplyManager); err != nil { return fmt.Errorf("cannot apply module-level kustomize: %w", err) } } @@ -92,7 +102,7 @@ var ModuleCmd = &cobra.Command{ // deployOneService runs a one-shot deploy Flow for a single service. // Mirrors initDeployService in service.go but inline-built so the // module loop can iterate without goroutine indirection. -func deployOneService(ctx context.Context, workspace *resources.Workspace, module *resources.Module, name string, env *resources.Environment) error { +func deployOneService(ctx context.Context, workspace *resources.Workspace, module *resources.Module, name string, env *resources.Environment, deploymentManager deployments.Manager) error { w := wool.Get(ctx).In("deployModule.deployOneService", wool.NameField(name)) service, err := module.LoadServiceFromName(ctx, name) @@ -100,8 +110,6 @@ func deployOneService(ctx context.Context, workspace *resources.Workspace, modul return w.Wrapf(err, "cannot load service") } - orchestration.SetDryRun(dryRun) - flow, err := orchestration.NewFlow(ctx, workspace, module, service, env, orchestration.DeployMode) if err != nil { return w.Wrap(err) @@ -123,9 +131,7 @@ func deployOneService(ctx context.Context, workspace *resources.Workspace, modul if err := flow.Load(ctx); err != nil { return w.Wrap(err) } - if !renderOnly { - flow.WithDeploymentManager(deployments.NewLocalApplyManager(ctx, workspace, env)) - } + flow.WithDeploymentManager(deploymentManager) if err := flow.Deploy(ctx); err != nil { return w.Wrapf(err, "deploy failed") @@ -139,7 +145,7 @@ func deployOneService(ctx context.Context, workspace *resources.Workspace, modul // for the module-level overlay matching the env. Silently no-ops when // the module hasn't scaffolded module/deployment/kustomize/ yet — // services-only modules are valid. -func applyModuleKustomize(ctx context.Context, module *resources.Module, env *resources.Environment) error { +func applyModuleKustomize(ctx context.Context, module *resources.Module, env *resources.Environment, manager *deployments.LocalApplyManager) error { w := wool.Get(ctx).In("applyModuleKustomize", wool.NameField(module.Name)) dir := path.Join(module.Dir(), "deployment", "kustomize", "overlays", env.Name) @@ -154,51 +160,7 @@ func applyModuleKustomize(ctx context.Context, module *resources.Module, env *re cli.Header(2, "Applying module-level kustomize at %s", dir) - // Render with kustomize. Fails loudly if the overlay is malformed — - // better here than mid-apply with a half-applied state. - build := exec.CommandContext(ctx, "kustomize", "build", dir) - var stdout, stderr bytes.Buffer - build.Stdout = &stdout - build.Stderr = &stderr - if err := build.Run(); err != nil { - return w.Wrapf(err, "kustomize build failed: %s", stderr.String()) - } - - // Resolve kubeconfig the same way the per-service path does, so a - // single workspace.codefly.yaml env declaration drives both layers. - configPath, err := deployments.GetK8sConfig(ctx, env) - if err != nil { - return w.Wrapf(err, "cannot get k8s config") - } - - // Apply via kubectl. Splitting on `---` keeps the per-resource log - // shape consistent with how cli/pkg/deployments/kubernetes.go - // applies per-service manifests. - objs := strings.Split(stdout.String(), "---") - for _, obj := range objs { - if strings.TrimSpace(obj) == "" { - continue - } - args := []string{"--kubeconfig", configPath} - if env.Cluster != nil && env.Cluster.Context != "" { - args = append(args, "--context", env.Cluster.Context) - } - args = append(args, "apply", "-f", "-") - - apply := exec.CommandContext(ctx, "kubectl", args...) - apply.Stdin = strings.NewReader(obj) - var out, errOut bytes.Buffer - apply.Stdout = &out - apply.Stderr = &errOut - if err := apply.Run(); err != nil { - return w.Wrapf(err, "kubectl apply failed: %s", errOut.String()) - } - line := strings.TrimSpace(out.String()) - if line != "" && !strings.Contains(line, "unchanged") { - cli.Info("%s", line) - } - } - return nil + return manager.ApplyModuleKustomize(ctx, module, dir) } func init() { @@ -207,6 +169,6 @@ func init() { // Go variable being the receiver for two different commands' flags // — each command has its own pflag set. ModuleCmd.Flags().StringVar(&envInput, "env", "local", "Environment to deploy the module") - ModuleCmd.Flags().BoolVar(&dryRun, "dry-run", false, "Dry run the deployment") + ModuleCmd.Flags().BoolVar(&dryRun, "dry-run", false, "Render the deployment without applying it") ModuleCmd.Flags().BoolVar(&renderOnly, "render-only", false, "Render kustomize manifests to disk without applying. Used for gitops flows where ArgoCD/Flux syncs from the rendered tree.") } diff --git a/cmd/deploy/service.go b/cmd/deploy/service.go index 196bf419..45b5f235 100644 --- a/cmd/deploy/service.go +++ b/cmd/deploy/service.go @@ -77,11 +77,19 @@ var ServiceCmd = &cobra.Command{ func initDeployService(ctx context.Context, workspace *resources.Workspace, module *resources.Module, service *resources.Service, standAlone bool) (*orchestration.Flow, error) { w := wool.Get(ctx).In("deployService", wool.ThisField(resources.WithUnique(service))) - orchestration.SetDryRun(dryRun) env, err := orchestration.SelectEnvironment(workspace, envInput) if err != nil { return nil, w.Wrap(err) } + var deploymentManager deployments.Manager + if directApplyRequested() { + deploymentManager, err = deployments.NewLocalApplyManager(ctx, workspace, env) + if err != nil { + return nil, w.Wrap(err) + } + } else { + deploymentManager = deployments.NewRenderManager(workspace, env) + } flow, err := orchestration.NewFlow(ctx, workspace, module, service, env, orchestration.DeployMode) if err != nil { @@ -100,18 +108,7 @@ func initDeployService(ctx context.Context, workspace *resources.Workspace, modu return nil, w.Wrap(err) } - // Apply mode (default): the LocalApplyManager runs `kustomize build - // | kubectl apply` after each agent renders its manifests, plus - // imports built images into k3d when the env declares it. - // - // Render-only mode (--render-only): skip the manager wiring. Agents - // still write the rendered kustomize tree to disk via KustomizeDeploy - // (in builder_deploy.go), but no kubectl apply runs. ArgoCD or a - // separate gitops sync picks the rendered tree up from the workspace - // once it's committed. - if !renderOnly { - flow.WithDeploymentManager(deployments.NewLocalApplyManager(ctx, workspace, env)) - } + flow.WithDeploymentManager(deploymentManager) return flow, nil } @@ -135,9 +132,13 @@ var envInput string var dryRun bool var renderOnly bool +func directApplyRequested() bool { + return !renderOnly && !dryRun +} + func init() { ServiceCmd.Flags().StringVar(&envInput, "env", "local", "Environment to deploy the service") ServiceCmd.Flags().BoolVar(&standAlone, "stand-alone", false, "Begin service as standalone, i.e. without its dependencies") - ServiceCmd.Flags().BoolVar(&dryRun, "dry-run", false, "Dry run the deployment") + ServiceCmd.Flags().BoolVar(&dryRun, "dry-run", false, "Render the deployment without applying it") ServiceCmd.Flags().BoolVar(&renderOnly, "render-only", false, "Render kustomize manifests to disk without applying. Used for gitops flows where ArgoCD/Flux syncs from the rendered tree.") } diff --git a/pkg/control/deploy.go b/pkg/control/deploy.go index 494cc68a..d17cf3e9 100644 --- a/pkg/control/deploy.go +++ b/pkg/control/deploy.go @@ -19,11 +19,6 @@ func (p *planeImpl) Deploy(ctx context.Context, req DeployRequest) (DeployResult return p.runDeploy(ctx, req) } -// runDeploy builds a DeployMode flow and drives it, mirroring `codefly deploy -// service`. It constructs the flow inline (rather than via buildFlow) because it -// needs the workspace + environment to wire the apply manager. DryRun renders -// manifests without applying (no deployment manager); otherwise the local apply -// manager runs kustomize build | kubectl apply. func (p *planeImpl) runDeploy(ctx context.Context, req DeployRequest) (DeployResult, error) { if req.Module != "" && req.Service == "" { return DeployResult{}, fmt.Errorf("module-wide deploy is not yet supported via the control plane; specify a service") @@ -40,9 +35,20 @@ func (p *planeImpl) runDeploy(ctx context.Context, req DeployRequest) (DeployRes if err != nil { return DeployResult{}, fmt.Errorf("select environment %q: %w", envName, err) } - // Set unconditionally to the requested value so a prior deploy's dry-run - // toggle (a package-level global in orchestration) never leaks into this one. - orchestration.SetDryRun(req.DryRun) + var deploymentManager deployments.Manager + var evidenceProvider deployments.EvidenceProvider + if req.DryRun { + manager := deployments.NewRenderManager(ws, env) + deploymentManager = manager + evidenceProvider = manager + } else { + manager, managerErr := deployments.NewLocalApplyManager(ctx, ws, env) + if managerErr != nil { + return DeployResult{}, managerErr + } + deploymentManager = manager + evidenceProvider = manager + } flow, err := orchestration.NewFlow(ctx, ws, module, service, env, orchestration.DeployMode) if err != nil { @@ -54,15 +60,45 @@ func (p *planeImpl) runDeploy(ctx context.Context, req DeployRequest) (DeployRes if err := flow.Load(ctx); err != nil { return DeployResult{}, fmt.Errorf("load flow: %w", err) } - // Wire the apply manager AFTER Load (matching the deploy command). Skipped on - // dry-run so agents only render manifests to disk. - if !req.DryRun { - flow.WithDeploymentManager(deployments.NewLocalApplyManager(ctx, ws, env)) - } + flow.WithDeploymentManager(deploymentManager) defer stopFlow(flow) if err := flow.Deploy(ctx); err != nil { - return DeployResult{Succeeded: false}, fmt.Errorf("deploy %s: %w", req.Service, err) + result, _ := deployResult(false, evidenceProvider) + return result, fmt.Errorf("deploy %s: %w", req.Service, err) + } + result, err := deployResult(true, evidenceProvider) + if err != nil { + return result, fmt.Errorf("collect deployment evidence for %s: %w", req.Service, err) + } + return result, nil +} + +func deployResult(succeeded bool, provider deployments.EvidenceProvider) (DeployResult, error) { + result := DeployResult{Succeeded: succeeded} + evidence := provider.Evidence() + if evidence.Target != nil { + target := evidence.Target + result.Target = &DeployTarget{ + Kind: target.Kind, + Kubeconfig: target.Kubeconfig, + Context: target.Context, + Cluster: target.Cluster, + APIServer: target.APIServer, + K3dCluster: target.K3dCluster, + ClusterIdentity: target.ClusterIdentity, + } + } + for _, tree := range evidence.RenderedTrees { + result.RenderedTrees = append(result.RenderedTrees, RenderedTree{ + Module: tree.Module, + Service: tree.Service, + Digest: tree.Digest, + Manifests: tree.Manifests, + }) + } + if len(result.RenderedTrees) == 0 { + return result, fmt.Errorf("rendered-tree evidence is unavailable") } - return DeployResult{Succeeded: true}, nil + return result, nil } diff --git a/pkg/control/deploy_test.go b/pkg/control/deploy_test.go new file mode 100644 index 00000000..21713430 --- /dev/null +++ b/pkg/control/deploy_test.go @@ -0,0 +1,105 @@ +package control + +import ( + "context" + "os" + "path/filepath" + "strings" + "testing" + + "github.com/codefly-dev/cli/pkg/deployments" + "github.com/stretchr/testify/require" +) + +func TestRunDeployRejectsRemoteTargetBeforeStartingFlow(t *testing.T) { + root := writeWorkspace(t) + workspace := fixtureWorkspaceYAML + `environments: + - name: production + cluster: + kind: eks + kubeconfig: /does/not/exist + context: k3d-production +` + if err := os.WriteFile(filepath.Join(root, "workspace.codefly.yaml"), []byte(workspace), 0o600); err != nil { + t.Fatal(err) + } + plane, err := NewAt(root) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = plane.Close() }) + + _, err = plane.Deploy(context.Background(), DeployRequest{ + Service: "backend/api", + Env: "production", + }) + + if err == nil { + t.Fatal("remote direct deploy succeeded") + } + if got := err.Error(); !containsAll(got, "exact local k3d target", "--render-only", "GitOps") { + t.Fatalf("unexpected error: %v", err) + } +} + +func TestDeployResultIncludesEveryRenderedTreeAndExactTarget(t *testing.T) { + provider := staticEvidenceProvider{evidence: deployments.DeploymentEvidence{ + Target: &deployments.VerifiedKubernetesTarget{ + Kind: "k3d", + Kubeconfig: "/tmp/kubeconfig", + Context: "k3d-dev", + Cluster: "k3d-dev", + APIServer: "https://127.0.0.1:6443", + K3dCluster: "dev", + ClusterIdentity: "sha256:cluster", + }, + RenderedTrees: []deployments.RenderedTreeEvidence{ + {Module: "backend", Service: "api", Digest: "sha256:api", Manifests: "kind: Deployment\n"}, + {Module: "shared", Service: "database", Digest: "sha256:database", Manifests: "kind: StatefulSet\n"}, + }, + }} + + result, err := deployResult(true, provider) + + require.NoError(t, err) + require.True(t, result.Succeeded) + require.Equal(t, "sha256:cluster", result.Target.ClusterIdentity) + require.Equal(t, []RenderedTree{ + {Module: "backend", Service: "api", Digest: "sha256:api", Manifests: "kind: Deployment\n"}, + {Module: "shared", Service: "database", Digest: "sha256:database", Manifests: "kind: StatefulSet\n"}, + }, result.RenderedTrees) +} + +func TestDeployResultCarriesRenderedManifestsWithoutMutationTarget(t *testing.T) { + provider := staticEvidenceProvider{evidence: deployments.DeploymentEvidence{ + RenderedTrees: []deployments.RenderedTreeEvidence{{ + Module: "backend", + Service: "api", + Digest: "sha256:api", + Manifests: "kind: Deployment\n", + }}, + }} + + result, err := deployResult(true, provider) + + require.NoError(t, err) + require.Nil(t, result.Target) + require.Equal(t, "kind: Deployment\n", result.RenderedTrees[0].Manifests) +} + +func containsAll(value string, fragments ...string) bool { + for _, fragment := range fragments { + if !strings.Contains(value, fragment) { + return false + } + } + return true +} + +type staticEvidenceProvider struct { + evidence deployments.DeploymentEvidence +} + +func (p staticEvidenceProvider) Evidence() deployments.DeploymentEvidence { + return p.evidence +} diff --git a/pkg/control/mutation.go b/pkg/control/mutation.go index a76b513c..49b115aa 100644 --- a/pkg/control/mutation.go +++ b/pkg/control/mutation.go @@ -73,7 +73,7 @@ func (p *planeImpl) PrepareMutation(ctx context.Context, m Mutation) (PreparedMu return PreparedMutation{Token: token}, nil } -func (p *planeImpl) ApplyPreparedMutation(ctx context.Context, token PreparedMutation) error { +func (p *planeImpl) ApplyPreparedMutation(ctx context.Context, token PreparedMutation) (MutationResult, error) { p.gate.mu.Lock() pending, ok := p.gate.pending[token.Token] if ok { @@ -84,32 +84,36 @@ func (p *planeImpl) ApplyPreparedMutation(ctx context.Context, token PreparedMut p.gate.mu.Unlock() if !ok { - return fmt.Errorf("unknown or already-consumed prepared mutation") + return MutationResult{}, fmt.Errorf("unknown or already-consumed prepared mutation") } if time.Now().After(pending.expiresAt) { - return fmt.Errorf("prepared mutation expired") + return MutationResult{}, fmt.Errorf("prepared mutation expired") } - return p.executeMutation(ctx, pending.mutation) + return executeMutation(ctx, p, pending.mutation) } -// executeMutation dispatches a prepared mutation to the underlying operation. -func (p *planeImpl) executeMutation(ctx context.Context, m Mutation) error { +type mutationExecutor interface { + ApplyEdit(context.Context, Edit) error + runDeploy(context.Context, DeployRequest) (DeployResult, error) +} + +func executeMutation(ctx context.Context, executor mutationExecutor, m Mutation) (MutationResult, error) { switch m.Kind { case MutationFile: edit, ok := m.Payload.(Edit) if !ok { - return fmt.Errorf("file mutation payload must be an Edit, got %T", m.Payload) + return MutationResult{}, fmt.Errorf("file mutation payload must be an Edit, got %T", m.Payload) } - return p.ApplyEdit(ctx, edit) + return MutationResult{}, executor.ApplyEdit(ctx, edit) case MutationDeploy: req, ok := m.Payload.(DeployRequest) if !ok { - return fmt.Errorf("deploy mutation payload must be a DeployRequest, got %T", m.Payload) + return MutationResult{}, fmt.Errorf("deploy mutation payload must be a DeployRequest, got %T", m.Payload) } - _, err := p.runDeploy(ctx, req) - return err + result, err := executor.runDeploy(ctx, req) + return MutationResult{Deploy: &result}, err default: - return fmt.Errorf("unsupported mutation kind %q", m.Kind) + return MutationResult{}, fmt.Errorf("unsupported mutation kind %q", m.Kind) } } diff --git a/pkg/control/mutation_test.go b/pkg/control/mutation_test.go index b6c7296c..564bab05 100644 --- a/pkg/control/mutation_test.go +++ b/pkg/control/mutation_test.go @@ -39,7 +39,7 @@ func TestPreparedFileMutationAppliesOnceThenIsConsumed(t *testing.T) { } // First apply performs the edit. - if err := p.ApplyPreparedMutation(ctx, token); err != nil { + if _, err := p.ApplyPreparedMutation(ctx, token); err != nil { t.Fatal(err) } data, _ := p.ReadFile(ctx, fixtureFile) @@ -48,17 +48,45 @@ func TestPreparedFileMutationAppliesOnceThenIsConsumed(t *testing.T) { } // The token is single-use — a second apply must fail. - if err := p.ApplyPreparedMutation(ctx, token); err == nil { + if _, err := p.ApplyPreparedMutation(ctx, token); err == nil { t.Error("prepared mutation token should be single-use") } } func TestApplyUnknownTokenFails(t *testing.T) { - if err := New().ApplyPreparedMutation(context.Background(), PreparedMutation{Token: "deadbeef"}); err == nil { + if _, err := New().ApplyPreparedMutation(context.Background(), PreparedMutation{Token: "deadbeef"}); err == nil { t.Error("applying an unknown token should fail") } } +func TestExecuteDeployMutationReturnsDeploymentEvidence(t *testing.T) { + want := DeployResult{ + Succeeded: true, + RenderedTrees: []RenderedTree{{ + Module: "backend", + Service: "api", + Digest: "sha256:rendered", + Manifests: "kind: Deployment\n", + }}, + } + executor := mutationExecutorStub{deployResult: want} + + result, err := executeMutation(context.Background(), executor, Mutation{ + Kind: MutationDeploy, + Payload: DeployRequest{Service: "backend/api"}, + }) + + if err != nil { + t.Fatal(err) + } + if result.Deploy == nil { + t.Fatal("prepared deploy returned no deployment result") + } + if result.Deploy.RenderedTrees[0].Digest != want.RenderedTrees[0].Digest { + t.Fatalf("prepared deploy digest = %q, want %q", result.Deploy.RenderedTrees[0].Digest, want.RenderedTrees[0].Digest) + } +} + func TestDeployRefusedUnderPreparedAuthority(t *testing.T) { ctx := context.Background() p := New() @@ -71,3 +99,15 @@ func TestDeployRefusedUnderPreparedAuthority(t *testing.T) { t.Error("direct Deploy should be refused under prepared authority") } } + +type mutationExecutorStub struct { + deployResult DeployResult +} + +func (mutationExecutorStub) ApplyEdit(context.Context, Edit) error { + return nil +} + +func (s mutationExecutorStub) runDeploy(context.Context, DeployRequest) (DeployResult, error) { + return s.deployResult, nil +} diff --git a/pkg/control/plane.go b/pkg/control/plane.go index 33bfe2a1..e232c08b 100644 --- a/pkg/control/plane.go +++ b/pkg/control/plane.go @@ -164,5 +164,5 @@ type TerminalController interface { type MutationAuthority interface { ConfigureMutationAuthority(ctx context.Context, cfg AuthorityConfig) error PrepareMutation(ctx context.Context, m Mutation) (PreparedMutation, error) - ApplyPreparedMutation(ctx context.Context, token PreparedMutation) error + ApplyPreparedMutation(ctx context.Context, token PreparedMutation) (MutationResult, error) } diff --git a/pkg/control/types.go b/pkg/control/types.go index 502a17d2..281bace8 100644 --- a/pkg/control/types.go +++ b/pkg/control/types.go @@ -179,9 +179,27 @@ type DeployRequest struct { // DeployResult is the outcome of a deploy. type DeployResult struct { - Succeeded bool - Rendered string // manifests when DryRun/render-only - Output string + Succeeded bool + RenderedTrees []RenderedTree + Target *DeployTarget + Output string +} + +type RenderedTree struct { + Module string + Service string + Digest string + Manifests string +} + +type DeployTarget struct { + Kind string + Kubeconfig string + Context string + Cluster string + APIServer string + K3dCluster string + ClusterIdentity string } // --- Source --- @@ -439,3 +457,7 @@ type Mutation struct { type PreparedMutation struct { Token string } + +type MutationResult struct { + Deploy *DeployResult +} diff --git a/pkg/deployments/k3d.go b/pkg/deployments/k3d.go index 24ac0f9c..ec60faad 100644 --- a/pkg/deployments/k3d.go +++ b/pkg/deployments/k3d.go @@ -5,75 +5,45 @@ import ( "context" "fmt" "os/exec" - "strings" "github.com/codefly-dev/core/wool" ) -// DetectK3dCluster returns the k3d cluster name if the current kubeconfig -// points to a k3d cluster, or "" if not k3d. -func DetectK3dCluster(ctx context.Context) string { - w := wool.Get(ctx).In("DetectK3dCluster") - - // Check current kubectl context name — k3d contexts are named "k3d-" - cmd := exec.CommandContext(ctx, "kubectl", "config", "current-context") - var out bytes.Buffer - cmd.Stdout = &out - if err := cmd.Run(); err != nil { - return "" - } - ctxName := strings.TrimSpace(out.String()) - if !strings.HasPrefix(ctxName, "k3d-") { - return "" - } - cluster := strings.TrimPrefix(ctxName, "k3d-") - w.Debug("detected k3d cluster", wool.Field("cluster", cluster)) - return cluster -} - -// K3dImportImages imports Docker images into a k3d cluster so pods -// can use them without a registry. For images not present locally, -// attempts a docker pull first (handles stock images like neo4j:5). -func K3dImportImages(ctx context.Context, cluster string, images []string) error { - w := wool.Get(ctx).In("K3dImportImages") - - if len(images) == 0 { - return nil - } - - // Ensure all images exist locally — pull if needed. - var available []string +func EnsureImagesAvailable(ctx context.Context, images []string) error { + w := wool.Get(ctx).In("EnsureImagesAvailable") for _, img := range images { if imageExistsLocally(ctx, img) { - available = append(available, img) continue } - // Try pulling stock/external images. w.Info(fmt.Sprintf("pulling %s...", img)) pullCmd := exec.CommandContext(ctx, "docker", "pull", img) + var stderr bytes.Buffer + pullCmd.Stderr = &stderr if err := pullCmd.Run(); err != nil { - w.Warn(fmt.Sprintf("cannot pull %s, skipping", img)) - continue + return w.Wrapf(err, "cannot pull %s: %s", img, stderr.String()) } - available = append(available, img) } + return nil +} - if len(available) == 0 { +func K3dImportImages(ctx context.Context, cluster string, images []string) error { + w := wool.Get(ctx).In("K3dImportImages") + if len(images) == 0 { return nil } - args := append([]string{"image", "import", "-c", cluster}, available...) + args := append([]string{"image", "import", "-c", cluster}, images...) cmd := exec.CommandContext(ctx, "k3d", args...) var stderr bytes.Buffer cmd.Stderr = &stderr - w.Info(fmt.Sprintf("importing %d image(s) into k3d cluster %q", len(available), cluster)) + w.Info(fmt.Sprintf("importing %d image(s) into k3d cluster %q", len(images), cluster)) if err := cmd.Run(); err != nil { return w.Wrapf(err, "k3d image import failed: %s", stderr.String()) } - for _, img := range available { + for _, img := range images { w.Info(fmt.Sprintf("imported %s", img)) } return nil diff --git a/pkg/deployments/kubernetes.go b/pkg/deployments/kubernetes.go index ba1470fe..d8135737 100644 --- a/pkg/deployments/kubernetes.go +++ b/pkg/deployments/kubernetes.go @@ -3,15 +3,46 @@ package deployments import ( "bytes" "context" + "crypto/sha256" + "encoding/json" + "fmt" "os" "os/exec" "path" + "path/filepath" "strings" "github.com/codefly-dev/core/resources" "github.com/codefly-dev/core/wool" + "gopkg.in/yaml.v3" ) +type VerifiedKubernetesTarget struct { + Kind string + Kubeconfig string + Context string + Cluster string + APIServer string + K3dCluster string + // ClusterIdentity is the digest of the complete kubeconfig cluster entry, + // including certificate and routing settings. + ClusterIdentity string +} + +type kubeconfigView struct { + CurrentContext string `json:"current-context" yaml:"current-context"` + Clusters []struct { + Name string `json:"name" yaml:"name"` + Cluster map[string]any `json:"cluster" yaml:"cluster"` + } `json:"clusters" yaml:"clusters"` + Contexts []struct { + Name string `json:"name" yaml:"name"` + Context struct { + Cluster string `json:"cluster" yaml:"cluster"` + } `json:"context" yaml:"context"` + } `json:"contexts" yaml:"contexts"` +} + // GetK8sConfig resolves the kubeconfig path for an environment. // Lookup order: // 1. env.Cluster.Kubeconfig (declared in workspace.codefly.yaml). Tilde @@ -52,13 +83,221 @@ func GetK8sConfig(ctx context.Context, env *resources.Environment) (string, erro return path.Join(home, ".kube/config"), nil } -func kubectlApply(ctx context.Context, configPath, kubeContext, resource string) error { +func VerifyLocalK3dTarget(ctx context.Context, env *resources.Environment) (VerifiedKubernetesTarget, error) { + target, _, err := verifyLocalK3dTarget(ctx, env) + return target, err +} + +func verifyLocalK3dTarget(ctx context.Context, env *resources.Environment) (VerifiedKubernetesTarget, []byte, error) { + if env == nil || env.Cluster == nil || env.Cluster.Kind != "k3d" { + envName := "" + kind := "" + if env != nil { + envName = env.Name + if env.Cluster != nil { + kind = env.Cluster.Kind + } + } + return VerifiedKubernetesTarget{}, nil, fmt.Errorf( + "direct Kubernetes apply is allowed only for an exact local k3d target; environment %q declares cluster kind %q; use --render-only and publish the rendered manifests through GitOps", + envName, + kind, + ) + } + + configPath, err := GetK8sConfig(ctx, env) + if err != nil { + return VerifiedKubernetesTarget{}, nil, fmt.Errorf("resolve kubeconfig: %w", err) + } + if len(filepath.SplitList(configPath)) != 1 { + return VerifiedKubernetesTarget{}, nil, fmt.Errorf("direct Kubernetes apply requires exactly one kubeconfig, got %q", configPath) + } + configPath, err = filepath.Abs(configPath) + if err != nil { + return VerifiedKubernetesTarget{}, nil, fmt.Errorf("resolve kubeconfig %q: %w", configPath, err) + } + kubeContext := env.Cluster.Context + if kubeContext == "" { + return VerifiedKubernetesTarget{}, nil, fmt.Errorf( + "environment %q must declare cluster.context to permit direct apply; use --render-only and publish the rendered manifests through GitOps", + env.Name, + ) + } + + selected, snapshot, err := readSelectedKubeconfig(ctx, configPath) + if err != nil { + return VerifiedKubernetesTarget{}, nil, err + } + currentContext := selected.CurrentContext + if currentContext == "" { + return VerifiedKubernetesTarget{}, nil, fmt.Errorf("kubeconfig %q has no current context; refusing direct apply", configPath) + } + if kubeContext != currentContext { + return VerifiedKubernetesTarget{}, nil, fmt.Errorf( + "declared context %q does not match kubeconfig %q current context %q; refusing direct apply", + kubeContext, + configPath, + currentContext, + ) + } + if !strings.HasPrefix(kubeContext, "k3d-") || kubeContext == "k3d-" { + return VerifiedKubernetesTarget{}, nil, fmt.Errorf("context %q is not a k3d context; refusing direct apply", kubeContext) + } + + clusterName, apiServer, clusterIdentity, err := selected.target(kubeContext) + if err != nil { + return VerifiedKubernetesTarget{}, nil, fmt.Errorf("resolve context %q from kubeconfig %q: %w", kubeContext, configPath, err) + } + k3dCluster := strings.TrimPrefix(kubeContext, "k3d-") + owned, err := readK3dKubeconfig(ctx, k3dCluster) + if err != nil { + return VerifiedKubernetesTarget{}, nil, err + } + ownedCluster, ownedServer, ownedIdentity, err := owned.target(kubeContext) + if err != nil { + return VerifiedKubernetesTarget{}, nil, fmt.Errorf("resolve k3d-owned cluster %q: %w", k3dCluster, err) + } + if owned.CurrentContext != kubeContext || + ownedCluster != clusterName || + ownedServer != apiServer || + ownedIdentity != clusterIdentity { + return VerifiedKubernetesTarget{}, nil, fmt.Errorf( + "context %q resolves to cluster %q at %q, which does not match k3d-owned cluster %q at %q; refusing direct apply", + kubeContext, + clusterName, + apiServer, + ownedCluster, + ownedServer, + ) + } + + return VerifiedKubernetesTarget{ + Kind: env.Cluster.Kind, + Kubeconfig: filepath.Clean(configPath), + Context: kubeContext, + Cluster: clusterName, + APIServer: apiServer, + K3dCluster: k3dCluster, + ClusterIdentity: clusterIdentity, + }, snapshot, nil +} + +func VerifyLocalK3dTargetUnchanged(ctx context.Context, env *resources.Environment, planned *VerifiedKubernetesTarget) error { + _, err := verifiedKubeconfigSnapshot(ctx, env, planned) + return err +} + +func verifiedKubeconfigSnapshot(ctx context.Context, env *resources.Environment, planned *VerifiedKubernetesTarget) ([]byte, error) { + current, snapshot, err := verifyLocalK3dTarget(ctx, env) + if err != nil { + return nil, err + } + if current != *planned { + return nil, fmt.Errorf( + "Kubernetes target changed after validation (planned context %q, cluster %q, server %q; current context %q, cluster %q, server %q); refusing direct apply", + planned.Context, + planned.Cluster, + planned.APIServer, + current.Context, + current.Cluster, + current.APIServer, + ) + } + return snapshot, nil +} + +func readSelectedKubeconfig(ctx context.Context, configPath string) (kubeconfigView, []byte, error) { + cmd := exec.CommandContext(ctx, "kubectl", "--kubeconfig", configPath, "config", "view", "--raw", "--flatten", "--minify", "-o", "json") + var stdout, stderr bytes.Buffer + cmd.Stdout = &stdout + cmd.Stderr = &stderr + if err := cmd.Run(); err != nil { + return kubeconfigView{}, nil, fmt.Errorf("read kubeconfig %q: %w: %s", configPath, err, strings.TrimSpace(stderr.String())) + } + var config kubeconfigView + if err := json.Unmarshal(stdout.Bytes(), &config); err != nil { + return kubeconfigView{}, nil, fmt.Errorf("decode kubeconfig %q: %w", configPath, err) + } + return config, bytes.Clone(stdout.Bytes()), nil +} + +func readK3dKubeconfig(ctx context.Context, cluster string) (kubeconfigView, error) { + cmd := exec.CommandContext(ctx, "k3d", "kubeconfig", "get", cluster) + var stdout, stderr bytes.Buffer + cmd.Stdout = &stdout + cmd.Stderr = &stderr + if err := cmd.Run(); err != nil { + return kubeconfigView{}, fmt.Errorf("verify k3d-owned cluster %q: %w: %s", cluster, err, strings.TrimSpace(stderr.String())) + } + var config kubeconfigView + if err := yaml.Unmarshal(stdout.Bytes(), &config); err != nil { + return kubeconfigView{}, fmt.Errorf("decode k3d-owned cluster %q kubeconfig: %w", cluster, err) + } + return config, nil +} + +func (config kubeconfigView) target(contextName string) (string, string, string, error) { + clusterName := "" + contextMatches := 0 + for _, candidate := range config.Contexts { + if candidate.Name != contextName { + continue + } + contextMatches++ + if contextMatches > 1 { + return "", "", "", fmt.Errorf("context %q is declared more than once", contextName) + } + clusterName = candidate.Context.Cluster + } + if clusterName == "" { + return "", "", "", fmt.Errorf("context %q does not select a cluster", contextName) + } + + apiServer := "" + clusterIdentity := "" + clusterMatches := 0 + for _, candidate := range config.Clusters { + if candidate.Name != clusterName { + continue + } + clusterMatches++ + if clusterMatches > 1 { + return "", "", "", fmt.Errorf("cluster %q is declared more than once", clusterName) + } + var ok bool + apiServer, ok = candidate.Cluster["server"].(string) + if !ok || apiServer == "" { + return "", "", "", fmt.Errorf("cluster %q has no API server", clusterName) + } + encoded, err := json.Marshal(candidate.Cluster) + if err != nil { + return "", "", "", fmt.Errorf("encode cluster %q identity: %w", clusterName, err) + } + clusterIdentity = fmt.Sprintf("sha256:%x", sha256.Sum256(encoded)) + } + if apiServer == "" { + return "", "", "", fmt.Errorf("cluster %q has no API server", clusterName) + } + return clusterName, apiServer, clusterIdentity, nil +} + +func kubectlApply(ctx context.Context, target *VerifiedKubernetesTarget, kubeconfig []byte, resource string) error { w := wool.Get(ctx).In("kubectlApply") - // Prepare the kubectl command with stdin from the resource string - args := []string{"--kubeconfig", configPath} - if kubeContext != "" { - args = append(args, "--context", kubeContext) + snapshot, err := os.CreateTemp("", "codefly-verified-kubeconfig-*") + if err != nil { + return w.Wrapf(err, "cannot create verified kubeconfig snapshot") } + snapshotPath := snapshot.Name() + defer os.Remove(snapshotPath) + if _, err := snapshot.Write(kubeconfig); err != nil { + _ = snapshot.Close() + return w.Wrapf(err, "cannot write verified kubeconfig snapshot") + } + if err := snapshot.Close(); err != nil { + return w.Wrapf(err, "cannot close verified kubeconfig snapshot") + } + + args := []string{"--kubeconfig", snapshotPath, "--context", target.Context} args = append(args, "apply", "-f", "-") cmd := exec.CommandContext(ctx, "kubectl", args...) cmd.Stdin = bytes.NewBufferString(resource) @@ -70,7 +309,7 @@ func kubectlApply(ctx context.Context, configPath, kubeContext, resource string) cmd.Stderr = &stderr // Execute the command - err := cmd.Run() + err = cmd.Run() if err != nil { return w.Wrapf(err, "cannot run kubectl apply: %s", stderr.String()) } @@ -82,21 +321,14 @@ func kubectlApply(ctx context.Context, configPath, kubeContext, resource string) return nil } -func KubernetesApply(ctx context.Context, service *resources.Service, env *resources.Environment, sources ...string) error { - w := wool.Get(ctx).In("KubernetesApply", wool.ThisField(resources.WithUnique(service))) - // Create the Kubernetes client - configPath, err := GetK8sConfig(ctx, env) - if err != nil { - return w.Wrapf(err, "cannot get k8s client") - } - var kubeContext string - if env != nil && env.Cluster != nil { - kubeContext = env.Cluster.Context - } - +func KubernetesApply(ctx context.Context, env *resources.Environment, target *VerifiedKubernetesTarget, sources ...string) error { + w := wool.Get(ctx).In("KubernetesApply") for _, r := range sources { - err := kubectlApply(ctx, configPath, kubeContext, r) + snapshot, err := verifiedKubeconfigSnapshot(ctx, env, target) if err != nil { + return w.Wrapf(err, "cannot verify Kubernetes target before apply") + } + if err := kubectlApply(ctx, target, snapshot, r); err != nil { return w.Wrapf(err, "cannot apply resource") } } diff --git a/pkg/deployments/kubernetes_test.go b/pkg/deployments/kubernetes_test.go new file mode 100644 index 00000000..eefd6972 --- /dev/null +++ b/pkg/deployments/kubernetes_test.go @@ -0,0 +1,552 @@ +package deployments + +import ( + "context" + "encoding/json" + "fmt" + "os" + "path/filepath" + "sync" + "testing" + + builderv0 "github.com/codefly-dev/core/generated/go/codefly/services/builder/v0" + "github.com/codefly-dev/core/resources" + "github.com/stretchr/testify/require" +) + +func TestVerifyLocalK3dTargetRejectsRemoteKindsBeforeInspectingKubeconfig(t *testing.T) { + for _, kind := range []string{"eks", "gke", "aks", "external"} { + t.Run(kind, func(t *testing.T) { + env := &resources.Environment{ + Name: "production", + Cluster: &resources.EnvironmentCluster{ + Kind: kind, + Kubeconfig: filepath.Join(t.TempDir(), "config"), + Context: "k3d-production", + }, + } + + _, err := VerifyLocalK3dTarget(context.Background(), env) + require.Error(t, err) + require.Contains(t, err.Error(), "exact local k3d target") + require.Contains(t, err.Error(), "--render-only") + require.Contains(t, err.Error(), "GitOps") + }) + } +} + +func TestVerifyLocalK3dTargetRejectsStaleCurrentContext(t *testing.T) { + harness := newKubernetesCommandHarness(t) + harness.writeSelected(kubeconfigDocument("eks-production", "eks-production", "production", "https://eks.example.com")) + harness.writeOwned(kubeconfigDocument("k3d-dev", "k3d-dev", "k3d-dev", "https://127.0.0.1:6443")) + + env := harness.environment("k3d-dev") + _, err := VerifyLocalK3dTarget(context.Background(), env) + + require.Error(t, err) + require.Contains(t, err.Error(), "does not match") + require.Contains(t, err.Error(), "eks-production") +} + +func TestVerifyLocalK3dTargetRejectsMismatchedKubeconfigServer(t *testing.T) { + harness := newKubernetesCommandHarness(t) + harness.writeSelected(kubeconfigDocument("k3d-dev", "k3d-dev", "k3d-dev", "https://127.0.0.1:6443")) + harness.writeOwned(kubeconfigDocument("k3d-dev", "k3d-dev", "k3d-dev", "https://127.0.0.1:6553")) + + _, err := VerifyLocalK3dTarget(context.Background(), harness.environment("k3d-dev")) + + require.Error(t, err) + require.Contains(t, err.Error(), "does not match k3d-owned cluster") +} + +func TestVerifyLocalK3dTargetRejectsMismatchedClusterConnectionIdentity(t *testing.T) { + for _, test := range []struct { + name string + selected map[string]any + owned map[string]any + }{ + { + name: "certificate authority", + selected: map[string]any{ + "server": "https://127.0.0.1:6443", + "certificate-authority-data": "selected-ca", + }, + owned: map[string]any{ + "server": "https://127.0.0.1:6443", + "certificate-authority-data": "owned-ca", + }, + }, + { + name: "proxy", + selected: map[string]any{ + "server": "https://127.0.0.1:6443", + "proxy-url": "https://remote.example.com", + }, + owned: map[string]any{ + "server": "https://127.0.0.1:6443", + }, + }, + { + name: "insecure tls", + selected: map[string]any{ + "server": "https://127.0.0.1:6443", + "insecure-skip-tls-verify": true, + }, + owned: map[string]any{ + "server": "https://127.0.0.1:6443", + }, + }, + { + name: "tls server name", + selected: map[string]any{ + "server": "https://127.0.0.1:6443", + "tls-server-name": "remote.example.com", + }, + owned: map[string]any{ + "server": "https://127.0.0.1:6443", + }, + }, + } { + t.Run(test.name, func(t *testing.T) { + harness := newKubernetesCommandHarness(t) + harness.writeSelected(kubeconfigDocumentWithCluster("k3d-dev", "k3d-dev", "k3d-dev", test.selected)) + harness.writeOwned(kubeconfigDocumentWithCluster("k3d-dev", "k3d-dev", "k3d-dev", test.owned)) + + _, err := VerifyLocalK3dTarget(context.Background(), harness.environment("k3d-dev")) + + require.Error(t, err) + require.Contains(t, err.Error(), "does not match k3d-owned cluster") + }) + } +} + +func TestVerifyLocalK3dTargetRequiresDeclaredContext(t *testing.T) { + harness := newKubernetesCommandHarness(t) + config := kubeconfigDocument("k3d-other", "k3d-other", "k3d-other", "https://127.0.0.1:7443") + harness.writeSelected(config) + harness.writeOwned(config) + + _, err := VerifyLocalK3dTarget(context.Background(), harness.environment("")) + + require.Error(t, err) + require.Contains(t, err.Error(), "must declare cluster.context") +} + +func TestVerifyLocalK3dTargetRejectsRenamedNonK3dContext(t *testing.T) { + harness := newKubernetesCommandHarness(t) + harness.writeSelected(kubeconfigDocument("k3d-production", "k3d-production", "production", "https://eks.example.com")) + harness.writeOwned(kubeconfigDocument("k3d-production", "k3d-production", "k3d-production", "https://127.0.0.1:6443")) + + _, err := VerifyLocalK3dTarget(context.Background(), harness.environment("k3d-production")) + + require.Error(t, err) + require.Contains(t, err.Error(), "does not match k3d-owned cluster") +} + +func TestVerifyLocalK3dTargetBindsExactIdentity(t *testing.T) { + harness := newKubernetesCommandHarness(t) + config := kubeconfigDocument("k3d-dev", "k3d-dev", "k3d-dev", "https://127.0.0.1:6443") + harness.writeSelected(config) + harness.writeOwned(config) + + target, err := VerifyLocalK3dTarget(context.Background(), harness.environment("k3d-dev")) + + require.NoError(t, err) + require.Regexp(t, `^sha256:[0-9a-f]{64}$`, target.ClusterIdentity) + target.ClusterIdentity = "" + require.Equal(t, VerifiedKubernetesTarget{ + Kind: "k3d", + Kubeconfig: harness.kubeconfig, + Context: "k3d-dev", + Cluster: "k3d-dev", + APIServer: "https://127.0.0.1:6443", + K3dCluster: "dev", + }, target) +} + +func TestKubernetesApplyRejectsContextSwapAfterValidation(t *testing.T) { + harness := newKubernetesCommandHarness(t) + dev := kubeconfigDocument("k3d-dev", "k3d-dev", "k3d-dev", "https://127.0.0.1:6443") + harness.writeSelected(dev) + harness.writeOwned(dev) + env := harness.environment("k3d-dev") + + target, err := VerifyLocalK3dTarget(context.Background(), env) + require.NoError(t, err) + + other := kubeconfigDocument("k3d-other", "k3d-other", "k3d-other", "https://127.0.0.1:7443") + harness.writeSelected(other) + harness.writeOwned(other) + + err = KubernetesApply(context.Background(), env, &target, "apiVersion: v1\nkind: Namespace\nmetadata:\n name: blocked\n") + require.Error(t, err) + require.Contains(t, err.Error(), "does not match") + require.Contains(t, err.Error(), "k3d-other") + require.NoFileExists(t, harness.applyLog) +} + +func TestLocalApplyManagerFailsClosedBeforeApplyWhenImagePreparationFails(t *testing.T) { + for _, test := range []struct { + name string + pullFails bool + importFails bool + }{ + {name: "pull failure", pullFails: true}, + {name: "import failure", importFails: true}, + } { + t.Run(test.name, func(t *testing.T) { + harness := newKubernetesCommandHarness(t) + config := kubeconfigDocument("k3d-dev", "k3d-dev", "k3d-dev", "https://127.0.0.1:6443") + harness.writeSelected(config) + harness.writeOwned(config) + if test.pullFails { + t.Setenv("FAKE_DOCKER_INSPECT_FAIL", "1") + t.Setenv("FAKE_DOCKER_PULL_FAIL", "1") + } + if test.importFails { + t.Setenv("FAKE_K3D_IMPORT_FAIL", "1") + } + + workspace, module, service := deploymentFixture(t) + manager, err := NewLocalApplyManager(context.Background(), workspace, harness.environment("k3d-dev")) + require.NoError(t, err) + + err = manager.Handle(context.Background(), service, module, kubernetesDeploymentOutput()) + require.Error(t, err) + require.NoFileExists(t, harness.applyLog) + require.NoFileExists(t, harness.kustomizeLog) + }) + } +} + +func TestLocalApplyManagerRecordsTargetAndRenderedTreeDigest(t *testing.T) { + harness := newKubernetesCommandHarness(t) + config := kubeconfigDocument("k3d-dev", "k3d-dev", "k3d-dev", "https://127.0.0.1:6443") + harness.writeSelected(config) + harness.writeOwned(config) + workspace, module, service := deploymentFixture(t) + manager, err := NewLocalApplyManager(context.Background(), workspace, harness.environment("k3d-dev")) + require.NoError(t, err) + + require.NoError(t, manager.Handle(context.Background(), service, module, kubernetesDeploymentOutput())) + + evidence := manager.Evidence() + require.NotNil(t, evidence.Target) + require.Equal(t, "https://127.0.0.1:6443", evidence.Target.APIServer) + require.Len(t, evidence.RenderedTrees, 1) + require.Equal(t, "backend", evidence.RenderedTrees[0].Module) + require.Equal(t, "api", evidence.RenderedTrees[0].Service) + require.Regexp(t, `^sha256:[0-9a-f]{64}$`, evidence.RenderedTrees[0].Digest) + require.Contains(t, evidence.RenderedTrees[0].Manifests, "kind: Namespace") + require.FileExists(t, harness.applyLog) +} + +func TestRenderManagerReturnsRenderedEvidenceWithoutApplying(t *testing.T) { + harness := newKubernetesCommandHarness(t) + workspace, module, service := deploymentFixture(t) + manager := NewRenderManager(workspace, harness.environment("k3d-dev")) + + require.NoError(t, manager.Handle(context.Background(), service, module, kubernetesDeploymentOutput())) + + evidence := manager.Evidence() + require.Nil(t, evidence.Target) + require.Equal(t, []RenderedTreeEvidence{{ + Module: "backend", + Service: "api", + Digest: evidence.RenderedTrees[0].Digest, + Manifests: "apiVersion: v1\nkind: Namespace\nmetadata:\n name: applied\n", + }}, evidence.RenderedTrees) + require.Regexp(t, `^sha256:[0-9a-f]{64}$`, evidence.RenderedTrees[0].Digest) + require.NoFileExists(t, harness.applyLog) +} + +func TestRenderManagersKeepConcurrentRequestEvidenceIsolated(t *testing.T) { + harness := newKubernetesCommandHarness(t) + firstWorkspace, firstModule, firstService := deploymentFixture(t) + secondWorkspace, secondModule, secondService := deploymentFixture(t) + first := NewRenderManager(firstWorkspace, harness.environment("k3d-dev")) + second := NewRenderManager(secondWorkspace, harness.environment("k3d-dev")) + + var wait sync.WaitGroup + errorsOut := make(chan error, 2) + for _, deploy := range []struct { + manager *RenderManager + module *resources.Module + service *resources.Service + }{ + {manager: first, module: firstModule, service: firstService}, + {manager: second, module: secondModule, service: secondService}, + } { + wait.Add(1) + go func() { + defer wait.Done() + errorsOut <- deploy.manager.Handle(context.Background(), deploy.service, deploy.module, kubernetesDeploymentOutput()) + }() + } + wait.Wait() + close(errorsOut) + for err := range errorsOut { + require.NoError(t, err) + } + + require.Len(t, first.Evidence().RenderedTrees, 1) + require.Len(t, second.Evidence().RenderedTrees, 1) + require.NoFileExists(t, harness.applyLog) +} + +func TestLocalApplyManagerRecordsModuleTreeEvidence(t *testing.T) { + harness := newKubernetesCommandHarness(t) + config := kubeconfigDocument("k3d-dev", "k3d-dev", "k3d-dev", "https://127.0.0.1:6443") + harness.writeSelected(config) + harness.writeOwned(config) + workspace, module, _ := deploymentFixture(t) + dir := filepath.Join(module.Dir(), "deployment", "kustomize", "overlays", "local") + writeTestFile(t, filepath.Join(dir, "kustomization.yaml"), "resources: []\n") + writeTestFile(t, filepath.Join(module.Dir(), "deployment", "kustomize", "base", "namespace.yaml"), "kind: Namespace\n") + manager, err := NewLocalApplyManager(context.Background(), workspace, harness.environment("k3d-dev")) + require.NoError(t, err) + + require.NoError(t, manager.ApplyModuleKustomize(context.Background(), module, dir)) + + evidence := manager.Evidence() + require.Len(t, evidence.RenderedTrees, 1) + require.Equal(t, "backend", evidence.RenderedTrees[0].Module) + require.Empty(t, evidence.RenderedTrees[0].Service) + require.Regexp(t, `^sha256:[0-9a-f]{64}$`, evidence.RenderedTrees[0].Digest) + wantDigest, err := RenderedTreeDigest(filepath.Join(module.Dir(), "deployment", "kustomize")) + require.NoError(t, err) + require.Equal(t, wantDigest, evidence.RenderedTrees[0].Digest) + require.Contains(t, evidence.RenderedTrees[0].Manifests, "kind: Namespace") + require.FileExists(t, harness.applyLog) +} + +func TestRenderedTreeDigestBindsPathsAndContents(t *testing.T) { + first := t.TempDir() + second := t.TempDir() + writeTestFile(t, filepath.Join(first, "base", "deployment.yaml"), "kind: Deployment\n") + writeTestFile(t, filepath.Join(first, "overlays", "local", "kustomization.yaml"), "resources: []\n") + writeTestFile(t, filepath.Join(second, "overlays", "local", "kustomization.yaml"), "resources: []\n") + writeTestFile(t, filepath.Join(second, "base", "deployment.yaml"), "kind: Deployment\n") + + firstDigest, err := RenderedTreeDigest(first) + require.NoError(t, err) + secondDigest, err := RenderedTreeDigest(second) + require.NoError(t, err) + require.Equal(t, firstDigest, secondDigest) + + writeTestFile(t, filepath.Join(second, "base", "deployment.yaml"), "kind: StatefulSet\n") + changedDigest, err := RenderedTreeDigest(second) + require.NoError(t, err) + require.NotEqual(t, firstDigest, changedDigest) +} + +func TestExtractKustomizeImagesUsesOriginalNameWhenNewNameIsOmitted(t *testing.T) { + dir := t.TempDir() + writeTestFile(t, filepath.Join(dir, "kustomization.yaml"), `images: + - name: app + newTag: test + - name: registry.example.com:5000/team/worker:old + newTag: current + - name: original + newName: replacement + newTag: latest +`) + + images, err := extractKustomizeImages(dir) + + require.NoError(t, err) + require.Equal(t, []string{ + "app:test", + "registry.example.com:5000/team/worker:current", + "replacement:latest", + }, images) +} + +type kubernetesCommandHarness struct { + t *testing.T + selected string + owned string + kubeconfig string + applyLog string + kustomizeLog string +} + +func newKubernetesCommandHarness(t *testing.T) kubernetesCommandHarness { + t.Helper() + root := t.TempDir() + bin := filepath.Join(root, "bin") + require.NoError(t, os.MkdirAll(bin, 0o755)) + harness := kubernetesCommandHarness{ + t: t, + selected: filepath.Join(root, "selected.json"), + owned: filepath.Join(root, "owned.yaml"), + kubeconfig: filepath.Join(root, "declared-kubeconfig"), + applyLog: filepath.Join(root, "apply.log"), + kustomizeLog: filepath.Join(root, "kustomize.log"), + } + writeTestFile(t, harness.kubeconfig, "fixture") + writeExecutable(t, filepath.Join(bin, "kubectl"), `#!/bin/sh +case " $* " in + *" config view "*) + cat "$FAKE_SELECTED_KUBECONFIG" + ;; + *" apply "*) + if [ "$2" = "$FAKE_DECLARED_KUBECONFIG" ] || [ ! -f "$2" ]; then + exit 93 + fi + printf '%s\n' "$*" >> "$FAKE_APPLY_LOG" + cat >/dev/null + printf 'resource configured\n' + ;; + *) + exit 90 + ;; +esac +`) + writeExecutable(t, filepath.Join(bin, "k3d"), `#!/bin/sh +if [ "$1" = "kubeconfig" ] && [ "$2" = "get" ]; then + cat "$FAKE_K3D_KUBECONFIG" + exit 0 +fi +if [ "$1" = "image" ] && [ "$2" = "import" ]; then + if [ "$FAKE_K3D_IMPORT_FAIL" = "1" ]; then + exit 41 + fi + exit 0 +fi +exit 91 +`) + writeExecutable(t, filepath.Join(bin, "docker"), `#!/bin/sh +if [ "$1" = "image" ] && [ "$2" = "inspect" ]; then + if [ "$FAKE_DOCKER_INSPECT_FAIL" = "1" ]; then + exit 42 + fi + exit 0 +fi +if [ "$1" = "pull" ]; then + if [ "$FAKE_DOCKER_PULL_FAIL" = "1" ]; then + exit 43 + fi + exit 0 +fi +exit 92 +`) + writeExecutable(t, filepath.Join(bin, "kustomize"), `#!/bin/sh +printf '%s\n' "$*" >> "$FAKE_KUSTOMIZE_LOG" +printf 'apiVersion: v1\nkind: Namespace\nmetadata:\n name: applied\n' +`) + t.Setenv("PATH", bin+string(os.PathListSeparator)+os.Getenv("PATH")) + t.Setenv("FAKE_SELECTED_KUBECONFIG", harness.selected) + t.Setenv("FAKE_K3D_KUBECONFIG", harness.owned) + t.Setenv("FAKE_DECLARED_KUBECONFIG", harness.kubeconfig) + t.Setenv("FAKE_APPLY_LOG", harness.applyLog) + t.Setenv("FAKE_KUSTOMIZE_LOG", harness.kustomizeLog) + return harness +} + +func (h kubernetesCommandHarness) environment(contextName string) *resources.Environment { + return &resources.Environment{ + Name: "local", + Cluster: &resources.EnvironmentCluster{ + Kind: "k3d", + Kubeconfig: h.kubeconfig, + Context: contextName, + }, + } +} + +func (h kubernetesCommandHarness) writeSelected(content string) { + h.t.Helper() + require.NoError(h.t, os.WriteFile(h.selected, []byte(content), 0o600)) +} + +func (h kubernetesCommandHarness) writeOwned(content string) { + h.t.Helper() + require.NoError(h.t, os.WriteFile(h.owned, []byte(content), 0o600)) +} + +func kubeconfigDocument(currentContext, contextName, clusterName, apiServer string) string { + return kubeconfigDocumentWithCluster(currentContext, contextName, clusterName, map[string]any{ + "server": apiServer, + }) +} + +func kubeconfigDocumentWithCluster(currentContext, contextName, clusterName string, cluster map[string]any) string { + document := map[string]any{ + "apiVersion": "v1", + "current-context": currentContext, + "contexts": []any{map[string]any{ + "name": contextName, + "context": map[string]any{ + "cluster": clusterName, + }, + }}, + "clusters": []any{map[string]any{ + "name": clusterName, + "cluster": cluster, + }}, + } + encoded, err := json.Marshal(document) + if err != nil { + panic(err) + } + return string(encoded) +} + +func deploymentFixture(t *testing.T) (*resources.Workspace, *resources.Module, *resources.Service) { + t.Helper() + root := t.TempDir() + writeTestFile(t, filepath.Join(root, "workspace.codefly.yaml"), `name: target-test +layout: modules +modules: + - name: backend +`) + writeTestFile(t, filepath.Join(root, "modules", "backend", "module.codefly.yaml"), `kind: module +name: backend +services: + - name: api +`) + writeTestFile(t, filepath.Join(root, "modules", "backend", "services", "api", "service.codefly.yaml"), `kind: service +name: api +version: 0.0.0 +module: backend +agent: + kind: runtime::service + name: test + version: 0.0.0 + publisher: codefly.dev +`) + writeTestFile(t, filepath.Join(root, "deployments", "modules", "backend", "services", "api", "overlays", "local", "kustomization.yaml"), `images: + - name: app + newName: app + newTag: test +`) + workspace, err := resources.LoadWorkspaceFromDir(context.Background(), root) + require.NoError(t, err) + module, err := workspace.LoadModuleFromName(context.Background(), "backend") + require.NoError(t, err) + service, err := module.LoadServiceFromName(context.Background(), "api") + require.NoError(t, err) + return workspace, module, service +} + +func kubernetesDeploymentOutput() *builderv0.DeploymentOutput { + return &builderv0.DeploymentOutput{ + Kind: &builderv0.DeploymentOutput_Kubernetes{ + Kubernetes: &builderv0.KubernetesDeploymentOutput{ + Kind: builderv0.KubernetesDeploymentOutput_KUSTOMIZE, + }, + }, + } +} + +func writeExecutable(t *testing.T, filePath, content string) { + t.Helper() + require.NoError(t, os.WriteFile(filePath, []byte(content), 0o700)) +} + +func writeTestFile(t *testing.T, filePath, content string) { + t.Helper() + require.NoError(t, os.MkdirAll(filepath.Dir(filePath), 0o755)) + require.NoError(t, os.WriteFile(filePath, []byte(content), 0o600), fmt.Sprintf("write %s", filePath)) +} diff --git a/pkg/deployments/kustomize.go b/pkg/deployments/kustomize.go index cf0d8c00..c10be0c6 100644 --- a/pkg/deployments/kustomize.go +++ b/pkg/deployments/kustomize.go @@ -3,13 +3,17 @@ package deployments import ( "bytes" "context" + "crypto/sha256" + "encoding/binary" "fmt" + "io/fs" + "os" "os/exec" "path" + "path/filepath" "strings" "github.com/codefly-dev/core/resources" - "github.com/codefly-dev/core/wool" ) func KustomizeDir(ctx context.Context, workspace *resources.Workspace, module *resources.Module, service *resources.Service) string { @@ -20,19 +24,76 @@ func KustomizeDirForEnv(ctx context.Context, workspace *resources.Workspace, mod return path.Join(KustomizeDir(ctx, workspace, module, service), "overlays", env.Name) } -func KustomizeApply(ctx context.Context, service *resources.Service, env *resources.Environment, dir string) error { - w := wool.Get(ctx).In("Builder", wool.ThisField(resources.WithUnique(service))) - w.Debug("applying kustomize", wool.DirField(dir)) - cmd := exec.Command("kustomize", "build", dir) +func renderKustomize(ctx context.Context, tree, treeDigest, dir string) (string, []string, error) { + cmd := exec.CommandContext(ctx, "kustomize", "build", dir) var stdout, stderr bytes.Buffer cmd.Stdout = &stdout cmd.Stderr = &stderr err := cmd.Run() if err != nil { - return w.Wrapf(err, "cannot run kustomize build: %s", stderr.String()) + return "", nil, fmt.Errorf("cannot run kustomize build: %w: %s", err, stderr.String()) } - // Split the output into individual objs - objs := strings.Split(stdout.String(), "---") - w.Info(fmt.Sprintf("Found %d resources to apply", len(objs))) - return KubernetesApply(ctx, service, env, objs...) + currentDigest, err := RenderedTreeDigest(tree) + if err != nil { + return "", nil, fmt.Errorf("cannot verify rendered deployment tree: %w", err) + } + if currentDigest != treeDigest { + return "", nil, fmt.Errorf("rendered deployment tree changed after validation; refusing direct apply") + } + manifests := stdout.String() + return manifests, strings.Split(manifests, "---"), nil +} + +func RenderedTreeDigest(root string) (string, error) { + digest := sha256.New() + err := filepath.WalkDir(root, func(filePath string, entry fs.DirEntry, walkErr error) error { + if walkErr != nil { + return walkErr + } + if entry.IsDir() { + return nil + } + relative, err := filepath.Rel(root, filePath) + if err != nil { + return err + } + relative = filepath.ToSlash(relative) + var kind byte + var content []byte + switch { + case entry.Type().IsRegular(): + kind = 1 + content, err = os.ReadFile(filePath) + if err != nil { + return err + } + case entry.Type()&os.ModeSymlink != 0: + target, err := os.Readlink(filePath) + if err != nil { + return err + } + kind = 2 + content = []byte(target) + default: + return fmt.Errorf("unsupported file type in rendered deployment tree: %s", filePath) + } + if _, err := digest.Write([]byte{kind}); err != nil { + return err + } + if err := binary.Write(digest, binary.BigEndian, uint64(len(relative))); err != nil { + return err + } + if _, err := digest.Write([]byte(relative)); err != nil { + return err + } + if err := binary.Write(digest, binary.BigEndian, uint64(len(content))); err != nil { + return err + } + _, err = digest.Write(content) + return err + }) + if err != nil { + return "", err + } + return fmt.Sprintf("sha256:%x", digest.Sum(nil)), nil } diff --git a/pkg/deployments/manager.go b/pkg/deployments/manager.go index e1654347..073c0e06 100644 --- a/pkg/deployments/manager.go +++ b/pkg/deployments/manager.go @@ -4,6 +4,9 @@ import ( "context" "os" "path" + "sort" + "strings" + "sync" builderv0 "github.com/codefly-dev/core/generated/go/codefly/services/builder/v0" "github.com/codefly-dev/core/resources" @@ -15,6 +18,58 @@ type Manager interface { Handle(ctx context.Context, service *resources.Service, module *resources.Module, deploy *builderv0.DeploymentOutput) error } +type RenderedTreeEvidence struct { + Module string + Service string + Digest string + Manifests string +} + +type DeploymentEvidence struct { + Target *VerifiedKubernetesTarget + RenderedTrees []RenderedTreeEvidence +} + +type EvidenceProvider interface { + Evidence() DeploymentEvidence +} + +type evidenceKey struct { + module string + service string +} + +type evidenceRecorder struct { + mu sync.Mutex + trees map[evidenceKey]RenderedTreeEvidence +} + +func newEvidenceRecorder() evidenceRecorder { + return evidenceRecorder{trees: map[evidenceKey]RenderedTreeEvidence{}} +} + +func (r *evidenceRecorder) record(tree RenderedTreeEvidence) { + r.mu.Lock() + defer r.mu.Unlock() + r.trees[evidenceKey{module: tree.Module, service: tree.Service}] = tree +} + +func (r *evidenceRecorder) renderedTrees() []RenderedTreeEvidence { + r.mu.Lock() + defer r.mu.Unlock() + trees := make([]RenderedTreeEvidence, 0, len(r.trees)) + for _, tree := range r.trees { + trees = append(trees, tree) + } + sort.Slice(trees, func(i, j int) bool { + if trees[i].Module == trees[j].Module { + return trees[i].Service < trees[j].Service + } + return trees[i].Module < trees[j].Module + }) + return trees +} + func GetKubernetesDeployment(ctx context.Context, dockerBuildContext *builderv0.DockerBuildContext, workspace *resources.Workspace, module *resources.Module, service *resources.Service, env *resources.Environment, namespace string) (*builderv0.Deployment, error) { return &builderv0.Deployment{ Kind: &builderv0.Deployment_Kubernetes{ @@ -27,16 +82,24 @@ func GetKubernetesDeployment(ctx context.Context, dockerBuildContext *builderv0. }, nil } -func NewLocalApplyManager(ctx context.Context, workspace *resources.Workspace, env *resources.Environment) *LocalApplyManager { +func NewLocalApplyManager(ctx context.Context, workspace *resources.Workspace, env *resources.Environment) (*LocalApplyManager, error) { + target, err := VerifyLocalK3dTarget(ctx, env) + if err != nil { + return nil, err + } return &LocalApplyManager{ Workspace: workspace, Env: env, - } + target: target, + evidence: newEvidenceRecorder(), + }, nil } type LocalApplyManager struct { Workspace *resources.Workspace Env *resources.Environment + target VerifiedKubernetesTarget + evidence evidenceRecorder } func (l *LocalApplyManager) Handle(ctx context.Context, service *resources.Service, module *resources.Module, deploy *builderv0.DeploymentOutput) error { @@ -44,13 +107,20 @@ func (l *LocalApplyManager) Handle(ctx context.Context, service *resources.Servi switch v := deploy.Kind.(type) { case *builderv0.DeploymentOutput_Kubernetes: if v.Kubernetes.Kind == builderv0.KubernetesDeploymentOutput_KUSTOMIZE { - // Import images into k3d if applicable. - if err := l.importImagesIfK3d(ctx, module, service); err != nil { - w.Warn("k3d image import failed (continuing)", wool.Field("error", err.Error())) + if err := VerifyLocalK3dTargetUnchanged(ctx, l.Env, &l.target); err != nil { + return w.Wrapf(err, "cannot verify Kubernetes target before image import") } - err := l.KustomizeApply(ctx, module, service) + tree := KustomizeDir(ctx, l.Workspace, module, service) + digest, err := RenderedTreeDigest(tree) if err != nil { + return w.Wrapf(err, "cannot digest rendered deployment tree") + } + + if err := l.importImages(ctx, module, service); err != nil { + return w.Wrapf(err, "cannot import images into verified k3d cluster") + } + if err := l.applyTree(ctx, module.Name, service.Name, tree, digest, KustomizeDirForEnv(ctx, l.Workspace, module, service, l.Env)); err != nil { return w.Wrapf(err, "cannot apply kustomize") } } @@ -61,50 +131,109 @@ func (l *LocalApplyManager) Handle(ctx context.Context, service *resources.Servi } var _ Manager = &LocalApplyManager{} +var _ EvidenceProvider = &LocalApplyManager{} -func (l *LocalApplyManager) KustomizeApply(ctx context.Context, module *resources.Module, service *resources.Service) error { - w := wool.Get(ctx).In("Builder", wool.ThisField(resources.WithUnique(service))) - dir := KustomizeDirForEnv(ctx, l.Workspace, module, service, l.Env) - - err := KustomizeApply(ctx, service, l.Env, dir) +func (l *LocalApplyManager) applyTree(ctx context.Context, module, service, tree, digest, dir string) error { + manifests, resourcesToApply, err := renderKustomize(ctx, tree, digest, dir) if err != nil { - return w.Wrapf(err, "cannot apply kustomize") + return err + } + if err := KubernetesApply(ctx, l.Env, &l.target, resourcesToApply...); err != nil { + return err } + l.evidence.record(RenderedTreeEvidence{ + Module: module, + Service: service, + Digest: digest, + Manifests: manifests, + }) return nil } -// importImagesIfK3d imports freshly-built Docker images into the k3d -// cluster so pods can use them without going through a registry. -// -// Cluster decision order: -// 1. If env declares a non-k3d Cluster.Kind (eks, gke, …), skip — those -// clusters pull from a registry and image-import is a no-op that just -// wastes a `kubectl config current-context` exec. -// 2. If env declares Cluster.Kind == "k3d" (or no Cluster block at all, -// which we treat as legacy = local k3d), run the runtime detection -// to discover the actual cluster name. -// 3. If detection returns "" (no k3d cluster at the current context), -// skip silently — user might be on minikube / kind / a real cluster. -func (l *LocalApplyManager) importImagesIfK3d(ctx context.Context, module *resources.Module, service *resources.Service) error { - if l.Env != nil && !l.Env.IsK3d() { - return nil +func (l *LocalApplyManager) ApplyModuleKustomize(ctx context.Context, module *resources.Module, dir string) error { + if err := VerifyLocalK3dTargetUnchanged(ctx, l.Env, &l.target); err != nil { + return err } - - cluster := DetectK3dCluster(ctx) - if cluster == "" { - return nil + tree := path.Join(module.Dir(), "deployment", "kustomize") + digest, err := RenderedTreeDigest(tree) + if err != nil { + return err } + return l.applyTree(ctx, module.Name, "", tree, digest, dir) +} - // Read the kustomize overlay to find image references. +func (l *LocalApplyManager) importImages(ctx context.Context, module *resources.Module, service *resources.Service) error { dir := KustomizeDirForEnv(ctx, l.Workspace, module, service, l.Env) images, err := extractKustomizeImages(dir) if err != nil { return err } + if err := EnsureImagesAvailable(ctx, images); err != nil { + return err + } + if err := VerifyLocalK3dTargetUnchanged(ctx, l.Env, &l.target); err != nil { + return err + } + return K3dImportImages(ctx, l.target.K3dCluster, images) +} + +func (l *LocalApplyManager) Evidence() DeploymentEvidence { + target := l.target + return DeploymentEvidence{ + Target: &target, + RenderedTrees: l.evidence.renderedTrees(), + } +} + +type RenderManager struct { + Workspace *resources.Workspace + Env *resources.Environment + evidence evidenceRecorder +} + +func NewRenderManager(workspace *resources.Workspace, env *resources.Environment) *RenderManager { + return &RenderManager{ + Workspace: workspace, + Env: env, + evidence: newEvidenceRecorder(), + } +} - return K3dImportImages(ctx, cluster, images) +func (r *RenderManager) Handle(ctx context.Context, service *resources.Service, module *resources.Module, deploy *builderv0.DeploymentOutput) error { + w := wool.Get(ctx).In("Builder") + switch v := deploy.Kind.(type) { + case *builderv0.DeploymentOutput_Kubernetes: + if v.Kubernetes.Kind != builderv0.KubernetesDeploymentOutput_KUSTOMIZE { + return nil + } + tree := KustomizeDir(ctx, r.Workspace, module, service) + digest, err := RenderedTreeDigest(tree) + if err != nil { + return w.Wrapf(err, "cannot digest rendered deployment tree") + } + manifests, _, err := renderKustomize(ctx, tree, digest, KustomizeDirForEnv(ctx, r.Workspace, module, service, r.Env)) + if err != nil { + return w.Wrapf(err, "cannot render kustomize") + } + r.evidence.record(RenderedTreeEvidence{ + Module: module.Name, + Service: service.Name, + Digest: digest, + Manifests: manifests, + }) + default: + return w.NewError("unsupported deployment kind %T", deploy.Kind) + } + return nil +} + +func (r *RenderManager) Evidence() DeploymentEvidence { + return DeploymentEvidence{RenderedTrees: r.evidence.renderedTrees()} } +var _ Manager = &RenderManager{} +var _ EvidenceProvider = &RenderManager{} + // kustomization is a minimal representation of kustomization.yaml for image extraction. type kustomization struct { Images []kustomizeImage `yaml:"images"` @@ -131,7 +260,16 @@ func extractKustomizeImages(dir string) ([]string, error) { var images []string for _, img := range k.Images { ref := img.NewName + if ref == "" { + ref = img.Name + } if img.NewTag != "" { + if digest := strings.IndexByte(ref, '@'); digest >= 0 { + ref = ref[:digest] + } + if tag := strings.LastIndexByte(ref, ':'); tag > strings.LastIndexByte(ref, '/') { + ref = ref[:tag] + } ref += ":" + img.NewTag } if ref != "" { diff --git a/pkg/orchestration/builder.go b/pkg/orchestration/builder.go index 33971414..2b8f0a23 100644 --- a/pkg/orchestration/builder.go +++ b/pkg/orchestration/builder.go @@ -271,9 +271,3 @@ func SetBuilderPush() { func (b *Builder) Unique() string { return b.instance.Unique() } - -var dryRun bool - -func SetDryRun(d bool) { - dryRun = d -} diff --git a/pkg/orchestration/builder_deploy.go b/pkg/orchestration/builder_deploy.go index 8613af14..93feac7d 100644 --- a/pkg/orchestration/builder_deploy.go +++ b/pkg/orchestration/builder_deploy.go @@ -93,10 +93,6 @@ func (b *Builder) Deploy(ctx context.Context) (*OutputProperty, error) { return nil, w.Wrapf(err, "cannot process outputProperty for deploy") } - // Handle - if dryRun { - return outputProperty, nil - } if resp.Deployment == nil { return outputProperty, nil }