diff --git a/internal/declarative/planner/event_gateway_cluster_policy_planner.go b/internal/declarative/planner/event_gateway_cluster_policy_planner.go index 78062be4c..8e6d076af 100644 --- a/internal/declarative/planner/event_gateway_cluster_policy_planner.go +++ b/internal/declarative/planner/event_gateway_cluster_policy_planner.go @@ -510,24 +510,31 @@ func (p *Planner) extractClusterPolicyConfig( func configFieldsMatch(current, desired map[string]any) bool { for key, desiredVal := range desired { currentVal, exists := current[key] - if !exists { + if !exists || !configValuesMatch(currentVal, desiredVal) { return false } + } + return true +} - // Recursive comparison for nested maps - desiredMap, desiredIsMap := desiredVal.(map[string]any) - currentMap, currentIsMap := currentVal.(map[string]any) - if desiredIsMap && currentIsMap { - if !configFieldsMatch(currentMap, desiredMap) { +func configValuesMatch(current, desired any) bool { + switch desired := desired.(type) { + case map[string]any: + current, ok := current.(map[string]any) + return ok && configFieldsMatch(current, desired) + case []any: + current, ok := current.([]any) + if !ok || len(current) != len(desired) || (current == nil) != (desired == nil) { + return false + } + // Array order matters for policy rules. + for i, desiredElement := range desired { + if !configValuesMatch(current[i], desiredElement) { return false } - continue - } - - // For slices, use DeepEqual (order matters for rules) - if !reflect.DeepEqual(currentVal, desiredVal) { - return false } + return true + default: + return reflect.DeepEqual(current, desired) } - return true } diff --git a/internal/declarative/planner/event_gateway_cluster_policy_planner_test.go b/internal/declarative/planner/event_gateway_cluster_policy_planner_test.go index 6b6ef266e..86798b595 100644 --- a/internal/declarative/planner/event_gateway_cluster_policy_planner_test.go +++ b/internal/declarative/planner/event_gateway_cluster_policy_planner_test.go @@ -210,3 +210,43 @@ func TestShouldUpdateClusterPolicy_ConfigChangedNestedField(t *testing.T) { assert.NotNil(t, updateFields, "updateFields should contain the new config") require.Contains(t, changedFields, "config", "config should be in changed fields") } + +func TestConfigFieldsMatchArrays(t *testing.T) { + for _, tt := range []struct { + name string + current any + desired any + want bool + }{ + { + name: "API adds fields inside array objects", + current: []any{map[string]any{FieldID: "key-id", FieldName: "key-name"}}, + desired: []any{map[string]any{FieldID: "key-id"}}, + want: true, + }, + { + name: "changed nested field", + current: []any{map[string]any{FieldID: "old-key"}}, + desired: []any{map[string]any{FieldID: "new-key"}}, + }, + { + name: "missing nested field", + current: []any{map[string]any{FieldName: "key-name"}}, + desired: []any{map[string]any{FieldID: "key-id"}}, + }, + {name: "order matters", current: []any{"a", "b"}, desired: []any{"b", "a"}}, + {name: "extra element", current: []any{"a", "b"}, desired: []any{"a"}}, + {name: "missing element", current: []any{"a"}, desired: []any{"a", "b"}}, + {name: "different type", current: "a", desired: []any{"a"}}, + {name: "nil versus empty", current: []any(nil), desired: []any{}}, + {name: "equal scalar elements", current: []any{"a", true}, desired: []any{"a", true}, want: true}, + {name: "nested arrays", current: []any{[]any{"a"}}, desired: []any{[]any{"a"}}, want: true}, + } { + t.Run(tt.name, func(t *testing.T) { + assert.Equal(t, tt.want, configFieldsMatch( + map[string]any{FieldConfig: tt.current}, + map[string]any{FieldConfig: tt.desired}, + )) + }) + } +} diff --git a/internal/declarative/planner/event_gateway_produce_policy_planner_test.go b/internal/declarative/planner/event_gateway_produce_policy_planner_test.go index 8528681c9..76c532229 100644 --- a/internal/declarative/planner/event_gateway_produce_policy_planner_test.go +++ b/internal/declarative/planner/event_gateway_produce_policy_planner_test.go @@ -529,3 +529,73 @@ func TestPrepareProducePolicyParentRefsResolvesModifyHeadersParent(t *testing.T) *prepared[1].EventGatewayModifyHeadersPolicyCreate.ParentPolicyID, ) } + +func TestShouldUpdateProducePolicy_EncryptFieldsReference(t *testing.T) { + for _, tt := range []struct { + name string + keyName string + wantUpdate bool + }{ + {name: "same key", keyName: "static-key-name"}, + {name: "different key", keyName: "another-key-name", wantUpdate: true}, + } { + t.Run(tt.name, func(t *testing.T) { + desired := producePolicyResourceFromJSON(t, `{ + "ref": "encrypt-fields", + "type": "encrypt_fields", + "name": "encrypt-fields", + "parent_policy_id": "parent-id", + "config": { + "failure_mode": "reject", + "encrypt_fields": [{ + "paths": "record.value.content.ssn", + "encryption_key": { + "type": "static", + "key": {"id": "__REF__:static-key-ref#id"} + } + }] + } + }`) + p := newTestPlanner() + p.resolver = NewReferenceResolver(nil, &resources.ResourceSet{ + EventGatewayStaticKeys: []resources.EventGatewayStaticKeyResource{{ + Ref: "static-key-ref", + EventGatewayStaticKeyCreate: kkComps.EventGatewayStaticKeyCreate{Name: tt.keyName}, + }}, + }) + current := state.EventGatewayVirtualClusterProducePolicyInfo{ + EventGatewayPolicy: kkComps.EventGatewayPolicy{ + ID: "policy-id", + Name: new("encrypt-fields"), + Type: "encrypt_fields", + ParentPolicyID: new("parent-id"), + }, + RawConfig: map[string]any{ + "failure_mode": "reject", + "encrypt_fields": []any{map[string]any{ + "paths": "record.value.content.ssn", + "encryption_key": map[string]any{ + "type": "static", + "key": map[string]any{FieldID: "key-id", FieldName: "static-key-name"}, + }, + }}, + }, + } + + needsUpdate, updateFields, changedFields := p.shouldUpdateProducePolicy(current, desired) + assert.Equal(t, tt.wantUpdate, needsUpdate) + if tt.wantUpdate { + assert.Contains(t, changedFields, FieldConfig) + assert.NotNil(t, updateFields) + } else { + assert.Empty(t, changedFields) + assert.Nil(t, updateFields) + } + // Comparison must not replace the desired reference in the request payload. + fields := p.producePolicyToFields(desired) + id, ok := stringValueAtFieldPath(fields, "config.encrypt_fields.0.encryption_key.key.id") + require.True(t, ok) + assert.Equal(t, "__REF__:static-key-ref#id", id) + }) + } +} diff --git a/test/e2e/scenarios/event-gateway/produce-policy/overlays/009-encrypt-fields-policy/config.yaml b/test/e2e/scenarios/event-gateway/produce-policy/overlays/009-encrypt-fields-policy/config.yaml new file mode 100644 index 000000000..70429c744 --- /dev/null +++ b/test/e2e/scenarios/event-gateway/produce-policy/overlays/009-encrypt-fields-policy/config.yaml @@ -0,0 +1,79 @@ +_defaults: + kongctl: + namespace: event-gateways-vc-produce-policy-comprehensive-test + +event_gateways: + - ref: egw-cp-for-vc-produce-policy-test + name: egw-cp-for-vc-produce-policy-test + description: "EGW CP for Virtual Cluster Produce Policy Test" + backend_clusters: + - ref: default-backend-cluster + name: default-backend-cluster + description: "Backend Cluster for Test" + bootstrap_servers: + - "egw-backend-1.example.com:9092" + - "egw-backend-2.example.com:9092" + authentication: + type: anonymous + tls: + enabled: true + insecure_skip_verify: false + tls_versions: + - tls12 + - tls13 + virtual_clusters: + - ref: default-virtual-cluster-for-produce-policy-test + name: default-virtual-cluster-name-for-produce-policy-test + description: "Virtual Cluster for Produce Policy Test" + destination: + id: !ref default-backend-cluster#id + authentication: + - type: anonymous + acl_mode: enforce_on_gateway + dns_label: vc-default + produce_policies: + - ref: schema-validation-produce-policy + name: schema-validation-produce-policy-name-for-test + description: Schema Validation Produce Policy description updated + type: schema_validation + enabled: false + labels: + env: production + version: v2 + config: + key_validation_action: mark + value_validation_action: mark + type: confluent_schema_registry + schema_registry: + id: !ref schema-registry-for-produce-policy#id + - ref: encrypt-fields-produce-policy + name: encrypt-fields-produce-policy-name-for-test + type: encrypt_fields + enabled: true + parent_policy_id: !ref schema-validation-produce-policy#id + config: + failure_mode: reject + encrypt_fields: + - paths: 'record.value.content["customer.ssn"]' + encryption_key: + type: static + key: + id: !ref static-key-for-produce-policy#id + schema_registries: + - ref: schema-registry-for-produce-policy + name: schema-registry-for-produce-policy-name + description: "Schema Registry for Produce Policy Test" + type: confluent + config: + schema_type: json + endpoint: https://schema-registry.example.com + timeout_seconds: 30 + authentication: + type: basic + username: testuser + password: !secret {source: !env KONGCTL_E2E_SCHEMA_REGISTRY_PASSWORD} + static_keys: + - ref: static-key-for-produce-policy + name: static-key-for-produce-policy-name + description: "Static Key for Produce Policy Test" + value: "YXNkZmdoamthc2RmZ2hqa2FzZGZnaGprYXNkZmdoams=" diff --git a/test/e2e/scenarios/event-gateway/produce-policy/scenario.yaml b/test/e2e/scenarios/event-gateway/produce-policy/scenario.yaml index a6815f6de..f57ed42f1 100644 --- a/test/e2e/scenarios/event-gateway/produce-policy/scenario.yaml +++ b/test/e2e/scenarios/event-gateway/produce-policy/scenario.yaml @@ -17,6 +17,8 @@ vars: producePolicyDescription: Default Produce Policy description svPolicyRef: schema-validation-produce-policy svPolicyName: schema-validation-produce-policy-name-for-test + efPolicyRef: encrypt-fields-produce-policy + efPolicyName: encrypt-fields-produce-policy-name-for-test staticKeyRef: static-key-for-produce-policy staticKeyName: static-key-for-produce-policy-name schemaRegistryRef: schema-registry-for-produce-policy @@ -829,7 +831,201 @@ steps: name: "{{ .vars.svPolicyName }}" type: schema_validation - - name: 011-schema-validation-sync-delete + - name: 011-encrypt-fields-create + inputOverlayDirs: + - overlays/009-encrypt-fields-policy + commands: + - name: 001-plan-encrypt-fields-create + outputFormat: disable + stdoutFile: "{{ .workdir }}/plan-encrypt-fields-create.json" + run: + - plan + - -f + - "{{ .workdir }}/config.yaml" + - --mode + - apply + assertions: + - select: metadata + expect: + fields: + mode: apply + - select: summary + expect: + fields: + total_changes: 2 + by_action.CREATE: 2 + - select: >- + changes[?resource_type=='event_gateway_static_key' && resource_ref=='{{ .vars.staticKeyRef }}'] | [0] + expect: + fields: + action: CREATE + resource_ref: "{{ .vars.staticKeyRef }}" + - select: >- + changes[?resource_type=='event_gateway_virtual_cluster_produce_policy' && resource_ref=='{{ .vars.efPolicyRef }}'] | [0] + expect: + fields: + action: CREATE + resource_ref: "{{ .vars.efPolicyRef }}" + - name: 002-apply-encrypt-fields-create + outputFormat: json + run: + - apply + - --plan + - "{{ .workdir }}/plan-encrypt-fields-create.json" + - --auto-approve + assertions: + - select: summary + expect: + fields: + applied: 2 + failed: 0 + - name: 003-get-static-key + run: + - get + - event-gateway + - static-keys + - --gateway-name + - "{{ .vars.egwName }}" + - -o + - json + recordVar: + name: encryptFieldsKeyID + responsePath: "[?name=='{{ .vars.staticKeyName }}'] | [0].id" + assertions: + - select: "length([?name=='{{ .vars.staticKeyName }}'] | [0].id) > `0`" + expect: + fields: + "@": true + - name: 004-verify-encrypt-fields-create + run: + - get + - event-gateway + - virtual-cluster + - produce-policies + - --gateway-name + - "{{ .vars.egwName }}" + - --virtual-cluster-name + - "{{ .vars.virtualClusterName }}" + - -o + - json + recordVar: + name: schemaValidationPolicyID + responsePath: "[?name=='{{ .vars.svPolicyName }}'] | [0].id" + assertions: + - select: "length(@)" + expect: + fields: + "@": 2 + - select: "length([?name=='{{ .vars.svPolicyName }}'] | [0].id) > `0`" + expect: + fields: + "@": true + - select: "[?name=='{{ .vars.efPolicyName }}'] | [0]" + expect: + fields: + name: "{{ .vars.efPolicyName }}" + type: encrypt_fields + enabled: true + parent_policy_id: "{{ .vars.schemaValidationPolicyID }}" + config.failure_mode: reject + length(config.encrypt_fields): 1 + config.encrypt_fields[0].paths: 'record.value.content["customer.ssn"]' + config.encrypt_fields[0].encryption_key.type: static + - select: "[?name=='{{ .vars.efPolicyName }}'] | [0].config.encrypt_fields[0].encryption_key.key.id" + expect: + fields: + "@": "{{ .vars.encryptFieldsKeyID }}" + - name: 005-idempotency-check + outputFormat: disable + run: + - plan + - -f + - "{{ .workdir }}/config.yaml" + - --mode + - apply + assertions: + - select: summary + expect: + fields: + total_changes: 0 + + - name: 012-encrypt-fields-sync-delete + inputOverlayDirs: + - overlays/006-update-schema-validation + inputOverlayOps: + - file: config.yaml + match: "event_gateways[?ref=='egw-cp-for-vc-produce-policy-test'] | [0]" + set: + static_keys: [] + commands: + - name: 001-plan-encrypt-fields-sync-delete + outputFormat: disable + stdoutFile: "{{ .workdir }}/plan-encrypt-fields-sync-delete.json" + run: + - plan + - -f + - "{{ .workdir }}/config.yaml" + - --mode + - sync + assertions: + - select: metadata + expect: + fields: + mode: sync + - select: summary + expect: + fields: + total_changes: 2 + by_action.DELETE: 2 + - select: >- + changes[?resource_type=='event_gateway_virtual_cluster_produce_policy' && resource_ref=='{{ .vars.efPolicyName }}'] | [0] + expect: + fields: + action: DELETE + resource_ref: "{{ .vars.efPolicyName }}" + - select: >- + changes[?resource_type=='event_gateway_static_key' && resource_ref=='{{ .vars.staticKeyName }}'] | [0] + expect: + fields: + action: DELETE + resource_ref: "{{ .vars.staticKeyName }}" + - name: 002-sync-delete-encrypt-fields + outputFormat: json + run: + - sync + - --plan + - "{{ .workdir }}/plan-encrypt-fields-sync-delete.json" + - --auto-approve + assertions: + - select: summary + expect: + fields: + applied: 2 + failed: 0 + - name: 003-verify-encrypt-fields-sync-delete + run: + - get + - event-gateway + - virtual-cluster + - produce-policies + - --gateway-name + - "{{ .vars.egwName }}" + - --virtual-cluster-name + - "{{ .vars.virtualClusterName }}" + - -o + - json + assertions: + - select: "length(@)" + expect: + fields: + "@": 1 + - select: "[0]" + expect: + fields: + name: "{{ .vars.svPolicyName }}" + type: schema_validation + + - name: 013-schema-validation-sync-delete inputOverlayDirs: - overlays/007-sync-delete-schema-validation commands: