diff --git a/CHANGELOG.md b/CHANGELOG.md index 75aea93a..51600b11 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -227,6 +227,44 @@ Both gates use the full resolver chain (`.spec`, `.metadata`, serve intent, prof `crdFiles` / `crFiles` added to `E2ESpec`. `tests/simulate-envtest/04-conditional-reconciliation/` covers gate-pass and gate-discard via envtest simulate. `examples/intermediate/05-when-conditions/conditional-reconciliation/` — App (reconcileGate) + Route (unconditional) pack. +### Per-target `operatorBox` — surface-specific reconciliation + +`serve.target..operatorBox` overrides the CRD-level `operatorBox` for CRs routed through that surface. The gateway stamps `orkestra.orkspace.io/serve-target` on every applied CR; the runtime reads that annotation at reconcile time and uses the matching target's templates instead of the shared CRD-level ones. CRs applied via `kubectl apply` (no annotation) fall back to the CRD-level `operatorBox`. + +```yaml +operatorBox: + onCreate: + deployments: + - name: "{{ .metadata.name }}" + services: + - name: "{{ .metadata.name }}-svc" + +serve: + enabled: true + target: + web: + primary: true + operatorBox: + onCreate: + deployments: + - name: "{{ .metadata.name }}-web" + apifixture: + operatorBox: + onCreate: + deployments: + - name: "{{ .metadata.name }}-apifixture" +``` + +`preReconcile` and `status` follow the same pattern — a target may declare its own gate conditions or status fields, with the CRD-level config as the fallback when absent. Reconciler-level settings (`reconciler.workers`, `reconciler.resync`, `reconciler.queue`, `reconciler.profile`, `autoscale`, `rollback`) are rejected by `ork validate` on target entries — the worker pool is fixed at CRD level. + +Cleanup on target change is handled automatically: `DeleteIfOwned` removes resources declared by the previous target's operatorBox that are no longer present in the new one. + +**`ork simulate --target `** — simulates a specific target's operatorBox. Also declarable in `simulate.yaml` via `spec.target:`. CLI flag takes precedence over the spec field. + +**`simulate.yaml` `spec.target:`** — new field. Pins the simulated reconciliation to a named target's operatorBox, equivalent to passing `--target` on the CLI. + +**`ork simulate` refactored** — CLI simulate helpers now take a `cliSimulateOptions` struct instead of a flat parameter list, reducing signature length across `runSimulate`, `runSimulateFromSpec`, `runSimulateDiscovery`, and `simulateOne`. + ### Serve modes, apply-time controls, and field selectors Three new blocks under `serve` and per target give platform teams granular control over the Gateway API surface, override behaviour, and full CR routing. diff --git a/cmd/cli/play_chain.go b/cmd/cli/play_chain.go index 568875a2..a4bf75f0 100644 --- a/cmd/cli/play_chain.go +++ b/cmd/cli/play_chain.go @@ -187,7 +187,7 @@ func playRunSimulate(ctx context.Context, katalogFile string, obj *unstructured. if simulateConfig != "" { return runSimulateWithCR(ctx, simulateConfig, tmp.Name()) } - return runSimulate(ctx, katalogFile, tmp.Name(), "", 10, simulate.RunOptions{SkipExternal: true}, false, false, "") + return runSimulate(ctx, katalogFile, tmp.Name(), cliSimulateOptions{MaxCycles: 10, SkipExternal: true}) } // runSimulateWithCR runs simulate using the spec file for katalog/cycles/expect @@ -252,7 +252,7 @@ func runSimulateWithCR(ctx context.Context, specPath, crFile string) error { crdOpts.Peers = in.peers crdOpts.ExistingInstances = in.existing expect := simulate.ExpectForCRD(doc.Spec.Expect, name) - if err := simulateOne(ctx, kat, name, in.cr, cycles, crdOpts, false, false, nil, "", expect); err != nil { + if err := simulateOne(ctx, kat, name, in.cr, cycles, crdOpts, cliSimulateOptions{}, nil, expect); err != nil { failed = append(failed, name) } } diff --git a/cmd/cli/push.go b/cmd/cli/push.go index 412baa05..cd7b8a4d 100644 --- a/cmd/cli/push.go +++ b/cmd/cli/push.go @@ -201,7 +201,7 @@ var pushCmd = &cobra.Command{ } else { fmt.Printf("\nRunning simulate gate (%s)...\n", registry.FileSimulate) start := time.Now() - if err := runSimulateFromSpec(cmd.Context(), simFile, "", 10, false, false, ""); err != nil { + if err := runSimulateFromSpec(cmd.Context(), simFile, cliSimulateOptions{MaxCycles: 10}); err != nil { return fmt.Errorf("✗ Simulate gate failed — push blocked\n Run 'ork simulate' to see the failures\n Use --force to override (recorded in the artifact)\n\n%w", err) } dur := time.Since(start).Round(time.Millisecond).String() diff --git a/cmd/cli/simulate.go b/cmd/cli/simulate.go index 2705fbbf..62494229 100644 --- a/cmd/cli/simulate.go +++ b/cmd/cli/simulate.go @@ -26,6 +26,17 @@ import ( sigsyaml "sigs.k8s.io/yaml" ) +// cliSimulateOptions groups the CLI-level flags that flow through all simulate helpers. +type cliSimulateOptions struct { + CRDName string + MaxCycles int + Target string + SkipExternal bool + DebugOps bool + UseEnvtest bool + K8sVersion string +} + var simulateCmd = &cobra.Command{ Use: "simulate", Short: "Simulate operator reconciliation in memory — no cluster required", @@ -40,18 +51,16 @@ should produce so the run is repeatable and verifiable: ork simulate -f katalog.yaml --cr cr.yaml # direct flags; op-print only ork simulate ./... # discovers all simulate.yaml files recursively`, RunE: func(cmd *cobra.Command, args []string) error { - crdName, _ := cmd.Flags().GetString("crd") - maxCycles, _ := cmd.Flags().GetInt("cycles") - target, _ := cmd.Flags().GetString("target") - - skipExternal, _ := cmd.Flags().GetBool("skip-external") - debugOps, _ := cmd.Flags().GetBool("debug-ops") - devServer, _ := cmd.Flags().GetBool("dev-server") - useEnvtest, _ := cmd.Flags().GetBool("envtest") - k8sVersion, _ := cmd.Flags().GetString("k8s-version") - opts := simulate.RunOptions{SkipExternal: skipExternal, Target: target} - - if devServer { + cliOpts := cliSimulateOptions{} + cliOpts.CRDName, _ = cmd.Flags().GetString("crd") + cliOpts.MaxCycles, _ = cmd.Flags().GetInt("cycles") + cliOpts.Target, _ = cmd.Flags().GetString("target") + cliOpts.SkipExternal, _ = cmd.Flags().GetBool("skip-external") + cliOpts.DebugOps, _ = cmd.Flags().GetBool("debug-ops") + cliOpts.UseEnvtest, _ = cmd.Flags().GetBool("envtest") + cliOpts.K8sVersion, _ = cmd.Flags().GetString("k8s-version") + + if devServer, _ := cmd.Flags().GetBool("dev-server"); devServer { devServerPort, _ := cmd.Flags().GetInt("dev-server-port") if err := devserver.Start(devServerPort); err != nil { return fmt.Errorf("starting dev server: %w", err) @@ -61,8 +70,7 @@ should produce so the run is repeatable and verifiable: // Discovery mode: ork simulate ./... if len(args) > 0 && args[0] == "./..." { skipRaw, _ := cmd.Flags().GetStringSlice("skip") - root := "." - return runSimulateDiscovery(cmd.Context(), root, crdName, maxCycles, skipRaw, debugOps, useEnvtest, k8sVersion) + return runSimulateDiscovery(cmd.Context(), ".", skipRaw, cliOpts) } katalogFile, _ := cmd.Flags().GetString("file") @@ -83,7 +91,7 @@ should produce so the run is repeatable and verifiable: // Simulate kind: assert mode if isSimulateDoc(katalogFile) { - return runSimulateFromSpec(cmd.Context(), katalogFile, crdName, maxCycles, debugOps, useEnvtest, k8sVersion) + return runSimulateFromSpec(cmd.Context(), katalogFile, cliOpts) } // Reject E2E files with a clear message @@ -99,11 +107,12 @@ should produce so the run is repeatable and verifiable: return fmt.Errorf("--cr is required") } - return runSimulate(cmd.Context(), katalogFile, crFile, crdName, maxCycles, opts, debugOps, useEnvtest, k8sVersion) + return runSimulate(cmd.Context(), katalogFile, crFile, cliOpts) }, } -func runSimulate(ctx context.Context, katalogFile, crFile, crdName string, maxCycles int, opts simulate.RunOptions, debugOps, useEnvtest bool, k8sVersion string) error { +func runSimulate(ctx context.Context, katalogFile, crFile string, cliOpts cliSimulateOptions) error { + maxCycles := cliOpts.MaxCycles if maxCycles <= 0 { maxCycles = 10 } @@ -134,12 +143,14 @@ func runSimulate(ctx context.Context, katalogFile, crFile, crdName string, maxCy // If --crd is given, simulate that CRD only. Otherwise simulate all. var targets []string - if crdName != "" { - targets = []string{crdName} + if cliOpts.CRDName != "" { + targets = []string{cliOpts.CRDName} } else { targets = kat.CRDNames() } + baseOpts := simulate.RunOptions{SkipExternal: cliOpts.SkipExternal, Target: cliOpts.Target} + for _, name := range targets { crdEntry, ok := kat.CRDEntry(name) if !ok { @@ -153,17 +164,17 @@ func runSimulate(ctx context.Context, katalogFile, crFile, crdName string, maxCy } return fmt.Errorf("no CR found for CRD %q (kind: %s) in %s", name, crdEntry.APITypes.Kind, crFile) } - crdOpts := opts + crdOpts := baseOpts crdOpts.Peers = in.peers crdOpts.ExistingInstances = in.existing - if err := simulateOne(ctx, kat, name, in.cr, maxCycles, crdOpts, debugOps, useEnvtest, nil, k8sVersion, nil); err != nil { + if err := simulateOne(ctx, kat, name, in.cr, maxCycles, crdOpts, cliOpts, nil, nil); err != nil { return err } } return nil } -func simulateOne(ctx context.Context, kat *katalog.Katalog, crdName string, cr *unstructured.Unstructured, maxCycles int, opts simulate.RunOptions, debugOps, useEnvtest bool, crdPaths []string, k8sVersion string, expect *orktypes.SimulateExpect) error { +func simulateOne(ctx context.Context, kat *katalog.Katalog, crdName string, cr *unstructured.Unstructured, maxCycles int, opts simulate.RunOptions, cliOpts cliSimulateOptions, crdPaths []string, expect *orktypes.SimulateExpect) error { fmt.Printf("Simulating %s/%s\n", crdName, cr.GetName()) // Emit notes for operatorBox blocks that cannot execute in the fake cluster. @@ -185,11 +196,11 @@ func simulateOne(ctx context.Context, kat *katalog.Katalog, crdName string, cr * start := time.Now() var result *simulate.Result var err error - if useEnvtest { + if cliOpts.UseEnvtest { if len(crdPaths) == 0 { return fmt.Errorf("%s --envtest requires spec.crd or spec.crdFiles to be set", failureMark()) } - result, err = simulate.RunWithEnvtest(ctx, kat, crdName, cr, maxCycles, opts, crdPaths, k8sVersion) + result, err = simulate.RunWithEnvtest(ctx, kat, crdName, cr, maxCycles, opts, crdPaths, cliOpts.K8sVersion) } else { result, err = simulate.Run(ctx, kat, crdName, cr, maxCycles, opts) } @@ -201,7 +212,7 @@ func simulateOne(ctx context.Context, kat *katalog.Katalog, crdName string, cr * spin.Stop() elapsed := time.Since(start) - if debugOps { + if cliOpts.DebugOps { fmt.Printf(" [debug-ops] %d total ops recorded across all cycles:\n", len(result.AllOps)) for _, op := range result.AllOps { fmt.Printf(" [debug-ops] cycle=%-2d verb=%-8s resource=%-20s name=%s\n", @@ -551,7 +562,7 @@ func isSimulateDoc(path string) bool { // runSimulateFromSpec loads a simulate.yaml and runs it in assert mode. // Aggregator form (imports, no spec) expands each imported file in order. -func runSimulateFromSpec(ctx context.Context, path string, crdName string, maxCycles int, debugOps, useEnvtest bool, k8sVersion string) error { +func runSimulateFromSpec(ctx context.Context, path string, cliOpts cliSimulateOptions) error { if abs, err := filepath.Abs(path); err == nil { path = abs } @@ -574,7 +585,7 @@ func runSimulateFromSpec(ctx context.Context, path string, crdName string, maxCy if !filepath.IsAbs(impPath) { impPath = filepath.Join(dir, impPath) } - if err := runSimulateFromSpec(ctx, impPath, crdName, maxCycles, debugOps, useEnvtest, k8sVersion); err != nil { + if err := runSimulateFromSpec(ctx, impPath, cliOpts); err != nil { return err } } @@ -595,10 +606,18 @@ func runSimulateFromSpec(ctx context.Context, path string, crdName string, maxCy cycles := doc.Spec.Cycles if cycles <= 0 { - cycles = maxCycles + cycles = cliOpts.MaxCycles } - opts := simulate.RunOptions{SkipExternal: doc.Spec.SkipExternal} + // CLI flag wins over spec field for both target and skipExternal. + effectiveTarget := cliOpts.Target + if effectiveTarget == "" { + effectiveTarget = doc.Spec.Target + } + opts := simulate.RunOptions{ + SkipExternal: cliOpts.SkipExternal || doc.Spec.SkipExternal, + Target: effectiveTarget, + } katalogPath := filepath.Join(dir, doc.Spec.Katalog) @@ -651,8 +670,8 @@ func runSimulateFromSpec(ctx context.Context, path string, crdName string, maxCy } var targets []string - if crdName != "" { - targets = []string{crdName} + if cliOpts.CRDName != "" { + targets = []string{cliOpts.CRDName} } else { targets = kat.CRDNames() } @@ -674,7 +693,7 @@ func runSimulateFromSpec(ctx context.Context, path string, crdName string, maxCy crdOpts.Peers = in.peers crdOpts.ExistingInstances = in.existing expect := simulate.ExpectForCRD(doc.Spec.Expect, name) - if err := simulateOne(ctx, kat, name, in.cr, cycles, crdOpts, debugOps, useEnvtest, crdPaths, k8sVersion, expect); err != nil { + if err := simulateOne(ctx, kat, name, in.cr, cycles, crdOpts, cliOpts, crdPaths, expect); err != nil { failed = append(failed, name) } } @@ -709,9 +728,9 @@ type simulateFileResult struct { cycleErrs bool } -// runSimulateDiscovery finds all e2e.yaml files under root, simulates each, +// runSimulateDiscovery finds all simulate.yaml files under root, simulates each, // and prints an aggregate summary. -func runSimulateDiscovery(ctx context.Context, root, crdName string, maxCycles int, skip []string, debugOps, useEnvtest bool, k8sVersion string) error { +func runSimulateDiscovery(ctx context.Context, root string, skip []string, cliOpts cliSimulateOptions) error { var patterns []string for _, s := range skip { patterns = append(patterns, s) @@ -729,12 +748,17 @@ func runSimulateDiscovery(ctx context.Context, root, crdName string, maxCycles i absRoot, _ := filepath.Abs(root) + // In discovery mode each file declares its own target; don't let a CLI + // --target flag override every file in the suite. + fileOpts := cliOpts + fileOpts.Target = "" + var results []simulateFileResult for _, p := range paths { rel, _ := filepath.Rel(absRoot, p) start := time.Now() - err := runSimulateFromSpec(ctx, p, crdName, maxCycles, debugOps, useEnvtest, k8sVersion) + err := runSimulateFromSpec(ctx, p, fileOpts) elapsed := time.Since(start) var res simulateFileResult diff --git a/cmd/internal/runtime_konstructor.go b/cmd/internal/runtime_konstructor.go index 6e16d4f3..7d5a4b08 100644 --- a/cmd/internal/runtime_konstructor.go +++ b/cmd/internal/runtime_konstructor.go @@ -316,9 +316,9 @@ func konstructRuntime(kfg *konfig.Konfig, m *merger.Merger, ctx context.Context) } // ── Enqueue filter — Tier 2b (pre-enqueue condition gate) ───────────── - // Register when the CRD declares operatorBox.preReconcile.enqueueGate or - // preReconcile.external conditions. - if rc := crd.PreReconcileCheck(); rc.HasEnqueueGate() { + // Register when any operatorBox (CRD-level or per-target) declares an + // enqueueGate. EvaluateEnqueueFilter resolves the effective box at runtime. + if crd.HasAnyEnqueueGate() { crdNameForFilter := crd.Name katForFilter := kat cs := kube.Clientset() diff --git a/documentation/concepts/self-service/02-target-mode.md b/documentation/concepts/self-service/02-target-mode.md index 3b54044d..ba787eb3 100644 --- a/documentation/concepts/self-service/02-target-mode.md +++ b/documentation/concepts/self-service/02-target-mode.md @@ -220,6 +220,51 @@ At apply time, `.status` is not yet available. Callers should poll `pollUrl` to --- +!!! tip "The one-sentence version" + The Katalog declares operators; targets declare operational profiles; OperatorBoxes execute those profiles; the Gateway selects the profile; the Runtime provides the shared machinery. + +## Target as a unit of runtime execution + +A target is not only a routing identifier. Each named target in `serve.target` can carry its own `operatorBox` — a complete set of lifecycle hooks, resource templates, and `preReconcile` gates. When the CR is reconciled, the reconciler selects the `operatorBox` that matches the active target. + +This means the same CRD can provision different infrastructure depending on which surface submitted the intent: + +```yaml +serve: + target: + web: + primary: true + operatorBox: + preReconcile: + enqueueGate: + when: + - field: "{{ .spec.image }}" + notEquals: "" + onCreate: + deployments: + - name: "{{ .metadata.name }}-web" + + regional: + operatorBox: + preReconcile: + reconcileGate: + when: + - field: "{{ len .spec.regions }}" + notEquals: "0" + onCreate: + deployments: + - name: "{{ .metadata.name }}-{{ .item }}" + forEach: + field: spec.regions + as: item +``` + +When a CR switches from one target to another (re-submitted via a different surface), the reconciler detects the change and cleans up resources from the previous target before applying the new box. + +→ [Per-target operatorBox schema](../../reference/schema/02-katalog/26-serve-target-operatorbox.md) — gates, surface cleanup, `keepPreviousSurface` + +--- + ## See also → [`serve.target` schema reference](../../reference/schema/02-katalog/20-serve.md#servetarget) diff --git a/documentation/concepts/self-service/10-multi-cluster-routing.md b/documentation/concepts/self-service/10-multi-cluster-routing.md index 17e6f4a3..9c8a27c2 100644 --- a/documentation/concepts/self-service/10-multi-cluster-routing.md +++ b/documentation/concepts/self-service/10-multi-cluster-routing.md @@ -111,10 +111,21 @@ the local cluster — the one the gateway runs on. This is unchanged behaviour. ## Read path behaviour -When the gateway reads resources (GET requests), cluster templates are not -resolved — the intent fields are not available on the read path. The gateway -falls back to the local cluster for reads and lists. Writes (POST, PATCH, -DELETE) resolve the cluster expression against the submitted fields. +GET requests for resources and schema support a `?cluster=` query parameter +that routes the request to the named registered cluster: + +```bash +# Read a resource from a specific cluster +curl /api/v1/resources/AppRequest/default/payments-api?cluster=prod \ + -H "Authorization: Bearer $TOKEN" + +# Get the schema for a target on a specific cluster +curl /api/v1/schema?target=app&cluster=staging \ + -H "Authorization: Bearer $TOKEN" +``` + +When `?cluster` is omitted, the gateway reads from the local cluster. When the +cluster name is not registered, the gateway returns a 404. ## Onboarding a new cluster diff --git a/documentation/reference/schema/02-katalog/20-serve.md b/documentation/reference/schema/02-katalog/20-serve.md index 00183c04..08d51156 100644 --- a/documentation/reference/schema/02-katalog/20-serve.md +++ b/documentation/reference/schema/02-katalog/20-serve.md @@ -732,6 +732,68 @@ serve: | `include` | string | Path **relative to the katalog file** to a YAML file with `tokens:` and/or `config:` keys. Resolved at load time. Inline fields take precedence on merge. | | `tokens` | map | Per-entry token restrictions — same shape as `serve.tokens`. When set, only tokens listed here are checked for access to this surface. | | `config.response` | object | Same `default`, `payload`, `exclude`, and `poll` fields as `serve.config.response`. | +| `operatorBox` | object | Per-target operatorBox — overrides the CRD-level `operatorBox` for CRs routed through this surface. See [`serve.target..operatorBox`](#servetargetnameoperatorbox) below. | + +### `serve.target..operatorBox` + +When a CR is submitted through the gateway, the apply handler stamps `orkestra.orkspace.io/serve-target` on it. The runtime reconciler reads this annotation at reconcile time and uses the matching target's `operatorBox` instead of the CRD-level one. CRs applied via `kubectl apply` (no annotation) always fall back to the CRD-level `operatorBox`. + +This makes it possible to deploy different resources depending on which surface submitted the intent — without branching on `when:` conditions inside a single shared operatorBox: + +```yaml +spec: + crds: + website: + operatorBox: # fallback — used by kubectl apply / unknown targets + onCreate: + deployments: + - name: "{{ .metadata.name }}" + services: + - name: "{{ .metadata.name }}-svc" + + serve: + enabled: true + target: + web: + primary: true + operatorBox: # used when CR arrives via the "web" target + onCreate: + deployments: + - name: "{{ .metadata.name }}-web" + apifixture: + operatorBox: # used when CR arrives via the "apifixture" target + onCreate: + deployments: + - name: "{{ .metadata.name }}-apifixture" +``` + +**Resolution order** (most specific first): + +1. The target entry whose name matches `serve-alias` annotation (alias wins over primary) +2. The target entry whose name matches `serve-target` annotation +3. CRD-level `operatorBox` — fallback when no annotation or no matching target + +**Cleanup on target change** — when a CR moves between targets (e.g. re-submitted via a different surface), the previous target's resources are cleaned up automatically via a label-selector sweep on `orkestra-owner=.`. No manual cleanup is needed. To retain old-target resources deliberately, set `keepPreviousSurface: true` in `target..apply.overrides`. + +**What stays fixed at the CRD level** — worker counts, resync intervals, and autoscale config are always taken from the CRD-level `operatorBox`. Only templates (`onReconcile`, `onCreate`, `onDelete`), status, finalizers, and external/cross blocks are resolved per-target. + +→ [Full per-target operatorBox reference](26-serve-target-operatorbox.md) — preReconcile gates, surface switch cleanup, `keepPreviousSurface`, simulate patterns + +**Simulating a specific target** — use `spec.target` in the simulate file, or `--target` on the CLI: + +```yaml +# simulate-web.yaml +spec: + katalog: ./katalog.yaml + cr: ./cr.yaml + target: web + expect: + ops: + - cycle: 1 + verb: create + resource: deployments + name: my-app-web +``` ### `include:` files for target entries diff --git a/documentation/reference/schema/02-katalog/26-serve-target-operatorbox.md b/documentation/reference/schema/02-katalog/26-serve-target-operatorbox.md new file mode 100644 index 00000000..e7e69bd1 --- /dev/null +++ b/documentation/reference/schema/02-katalog/26-serve-target-operatorbox.md @@ -0,0 +1,195 @@ +# serve.target operatorBox + +Each named target in `serve.target` can carry its own `operatorBox`. When a CR is applied through a specific target, the reconciler uses that target's `operatorBox` instead of the CRD-level one — a different set of child resources, different lifecycle hooks, and optionally a different set of `preReconcile` gates. + +This makes the target the unit of runtime execution: the same CRD can behave differently depending on which surface delivered the intent, without branching on `when:` conditions inside one shared box. + +--- + +## Declaration + +```yaml +spec: + crds: + website: + operatorBox: # CRD-level fallback — used by kubectl apply / unknown targets + onCreate: + deployments: + - name: "{{ .metadata.name }}" + + serve: + enabled: true + target: + web: + primary: true + operatorBox: # used when CR arrives via the "web" target + preReconcile: + enqueueGate: + when: + - field: "{{ .spec.image }}" + notEquals: "" + onCreate: + deployments: + - name: "{{ .metadata.name }}-web" + image: "{{ .spec.image }}" + replicas: "{{ .spec.replicas }}" + + regional: + operatorBox: # used when CR arrives via the "regional" target + preReconcile: + reconcileGate: + when: + - field: "{{ len .spec.regions }}" + notEquals: "0" + onCreate: + deployments: + - name: "{{ .metadata.name }}-{{ .item }}" + image: "{{ .spec.image }}" + forEach: + field: spec.regions + as: item + namespaces: + - name: "{{ .metadata.name }}-{{ .item }}" + forEach: + field: spec.regions + as: item +``` + +--- + +## Resolution order + +When the reconciler is invoked, it selects the `operatorBox` by checking the CR's annotations (most specific first): + +1. Target entry whose name matches `orkestra.orkspace.io/serve-alias` (alias wins over primary) +2. Target entry whose name matches `orkestra.orkspace.io/serve-target` +3. CRD-level `operatorBox` — fallback when no annotation or no matching target + +CRs applied via `kubectl apply` carry no gateway annotations, so they always use the CRD-level box. + +--- + +## `preReconcile` gates at the target level + +Per-target `operatorBox` blocks support the same `preReconcile` shape as the CRD-level box: + +```yaml +operatorBox: + preReconcile: + enqueueGate: # evaluated before the item enters the work queue + when: + - field: "{{ .spec.image }}" + notEquals: "" + + reconcileGate: # evaluated after dequeue, before reconciler runs + when: + - field: "{{ len .spec.regions }}" + notEquals: "0" + anyOf: + - field: '{{ .spec.tier }}' + equals: premium +``` + +| Gate | Evaluated by | Effect when condition fails | +|------|--------------|-----------------------------| +| `enqueueGate` | Informer (before work queue) | CR is dropped from the queue — no reconcile cycle starts | +| `reconcileGate` | Kordinator (after dequeue) | Reconcile is skipped for this cycle — item is requeued | + +Gate semantics are identical to those at the CRD level — the only difference is they apply only when this target's box is active. See [operatorBox preReconcile](04-operatorbox.md#prereconcile) for the full gate reference including `external:` calls. + +--- + +## Surface switch and resource cleanup + +When a CR is re-submitted via a different target (a surface switch), the reconciler detects the change by comparing two annotations: + +| Annotation | Written by | Value | +|---|---|---| +| `orkestra.orkspace.io/serve-alias` | Apply handler (per-request) | The alias name, or `""` for the primary target | +| `orkestra.orkspace.io/last-surface` | Reconciler (after each cycle) | The effective target active at the end of the previous reconcile | + +A mismatch between `last-surface` and the current effective target triggers a cleanup sweep before the new target's `operatorBox` runs. The sweep finds resources stamped with `orkestra-owner=.` and deletes them — both namespaced and cluster-scoped types. + +**Why a sweep and not template expansion?** When the gateway routes a CR away from the old target, the spec fields that drove that target's `forEach` declarations may already be absent (e.g. `spec.regions` is not included in the new POST body). Template-based deletion would expand `forEach` to nothing and silently miss the orphans. The label-selector sweep is immune to spec changes. + +After cleanup, `last-surface` is updated to the current target and reconciliation proceeds with the new box. + +--- + +## `keepPreviousSurface` + +To retain old-target resources after a switch (e.g. a canary scenario where both targets run simultaneously), set `keepPreviousSurface: true`. It can be declared at the CRD level or per-target: + +```yaml +# CRD level — applies to all targets +serve: + apply: + overrides: + keepPreviousSurface: true + +# Per-target — applies only when this target becomes active +target: + canary: + apply: + overrides: + keepPreviousSurface: true + operatorBox: + onCreate: + deployments: + - name: "{{ .metadata.name }}-canary" +``` + +CRD-level wins if set; per-target applies otherwise. + +| Value | Behaviour | +|-------|-----------| +| `false` (default) | Previous-surface resources are deleted on the first reconcile after a target switch | +| `true` | Previous-surface resources are left alive — no sweep runs | + +--- + +## What stays fixed at the CRD level + +Per-target `operatorBox` overrides lifecycle templates and gates. These fields are always taken from the CRD-level `operatorBox`: + +- Worker counts, resync intervals, and autoscale config +- `finalizers` +- `rollBackOnError` +- `reconciler.constructor` and `reconciler.hooks` (custom reconciler wiring) + +Only `onCreate`, `onReconcile`, `onDelete`, `preReconcile`, `status`, and `external`/`cross` blocks are resolved per-target. + +--- + +## Simulating a specific target + +Use `spec.target` in the simulate file, or `--target` on the CLI, to route the simulation through a named target's box: + +```yaml +# simulate-regional.yaml +spec: + katalog: ./katalog.yaml + cr: ./cr.yaml + target: regional + expect: + ops: + - cycle: 1 + verb: create + resource: deployments + name: my-app-eu-west + - cycle: 1 + verb: create + resource: namespaces + name: my-app-eu-west +``` + +`preReconcile` gates are evaluated in simulation. A `reconcileGate` that blocks (e.g. `len .spec.regions != 0` when regions is empty) will cause the simulate cycle to produce no ops — use `expect.crds` with `steady: false` to assert the gate fired rather than asserting specific resources. + +--- + +## See also + +- [operatorBox](04-operatorbox.md) — full `preReconcile` gate reference, `when:` conditions, `external:` calls +- [serve.target](20-serve.md#servetarget) — target map declaration, tokens, response config +- [garbage collection](../../internal/runners/docs/02-garbage-collection.md) — two-path deletion model: template-based (CR deletion) vs sweep-based (surface switch) +- [simulate](../05-simulate/index.md) — multi-target simulation patterns diff --git a/documentation/reference/schema/05-simulate/index.md b/documentation/reference/schema/05-simulate/index.md index eb0de521..417936eb 100644 --- a/documentation/reference/schema/05-simulate/index.md +++ b/documentation/reference/schema/05-simulate/index.md @@ -36,6 +36,7 @@ spec: cr: ./cr.yaml cycles: 5 skipExternal: false + target: web # optional — simulate a specific serve target's operatorBox expect: noErrors: true @@ -91,6 +92,7 @@ No `spec:`. Each imported file runs independently. The suite fails if any file f | `cr` | yes | | Path to the CR YAML file. Multi-doc YAML supported — each doc matched to its CRD by `kind`. Extra docs of the CRD-under-test's OWN kind aren't reconciled — they're seeded as pre-existing instances so `operator: unique` (and similar reconcile-time checks that list other instances) can be simulated. Only the first doc of that kind is actually reconciled | | `cycles` | no | `10` | Maximum number of reconcile cycles to run | | `skipExternal` | no | `false` | Stub all `external:` HTTP calls with an empty 200 response instead of hitting the real network | +| `target` | no | — | Serve target whose `operatorBox` governs the simulated reconciliation. When set, the runtime uses that target's templates instead of the CRD-level `operatorBox`. Equivalent to passing `--target` on the CLI; the CLI flag takes precedence when both are set. | ### `spec.expect` diff --git a/pkg/gateway/api/apply.go b/pkg/gateway/api/apply.go index 338850d8..ab081f91 100644 --- a/pkg/gateway/api/apply.go +++ b/pkg/gateway/api/apply.go @@ -326,7 +326,11 @@ func applyHandler( gvr = crd.GVR() // ─── Resolve final effective target ────────────────────────────────────────── - if effectiveTarget := crd.EffectiveServeTargetForMap(obj.Object); effectiveTarget != "" { + // Only override alias when a fieldSelector explicitly matched — never fall + // back to the primary target here, since that would mask an explicitly + // named target (e.g. "web") with the primary (e.g. "apifixture") and + // bypass the routing surface conflict check. + if effectiveTarget := crd.ServeTargetForFieldSelector(obj.Object); effectiveTarget != "" { alias = effectiveTarget } diff --git a/pkg/gateway/api/helper.go b/pkg/gateway/api/helper.go index 910bec56..0a31da2b 100644 --- a/pkg/gateway/api/helper.go +++ b/pkg/gateway/api/helper.go @@ -37,7 +37,7 @@ var ( validateK8sName = utils.ValidKubernetesName // Others - resourceChecker = utils.NewResourceChecker + resourceChecker = utils.NewResourceChecker nestedSlice = utils.NestedSlice nestedMap = utils.NestedMap deleteNestedPath = utils.DeleteNestedPath diff --git a/pkg/katalog/pre_reconcile.go b/pkg/katalog/pre_reconcile.go index 1acc928e..dc04fb3c 100644 --- a/pkg/katalog/pre_reconcile.go +++ b/pkg/katalog/pre_reconcile.go @@ -26,8 +26,10 @@ func (k *Katalog) EvaluatePreReconcile(ctx context.Context, crdName string, obj if !ok { return true, "" } - rc := entry.PreReconcileCheck() - if !rc.HasReconcileGate() { + target := orktypes.ResolveTargetFromAnnotations(obj.GetAnnotations()) + box := entry.EffectiveOperatorBox(target) + rc := box.PreReconcile + if rc == nil || !rc.HasReconcileGate() { return true, "" } @@ -79,8 +81,10 @@ func (k *Katalog) EvaluateEnqueueFilter(ctx context.Context, crdName string, obj if !ok { return true } - rc := entry.PreReconcileCheck() - if !rc.HasEnqueueGate() { + target := orktypes.ResolveTargetFromAnnotations(obj.GetAnnotations()) + box := entry.EffectiveOperatorBox(target) + rc := box.PreReconcile + if rc == nil || !rc.HasEnqueueGate() { return true } diff --git a/pkg/katalog/validate_serve_target.go b/pkg/katalog/validate_serve_target.go index 70af94e6..b781a1a6 100644 --- a/pkg/katalog/validate_serve_target.go +++ b/pkg/katalog/validate_serve_target.go @@ -1,6 +1,11 @@ package katalog -import "fmt" +import ( + "fmt" + "strings" + + orktypes "github.com/orkspace/orkestra/pkg/types" +) // validateServeTarget is called from the main Validate() chain after the // existing Serve checks. It catches: @@ -12,12 +17,16 @@ import "fmt" // make lookups ambiguous. Most commonly caused by two CRDs with the same // lowercased kind when serve.target is not set explicitly on one of them. // +// 3. Reconciler-level fields on per-target operatorBox — workers, resync, queue, +// autoscale, rollback are fixed at CRD level and ignored at runtime if declared +// on a target entry. Reject early to make the mistake visible. +// // if err := k.validateServeTarget(); err != nil { return err } func (k *Katalog) validateServeTarget() error { // Track seen targets → CRD name for duplicate detection. seen := make(map[string]string) - for crdName, crd := range k.Spec.CRDs { + for crdName, crd := range k.enabledCRDs { if !crd.HasServeTarget() { continue } @@ -43,11 +52,63 @@ func (k *Katalog) validateServeTarget() error { ) } seen[target] = crdName + + // 3. Per-target operatorBox may not declare reconciler-level settings. + if crd.Serve != nil { + for entryName, entry := range crd.Serve.Target.Entries { + if entry.OperatorBox == nil { + continue + } + if err := validateTargetOperatorBox(crdName, entryName, entry.OperatorBox); err != nil { + return err + } + } + } } return nil } +// validateTargetOperatorBox rejects fields that configure the reconciler worker +// pool — these are fixed at CRD level and have no effect per target. +func validateTargetOperatorBox(crdName, targetName string, box *orktypes.OperatorBoxConfig) error { + var bad []string + + if box.Reconciler != nil { + r := box.Reconciler + if r.Workers != 0 { + bad = append(bad, "reconciler.workers") + } + if r.Resync.Duration != 0 { + bad = append(bad, "reconciler.resync") + } + if r.Queue != (orktypes.Queue{}) { + bad = append(bad, "reconciler.queue") + } + if r.Profile != "" { + bad = append(bad, "reconciler.profile") + } + } + if box.Autoscale != nil { + bad = append(bad, "autoscale") + } + if box.Rollback != nil { + bad = append(bad, "rollback") + } + if box.RollBackOnError { + bad = append(bad, "rollBackOnError") + } + if len(bad) == 0 { + return nil + } + return fmt.Errorf( + "%s crd %q: serve.target.%s.operatorBox declares reconciler-level field(s): %s\n"+ + " These settings govern the worker pool and are fixed at CRD level.\n"+ + " Move them to the top-level operatorBox.", + failureMark(), crdName, targetName, strings.Join(bad, ", "), + ) +} + // isValidServeTarget reports whether s is a valid target string. // Valid: lowercase letters, digits, hyphens. Must be non-empty. func isValidServeTarget(s string) bool { diff --git a/pkg/katalog/validate_serve_target_test.go b/pkg/katalog/validate_serve_target_test.go new file mode 100644 index 00000000..9d86d3b2 --- /dev/null +++ b/pkg/katalog/validate_serve_target_test.go @@ -0,0 +1,269 @@ +package katalog + +import ( + "strings" + "testing" + "time" + + orktypes "github.com/orkspace/orkestra/pkg/types" +) + +// ── helpers ─────────────────────────────────────────────────────────────────── + +func katalogWithTargetCRDs(crds map[string]orktypes.CRDEntry) *Katalog { + return &Katalog{enabledCRDs: crds} +} + +func crdWithTarget(targetName string) orktypes.CRDEntry { + return orktypes.CRDEntry{ + APITypes: orktypes.APITypes{Kind: "MyResource"}, + Serve: &orktypes.ServeConfig{ + Enabled: true, + Target: orktypes.ServeTargetValue{ + Entries: map[string]*orktypes.ServeTargetConfig{ + targetName: {Primary: true}, + }, + }, + }, + } +} + +func crdWithTargetAndBox(targetName string, box *orktypes.OperatorBoxConfig) orktypes.CRDEntry { + crd := crdWithTarget(targetName) + crd.Serve.Target.Entries[targetName].OperatorBox = box + return crd +} + +// ── isValidServeTarget ──────────────────────────────────────────────────────── + +func TestIsValidServeTarget_Valid(t *testing.T) { + cases := []string{"app", "my-app", "web123", "a", "apifixture", "v2-preview"} + for _, c := range cases { + if !isValidServeTarget(c) { + t.Errorf("%q should be valid", c) + } + } +} + +func TestIsValidServeTarget_Invalid(t *testing.T) { + cases := []string{"", "MyApp", "my_app", "app app", "app.v2", "APP"} + for _, c := range cases { + if isValidServeTarget(c) { + t.Errorf("%q should be invalid", c) + } + } +} + +// ── validateServeTarget — format ────────────────────────────────────────────── + +func TestValidateServeTarget_ValidFormat(t *testing.T) { + k := katalogWithTargetCRDs(map[string]orktypes.CRDEntry{ + "res": crdWithTarget("my-resource"), + }) + if err := k.validateServeTarget(); err != nil { + t.Fatalf("unexpected error: %v", err) + } +} + +func TestValidateServeTarget_InvalidFormat(t *testing.T) { + k := katalogWithTargetCRDs(map[string]orktypes.CRDEntry{ + "res": crdWithTarget("My_Resource"), + }) + err := k.validateServeTarget() + if err == nil { + t.Fatal("expected error for invalid target format") + } + if !strings.Contains(err.Error(), "invalid") { + t.Errorf("error should mention invalid format, got: %v", err) + } +} + +// ── validateServeTarget — uniqueness ────────────────────────────────────────── + +func TestValidateServeTarget_DuplicateTarget(t *testing.T) { + k := katalogWithTargetCRDs(map[string]orktypes.CRDEntry{ + "res1": crdWithTarget("app"), + "res2": crdWithTarget("app"), + }) + err := k.validateServeTarget() + if err == nil { + t.Fatal("expected error for duplicate target") + } + if !strings.Contains(err.Error(), "app") { + t.Errorf("error should mention the duplicate target name, got: %v", err) + } +} + +func TestValidateServeTarget_UniqueTargets(t *testing.T) { + k := katalogWithTargetCRDs(map[string]orktypes.CRDEntry{ + "res1": crdWithTarget("web"), + "res2": crdWithTarget("api"), + }) + if err := k.validateServeTarget(); err != nil { + t.Fatalf("unexpected error for unique targets: %v", err) + } +} + +func TestValidateServeTarget_NoServe(t *testing.T) { + k := katalogWithTargetCRDs(map[string]orktypes.CRDEntry{ + "res": {APITypes: orktypes.APITypes{Kind: "MyResource"}}, + }) + if err := k.validateServeTarget(); err != nil { + t.Fatalf("unexpected error for CRD without serve: %v", err) + } +} + +// ── validateTargetOperatorBox — valid boxes ─────────────────────────────────── + +func TestValidateServeTarget_TargetBoxTemplatesOnly(t *testing.T) { + box := &orktypes.OperatorBoxConfig{ + OnCreate: &orktypes.HookTemplates{}, + Finalizers: []string{"my.finalizer/cleanup"}, + } + k := katalogWithTargetCRDs(map[string]orktypes.CRDEntry{ + "res": crdWithTargetAndBox("web", box), + }) + if err := k.validateServeTarget(); err != nil { + t.Fatalf("unexpected error for valid target box: %v", err) + } +} + +func TestValidateServeTarget_NilTargetBox(t *testing.T) { + k := katalogWithTargetCRDs(map[string]orktypes.CRDEntry{ + "res": crdWithTargetAndBox("web", nil), + }) + if err := k.validateServeTarget(); err != nil { + t.Fatalf("unexpected error for nil target box: %v", err) + } +} + +// ── validateTargetOperatorBox — forbidden fields ────────────────────────────── + +func TestValidateServeTarget_ReconcilerWorkers(t *testing.T) { + box := &orktypes.OperatorBoxConfig{ + Reconciler: &orktypes.ReconcilerConfig{Workers: 3}, + } + err := katalogWithTargetCRDs(map[string]orktypes.CRDEntry{ + "res": crdWithTargetAndBox("web", box), + }).validateServeTarget() + if err == nil { + t.Fatal("expected error for reconciler.workers on target box") + } + if !strings.Contains(err.Error(), "reconciler.workers") { + t.Errorf("error should name reconciler.workers, got: %v", err) + } +} + +func TestValidateServeTarget_ReconcilerResync(t *testing.T) { + box := &orktypes.OperatorBoxConfig{ + Reconciler: &orktypes.ReconcilerConfig{ + Resync: orktypes.Duration{Duration: 30 * time.Second}, + }, + } + err := katalogWithTargetCRDs(map[string]orktypes.CRDEntry{ + "res": crdWithTargetAndBox("web", box), + }).validateServeTarget() + if err == nil { + t.Fatal("expected error for reconciler.resync on target box") + } + if !strings.Contains(err.Error(), "reconciler.resync") { + t.Errorf("error should name reconciler.resync, got: %v", err) + } +} + +func TestValidateServeTarget_ReconcilerQueue(t *testing.T) { + box := &orktypes.OperatorBoxConfig{ + Reconciler: &orktypes.ReconcilerConfig{ + Queue: orktypes.Queue{MaxDepth: 50}, + }, + } + err := katalogWithTargetCRDs(map[string]orktypes.CRDEntry{ + "res": crdWithTargetAndBox("web", box), + }).validateServeTarget() + if err == nil { + t.Fatal("expected error for reconciler.queue on target box") + } + if !strings.Contains(err.Error(), "reconciler.queue") { + t.Errorf("error should name reconciler.queue, got: %v", err) + } +} + +func TestValidateServeTarget_ReconcilerProfile(t *testing.T) { + box := &orktypes.OperatorBoxConfig{ + Reconciler: &orktypes.ReconcilerConfig{Profile: "high-throughput"}, + } + err := katalogWithTargetCRDs(map[string]orktypes.CRDEntry{ + "res": crdWithTargetAndBox("web", box), + }).validateServeTarget() + if err == nil { + t.Fatal("expected error for reconciler.profile on target box") + } + if !strings.Contains(err.Error(), "reconciler.profile") { + t.Errorf("error should name reconciler.profile, got: %v", err) + } +} + +func TestValidateServeTarget_Autoscale(t *testing.T) { + box := &orktypes.OperatorBoxConfig{ + Autoscale: &orktypes.AutoscaleSpec{}, + } + err := katalogWithTargetCRDs(map[string]orktypes.CRDEntry{ + "res": crdWithTargetAndBox("web", box), + }).validateServeTarget() + if err == nil { + t.Fatal("expected error for autoscale on target box") + } + if !strings.Contains(err.Error(), "autoscale") { + t.Errorf("error should name autoscale, got: %v", err) + } +} + +func TestValidateServeTarget_Rollback(t *testing.T) { + box := &orktypes.OperatorBoxConfig{ + Rollback: &orktypes.RollbackBlock{}, + } + err := katalogWithTargetCRDs(map[string]orktypes.CRDEntry{ + "res": crdWithTargetAndBox("web", box), + }).validateServeTarget() + if err == nil { + t.Fatal("expected error for rollback on target box") + } + if !strings.Contains(err.Error(), "rollback") { + t.Errorf("error should name rollback, got: %v", err) + } +} + +func TestValidateServeTarget_RollBackOnError(t *testing.T) { + box := &orktypes.OperatorBoxConfig{RollBackOnError: true} + err := katalogWithTargetCRDs(map[string]orktypes.CRDEntry{ + "res": crdWithTargetAndBox("web", box), + }).validateServeTarget() + if err == nil { + t.Fatal("expected error for rollBackOnError on target box") + } + if !strings.Contains(err.Error(), "rollBackOnError") { + t.Errorf("error should name rollBackOnError, got: %v", err) + } +} + +func TestValidateServeTarget_MultipleViolations(t *testing.T) { + box := &orktypes.OperatorBoxConfig{ + Reconciler: &orktypes.ReconcilerConfig{ + Workers: 5, + Profile: "high-throughput", + }, + Autoscale: &orktypes.AutoscaleSpec{}, + } + err := katalogWithTargetCRDs(map[string]orktypes.CRDEntry{ + "res": crdWithTargetAndBox("web", box), + }).validateServeTarget() + if err == nil { + t.Fatal("expected error for multiple violations") + } + msg := err.Error() + for _, field := range []string{"reconciler.workers", "reconciler.profile", "autoscale"} { + if !strings.Contains(msg, field) { + t.Errorf("error should name %q, got: %v", field, msg) + } + } +} diff --git a/pkg/katalog/validation_methods.go b/pkg/katalog/validation_methods.go index fa2e6be0..8bec4fd0 100644 --- a/pkg/katalog/validation_methods.go +++ b/pkg/katalog/validation_methods.go @@ -299,32 +299,6 @@ func (k *Katalog) setDefaults(kfg *konfig.Konfig) error { } crd.OperatorBox.Reconciler = rec - // Apply defaults for targets - if crd.HasServeTarget() { - for _, target := range crd.Serve.Target.Entries { - if target.OperatorBox.IsEmpty() { - continue - } - if target.OperatorBox.Reconciler == nil { - target.OperatorBox.Reconciler = &orktypes.ReconcilerConfig{} - } - rec := target.OperatorBox.Reconciler - if rec.Workers == 0 { - rec.Workers = kfg.Katalog().DefaultWorkers() - } - if rec.Resync.Duration == 0 { - rec.Resync.Duration = kfg.Katalog().DefaultResync() - } - if rec.Queue.MaxDepth == 0 { - rec.Queue.MaxDepth = kfg.Katalog().DefaultQueueDepth() - } - if rec.Queue.FailureThreshold == 0 { - rec.Queue.FailureThreshold = kfg.Katalog().DefaultFailureThreshold() - } - target.OperatorBox.Reconciler = rec - } - } - // Handle Notifications if k.IsEmailNotificationEnabled() || k.IsSlackNotificationEnabled() { enabled := true diff --git a/pkg/labels/labels.go b/pkg/labels/labels.go index 7bdc6922..cf750e12 100644 --- a/pkg/labels/labels.go +++ b/pkg/labels/labels.go @@ -125,6 +125,11 @@ const ( // reconcile loops to determine whether a resource should be updated or deleted. OrkestraOwner = "orkestra-owner" + // OrkestraServeTarget records the effective serve surface (alias → target) + // that was active when a child resource was created. Used to detect orphaned + // resources after a surface switch and clean them up on the next reconcile. + OrkestraServeTarget = "orkestra-serve-target" + // LabelCreatedBy identifies the creator of a resource. LabelCreatedBy = "app.kubernetes.io/createdBy" @@ -195,4 +200,52 @@ const ( // Value is a JSON-encoded map of field paths to values. // Example: '{"spec.mealPlan":"dinner","spec.kitchenConfig":"standard"}' AnnotationServeSelector = "orkestra.orkspace.io/serve-selector" + + // AnnotationLastSurface records the serve surface (target name) that was active + // on the last successful reconcile of a CR. Written by the reconciler after + // surface orphan cleanup so that the next reconcile can detect a surface switch + // and clean up resources from the previous surface. + AnnotationLastSurface = "orkestra.orkspace.io/last-surface" ) + +// EffectiveOwnerKey returns the ownership identity to stamp on child resources +// and to use in DeleteIfOwned checks. When the owner has an active serve target +// the identity encodes both CR name and target so resources from different +// surfaces have distinct labels and orphan detection is precise. +// +// Format: +// - target mode: "." e.g. "hello-website.web" +// - direct apply: "" e.g. "hello-website" +func EffectiveOwnerKey(ownerName string, ownerAnnotations map[string]string) string { + if ownerAnnotations != nil { + if t := ownerAnnotations[AnnotationServeAlias]; t != "" { + return ownerName + "." + t + } + if t := ownerAnnotations[AnnotationServeTarget]; t != "" { + return ownerName + "." + t + } + } + return ownerName +} + +// StampOrkestraLabels stamps all Orkestra system ownership labels onto lbls. +// Consolidates the managed, owner, and serve-target labels into one call so +// every child resource is stamped consistently from its build* function. +// ownerAnnotations may be nil (e.g. direct kubectl apply — no serve-target set). +// +// OrkestraOwner encodes the surface identity via EffectiveOwnerKey so that +// resources from different serve surfaces carry distinct labels. This enables +// precise orphan detection when a CR switches targets. +func StampOrkestraLabels(lbls map[string]string, ownerName string, ownerAnnotations map[string]string) { + lbls[ManagedKey] = ManagedValue + lbls[OrkestraOwner] = EffectiveOwnerKey(ownerName, ownerAnnotations) + if ownerAnnotations != nil { + target := ownerAnnotations[AnnotationServeAlias] + if target == "" { + target = ownerAnnotations[AnnotationServeTarget] + } + if target != "" { + lbls[OrkestraServeTarget] = target + } + } +} diff --git a/pkg/registry/simulate/harness.go b/pkg/registry/simulate/harness.go index 3ab94908..ece10afd 100644 --- a/pkg/registry/simulate/harness.go +++ b/pkg/registry/simulate/harness.go @@ -109,10 +109,12 @@ func Run(ctx context.Context, kat *katalog.Katalog, crdName string, cr *unstruct // output explains exactly what was omitted and why. result := &Result{} box := effectiveOperatorBox(crdEntry, cr, opts.Target) + // Copy the effective box so we don't mutate the original CRD entry. + boxCopy := *box for _, phase := range []*orktypes.HookTemplates{ - box.OnCreate, - box.OnReconcile, - box.OnDelete, + boxCopy.OnCreate, + boxCopy.OnReconcile, + boxCopy.OnDelete, } { if phase == nil { continue @@ -121,6 +123,9 @@ func Run(ctx context.Context, kat *katalog.Katalog, crdName string, cr *unstruct *phase = filtered result.Notes = append(result.Notes, skipped...) } + // Build effective CRD entry — reconciler uses target's operatorBox, not the CRD-level one. + effectiveCRDEntry := crdEntry + effectiveCRDEntry.OperatorBox = boxCopy scheme, err := kat.Scheme() if err != nil { @@ -226,7 +231,7 @@ func Run(ctx context.Context, kat *katalog.Katalog, crdName string, cr *unstruct r = factoryFn(fakeKube, informer, event.Discard()) } else { r = reconciler.NewGenericReconciler( - crdEntry, + effectiveCRDEntry, informer, nil, fakeKube, @@ -242,7 +247,7 @@ func Run(ctx context.Context, kat *katalog.Katalog, crdName string, cr *unstruct return nil, err } - r = wrapWithGate(r, box.PreReconcile, kat.Notes, getFromIndexerOrFallback(indexer, key, cr)) + r = wrapWithGate(r, boxCopy.PreReconcile, kat.Notes, getFromIndexerOrFallback(indexer, key, cr)) loopResult := runLoop(ctx, r, fakeKube, key, maxCycles) loopResult.Notes = result.Notes diff --git a/pkg/registry/simulate/helper.go b/pkg/registry/simulate/helper.go index 615fd7c1..9a93d4e9 100644 --- a/pkg/registry/simulate/helper.go +++ b/pkg/registry/simulate/helper.go @@ -19,7 +19,7 @@ func effectiveOperatorBox(entry orktypes.CRDEntry, cr *unstructured.Unstructured return entry.EffectiveOperatorBox(target) } - effectiveTarget := orktypes.ResolveTargetFromAnnotations(cr) + effectiveTarget := orktypes.ResolveTargetFromAnnotations(cr.GetAnnotations()) return entry.EffectiveOperatorBox(effectiveTarget) } diff --git a/pkg/resources/clusterrolebindings/clusterrolebinding.go b/pkg/resources/clusterrolebindings/clusterrolebinding.go index 7d8a03d1..af212a43 100644 --- a/pkg/resources/clusterrolebindings/clusterrolebinding.go +++ b/pkg/resources/clusterrolebindings/clusterrolebinding.go @@ -143,7 +143,7 @@ func DeleteIfOwned(ctx context.Context, kube kubeclient.Interface, } return err } - if existing.Labels[labels.OrkestraOwner] != owner.GetName() { + if existing.Labels[labels.OrkestraOwner] != labels.EffectiveOwnerKey(owner.GetName(), owner.GetAnnotations()) { return nil } return kube.Clientset().RbacV1().ClusterRoleBindings().Delete(ctx, name, metav1.DeleteOptions{}) @@ -165,8 +165,6 @@ func Resolve(src orktypes.ClusterRoleBindingTemplateSource, ownerName string) Re for k, v := range src.Labels { spec.Labels[k] = v } - spec.Labels[labels.ManagedKey] = labels.ManagedValue - spec.Labels[labels.OrkestraOwner] = ownerName kind := src.RoleRef.Kind if kind == "" { @@ -192,6 +190,7 @@ func Resolve(src orktypes.ClusterRoleBindingTemplateSource, ownerName string) Re // ── Internal helpers ────────────────────────────────────────────────────────── func buildClusterRoleBinding(owner domain.Object, spec ResolvedClusterRoleBindingSpec) *rbacv1.ClusterRoleBinding { + labels.StampOrkestraLabels(spec.Labels, owner.GetName(), owner.GetAnnotations()) crb := &rbacv1.ClusterRoleBinding{ ObjectMeta: metav1.ObjectMeta{ Name: spec.Name, diff --git a/pkg/resources/clusterroles/clusterrole.go b/pkg/resources/clusterroles/clusterrole.go index aa42cfa6..d70f1f8e 100644 --- a/pkg/resources/clusterroles/clusterrole.go +++ b/pkg/resources/clusterroles/clusterrole.go @@ -133,7 +133,7 @@ func DeleteIfOwned(ctx context.Context, kube kubeclient.Interface, } return err } - if existing.Labels[labels.OrkestraOwner] != owner.GetName() { + if existing.Labels[labels.OrkestraOwner] != labels.EffectiveOwnerKey(owner.GetName(), owner.GetAnnotations()) { return nil } return kube.Clientset().RbacV1().ClusterRoles().Delete(ctx, name, metav1.DeleteOptions{}) @@ -155,8 +155,6 @@ func Resolve(src orktypes.ClusterRoleTemplateSource, ownerName string) ResolvedC for k, v := range src.Labels { spec.Labels[k] = v } - spec.Labels[labels.ManagedKey] = labels.ManagedValue - spec.Labels[labels.OrkestraOwner] = ownerName for _, r := range src.Rules { spec.Rules = append(spec.Rules, rbacv1.PolicyRule{ @@ -173,6 +171,7 @@ func Resolve(src orktypes.ClusterRoleTemplateSource, ownerName string) ResolvedC // ── Internal helpers ────────────────────────────────────────────────────────── func buildClusterRole(owner domain.Object, spec ResolvedClusterRoleSpec) *rbacv1.ClusterRole { + labels.StampOrkestraLabels(spec.Labels, owner.GetName(), owner.GetAnnotations()) cr := &rbacv1.ClusterRole{ ObjectMeta: metav1.ObjectMeta{ Name: spec.Name, diff --git a/pkg/resources/configmaps/configmap.go b/pkg/resources/configmaps/configmap.go index 4063a511..657e40e6 100644 --- a/pkg/resources/configmaps/configmap.go +++ b/pkg/resources/configmaps/configmap.go @@ -252,7 +252,7 @@ func DeleteIfOwned(ctx context.Context, kube kubeclient.Interface, } // Only delete if we own it - if existing.Labels[labels.OrkestraOwner] != owner.GetName() { + if existing.Labels[labels.OrkestraOwner] != labels.EffectiveOwnerKey(owner.GetName(), owner.GetAnnotations()) { return nil } return kube.Clientset().CoreV1().ConfigMaps(namespace). @@ -280,9 +280,6 @@ func Resolve(src orktypes.ConfigMapTemplateSource, ownerName string) ResolvedCon spec.Labels[k] = v } - spec.Labels[labels.ManagedKey] = labels.ManagedValue - spec.Labels[labels.OrkestraOwner] = ownerName - return spec } @@ -333,6 +330,7 @@ func resolveData( } func buildConfigMap(owner domain.Object, spec ResolvedConfigMapSpec, namespace string, data map[string]string) *corev1.ConfigMap { + labels.StampOrkestraLabels(spec.Labels, owner.GetName(), owner.GetAnnotations()) return &corev1.ConfigMap{ ObjectMeta: metav1.ObjectMeta{ Name: spec.Name, diff --git a/pkg/resources/cronjobs/cronjob.go b/pkg/resources/cronjobs/cronjob.go index d5da69f8..d12bc196 100644 --- a/pkg/resources/cronjobs/cronjob.go +++ b/pkg/resources/cronjobs/cronjob.go @@ -210,7 +210,7 @@ func DeleteIfOwned(ctx context.Context, kube kubeclient.Interface, } return fmt.Errorf("cronjob.DeleteIfOwned: getting %q: %w", name, err) } - if existing.Labels[labels.OrkestraOwner] != owner.GetName() { + if existing.Labels[labels.OrkestraOwner] != labels.EffectiveOwnerKey(owner.GetName(), owner.GetAnnotations()) { return nil } return kube.Clientset().BatchV1().CronJobs(namespace). @@ -288,15 +288,13 @@ func Resolve(src orktypes.CronJobTemplateSource, ownerName string, reg orktypes. spec.Labels[k] = v } - spec.Labels[labels.ManagedKey] = labels.ManagedValue - spec.Labels[labels.OrkestraOwner] = ownerName - return spec } // ── Internal helpers ────────────────────────────────────────────────────────── func buildCronJob(owner domain.Object, spec ResolvedCronJobSpec, namespace string) *batchv1.CronJob { + labels.StampOrkestraLabels(spec.Labels, owner.GetName(), owner.GetAnnotations()) cj := &batchv1.CronJob{ ObjectMeta: metav1.ObjectMeta{ Name: spec.Name, diff --git a/pkg/resources/customresources/custom.go b/pkg/resources/customresources/custom.go index de8d11a0..7716deb3 100644 --- a/pkg/resources/customresources/custom.go +++ b/pkg/resources/customresources/custom.go @@ -311,7 +311,7 @@ func DeleteIfOwned(ctx context.Context, kube kubeclient.Interface, owner domain. } // Only delete if we own it - if labelsMap[orklabels.OrkestraOwner] != owner.GetName() { + if labelsMap[orklabels.OrkestraOwner] != orklabels.EffectiveOwnerKey(owner.GetName(), owner.GetAnnotations()) { return nil } @@ -344,12 +344,11 @@ func buildUnstructured(spec ResolvedCustomResourceSpec, owner domain.Object, gvk // Labels: merge user-declared + Orkestra managed orklabels. Copy rather than // mutate spec.Metadata.Labels directly — that map may be shared/reused // (e.g. across forEach expansions of the same source). - lbls := make(map[string]string, len(spec.Metadata.Labels)+2) + lbls := make(map[string]string, len(spec.Metadata.Labels)+3) for k, v := range spec.Metadata.Labels { lbls[k] = v } - lbls[orklabels.ManagedKey] = orklabels.ManagedValue - lbls[orklabels.OrkestraOwner] = owner.GetName() + orklabels.StampOrkestraLabels(lbls, owner.GetName(), owner.GetAnnotations()) u.SetLabels(lbls) // Annotations: copy for the same reason as Labels above. diff --git a/pkg/resources/deployments/deployment.go b/pkg/resources/deployments/deployment.go index 26450b27..7c0548ef 100644 --- a/pkg/resources/deployments/deployment.go +++ b/pkg/resources/deployments/deployment.go @@ -148,7 +148,7 @@ func DeleteIfOwned(ctx context.Context, kube kubeclient.Interface, return err } // Only delete if we own it - if existing.Labels[labels.OrkestraOwner] != owner.GetName() { + if existing.Labels[labels.OrkestraOwner] != labels.EffectiveOwnerKey(owner.GetName(), owner.GetAnnotations()) { return nil } return kube.Clientset().AppsV1().Deployments(namespace). @@ -219,8 +219,6 @@ func Resolve(src orktypes.DeploymentTemplateSource, ownerName string, reg orktyp } // Orkestra system labels — always added - spec.Labels[labels.ManagedKey] = labels.ManagedValue - spec.Labels[labels.OrkestraOwner] = ownerName return spec } @@ -228,6 +226,7 @@ func Resolve(src orktypes.DeploymentTemplateSource, ownerName string, reg orktyp // ── Internal helpers ────────────────────────────────────────────────────────── func buildDeployment(owner domain.Object, spec ResolvedDeploymentSpec, namespace string) *appsv1.Deployment { + labels.StampOrkestraLabels(spec.Labels, owner.GetName(), owner.GetAnnotations()) // Debug line logger.Debug(). Interface("env", spec.Env). diff --git a/pkg/resources/hpas/hpa.go b/pkg/resources/hpas/hpa.go index ead2cc7b..43e5eec5 100644 --- a/pkg/resources/hpas/hpa.go +++ b/pkg/resources/hpas/hpa.go @@ -145,7 +145,7 @@ func DeleteIfOwned(ctx context.Context, kube kubeclient.Interface, } return err } - if existing.Labels[labels.OrkestraOwner] != owner.GetName() { + if existing.Labels[labels.OrkestraOwner] != labels.EffectiveOwnerKey(owner.GetName(), owner.GetAnnotations()) { return nil } return kube.Clientset().AutoscalingV2().HorizontalPodAutoscalers(namespace). @@ -208,8 +208,6 @@ func Resolve(src orktypes.HPATemplateSource, ownerName string, reg orktypes.Prof } // System labels - spec.Labels[labels.ManagedKey] = labels.ManagedValue - spec.Labels[labels.OrkestraOwner] = ownerName return spec } @@ -217,6 +215,7 @@ func Resolve(src orktypes.HPATemplateSource, ownerName string, reg orktypes.Prof // ── Internal helpers ────────────────────────────────────────────────────────── func buildHPA(owner domain.Object, spec ResolvedHPASpec, namespace string) *autoscalingv2.HorizontalPodAutoscaler { + labels.StampOrkestraLabels(spec.Labels, owner.GetName(), owner.GetAnnotations()) apiVersion := "" kind := "" if u, ok := owner.(*unstructured.Unstructured); ok { diff --git a/pkg/resources/ingresses/ingress.go b/pkg/resources/ingresses/ingress.go index 1ab03334..8a79ac02 100644 --- a/pkg/resources/ingresses/ingress.go +++ b/pkg/resources/ingresses/ingress.go @@ -143,7 +143,7 @@ func DeleteIfOwned(ctx context.Context, kube kubeclient.Interface, } return err } - if existing.Labels[labels.OrkestraOwner] != owner.GetName() { + if existing.Labels[labels.OrkestraOwner] != labels.EffectiveOwnerKey(owner.GetName(), owner.GetAnnotations()) { return nil } return kube.Clientset().NetworkingV1().Ingresses(namespace). @@ -190,8 +190,6 @@ func Resolve(src orktypes.IngressTemplateSource, ownerName string) ResolvedIngre } // System labels - spec.Labels[labels.ManagedKey] = labels.ManagedValue - spec.Labels[labels.OrkestraOwner] = ownerName if src.TLS != nil && src.TLS.Create { secretName := src.TLS.SecretName @@ -212,6 +210,7 @@ func Resolve(src orktypes.IngressTemplateSource, ownerName string) ResolvedIngre // ── Internal helpers ────────────────────────────────────────────────────────── func buildIngress(owner domain.Object, spec ResolvedIngressSpec, namespace string) *networkingv1.Ingress { + labels.StampOrkestraLabels(spec.Labels, owner.GetName(), owner.GetAnnotations()) apiVersion := "" kind := "" if u, ok := owner.(*unstructured.Unstructured); ok { diff --git a/pkg/resources/jobs/job.go b/pkg/resources/jobs/job.go index 3b200bc3..79302043 100644 --- a/pkg/resources/jobs/job.go +++ b/pkg/resources/jobs/job.go @@ -174,8 +174,6 @@ func Resolve(src orktypes.JobTemplateSource, backoffLimit int, ownerName string, } // System labels - spec.Labels[labels.ManagedKey] = labels.ManagedValue - spec.Labels[labels.OrkestraOwner] = ownerName return spec } @@ -183,6 +181,7 @@ func Resolve(src orktypes.JobTemplateSource, backoffLimit int, ownerName string, // ── Internal helpers ────────────────────────────────────────────────────────── func buildJob(owner domain.Object, spec ResolvedJobSpec, namespace string) *batchv1.Job { + labels.StampOrkestraLabels(spec.Labels, owner.GetName(), owner.GetAnnotations()) backoffLimit := int32(spec.BackoffLimit) container := corev1.Container{ diff --git a/pkg/resources/limitranges/limitrange.go b/pkg/resources/limitranges/limitrange.go index b0ee9c69..8a240e14 100644 --- a/pkg/resources/limitranges/limitrange.go +++ b/pkg/resources/limitranges/limitrange.go @@ -153,7 +153,7 @@ func DeleteIfOwned(ctx context.Context, kube kubeclient.Interface, } return err } - if existing.Labels[labels.OrkestraOwner] != owner.GetName() { + if existing.Labels[labels.OrkestraOwner] != labels.EffectiveOwnerKey(owner.GetName(), owner.GetAnnotations()) { return nil } return kube.Clientset().CoreV1().LimitRanges(namespace). @@ -234,8 +234,6 @@ func Resolve(src orktypes.LimitRangeTemplateSource, ownerName string, reg orktyp for k, v := range src.Labels { spec.Labels[k] = v } - spec.Labels[labels.ManagedKey] = labels.ManagedValue - spec.Labels[labels.OrkestraOwner] = ownerName return spec } @@ -290,6 +288,7 @@ func buildLimitRange( namespace string, limits []orktypes.LimitRangeItem, ) *corev1.LimitRange { + labels.StampOrkestraLabels(spec.Labels, owner.GetName(), owner.GetAnnotations()) return &corev1.LimitRange{ ObjectMeta: metav1.ObjectMeta{ Name: spec.Name, diff --git a/pkg/resources/namespaces/namespace.go b/pkg/resources/namespaces/namespace.go index a87ee38d..53830cdb 100644 --- a/pkg/resources/namespaces/namespace.go +++ b/pkg/resources/namespaces/namespace.go @@ -159,7 +159,7 @@ func DeleteIfOwned(ctx context.Context, kube kubeclient.Interface, return err } // Only delete if we own it - if existing.Labels[labels.OrkestraOwner] != owner.GetName() { + if existing.Labels[labels.OrkestraOwner] != labels.EffectiveOwnerKey(owner.GetName(), owner.GetAnnotations()) { return nil } return kube.Clientset().CoreV1().Namespaces(). @@ -184,15 +184,13 @@ func Resolve(src orktypes.NamespaceTemplateSource, ownerName string) ResolvedNam spec.Labels[k] = v } - spec.Labels[labels.ManagedKey] = labels.ManagedValue - spec.Labels[labels.OrkestraOwner] = ownerName - return spec } // ── Internal helpers ────────────────────────────────────────────────────────── func buildNamespace(owner domain.Object, spec ResolvedNamespaceSpec) *corev1.Namespace { + labels.StampOrkestraLabels(spec.Labels, owner.GetName(), owner.GetAnnotations()) ns := &corev1.Namespace{ ObjectMeta: metav1.ObjectMeta{ Name: spec.Name, diff --git a/pkg/resources/networkpolicies/networkpolicy.go b/pkg/resources/networkpolicies/networkpolicy.go index 9a8f4f70..3094a704 100644 --- a/pkg/resources/networkpolicies/networkpolicy.go +++ b/pkg/resources/networkpolicies/networkpolicy.go @@ -155,7 +155,7 @@ func DeleteIfOwned(ctx context.Context, kube kubeclient.Interface, } return err } - if existing.Labels[labels.OrkestraOwner] != owner.GetName() { + if existing.Labels[labels.OrkestraOwner] != labels.EffectiveOwnerKey(owner.GetName(), owner.GetAnnotations()) { return nil } return kube.Clientset().NetworkingV1().NetworkPolicies(namespace). @@ -243,8 +243,6 @@ func Resolve(src orktypes.NetworkPolicyTemplateSource, ownerName string, reg ork for k, v := range src.Labels { spec.Labels[k] = v } - spec.Labels[labels.ManagedKey] = labels.ManagedValue - spec.Labels[labels.OrkestraOwner] = ownerName return spec } @@ -315,6 +313,7 @@ func buildNetworkPolicyFromSpec( namespace string, npSpec networkingv1.NetworkPolicySpec, ) *networkingv1.NetworkPolicy { + labels.StampOrkestraLabels(spec.Labels, owner.GetName(), owner.GetAnnotations()) return &networkingv1.NetworkPolicy{ ObjectMeta: metav1.ObjectMeta{ Name: spec.Name, diff --git a/pkg/resources/pdbs/pdb.go b/pkg/resources/pdbs/pdb.go index 10fc7a74..3e7bde5c 100644 --- a/pkg/resources/pdbs/pdb.go +++ b/pkg/resources/pdbs/pdb.go @@ -153,7 +153,7 @@ func DeleteIfOwned(ctx context.Context, kube kubeclient.Interface, } return err } - if existing.Labels[labels.OrkestraOwner] != owner.GetName() { + if existing.Labels[labels.OrkestraOwner] != labels.EffectiveOwnerKey(owner.GetName(), owner.GetAnnotations()) { return nil } return kube.Clientset().PolicyV1().PodDisruptionBudgets(namespace). @@ -200,8 +200,6 @@ func Resolve(src orktypes.PDBTemplateSource, ownerName string, reg orktypes.Prof } // System labels - spec.Labels[labels.ManagedKey] = labels.ManagedValue - spec.Labels[labels.OrkestraOwner] = ownerName return spec } @@ -209,6 +207,7 @@ func Resolve(src orktypes.PDBTemplateSource, ownerName string, reg orktypes.Prof // ── Internal helpers ────────────────────────────────────────────────────────── func buildPDB(owner domain.Object, spec ResolvedPDBSpec, namespace string) *policyv1.PodDisruptionBudget { + labels.StampOrkestraLabels(spec.Labels, owner.GetName(), owner.GetAnnotations()) apiVersion := "" kind := "" if u, ok := owner.(*unstructured.Unstructured); ok { diff --git a/pkg/resources/pods/pod.go b/pkg/resources/pods/pod.go index afef196b..2792ac7f 100644 --- a/pkg/resources/pods/pod.go +++ b/pkg/resources/pods/pod.go @@ -153,7 +153,7 @@ func DeleteIfOwned(ctx context.Context, kube kubeclient.Interface, return err } // Only delete if we own it - if existing.Labels[labels.OrkestraOwner] != owner.GetName() { + if existing.Labels[labels.OrkestraOwner] != labels.EffectiveOwnerKey(owner.GetName(), owner.GetAnnotations()) { return nil } return kube.Clientset().CoreV1().Pods(namespace). @@ -203,8 +203,6 @@ func Resolve(src orktypes.PodTemplateSource, ownerName string, reg orktypes.Prof } // System labels — always present - spec.Labels[labels.ManagedKey] = labels.ManagedValue - spec.Labels[labels.OrkestraOwner] = ownerName return spec } @@ -212,6 +210,7 @@ func Resolve(src orktypes.PodTemplateSource, ownerName string, reg orktypes.Prof // ── Internal helpers ────────────────────────────────────────────────────────── func buildPod(owner domain.Object, spec ResolvedPodSpec, namespace string) *corev1.Pod { + labels.StampOrkestraLabels(spec.Labels, owner.GetName(), owner.GetAnnotations()) pod := &corev1.Pod{ ObjectMeta: metav1.ObjectMeta{ Name: spec.Name, diff --git a/pkg/resources/pvcs/pvc.go b/pkg/resources/pvcs/pvc.go index bdb6f80c..36ab125d 100644 --- a/pkg/resources/pvcs/pvc.go +++ b/pkg/resources/pvcs/pvc.go @@ -112,7 +112,7 @@ func DeleteIfOwned(ctx context.Context, kube kubeclient.Interface, owner domain. } return err } - if existing.Labels[labels.OrkestraOwner] != owner.GetName() { + if existing.Labels[labels.OrkestraOwner] != labels.EffectiveOwnerKey(owner.GetName(), owner.GetAnnotations()) { return nil } return kube.Clientset().CoreV1().PersistentVolumeClaims(namespace).Delete(ctx, name, metav1.DeleteOptions{}) @@ -142,8 +142,6 @@ func Resolve(src orktypes.PVCTemplateSource, ownerName string) ResolvedPVCSpec { for k, v := range src.Labels { spec.Labels[k] = v } - spec.Labels[labels.ManagedKey] = labels.ManagedValue - spec.Labels[labels.OrkestraOwner] = ownerName return spec } @@ -151,6 +149,7 @@ func Resolve(src orktypes.PVCTemplateSource, ownerName string) ResolvedPVCSpec { // ── Internal helpers ────────────────────────────────────────────────────────── func buildPVC(owner domain.Object, spec ResolvedPVCSpec, ns string) *corev1.PersistentVolumeClaim { + labels.StampOrkestraLabels(spec.Labels, owner.GetName(), owner.GetAnnotations()) apiVersion := "" kind := "" if u, ok := owner.(*unstructured.Unstructured); ok { diff --git a/pkg/resources/pvs/pv.go b/pkg/resources/pvs/pv.go index 5223e6ad..6d1fc40a 100644 --- a/pkg/resources/pvs/pv.go +++ b/pkg/resources/pvs/pv.go @@ -103,7 +103,7 @@ func DeleteIfOwned(ctx context.Context, kube kubeclient.Interface, owner domain. } return err } - if existing.Labels[labels.OrkestraOwner] != owner.GetName() { + if existing.Labels[labels.OrkestraOwner] != labels.EffectiveOwnerKey(owner.GetName(), owner.GetAnnotations()) { return nil } return kube.Clientset().CoreV1().PersistentVolumes().Delete(ctx, name, metav1.DeleteOptions{}) @@ -134,8 +134,6 @@ func Resolve(src orktypes.PVTemplateSource, ownerName string) ResolvedPVSpec { for k, v := range src.Labels { spec.Labels[k] = v } - spec.Labels[labels.ManagedKey] = labels.ManagedValue - spec.Labels[labels.OrkestraOwner] = ownerName return spec } @@ -143,6 +141,7 @@ func Resolve(src orktypes.PVTemplateSource, ownerName string) ResolvedPVSpec { // ── Internal helpers ────────────────────────────────────────────────────────── func buildPV(owner domain.Object, spec ResolvedPVSpec) *corev1.PersistentVolume { + labels.StampOrkestraLabels(spec.Labels, owner.GetName(), owner.GetAnnotations()) capacityQty := resource.MustParse(spec.Capacity) var accessModes []corev1.PersistentVolumeAccessMode diff --git a/pkg/resources/replicasets/replicaset.go b/pkg/resources/replicasets/replicaset.go index ba5ee40e..8e7099c6 100644 --- a/pkg/resources/replicasets/replicaset.go +++ b/pkg/resources/replicasets/replicaset.go @@ -145,7 +145,7 @@ func DeleteIfOwned(ctx context.Context, kube kubeclient.Interface, return err } - if existing.Labels[labels.OrkestraOwner] != owner.GetName() { + if existing.Labels[labels.OrkestraOwner] != labels.EffectiveOwnerKey(owner.GetName(), owner.GetAnnotations()) { return nil } @@ -213,15 +213,13 @@ func Resolve(src orktypes.ReplicaSetTemplateSource, ownerName string, reg orktyp spec.RollingUpdate = &r } - spec.Labels[labels.ManagedKey] = labels.ManagedValue - spec.Labels[labels.OrkestraOwner] = ownerName - return spec } // ── Internal helpers ────────────────────────────────────────────────────────── func buildReplicaSet(owner domain.Object, spec ResolvedReplicaSetSpec, namespace string) *appsv1.ReplicaSet { + labels.StampOrkestraLabels(spec.Labels, owner.GetName(), owner.GetAnnotations()) logger.Debug(). Interface("env", spec.Env). Interface("envFrom", spec.EnvFrom). diff --git a/pkg/resources/resourcequotas/resourcequota.go b/pkg/resources/resourcequotas/resourcequota.go index 44baca73..191fafb8 100644 --- a/pkg/resources/resourcequotas/resourcequota.go +++ b/pkg/resources/resourcequotas/resourcequota.go @@ -153,7 +153,7 @@ func DeleteIfOwned(ctx context.Context, kube kubeclient.Interface, } return err } - if existing.Labels[labels.OrkestraOwner] != owner.GetName() { + if existing.Labels[labels.OrkestraOwner] != labels.EffectiveOwnerKey(owner.GetName(), owner.GetAnnotations()) { return nil } return kube.Clientset().CoreV1().ResourceQuotas(namespace). @@ -236,8 +236,6 @@ func Resolve(src orktypes.ResourceQuotaTemplateSource, ownerName string, reg ork for k, v := range src.Labels { spec.Labels[k] = v } - spec.Labels[labels.ManagedKey] = labels.ManagedValue - spec.Labels[labels.OrkestraOwner] = ownerName return spec } @@ -297,6 +295,7 @@ func buildResourceQuota( namespace string, hard map[string]string, ) *corev1.ResourceQuota { + labels.StampOrkestraLabels(spec.Labels, owner.GetName(), owner.GetAnnotations()) return &corev1.ResourceQuota{ ObjectMeta: metav1.ObjectMeta{ Name: spec.Name, diff --git a/pkg/resources/rolebindings/rolebinding.go b/pkg/resources/rolebindings/rolebinding.go index da868150..683b0b19 100644 --- a/pkg/resources/rolebindings/rolebinding.go +++ b/pkg/resources/rolebindings/rolebinding.go @@ -160,7 +160,7 @@ func DeleteIfOwned(ctx context.Context, kube kubeclient.Interface, } return err } - if existing.Labels[labels.OrkestraOwner] != owner.GetName() { + if existing.Labels[labels.OrkestraOwner] != labels.EffectiveOwnerKey(owner.GetName(), owner.GetAnnotations()) { return nil } return kube.Clientset().RbacV1().RoleBindings(namespace).Delete(ctx, name, metav1.DeleteOptions{}) @@ -184,9 +184,6 @@ func Resolve(src orktypes.RoleBindingTemplateSource, ownerName string) ResolvedR spec.Labels[k] = v } - spec.Labels[labels.ManagedKey] = labels.ManagedValue - spec.Labels[labels.OrkestraOwner] = ownerName - kind := src.RoleRef.Kind if kind == "" { kind = "Role" @@ -211,6 +208,7 @@ func Resolve(src orktypes.RoleBindingTemplateSource, ownerName string) ResolvedR // ── Internal helpers ────────────────────────────────────────────────────────── func buildRoleBinding(owner domain.Object, spec ResolvedRoleBindingSpec, namespace string) *rbacv1.RoleBinding { + labels.StampOrkestraLabels(spec.Labels, owner.GetName(), owner.GetAnnotations()) return &rbacv1.RoleBinding{ ObjectMeta: metav1.ObjectMeta{ Name: spec.Name, diff --git a/pkg/resources/roles/role.go b/pkg/resources/roles/role.go index 47760383..1ece8ba5 100644 --- a/pkg/resources/roles/role.go +++ b/pkg/resources/roles/role.go @@ -150,7 +150,7 @@ func DeleteIfOwned(ctx context.Context, kube kubeclient.Interface, } return err } - if existing.Labels[labels.OrkestraOwner] != owner.GetName() { + if existing.Labels[labels.OrkestraOwner] != labels.EffectiveOwnerKey(owner.GetName(), owner.GetAnnotations()) { return nil } return kube.Clientset().RbacV1().Roles(namespace).Delete(ctx, name, metav1.DeleteOptions{}) @@ -174,9 +174,6 @@ func Resolve(src orktypes.RoleTemplateSource, ownerName string) ResolvedRoleSpec spec.Labels[k] = v } - spec.Labels[labels.ManagedKey] = labels.ManagedValue - spec.Labels[labels.OrkestraOwner] = ownerName - for _, r := range src.Rules { spec.Rules = append(spec.Rules, rbacv1.PolicyRule{ APIGroups: r.APIGroups, @@ -192,6 +189,7 @@ func Resolve(src orktypes.RoleTemplateSource, ownerName string) ResolvedRoleSpec // ── Internal helpers ────────────────────────────────────────────────────────── func buildRole(owner domain.Object, spec ResolvedRoleSpec, namespace string) *rbacv1.Role { + labels.StampOrkestraLabels(spec.Labels, owner.GetName(), owner.GetAnnotations()) return &rbacv1.Role{ ObjectMeta: metav1.ObjectMeta{ Name: spec.Name, diff --git a/pkg/resources/secrets/secret.go b/pkg/resources/secrets/secret.go index 7258f368..711714b4 100644 --- a/pkg/resources/secrets/secret.go +++ b/pkg/resources/secrets/secret.go @@ -262,7 +262,7 @@ func DeleteIfOwned(ctx context.Context, kube kubeclient.Interface, return err } // Only delete if we own it - if existing.Labels[labels.OrkestraOwner] != owner.GetName() { + if existing.Labels[labels.OrkestraOwner] != labels.EffectiveOwnerKey(owner.GetName(), owner.GetAnnotations()) { return nil } return kube.Clientset().CoreV1().Secrets(namespace). @@ -295,8 +295,6 @@ func Resolve(src orktypes.SecretTemplateSource, ownerName string) ResolvedSecret } // System labels - spec.Labels[labels.ManagedKey] = labels.ManagedValue - spec.Labels[labels.OrkestraOwner] = ownerName return spec } @@ -339,6 +337,7 @@ func resolveData( } func buildSecret(owner domain.Object, spec ResolvedSecretSpec, namespace string, data map[string][]byte, stringData map[string]string) *corev1.Secret { + labels.StampOrkestraLabels(spec.Labels, owner.GetName(), owner.GetAnnotations()) secretType := corev1.SecretTypeOpaque switch strings.ToLower(spec.Type) { case "kubernetes.io/tls": diff --git a/pkg/resources/serviceaccounts/serviceaccount.go b/pkg/resources/serviceaccounts/serviceaccount.go index c97b2810..22dfd449 100644 --- a/pkg/resources/serviceaccounts/serviceaccount.go +++ b/pkg/resources/serviceaccounts/serviceaccount.go @@ -122,7 +122,7 @@ func DeleteIfOwned(ctx context.Context, kube kubeclient.Interface, return err } // Only delete if we own it - if existing.Labels[labels.OrkestraOwner] != owner.GetName() { + if existing.Labels[labels.OrkestraOwner] != labels.EffectiveOwnerKey(owner.GetName(), owner.GetAnnotations()) { return nil } return kube.Clientset().CoreV1().ServiceAccounts(namespace). @@ -147,15 +147,13 @@ func Resolve(src orktypes.ServiceAccountTemplateSource, ownerName string) Resolv spec.Labels[k] = v } - spec.Labels[labels.ManagedKey] = labels.ManagedValue - spec.Labels[labels.OrkestraOwner] = ownerName - return spec } // ── Internal helpers ────────────────────────────────────────────────────────── func buildServiceAccount(owner domain.Object, spec ResolvedServiceAccountSpec, namespace string) *corev1.ServiceAccount { + labels.StampOrkestraLabels(spec.Labels, owner.GetName(), owner.GetAnnotations()) return &corev1.ServiceAccount{ ObjectMeta: metav1.ObjectMeta{ Name: spec.Name, diff --git a/pkg/resources/services/services.go b/pkg/resources/services/services.go index 489691cb..9d4fae99 100644 --- a/pkg/resources/services/services.go +++ b/pkg/resources/services/services.go @@ -147,7 +147,7 @@ func DeleteIfOwned(ctx context.Context, kube kubeclient.Interface, return err } // Only delete if we own it - if existing.Labels[labels.OrkestraOwner] != owner.GetName() { + if existing.Labels[labels.OrkestraOwner] != labels.EffectiveOwnerKey(owner.GetName(), owner.GetAnnotations()) { return nil } return kube.Clientset().CoreV1().Services(namespace). @@ -202,8 +202,6 @@ func Resolve(src orktypes.ServiceTemplateSource, ownerName string) ResolvedServi } // System labels - spec.Labels[labels.ManagedKey] = labels.ManagedValue - spec.Labels[labels.OrkestraOwner] = ownerName return spec } @@ -211,6 +209,7 @@ func Resolve(src orktypes.ServiceTemplateSource, ownerName string) ResolvedServi // ── Internal helpers ────────────────────────────────────────────────────────── func buildService(owner domain.Object, spec ResolvedServiceSpec, namespace string) *corev1.Service { + labels.StampOrkestraLabels(spec.Labels, owner.GetName(), owner.GetAnnotations()) svcType := corev1.ServiceTypeClusterIP switch spec.Type { case "NodePort": diff --git a/pkg/resources/statefulsets/statefulset.go b/pkg/resources/statefulsets/statefulset.go index d9cf2a50..b0ddbb1a 100644 --- a/pkg/resources/statefulsets/statefulset.go +++ b/pkg/resources/statefulsets/statefulset.go @@ -122,7 +122,7 @@ func DeleteIfOwned(ctx context.Context, kube kubeclient.Interface, owner domain. } return err } - if existing.Labels[labels.OrkestraOwner] != owner.GetName() { + if existing.Labels[labels.OrkestraOwner] != labels.EffectiveOwnerKey(owner.GetName(), owner.GetAnnotations()) { return nil } return kube.Clientset().AppsV1().StatefulSets(namespace).Delete(ctx, name, metav1.DeleteOptions{}) @@ -190,9 +190,6 @@ func Resolve(src orktypes.StatefulSetTemplateSource, ownerName string, reg orkty spec.Annotations[k] = v } - spec.Labels[labels.ManagedKey] = labels.ManagedValue - spec.Labels[labels.OrkestraOwner] = ownerName - if src.RollingUpdate != nil && src.RollingUpdate.Profile != "" { expansion, err := profiles.ApplyRollingUpdateProfile(src.RollingUpdate.Profile, reg) if err != nil { @@ -234,6 +231,7 @@ func resolveAccessModes(modes []string) []corev1.PersistentVolumeAccessMode { } func buildStatefulSet(owner domain.Object, spec ResolvedStatefulSetSpec, ns string) *appsv1.StatefulSet { + labels.StampOrkestraLabels(spec.Labels, owner.GetName(), owner.GetAnnotations()) apiVersion := "" kind := "" if u, ok := owner.(*unstructured.Unstructured); ok { diff --git a/pkg/runtime/kordinator/worker.go b/pkg/runtime/kordinator/worker.go index 185ad9b4..b8a50ff9 100644 --- a/pkg/runtime/kordinator/worker.go +++ b/pkg/runtime/kordinator/worker.go @@ -136,7 +136,7 @@ func (k *Kontroller) processItemForGVK(ctx context.Context, gvk string, item que // The reconciler is never called when conditions are not met — gated state // is idle, not failure; error rate and health state are unaffected. if entry, ok := k.katalog.Get(gvk); ok { - if entry.CRD.PreReconcileCheck().HasReconcileGate() { + if entry.CRD.HasAnyReconcileGate() { obj := k.objectFromCache(entry, item.Key) if gated, reason := k.evaluatePreReconcileCheck(ctx, obj, entry.CRD.Name); gated { k.crdHealthMap[gvk].RecordGated(reason) diff --git a/pkg/runtime/reconciler/generic.go b/pkg/runtime/reconciler/generic.go index cdfc6826..28429d53 100644 --- a/pkg/runtime/reconciler/generic.go +++ b/pkg/runtime/reconciler/generic.go @@ -240,6 +240,20 @@ func NewGenericReconciler[PTR domain.Object]( var _ domain.Reconciler = (*GenericReconciler[domain.Object])(nil) +// effectiveBox returns the operatorBox that governs this reconcile cycle. +// It reads the serve-target annotation (alias > target > empty) and delegates +// to CRDEntry.EffectiveOperatorBox. Falls back to the CRD-level box when the +// CR has no target annotation (e.g. direct kubectl apply). +// The system CleanupFinalizer is always included in the returned box. +func (r *GenericReconciler[PTR]) effectiveBox(obj PTR) orktypes.OperatorBoxConfig { + target := orktypes.ResolveTargetFromAnnotations(obj.GetAnnotations()) + box := *r.crd.EffectiveOperatorBox(target) + if !slices.Contains(box.Finalizers, labels.CleanupFinalizer) { + box.Finalizers = append(box.Finalizers, labels.CleanupFinalizer) + } + return box +} + // Reconcile dispatches to the correct reconcile implementation. // Order: // 1. Conditional provisioning (when blocks) — handled by runTemplateReconcile @@ -293,6 +307,11 @@ func (r *GenericReconciler[PTR]) reconcileCore(ctx context.Context, key string) } rawObj := obj.DeepCopyObject().(PTR) + // Resolve the effective operatorBox for this CR. CRs routed through the gateway + // carry a serve-target annotation; the box for that target governs this cycle. + // Falls back to the CRD-level box for direct kubectl applies (no annotation). + box := r.effectiveBox(rawObj) + // Normalize before mutation/validation/template rendering ───────────── // Normalize + base resolver obj, resolver, normalizeChanges, err := r.applyNormalize(ctx, rawObj) @@ -306,7 +325,7 @@ func (r *GenericReconciler[PTR]) reconcileCore(ctx context.Context, key string) resolver = resolver.WithProfiles(r.kat.Profiles) } if r.kat != nil && !r.kat.Notes.IsEmpty() { - resolver = resolver.WithUserNotes(r.kat.Notes) + resolver = resolver.WithUserNotes(r.kat.UserNotes()) } // Inject raw serve intent as .request. so operatorBox templates, // mutation rules, and validation rules can all read the caller's vocabulary. @@ -366,7 +385,7 @@ func (r *GenericReconciler[PTR]) reconcileCore(ctx context.Context, key string) r.event.Eventf(obj, corev1.EventTypeNormal, "Deleting", fmt.Sprintf("Deleting %s %s/%s", r.crd.GVKString(), obj.GetNamespace(), obj.GetName())) - return r.handleDeletion(ctx, resolver, obj) + return r.handleDeletion(ctx, resolver, obj, box) } // Namespace guard — skip reconcile for CRs in restricted or non-allowed namespaces. @@ -389,14 +408,14 @@ func (r *GenericReconciler[PTR]) reconcileCore(ctx context.Context, key string) // If RemoveFinalizers is false (normal operation), ensure required finalizers exist. // If RemoveFinalizers is true (e.g., for testing or forced cleanup), remove them. if !r.crd.RemoveFinalizers { - if err := r.ensureFinalizers(ctx, obj); err != nil { + if err := r.ensureFinalizers(ctx, obj, box); err != nil { r.event.Eventf(obj, corev1.EventTypeWarning, r.crd.APITypes.Kind+"FinalizerError", fmt.Sprintf("Failed to add finalizers: %v", err)) return err } } else { logger.FromContext(ctx).Debug().Msgf("removing finalizers for %s", obj.GetName()) - if err := r.removeFinalizers(ctx, obj); err != nil { + if err := r.removeFinalizers(ctx, obj, box); err != nil { r.event.Eventf(obj, corev1.EventTypeWarning, r.crd.APITypes.Kind+"FinalizerRemovalError", fmt.Sprintf("Failed to remove finalizers: %v", err)) return err @@ -466,7 +485,12 @@ func (r *GenericReconciler[PTR]) reconcileCore(ctx context.Context, key string) } // ── Step 5: Reconcile implementation ────────────────────────────────────── - return r.reconcileImpl(ctx, resolver, obj) + if err := r.reconcileImpl(ctx, resolver, obj, box); err != nil { + return err + } + + // ── Step 6: Surface orphan cleanup ──────────────────────────────────────── + return r.cleanupPreviousSurface(ctx, rawObj) } // reconcileImpl dispatches to the correct reconcile implementation. @@ -479,7 +503,7 @@ func (r *GenericReconciler[PTR]) reconcileCore(ctx context.Context, key string) // 4. Reconcile dispatch // 5. Failure trigger check — record failure; trigger rollback if threshold met // 6. Status patch -func (r *GenericReconciler[PTR]) reconcileImpl(ctx context.Context, resolver *orktmpl.Resolver, obj PTR) error { +func (r *GenericReconciler[PTR]) reconcileImpl(ctx context.Context, resolver *orktmpl.Resolver, obj PTR, box orktypes.OperatorBoxConfig) error { var err error // ── Phase 1: Rollback gate ──────────────────────────────────────────────── @@ -523,7 +547,7 @@ func (r *GenericReconciler[PTR]) reconcileImpl(ctx context.Context, resolver *or var valErr error resolver, lastValResult, valErr = r.applyReconcileTimeValidation(ctx, resolver, obj) if valErr != nil { - r.patchStatusWithChildren(ctx, obj, resolver, valErr, lastValResult) + r.patchStatusWithChildren(ctx, obj, resolver, valErr, box, lastValResult) return valErr } } @@ -537,7 +561,7 @@ func (r *GenericReconciler[PTR]) reconcileImpl(ctx context.Context, resolver *or Msg("reconcile mutation failed — continuing") } } - hasTemplates := r.operatorBox.OnCreate != nil || r.operatorBox.OnReconcile != nil + hasTemplates := box.OnCreate != nil || box.OnReconcile != nil switch { case r.hooks.OnReconcile != nil: // Go hooks — user-provided, full type-safe access. @@ -546,20 +570,20 @@ func (r *GenericReconciler[PTR]) reconcileImpl(ctx context.Context, resolver *or // Order: by default declared templates run first (hybrid 90/10 pattern). // Set hooks.runHooksFirst: true in the Katalog to run the hook first. if !r.crd.RunHooksFirst() && hasTemplates { - resolver, err = r.runTemplateReconcile(ctx, resolver, obj) + resolver, err = r.runTemplateReconcile(ctx, resolver, obj, box) } if err == nil { err = r.hooks.OnReconcile(ctx, obj) } if err == nil && r.crd.RunHooksFirst() && hasTemplates { - resolver, err = r.runTemplateReconcile(ctx, resolver, obj) + resolver, err = r.runTemplateReconcile(ctx, resolver, obj, box) } - case r.operatorBox.OnCreate != nil || r.operatorBox.OnReconcile != nil: + case box.OnCreate != nil || box.OnReconcile != nil: // Declarative templates — interpreted at runtime. // Requires: nothing. ork generate registry NOT needed. // The returned resolver carries cross/external/git data for status evaluation. - resolver, err = r.runTemplateReconcile(ctx, resolver, obj) + resolver, err = r.runTemplateReconcile(ctx, resolver, obj, box) default: // No-op — finalizers, events, metrics still handled above. @@ -634,7 +658,7 @@ func (r *GenericReconciler[PTR]) reconcileImpl(ctx context.Context, resolver *or // Always patch status — best-effort, never fails reconcile. // Called with the outcome so Ready condition reflects reality. // Must run before the error return so Ready=False is written on failure. - r.patchStatusWithChildren(ctx, obj, resolver, err, lastValResult) + r.patchStatusWithChildren(ctx, obj, resolver, err, box, lastValResult) if err != nil { logger.FromContext(ctx).Error().Err(err). @@ -680,7 +704,7 @@ func (r *GenericReconciler[PTR]) namespaceAllowed( // handleDeletion runs cleanup then removes our finalizers. // Finalizers are never removed on error — object stays protected until // cleanup succeeds. -func (r *GenericReconciler[PTR]) handleDeletion(ctx context.Context, resolver *orktmpl.Resolver, obj PTR) error { +func (r *GenericReconciler[PTR]) handleDeletion(ctx context.Context, resolver *orktmpl.Resolver, obj PTR, box orktypes.OperatorBoxConfig) error { switch { case r.hooks.OnDelete != nil: if err := r.hooks.OnDelete(ctx, obj); err != nil { @@ -689,8 +713,8 @@ func (r *GenericReconciler[PTR]) handleDeletion(ctx context.Context, resolver *o return fmt.Errorf("deletion hook: %w", err) } - case r.operatorBox.OnDelete != nil: - if err := r.runTemplateOnDelete(ctx, resolver, obj); err != nil { + case box.OnDelete != nil: + if err := r.runTemplateOnDelete(ctx, resolver, obj, box); err != nil { r.event.Eventf(obj, corev1.EventTypeWarning, r.crd.APITypes.Kind+"DeleteError", fmt.Sprintf("Template deletion failed: %v", err)) return fmt.Errorf("template deletion: %w", err) @@ -701,15 +725,15 @@ func (r *GenericReconciler[PTR]) handleDeletion(ctx context.Context, resolver *o // onDelete block exists: the GC does not cascade through owner references on them. // runTemplateOnDelete already handles this when OnDelete is set; run it here for all // other cases (no onDelete block, or Go hook path). - if r.operatorBox.OnDelete == nil { + if box.OnDelete == nil { if kube, ok := kubeclient.FromContext(ctx); ok { - if err := runners.DeleteOwnedClusterScopedResources(ctx, kube, resolver, obj, r.operatorBox); err != nil { + if err := runners.DeleteOwnedClusterScopedResources(ctx, kube, resolver, obj, box); err != nil { return fmt.Errorf("namespace cleanup: %w", err) } } } - if err := r.removeFinalizers(ctx, obj); err != nil { + if err := r.removeFinalizers(ctx, obj, box); err != nil { r.event.Eventf(obj, corev1.EventTypeWarning, r.crd.APITypes.Kind+"FinalizerRemovalError", fmt.Sprintf("Failed to remove finalizers: %v", err)) return err diff --git a/pkg/runtime/reconciler/helper.go b/pkg/runtime/reconciler/helper.go index 094305f6..a3046067 100644 --- a/pkg/runtime/reconciler/helper.go +++ b/pkg/runtime/reconciler/helper.go @@ -7,6 +7,7 @@ import ( "github.com/orkspace/orkestra/domain" "github.com/orkspace/orkestra/pkg/logger" + orktypes "github.com/orkspace/orkestra/pkg/types" corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" @@ -71,18 +72,18 @@ func (r *GenericReconciler[PTR]) getLatestObject(ctx context.Context, namespace, // // Configured finalizers: ["protection.orkestra.io/finalizer"] // After calling ensureFinalizers, the resource's metadata.finalizers will include it. -func (r *GenericReconciler[PTR]) ensureFinalizers(ctx context.Context, obj PTR) error { - if len(r.operatorBox.Finalizers) == 0 { +func (r *GenericReconciler[PTR]) ensureFinalizers(ctx context.Context, obj PTR, box orktypes.OperatorBoxConfig) error { + if len(box.Finalizers) == 0 { return nil } logger.Debug(). Str("name", obj.GetName()). - Any("crd finalizers", r.operatorBox.Finalizers). + Any("crd finalizers", box.Finalizers). Msgf("checking finalizers: %v", obj.GetFinalizers()) needsUpdate := false - for _, f := range r.operatorBox.Finalizers { + for _, f := range box.Finalizers { if !ContainsFinalizer(obj, f) { needsUpdate = true break @@ -93,7 +94,7 @@ func (r *GenericReconciler[PTR]) ensureFinalizers(ctx context.Context, obj PTR) } newFinalizers := obj.GetFinalizers() - for _, f := range r.operatorBox.Finalizers { + for _, f := range box.Finalizers { if !ContainsFinalizer(obj, f) { newFinalizers = append(newFinalizers, f) } @@ -109,14 +110,14 @@ func (r *GenericReconciler[PTR]) ensureFinalizers(ctx context.Context, obj PTR) return r.kube.PatchFinalizers(ctx, obj, newFinalizers) } -func (r *GenericReconciler[PTR]) removeFinalizers(ctx context.Context, obj PTR) error { +func (r *GenericReconciler[PTR]) removeFinalizers(ctx context.Context, obj PTR, box orktypes.OperatorBoxConfig) error { if len(obj.GetFinalizers()) == 0 { return nil } newFinalizers := make([]string, 0, len(obj.GetFinalizers())) for _, f := range obj.GetFinalizers() { - if !slices.Contains(r.operatorBox.Finalizers, f) { + if !slices.Contains(box.Finalizers, f) { newFinalizers = append(newFinalizers, f) } } diff --git a/pkg/runtime/reconciler/run_status.go b/pkg/runtime/reconciler/run_status.go index dd0894ea..68e6d15a 100644 --- a/pkg/runtime/reconciler/run_status.go +++ b/pkg/runtime/reconciler/run_status.go @@ -51,6 +51,7 @@ func (r *GenericReconciler[PTR]) patchStatusWithChildren( obj PTR, resolver *orktmpl.Resolver, reconcileErr error, + box orktypes.OperatorBoxConfig, valResult ...*ValidationResult, ) { // ── Layer 3: extend resolver with child resource state ───────────────── @@ -58,7 +59,7 @@ func (r *GenericReconciler[PTR]) patchStatusWithChildren( // The new resolver's data map includes a "children" key so that // status field expressions can reference child status: // {{ .children.cronjob.status.lastScheduleTime }} - if reconcileErr == nil && (r.operatorBox.OnCreate != nil || r.operatorBox.OnReconcile != nil) { + if reconcileErr == nil && (box.OnCreate != nil || box.OnReconcile != nil) { children := children.ReadChildren(ctx, r.kube, obj, resolver, r.crd) resolver = resolver.WithChildren(children) // ← reassign — WithChildren returns new resolver } @@ -68,7 +69,7 @@ func (r *GenericReconciler[PTR]) patchStatusWithChildren( if len(valResult) > 0 { vr = valResult[0] } - if err := runStatusPatch(ctx, r, obj, resolver, reconcileErr, vr); err != nil { + if err := runStatusPatch(ctx, r, obj, resolver, reconcileErr, vr, box); err != nil { logger.FromContext(ctx).Warn().Err(err). Str("name", obj.GetName()). Msg("status: patch failed — continuing") @@ -89,6 +90,7 @@ func runStatusPatch[PTR domain.Object]( resolver *orktmpl.Resolver, reconcileErr error, valResult *ValidationResult, + box orktypes.OperatorBoxConfig, ) error { // ── Layer 1: Ready condition ─────────────────────────────────────────── // Always written — on success and failure — so operators can monitor @@ -120,13 +122,13 @@ func runStatusPatch[PTR domain.Object]( // Errors in field resolution are logged as warnings and do not fail the reconcile. logger.FromContext(ctx).Debug(). Str("name", obj.GetName()). - Bool("has_status_config", r.operatorBox.Status != nil && r.operatorBox.Status.HasFields()). + Bool("has_status_config", box.Status != nil && box.Status.HasFields()). Bool("reconcile_error", reconcileErr != nil). AnErr("reconcile_err", reconcileErr). Msg("status: layer2 evaluation") - if r.operatorBox.Status != nil && r.operatorBox.Status.HasFields() { - fields := r.operatorBox.Status.Fields + if box.Status != nil && box.Status.HasFields() { + fields := box.Status.Fields if reconcileErr != nil { var conditional []orktypes.StatusFieldSpec for _, f := range fields { diff --git a/pkg/runtime/reconciler/run_surface_cleanup.go b/pkg/runtime/reconciler/run_surface_cleanup.go new file mode 100644 index 00000000..7e5af5db --- /dev/null +++ b/pkg/runtime/reconciler/run_surface_cleanup.go @@ -0,0 +1,56 @@ +package reconciler + +import ( + "context" + + "github.com/orkspace/orkestra/pkg/labels" + "github.com/orkspace/orkestra/pkg/logger" + "github.com/orkspace/orkestra/pkg/runtime/runners" + orktypes "github.com/orkspace/orkestra/pkg/types" +) + +// cleanupPreviousSurface deletes all resources belonging to the surface the CR +// was on before the current reconcile cycle, then stamps AnnotationLastSurface +// with the new target so subsequent reconciles skip this step. +// +// Detection: compares AnnotationLastSurface (previous target) with the current +// effective target derived from serve-alias / serve-target annotations. A mismatch +// means the CR was re-routed via the gateway and the old surface must be cleaned up. +// +// Deletion: label-selector sweep over all known resource types for the previous +// owner key ("."). Template-based deletion is not used because +// forEach spec fields may have been cleared before this runs (e.g. spec.regions +// removed when switching away from a regional target). +func (r *GenericReconciler[PTR]) cleanupPreviousSurface( + ctx context.Context, + rawObj PTR, +) error { + target := orktypes.ResolveTargetFromAnnotations(rawObj.GetAnnotations()) + if target == "" || r.crd.KeepPreviousSurface(target) { + return nil + } + + prevTarget := rawObj.GetAnnotations()[labels.AnnotationLastSurface] + if prevTarget != "" && prevTarget != target { + prevOwnerKey := labels.EffectiveOwnerKey(rawObj.GetName(), map[string]string{ + labels.AnnotationServeAlias: prevTarget, + }) + ns := rawObj.GetNamespace() + if err := runners.SweepOwnedNamespacedResources(ctx, r.kube, prevOwnerKey, ns); err != nil { + return err + } + if err := runners.SweepOwnedClusterScopedResources(ctx, r.kube, prevOwnerKey); err != nil { + return err + } + } + + if err := r.kube.PatchAnnotations(ctx, rawObj, map[string]string{ + labels.AnnotationLastSurface: target, + }); err != nil { + logger.FromContext(ctx).Warn().Err(err). + Str("name", rawObj.GetName()). + Msg("surface cleanup: failed to update last-surface annotation") + } + + return nil +} diff --git a/pkg/runtime/reconciler/run_template_reconcile.go b/pkg/runtime/reconciler/run_template_reconcile.go index 3e5d382e..915d5f87 100644 --- a/pkg/runtime/reconciler/run_template_reconcile.go +++ b/pkg/runtime/reconciler/run_template_reconcile.go @@ -28,7 +28,7 @@ import ( // runTemplateReconcile interprets the Katalog's onCreate and onReconcile blocks. // Returns the enriched resolver so callers (reconcileImpl) can pass cross/external // data into patchStatusWithChildren for status field evaluation. -func (r *GenericReconciler[PTR]) runTemplateReconcile(ctx context.Context, resolver *orktmpl.Resolver, obj domain.Object) (*orktmpl.Resolver, error) { +func (r *GenericReconciler[PTR]) runTemplateReconcile(ctx context.Context, resolver *orktmpl.Resolver, obj domain.Object, box orktypes.OperatorBoxConfig) (*orktmpl.Resolver, error) { kube, ok := kubeclient.FromContext(ctx) if !ok { return resolver, fmt.Errorf("kubeclient not found in context") @@ -42,8 +42,8 @@ func (r *GenericReconciler[PTR]) runTemplateReconcile(ctx context.Context, resol // Step 2: cross-CRD observation // Reads from sibling CRD informer caches via r.katalogRegistry — zero API calls. // Must run first so git, docker, external calls, and resources can reference .cross.* - if len(r.operatorBox.Cross) > 0 { - crossData := r.readCross(ctx, obj, r.operatorBox.Cross, resolver) + if len(box.Cross) > 0 { + crossData := r.readCross(ctx, obj, box.Cross, resolver) logger.FromContext(ctx).Debug(). Str("observer", obj.GetName()). Int("cross_entries", len(crossData)). @@ -57,13 +57,13 @@ func (r *GenericReconciler[PTR]) runTemplateReconcile(ctx context.Context, resol // Step 3: Git hook // Runs before external calls so URLs, tokens, and payloads can reference .git.commit, // .git.changed, and .git.path. Git is a declarative precondition for pipelines. - if t := r.operatorBox.OnReconcile; t != nil && t.Git != nil { + if t := box.OnReconcile; t != nil && t.Git != nil { resolver, err = runGit(ctx, r.crd.GVKString(), resolver, kube, obj, r.crd.GVR(), t.Git) if err != nil { return resolver, fmt.Errorf("git hook: %w", err) } } - if t := r.operatorBox.OnCreate; t != nil && t.Git != nil { + if t := box.OnCreate; t != nil && t.Git != nil { resolver, err = runGit(ctx, r.crd.GVKString(), resolver, kube, obj, r.crd.GVR(), t.Git) if err != nil { return resolver, fmt.Errorf("git hook: %w", err) @@ -72,13 +72,13 @@ func (r *GenericReconciler[PTR]) runTemplateReconcile(ctx context.Context, resol // Step 4: external HTTP calls // Runs after Git so external URLs can embed commit hashes or paths. - if t := r.operatorBox.OnReconcile; t != nil && len(t.External) > 0 { + if t := box.OnReconcile; t != nil && len(t.External) > 0 { resolver, err = runExternal(ctx, r.crd.GVKString(), resolver, t.External, r.kube.Clientset()) if err != nil { return resolver, fmt.Errorf("external calls: %w", err) } } - if t := r.operatorBox.OnCreate; t != nil && len(t.External) > 0 { + if t := box.OnCreate; t != nil && len(t.External) > 0 { resolver, err = runExternal(ctx, r.crd.GVKString(), resolver, t.External, r.kube.Clientset()) if err != nil { return resolver, fmt.Errorf("external calls: %w", err) @@ -87,13 +87,13 @@ func (r *GenericReconciler[PTR]) runTemplateReconcile(ctx context.Context, resol // Step 5: Docker hook // Runs after external so build/push can use tokens or metadata from external calls. - if t := r.operatorBox.OnReconcile; t != nil && t.Docker != nil { + if t := box.OnReconcile; t != nil && t.Docker != nil { resolver, err = runDocker(ctx, r.crd.GVKString(), resolver, t.Docker) if err != nil { return resolver, fmt.Errorf("docker hook: %w", err) } } - if t := r.operatorBox.OnCreate; t != nil && t.Docker != nil { + if t := box.OnCreate; t != nil && t.Docker != nil { resolver, err = runDocker(ctx, r.crd.GVKString(), resolver, t.Docker) if err != nil { return resolver, fmt.Errorf("docker hook: %w", err) @@ -101,23 +101,23 @@ func (r *GenericReconciler[PTR]) runTemplateReconcile(ctx context.Context, resol } // Step 6: onCreate resource groups (update=false) - if t := r.operatorBox.OnCreate; t != nil { + if t := box.OnCreate; t != nil { if err := r.runResourceGroup(ctx, kube, resolver, obj, t, false); err != nil { return resolver, err } } // Step 7: onReconcile resource groups (update=true) - if t := r.operatorBox.OnReconcile; t != nil { + if t := box.OnReconcile; t != nil { if err := r.runResourceGroup(ctx, kube, resolver, obj, t, true); err != nil { return resolver, err } } // Step 8: provider dispatch - if len(r.operatorBox.ProviderBlocks) > 0 && r.providerRegistry != nil && r.providerRegistry.Len() > 0 { + if len(box.ProviderBlocks) > 0 && r.providerRegistry != nil && r.providerRegistry.Len() > 0 { kubeReader := &kubeReaderAdapter{kube: kube} - if err := runProviders(ctx, obj, resolver, r.operatorBox.ProviderBlocks, r.providerRegistry, kubeReader, r.providerStats); err != nil { + if err := runProviders(ctx, obj, resolver, box.ProviderBlocks, r.providerRegistry, kubeReader, r.providerStats); err != nil { return resolver, fmt.Errorf("providers: %w", err) } } @@ -247,7 +247,7 @@ func (r *GenericReconciler[PTR]) runResourceGroup( } // runTemplateOnDelete interprets the onDelete block. -func (r *GenericReconciler[PTR]) runTemplateOnDelete(ctx context.Context, resolver *orktmpl.Resolver, obj domain.Object) error { +func (r *GenericReconciler[PTR]) runTemplateOnDelete(ctx context.Context, resolver *orktmpl.Resolver, obj domain.Object, box orktypes.OperatorBoxConfig) error { kube, ok := kubeclient.FromContext(ctx) if !ok { return fmt.Errorf("kubeclient not found in context") @@ -255,7 +255,7 @@ func (r *GenericReconciler[PTR]) runTemplateOnDelete(ctx context.Context, resolv guard := r.namespaceGuardFunc(ctx, obj) - if t := r.operatorBox.OnDelete; t != nil { + if t := box.OnDelete; t != nil { if t.Ordered { if err := r.runOrderedDelete(ctx, kube, resolver, obj, t, guard); err != nil { return err @@ -268,16 +268,16 @@ func (r *GenericReconciler[PTR]) runTemplateOnDelete(ctx context.Context, resolv } } - if len(r.operatorBox.ProviderBlocks) > 0 && r.providerRegistry != nil { + if len(box.ProviderBlocks) > 0 && r.providerRegistry != nil { kubeReader := &kubeReaderAdapter{kube: kube} - if err := runProviderDelete(ctx, obj, resolver, r.operatorBox.ProviderBlocks, r.providerRegistry, kubeReader, r.providerStats); err != nil { + if err := runProviderDelete(ctx, obj, resolver, box.ProviderBlocks, r.providerRegistry, kubeReader, r.providerStats); err != nil { return fmt.Errorf("provider cleanup: %w", err) } } // Cluster-scoped resources cannot have namespace-scoped owners, so GC never cleans them up. // Always run explicit cleanup regardless of ordered/unordered path. - if err := runners.DeleteOwnedClusterScopedResources(ctx, kube, resolver, obj, r.operatorBox); err != nil { + if err := runners.DeleteOwnedClusterScopedResources(ctx, kube, resolver, obj, box); err != nil { return fmt.Errorf("cluster-scoped resource cleanup: %w", err) } diff --git a/pkg/runtime/runners/cluster_scoped_deletion.go b/pkg/runtime/runners/cluster_scoped_deletion.go index 68218c16..63c441fc 100644 --- a/pkg/runtime/runners/cluster_scoped_deletion.go +++ b/pkg/runtime/runners/cluster_scoped_deletion.go @@ -6,6 +6,7 @@ import ( "fmt" "github.com/orkspace/orkestra/domain" + "github.com/orkspace/orkestra/pkg/children" "github.com/orkspace/orkestra/pkg/kubeclient" orkcrb "github.com/orkspace/orkestra/pkg/resources/clusterrolebindings" orkcr "github.com/orkspace/orkestra/pkg/resources/clusterroles" @@ -32,35 +33,21 @@ func DeleteOwnedClusterScopedResources( obj domain.Object, box orktypes.OperatorBoxConfig, ) error { - // Delete Namespaces if err := deleteOwnedNamespaces(ctx, kube, resolver, obj, box); err != nil { return err } - - // Delete ClusterRoles if err := deleteOwnedClusterRoles(ctx, kube, resolver, obj, box); err != nil { return err } - - // Delete ClusterRoleBindings if err := deleteOwnedClusterRoleBindings(ctx, kube, resolver, obj, box); err != nil { return err } - - // Delete PersistentVolumes if err := deleteOwnedPersistentVolumes(ctx, kube, resolver, obj, box); err != nil { return err } - - // Delete Cluster-scoped Custom Resources - if err := deleteOwnedCustomResources(ctx, kube, resolver, obj, box); err != nil { - return err - } - - return nil + return deleteOwnedCustomResources(ctx, kube, resolver, obj, box) } -// deleteOwnedNamespaces explicitly deletes all Namespaces owned by this CR. func deleteOwnedNamespaces( ctx context.Context, kube kubeclient.Interface, @@ -69,19 +56,14 @@ func deleteOwnedNamespaces( box orktypes.OperatorBoxConfig, ) error { var srcs []orktypes.NamespaceTemplateSource - if box.OnCreate != nil { - srcs = append(srcs, box.OnCreate.Namespaces...) - } - if box.OnReconcile != nil { - srcs = append(srcs, box.OnReconcile.Namespaces...) - } - if box.OnDelete != nil { - srcs = append(srcs, box.OnDelete.Namespaces...) + for _, hook := range allHooks(box) { + if hook != nil { + srcs = append(srcs, hook.Namespaces...) + } } - - for i, src := range srcs { - name, err := resolver.Resolve(src.Name) - if err != nil || name == "" { + for i, src := range children.ExpandForEachNamespaces(resolver, srcs) { + name, _ := resolver.Resolve(src.Name) + if name == "" { continue } if err := orkns.DeleteIfOwned(ctx, kube, obj, name); err != nil { @@ -91,7 +73,6 @@ func deleteOwnedNamespaces( return nil } -// deleteOwnedClusterRoles explicitly deletes all ClusterRoles owned by this CR. func deleteOwnedClusterRoles( ctx context.Context, kube kubeclient.Interface, @@ -100,19 +81,14 @@ func deleteOwnedClusterRoles( box orktypes.OperatorBoxConfig, ) error { var srcs []orktypes.ClusterRoleTemplateSource - if box.OnCreate != nil { - srcs = append(srcs, box.OnCreate.ClusterRoles...) - } - if box.OnReconcile != nil { - srcs = append(srcs, box.OnReconcile.ClusterRoles...) - } - if box.OnDelete != nil { - srcs = append(srcs, box.OnDelete.ClusterRoles...) + for _, hook := range allHooks(box) { + if hook != nil { + srcs = append(srcs, hook.ClusterRoles...) + } } - - for i, src := range srcs { - name, err := resolver.Resolve(src.Name) - if err != nil || name == "" { + for i, src := range children.ExpandForEachClusterRoles(resolver, srcs) { + name, _ := resolver.Resolve(src.Name) + if name == "" { continue } if err := orkcr.DeleteIfOwned(ctx, kube, obj, name); err != nil { @@ -122,7 +98,6 @@ func deleteOwnedClusterRoles( return nil } -// deleteOwnedClusterRoleBindings explicitly deletes all ClusterRoleBindings owned by this CR. func deleteOwnedClusterRoleBindings( ctx context.Context, kube kubeclient.Interface, @@ -131,19 +106,14 @@ func deleteOwnedClusterRoleBindings( box orktypes.OperatorBoxConfig, ) error { var srcs []orktypes.ClusterRoleBindingTemplateSource - if box.OnCreate != nil { - srcs = append(srcs, box.OnCreate.ClusterRoleBindings...) - } - if box.OnReconcile != nil { - srcs = append(srcs, box.OnReconcile.ClusterRoleBindings...) - } - if box.OnDelete != nil { - srcs = append(srcs, box.OnDelete.ClusterRoleBindings...) + for _, hook := range allHooks(box) { + if hook != nil { + srcs = append(srcs, hook.ClusterRoleBindings...) + } } - - for i, src := range srcs { - name, err := resolver.Resolve(src.Name) - if err != nil || name == "" { + for i, src := range children.ExpandForEachClusterRoleBindings(resolver, srcs) { + name, _ := resolver.Resolve(src.Name) + if name == "" { continue } if err := orkcrb.DeleteIfOwned(ctx, kube, obj, name); err != nil { @@ -153,7 +123,6 @@ func deleteOwnedClusterRoleBindings( return nil } -// deleteOwnedPersistentVolumes explicitly deletes all PersistentVolumes owned by this CR. func deleteOwnedPersistentVolumes( ctx context.Context, kube kubeclient.Interface, @@ -162,19 +131,14 @@ func deleteOwnedPersistentVolumes( box orktypes.OperatorBoxConfig, ) error { var srcs []orktypes.PVTemplateSource - if box.OnCreate != nil { - srcs = append(srcs, box.OnCreate.PersistentVolumes...) - } - if box.OnReconcile != nil { - srcs = append(srcs, box.OnReconcile.PersistentVolumes...) - } - if box.OnDelete != nil { - srcs = append(srcs, box.OnDelete.PersistentVolumes...) + for _, hook := range allHooks(box) { + if hook != nil { + srcs = append(srcs, hook.PersistentVolumes...) + } } - - for i, src := range srcs { - name, err := resolver.Resolve(src.Name) - if err != nil || name == "" { + for i, src := range children.ExpandForEachPVs(resolver, srcs) { + name, _ := resolver.Resolve(src.Name) + if name == "" { continue } if err := orkpv.DeleteIfOwned(ctx, kube, obj, name); err != nil { @@ -184,7 +148,6 @@ func deleteOwnedPersistentVolumes( return nil } -// deleteOwnedCustomResources explicitly deletes all Cluster-scoped custom resources owned by this CR. func deleteOwnedCustomResources( ctx context.Context, kube kubeclient.Interface, @@ -193,31 +156,27 @@ func deleteOwnedCustomResources( box orktypes.OperatorBoxConfig, ) error { var srcs []orktypes.CustomResourceTemplateSource - if box.OnCreate != nil { - srcs = append(srcs, box.OnCreate.CustomResource...) - } - if box.OnReconcile != nil { - srcs = append(srcs, box.OnReconcile.CustomResource...) - } - if box.OnDelete != nil { - srcs = append(srcs, box.OnDelete.CustomResource...) + for _, hook := range allHooks(box) { + if hook != nil { + srcs = append(srcs, hook.CustomResource...) + } } - - for i, src := range srcs { - // Only process cluster-scoped CRs + for i, src := range children.ExpandForEachCustomResources(resolver, srcs) { if src.IsNamespaced() { - continue // GC handles these + continue // GC handles namespace-scoped custom resources } - - name, err := resolver.Resolve(src.Metadata.Name) - if err != nil || name == "" { + name, _ := resolver.Resolve(src.Metadata.Name) + if name == "" { continue } - - // Cluster-scoped custom resources need explicit cleanup if err := orkcust.DeleteIfOwned(ctx, kube, obj, name, "", src.APIVersion, src.Kind); err != nil { return fmt.Errorf("customresource[%d] %q: %w", i, name, err) } } return nil } + +// allHooks returns all three lifecycle hook blocks in declaration order. +func allHooks(box orktypes.OperatorBoxConfig) []*orktypes.HookTemplates { + return []*orktypes.HookTemplates{box.OnCreate, box.OnReconcile, box.OnDelete} +} diff --git a/pkg/runtime/runners/docs/02-garbage-collection.md b/pkg/runtime/runners/docs/02-garbage-collection.md index c8836700..b4f861b7 100644 --- a/pkg/runtime/runners/docs/02-garbage-collection.md +++ b/pkg/runtime/runners/docs/02-garbage-collection.md @@ -9,20 +9,28 @@ Kubernetes GC does **not** cascade owner references from namespace-scoped resour | Namespace-scoped → | Namespace-scoped | ✅ | | Namespace-scoped → | Cluster-scoped | ❌ | -When a ReconcilerProbe CR (namespace-scoped) creates cluster-scoped resources (Namespaces, ClusterRoles, ClusterRoleBindings, PersistentVolumes), they become orphaned on CR deletion. +When a CR creates cluster-scoped resources (Namespaces, ClusterRoles, ClusterRoleBindings, PersistentVolumes), they orphan on CR deletion. Explicit deletion is required. -## The Solution +## Two deletion paths -[`cluster_scoped_deletion.go`](../cluster_scoped_deletion.go) provides `DeleteOwnedClusterScopedResources()` — called during CR deletion to explicitly clean up owned cluster-scoped resources. +| Path | Mechanism | File | +|------|-----------|------| +| **CR deletion** | Template-based — resolves names from the box declaration, expands `forEach` | [`cluster_scoped_deletion.go`](../cluster_scoped_deletion.go) | +| **Surface switch** | Label-selector sweep — finds resources by `orkestra-owner=` | [`surface_sweep.go`](../surface_sweep.go) | -### How it works +### Why two mechanisms? -1. **Collect sources** from `OnCreate`, `OnReconcile`, `OnDelete` -2. **Resolve names** through the template resolver -3. **Check ownership** — only delete if the CR is the owner -4. **Delete** the resource +Template-based deletion works at CR deletion time because the CR's spec is intact — names resolve, `forEach` expands correctly. -### Currently handled resources +For surface switches the spec may have already changed before cleanup runs (e.g. `spec.regions` is cleared when routing away from a `regional` target). Template-based deletion would expand `forEach` to nothing and silently miss the orphans. Label-selector sweep is immune to spec changes. + +Namespace-scoped resources on the CR deletion path do not need explicit cleanup — Kubernetes GC handles them via owner references. Only cluster-scoped resources require `DeleteOwnedClusterScopedResources`. + +## CR deletion — `DeleteOwnedClusterScopedResources` + +Covers `onCreate`, `onReconcile`, and `onDelete` blocks. Called from the reconciler deletion path. + +### Currently handled resource types | Resource | |----------| @@ -30,60 +38,54 @@ When a ReconcilerProbe CR (namespace-scoped) creates cluster-scoped resources (N | ClusterRole | | ClusterRoleBinding | | PersistentVolume | -| Custom Resources | +| Custom Resources (cluster-scoped) | + +### Adding a new cluster-scoped resource type + +1. Add `DeleteIfOwned` to the resource package (`pkg/resources//`) +2. Add a `deleteOwned` function to `cluster_scoped_deletion.go` following the existing pattern — collect sources from `allHooks(box)`, call `ExpandForEach*`, resolve name via `resolver.Resolve`, call `DeleteIfOwned` +3. Call it from `DeleteOwnedClusterScopedResources` +4. Update the table above + +## Surface switch — `SweepOwned*` -## Adding a new cluster-scoped resource +`SweepOwnedNamespacedResources` and `SweepOwnedClusterScopedResources` list all resources of each known type with `orkestra-owner=` and delete them. -1. **Add `DeleteIfOwned`** to the resource package (`pkg/resources//`) -2. **Add collection function** in `cluster_scoped_deletion.go` following the existing pattern -3. **Call it** from `DeleteOwnedClusterScopedResources()` -4. **Update the list** above +### Adding a new resource type to the sweep -### Pattern to follow +Add a list+delete block to the appropriate sweep function in `surface_sweep.go`: ```go -func deleteOwned( - ctx context.Context, - kube kubeclient.KubeClient, - resolver *orktmpl.Resolver, - obj domain.Object, - box orktypes.OperatorBoxConfig, -) error { - var srcs []orktypes.TemplateSource - if box.OnCreate != nil { - srcs = append(srcs, box.OnCreate....) - } - if box.OnReconcile != nil { - srcs = append(srcs, box.OnReconcile....) - } - if box.OnDelete != nil { - srcs = append(srcs, box.OnDelete....) +// namespaced +if list, err := cs.AppsV1().(ns).List(ctx, opts); err == nil { + for _, r := range list.Items { + collect("/"+r.Name, cs.AppsV1().(ns).Delete(ctx, r.Name, dopts)) } +} - for i, src := range srcs { - name, err := resolver.Resolve(src.Name) - if err != nil || name == "" { - continue - } - if err := ork.DeleteIfOwned(ctx, kube, obj, name); err != nil { - return fmt.Errorf("[%d] %q: %w", i, name, err) - } +// cluster-scoped +if list, err := cs.().().List(ctx, opts); err == nil { + for _, r := range list.Items { + collect("/"+r.Name, cs.().().Delete(ctx, r.Name, dopts)) } - return nil } ``` -## Integration point - -In the reconciler deletion path: +## Integration points +**CR deletion** (`handleDeletion`): ```go -// Cluster-scoped resources require explicit deletion: GC does not cascade owner references to them. -if err := runners.DeleteOwnedClusterScopedResources(ctx, kube, resolver, obj, r.operatorBox); err != nil { - return fmt.Errorf("cluster-scoped resource cleanup: %w", err) +if err := runners.DeleteOwnedClusterScopedResources(ctx, kube, resolver, obj, box); err != nil { + return err } ``` +**Surface switch** (`cleanupPreviousSurface`): +```go +runners.SweepOwnedNamespacedResources(ctx, kube, prevOwnerKey, ns) +runners.SweepOwnedClusterScopedResources(ctx, kube, prevOwnerKey) +``` + --- → Back: [README](../README.md) diff --git a/pkg/runtime/runners/surface_sweep.go b/pkg/runtime/runners/surface_sweep.go new file mode 100644 index 00000000..c17345bb --- /dev/null +++ b/pkg/runtime/runners/surface_sweep.go @@ -0,0 +1,197 @@ +// pkg/runners/surface_sweep.go +package runners + +import ( + "context" + "fmt" + + "github.com/orkspace/orkestra/pkg/kubeclient" + "github.com/orkspace/orkestra/pkg/labels" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +// SweepOwnedNamespacedResources deletes all namespace-scoped resources in ns +// whose orkestra-owner label matches ownerKey. Used for surface-switch cleanup +// where template-based deletion cannot work: when a spec field that drove a +// forEach declaration is removed before cleanup runs (e.g. spec.regions is +// cleared when switching away from a regional target), ExpandForEach produces +// nothing and the orphaned resources are never found by name. +// +// This is a label-selector sweep — it finds resources by ownership label +// rather than by resolving template names, so it is immune to spec changes. +func SweepOwnedNamespacedResources( + ctx context.Context, + kube kubeclient.Interface, + ownerKey string, + ns string, +) error { + cs := kube.Clientset() + sel := labels.OrkestraOwner + "=" + ownerKey + opts := metav1.ListOptions{LabelSelector: sel} + dopts := metav1.DeleteOptions{} + + var errs []error + collect := func(label string, err error) { + if err != nil { + errs = append(errs, fmt.Errorf("%s: %w", label, err)) + } + } + + // ── Workloads ───────────────────────────────────────────────────────────── + if list, err := cs.AppsV1().Deployments(ns).List(ctx, opts); err == nil { + for _, r := range list.Items { + collect("deployment/"+r.Name, cs.AppsV1().Deployments(ns).Delete(ctx, r.Name, dopts)) + } + } + if list, err := cs.AppsV1().ReplicaSets(ns).List(ctx, opts); err == nil { + for _, r := range list.Items { + collect("replicaset/"+r.Name, cs.AppsV1().ReplicaSets(ns).Delete(ctx, r.Name, dopts)) + } + } + if list, err := cs.AppsV1().StatefulSets(ns).List(ctx, opts); err == nil { + for _, r := range list.Items { + collect("statefulset/"+r.Name, cs.AppsV1().StatefulSets(ns).Delete(ctx, r.Name, dopts)) + } + } + + // ── Networking ──────────────────────────────────────────────────────────── + if list, err := cs.CoreV1().Services(ns).List(ctx, opts); err == nil { + for _, r := range list.Items { + collect("service/"+r.Name, cs.CoreV1().Services(ns).Delete(ctx, r.Name, dopts)) + } + } + if list, err := cs.NetworkingV1().Ingresses(ns).List(ctx, opts); err == nil { + for _, r := range list.Items { + collect("ingress/"+r.Name, cs.NetworkingV1().Ingresses(ns).Delete(ctx, r.Name, dopts)) + } + } + if list, err := cs.NetworkingV1().NetworkPolicies(ns).List(ctx, opts); err == nil { + for _, r := range list.Items { + collect("networkpolicy/"+r.Name, cs.NetworkingV1().NetworkPolicies(ns).Delete(ctx, r.Name, dopts)) + } + } + + // ── Config ──────────────────────────────────────────────────────────────── + if list, err := cs.CoreV1().ConfigMaps(ns).List(ctx, opts); err == nil { + for _, r := range list.Items { + collect("configmap/"+r.Name, cs.CoreV1().ConfigMaps(ns).Delete(ctx, r.Name, dopts)) + } + } + if list, err := cs.CoreV1().Secrets(ns).List(ctx, opts); err == nil { + for _, r := range list.Items { + collect("secret/"+r.Name, cs.CoreV1().Secrets(ns).Delete(ctx, r.Name, dopts)) + } + } + + // ── Storage ─────────────────────────────────────────────────────────────── + if list, err := cs.CoreV1().PersistentVolumeClaims(ns).List(ctx, opts); err == nil { + for _, r := range list.Items { + collect("pvc/"+r.Name, cs.CoreV1().PersistentVolumeClaims(ns).Delete(ctx, r.Name, dopts)) + } + } + + // ── Autoscaling / policy ────────────────────────────────────────────────── + if list, err := cs.AutoscalingV2().HorizontalPodAutoscalers(ns).List(ctx, opts); err == nil { + for _, r := range list.Items { + collect("hpa/"+r.Name, cs.AutoscalingV2().HorizontalPodAutoscalers(ns).Delete(ctx, r.Name, dopts)) + } + } + if list, err := cs.PolicyV1().PodDisruptionBudgets(ns).List(ctx, opts); err == nil { + for _, r := range list.Items { + collect("pdb/"+r.Name, cs.PolicyV1().PodDisruptionBudgets(ns).Delete(ctx, r.Name, dopts)) + } + } + + // ── RBAC ───────────────────────────────────────────────────────────────── + if list, err := cs.CoreV1().ServiceAccounts(ns).List(ctx, opts); err == nil { + for _, r := range list.Items { + collect("serviceaccount/"+r.Name, cs.CoreV1().ServiceAccounts(ns).Delete(ctx, r.Name, dopts)) + } + } + if list, err := cs.RbacV1().Roles(ns).List(ctx, opts); err == nil { + for _, r := range list.Items { + collect("role/"+r.Name, cs.RbacV1().Roles(ns).Delete(ctx, r.Name, dopts)) + } + } + if list, err := cs.RbacV1().RoleBindings(ns).List(ctx, opts); err == nil { + for _, r := range list.Items { + collect("rolebinding/"+r.Name, cs.RbacV1().RoleBindings(ns).Delete(ctx, r.Name, dopts)) + } + } + + // ── Quota / limits ──────────────────────────────────────────────────────── + if list, err := cs.CoreV1().ResourceQuotas(ns).List(ctx, opts); err == nil { + for _, r := range list.Items { + collect("resourcequota/"+r.Name, cs.CoreV1().ResourceQuotas(ns).Delete(ctx, r.Name, dopts)) + } + } + if list, err := cs.CoreV1().LimitRanges(ns).List(ctx, opts); err == nil { + for _, r := range list.Items { + collect("limitrange/"+r.Name, cs.CoreV1().LimitRanges(ns).Delete(ctx, r.Name, dopts)) + } + } + + // ── Pods / CronJobs ─────────────────────────────────────────────────────── + if list, err := cs.CoreV1().Pods(ns).List(ctx, opts); err == nil { + for _, r := range list.Items { + collect("pod/"+r.Name, cs.CoreV1().Pods(ns).Delete(ctx, r.Name, dopts)) + } + } + if list, err := cs.BatchV1().CronJobs(ns).List(ctx, opts); err == nil { + for _, r := range list.Items { + collect("cronjob/"+r.Name, cs.BatchV1().CronJobs(ns).Delete(ctx, r.Name, dopts)) + } + } + + if len(errs) > 0 { + return fmt.Errorf("surface sweep %q: %v", ownerKey, errs) + } + return nil +} + +// SweepOwnedClusterScopedResources deletes all cluster-scoped resources whose +// orkestra-owner label matches ownerKey. Mirrors SweepOwnedNamespacedResources +// but operates on cluster-scoped types that Kubernetes GC does not cascade. +func SweepOwnedClusterScopedResources( + ctx context.Context, + kube kubeclient.Interface, + ownerKey string, +) error { + cs := kube.Clientset() + sel := labels.OrkestraOwner + "=" + ownerKey + opts := metav1.ListOptions{LabelSelector: sel} + dopts := metav1.DeleteOptions{} + + var errs []error + collect := func(label string, err error) { + if err != nil { + errs = append(errs, fmt.Errorf("%s: %w", label, err)) + } + } + + if list, err := cs.CoreV1().Namespaces().List(ctx, opts); err == nil { + for _, r := range list.Items { + collect("namespace/"+r.Name, cs.CoreV1().Namespaces().Delete(ctx, r.Name, dopts)) + } + } + if list, err := cs.RbacV1().ClusterRoles().List(ctx, opts); err == nil { + for _, r := range list.Items { + collect("clusterrole/"+r.Name, cs.RbacV1().ClusterRoles().Delete(ctx, r.Name, dopts)) + } + } + if list, err := cs.RbacV1().ClusterRoleBindings().List(ctx, opts); err == nil { + for _, r := range list.Items { + collect("clusterrolebinding/"+r.Name, cs.RbacV1().ClusterRoleBindings().Delete(ctx, r.Name, dopts)) + } + } + if list, err := cs.CoreV1().PersistentVolumes().List(ctx, opts); err == nil { + for _, r := range list.Items { + collect("pv/"+r.Name, cs.CoreV1().PersistentVolumes().Delete(ctx, r.Name, dopts)) + } + } + + if len(errs) > 0 { + return fmt.Errorf("cluster-scoped surface sweep %q: %v", ownerKey, errs) + } + return nil +} diff --git a/pkg/types/methods.go b/pkg/types/methods.go index ed437104..a161057e 100644 --- a/pkg/types/methods.go +++ b/pkg/types/methods.go @@ -23,6 +23,44 @@ func (e CRDEntry) PreReconcileCheck() *PreReconcileConfig { return e.OperatorBox.PreReconcile } +// HasAnyEnqueueGate reports whether the CRD-level or any per-target operatorBox +// declares an enqueueGate. Used at startup to decide whether to register the +// informer enqueue filter — must register if ANY surface can gate enqueueing. +func (e CRDEntry) HasAnyEnqueueGate() bool { + if rc := e.OperatorBox.PreReconcile; rc != nil && rc.HasEnqueueGate() { + return true + } + if e.Serve != nil { + for _, cfg := range e.Serve.Target.Entries { + if cfg.OperatorBox != nil { + if rc := cfg.OperatorBox.PreReconcile; rc != nil && rc.HasEnqueueGate() { + return true + } + } + } + } + return false +} + +// HasAnyReconcileGate reports whether the CRD-level or any per-target operatorBox +// declares a reconcileGate. Used at dequeue time to decide whether to evaluate +// the gate before calling the reconciler. +func (e CRDEntry) HasAnyReconcileGate() bool { + if rc := e.OperatorBox.PreReconcile; rc != nil && rc.HasReconcileGate() { + return true + } + if e.Serve != nil { + for _, cfg := range e.Serve.Target.Entries { + if cfg.OperatorBox != nil { + if rc := cfg.OperatorBox.PreReconcile; rc != nil && rc.HasReconcileGate() { + return true + } + } + } + } + return false +} + // IsBuiltInType reports whether this CRD represents a built‑in Kubernetes resource. // Built‑ins rely on enrichment to populate group, version, plural, and scope. func (c *CRDEntry) IsBuiltInType() bool { diff --git a/pkg/types/types_crd_entry.go b/pkg/types/types_crd_entry.go index ac103f7b..b019198f 100644 --- a/pkg/types/types_crd_entry.go +++ b/pkg/types/types_crd_entry.go @@ -5,7 +5,6 @@ import ( "sort" "github.com/orkspace/orkestra/pkg/labels" - "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/runtime/schema" ) @@ -280,13 +279,27 @@ type CRDEntry struct { } // EffectiveOperatorBox returns the operatorBox for a given target. +// When a target entry declares its own operatorBox, non-template fields +// (preReconcile, status) fall back to the CRD-level values when absent +// on the target — so a CRD-level gate or status config applies to all +// surfaces unless a target explicitly overrides it. func (c *CRDEntry) EffectiveOperatorBox(target string) *OperatorBoxConfig { if target == "" { return &c.OperatorBox } if c.Serve != nil && c.Serve.Target.Entries != nil { if cfg, ok := c.Serve.Target.Entries[target]; ok && cfg.OperatorBox != nil { - return cfg.OperatorBox + box := *cfg.OperatorBox + if box.PreReconcile == nil { + box.PreReconcile = c.OperatorBox.PreReconcile + } + if box.Status == nil { + box.Status = c.OperatorBox.Status + } + if box.Reconciler == nil { + box.Reconciler = c.OperatorBox.Reconciler + } + return &box } } return &c.OperatorBox @@ -297,8 +310,7 @@ func (c *CRDEntry) EffectiveOperatorBox(target string) *OperatorBoxConfig { // 1. serve-alias annotation (most specific) // 2. serve-target annotation (primary target) // 3. Empty string (no target found) -func ResolveTargetFromAnnotations(obj *unstructured.Unstructured) string { - annotations := obj.GetAnnotations() +func ResolveTargetFromAnnotations(annotations map[string]string) string { if annotations == nil { return "" } diff --git a/pkg/types/types_serve_apply_config.go b/pkg/types/types_serve_apply_config.go index 04546640..3a633e8e 100644 --- a/pkg/types/types_serve_apply_config.go +++ b/pkg/types/types_serve_apply_config.go @@ -23,6 +23,12 @@ type ServeApplyOverrides struct { // Can be overridden per-request with ?overwrite=true regardless of this setting. // Default: false. ResourceConflict *bool `yaml:"resourceConflict,omitempty" json:"resourceConflict,omitempty"` + + // KeepPreviousSurface, when true, leaves resources created by a prior serve + // target alive when the CR switches to a new target. By default (false) the + // reconciler deletes orphaned child resources whose serve-target label no + // longer matches the CR's active surface. + KeepPreviousSurface *bool `yaml:"keepPreviousSurface,omitempty" json:"keepPreviousSurface,omitempty"` } // ServeApplyConfig configures apply-time behaviour for a CRD or target. @@ -39,7 +45,8 @@ func (s *ServeConfig) HasOverride() bool { if s == nil || s.Apply == nil || s.Apply.Overrides == nil { return false } - return s.Apply.Overrides.TargetConflict != nil || s.Apply.Overrides.ResourceConflict != nil + o := s.Apply.Overrides + return o.TargetConflict != nil || o.ResourceConflict != nil || o.KeepPreviousSurface != nil } // HasOverride reports whether the ServeTargetConfig has any override fields set. @@ -47,7 +54,8 @@ func (t *ServeTargetConfig) HasOverride() bool { if t == nil || t.Apply == nil || t.Apply.Overrides == nil { return false } - return t.Apply.Overrides.TargetConflict != nil || t.Apply.Overrides.ResourceConflict != nil + o := t.Apply.Overrides + return o.TargetConflict != nil || o.ResourceConflict != nil || o.KeepPreviousSurface != nil } // HasOverride reports whether the CRDEntry has any override fields set @@ -118,6 +126,26 @@ func (c *CRDEntry) HasResourceConflict() bool { return false } +// KeepPreviousSurface reports whether the CRDEntry or the given target +// has keepPreviousSurface enabled. CRD-level wins if set; otherwise the +// target entry is checked. +func (c *CRDEntry) KeepPreviousSurface(target string) bool { + if !c.ServeEnabled() { + return false + } + if c.Serve.Apply != nil && c.Serve.Apply.Overrides != nil && c.Serve.Apply.Overrides.KeepPreviousSurface != nil { + return *c.Serve.Apply.Overrides.KeepPreviousSurface + } + if c.Serve.Target.Entries != nil { + if cfg, ok := c.Serve.Target.Entries[target]; ok { + if cfg.Apply != nil && cfg.Apply.Overrides != nil && cfg.Apply.Overrides.KeepPreviousSurface != nil { + return *cfg.Apply.Overrides.KeepPreviousSurface + } + } + } + return false +} + // EffectiveServeTargetForCR returns the effective target for a given CR. // Resolution order: // 1. If the CR matches a target via fieldSelector, use that target. diff --git a/pkg/types/types_simulate.go b/pkg/types/types_simulate.go index 0ab26519..7d053316 100644 --- a/pkg/types/types_simulate.go +++ b/pkg/types/types_simulate.go @@ -29,6 +29,7 @@ type SimulateSpec struct { CRDFiles []string `yaml:"crdFiles,omitempty"` Cycles int `yaml:"cycles,omitempty"` // default 10 when unset SkipExternal bool `yaml:"skipExternal,omitempty"` + Target string `yaml:"target,omitempty"` // serve target to use for reconciliation Expect *SimulateExpect `yaml:"expect,omitempty"` // nil = op-print only }