diff --git a/data/k8s/dev/devteam/another-opensearch.yaml b/data/k8s/dev/devteam/another-opensearch.yaml index 2e952b75d..879c98ae0 100644 --- a/data/k8s/dev/devteam/another-opensearch.yaml +++ b/data/k8s/dev/devteam/another-opensearch.yaml @@ -37,6 +37,7 @@ spec: userConfig: opensearch_version: "2" status: + version: "2.19.3" conditions: - lastTransitionTime: "2023-11-08T10:36:06Z" message: Instance was created or update on Aiven side diff --git a/data/k8s/dev/devteam/opensearch.yaml b/data/k8s/dev/devteam/opensearch.yaml index 8c3fd294b..d1cd25aad 100644 --- a/data/k8s/dev/devteam/opensearch.yaml +++ b/data/k8s/dev/devteam/opensearch.yaml @@ -36,6 +36,7 @@ spec: tenant: nav terminationProtection: true status: + version: "2.19.3" conditions: - lastTransitionTime: "2023-11-08T10:24:54Z" message: Instance was created or update on Aiven side @@ -79,6 +80,7 @@ spec: tenant: nav terminationProtection: true status: + version: "2.19.3" conditions: - lastTransitionTime: "2023-11-08T10:24:54Z" message: Instance was created or update on Aiven side diff --git a/data/k8s/dev/devteam/valkey.yaml b/data/k8s/dev/devteam/valkey.yaml index c57eec938..3849b429e 100644 --- a/data/k8s/dev/devteam/valkey.yaml +++ b/data/k8s/dev/devteam/valkey.yaml @@ -30,6 +30,7 @@ spec: tenant: nav terminationProtection: true status: + version: "8.1.9" conditions: - lastTransitionTime: "2023-11-20T19:07:04Z" message: Instance was created or update on Aiven side @@ -68,6 +69,7 @@ spec: tenant: nav terminationProtection: true status: + version: "8.1.9" conditions: - lastTransitionTime: "2023-11-20T19:07:04Z" message: Instance was created or update on Aiven side diff --git a/flake.lock b/flake.lock index 8aea385fc..48705d8b1 100644 --- a/flake.lock +++ b/flake.lock @@ -1,61 +1,61 @@ { - "nodes": { - "flake-utils": { - "inputs": { - "systems": "systems" - }, - "locked": { - "lastModified": 1731533236, - "narHash": "sha256-l0KFg5HjrsfsO/JpG+r7fRrqm12kzFHyUHqHCVpMMbI=", - "owner": "numtide", - "repo": "flake-utils", - "rev": "11707dc2f618dd54ca8739b309ec4fc024de578b", - "type": "github" - }, - "original": { - "owner": "numtide", - "repo": "flake-utils", - "type": "github" - } - }, - "nixpkgs": { - "locked": { - "lastModified": 1746328495, - "narHash": "sha256-uKCfuDs7ZM3QpCE/jnfubTg459CnKnJG/LwqEVEdEiw=", - "owner": "nixos", - "repo": "nixpkgs", - "rev": "979daf34c8cacebcd917d540070b52a3c2b9b16e", - "type": "github" - }, - "original": { - "owner": "nixos", - "ref": "nixos-unstable", - "repo": "nixpkgs", - "type": "github" - } - }, - "root": { - "inputs": { - "flake-utils": "flake-utils", - "nixpkgs": "nixpkgs" - } - }, - "systems": { - "locked": { - "lastModified": 1681028828, - "narHash": "sha256-Vy1rq5AaRuLzOxct8nz4T6wlgyUR7zLU309k9mBC768=", - "owner": "nix-systems", - "repo": "default", - "rev": "da67096a3b9bf56a91d16901293e51ba5b49a27e", - "type": "github" - }, - "original": { - "owner": "nix-systems", - "repo": "default", - "type": "github" - } - } - }, - "root": "root", - "version": 7 + "nodes": { + "flake-utils": { + "inputs": { + "systems": "systems" + }, + "locked": { + "lastModified": 1731533236, + "narHash": "sha256-l0KFg5HjrsfsO/JpG+r7fRrqm12kzFHyUHqHCVpMMbI=", + "owner": "numtide", + "repo": "flake-utils", + "rev": "11707dc2f618dd54ca8739b309ec4fc024de578b", + "type": "github" + }, + "original": { + "owner": "numtide", + "repo": "flake-utils", + "type": "github" + } + }, + "nixpkgs": { + "locked": { + "lastModified": 1787736819, + "narHash": "sha256-cV5xEJJK3BvhU8rEd4mC9UsmDi5qscv/kzGPhBRC5WA=", + "owner": "nixos", + "repo": "nixpkgs", + "rev": "9fbb54b33e91ee4ca368e35a78e0613c720600b3", + "type": "github" + }, + "original": { + "owner": "nixos", + "ref": "nixos-unstable", + "repo": "nixpkgs", + "type": "github" + } + }, + "root": { + "inputs": { + "flake-utils": "flake-utils", + "nixpkgs": "nixpkgs" + } + }, + "systems": { + "locked": { + "lastModified": 1681028828, + "narHash": "sha256-Vy1rq5AaRuLzOxct8nz4T6wlgyUR7zLU309k9mBC768=", + "owner": "nix-systems", + "repo": "default", + "rev": "da67096a3b9bf56a91d16901293e51ba5b49a27e", + "type": "github" + }, + "original": { + "owner": "nix-systems", + "repo": "default", + "type": "github" + } + } + }, + "root": "root", + "version": 7 } diff --git a/flake.nix b/flake.nix index 48f06dcd8..cb8c9a845 100644 --- a/flake.nix +++ b/flake.nix @@ -20,16 +20,16 @@ ( final: prev: let - version = "1.24.2"; + version = "1.26.7"; newerGoVersion = prev.go.overrideAttrs (old: { inherit version; src = prev.fetchurl { url = "https://go.dev/dl/go${version}.src.tar.gz"; - hash = "sha256-ncd/+twW2DehvzLZnGJMtN8GR87nsRnt2eexvMBfLgA="; + hash = ""; }; }); nixpkgsVersion = prev.go.version; - newVersionNotInNixpkgs = -1 == builtins.compareVersions nixpkgsVersion version; + newVersionNotInNixpkgs = -1 == prev.lib.compareVersions nixpkgsVersion version; in { go = if newVersionNotInNixpkgs then newerGoVersion else prev.go; @@ -70,7 +70,7 @@ nodejs_22 mise - nodePackages.prettier + prettier ] ++ [ gqlgen ]; }; diff --git a/go.mod b/go.mod index fbb9d0ec3..01cfa6a18 100644 --- a/go.mod +++ b/go.mod @@ -21,6 +21,7 @@ require ( cloud.google.com/go/pubsub v1.50.1 github.com/99designs/gqlgen v0.17.90 github.com/GoogleCloudPlatform/k8s-config-connector v1.128.0 + github.com/Masterminds/semver/v3 v3.4.0 github.com/aiven/go-client-codegen v0.106.0 github.com/blevesearch/bleve/v2 v2.5.0 github.com/blevesearch/bleve_index_api v1.2.7 @@ -108,7 +109,6 @@ require ( github.com/BurntSushi/toml v1.6.0 // indirect github.com/DataDog/sketches-go v1.4.8 // indirect github.com/Masterminds/goutils v1.1.1 // indirect - github.com/Masterminds/semver/v3 v3.4.0 // indirect github.com/Masterminds/sprig/v3 v3.3.0 // indirect github.com/Microsoft/go-winio v0.6.2 // indirect github.com/RoaringBitmap/roaring/v2 v2.4.5 // indirect diff --git a/integration_tests/k8s_resources/opensearch_crud/dev/someteamname/opensearch_downgrade.yaml b/integration_tests/k8s_resources/opensearch_crud/dev/someteamname/opensearch_downgrade.yaml new file mode 100644 index 000000000..42104ffc3 --- /dev/null +++ b/integration_tests/k8s_resources/opensearch_crud/dev/someteamname/opensearch_downgrade.yaml @@ -0,0 +1,47 @@ +apiVersion: aiven.io/v1alpha1 +kind: OpenSearch +metadata: + annotations: + controllers.aiven.io/generation-was-processed: "2" + controllers.aiven.io/instance-is-running: "true" + nais.io/created_by: aiven-iac-migration + creationTimestamp: "2023-11-08T10:35:59Z" + finalizers: + - finalizers.aiven.io/delete-remote-resource + generation: 2 + labels: + team: teampam + nais.io/managed-by: console + name: opensearch-someteamname-downgrade + namespace: teampam + resourceVersion: "3990043291" + uid: 3b9e5e9b-3cf5-4c3c-a1fd-925fa7b04fe1 +spec: + cloudName: google-europe-north1 + connInfoSecretTarget: + name: "" + disk_space: 525G + plan: hobbyist + project: nav-prod + projectVpcId: fff21e17-95d5-408b-8df5-15aacf38f5de + tags: + environment: prod + team: teampam + tenant: nav + terminationProtection: true + userConfig: + opensearch_version: "3.6" +status: + version: "3.6.0" + conditions: + - lastTransitionTime: "2023-11-08T10:36:06Z" + message: Instance was created or update on Aiven side + reason: Updated + status: "True" + type: Initialized + - lastTransitionTime: "2024-01-10T09:40:58Z" + message: Instance is running on Aiven side + reason: CheckRunning + status: "True" + type: Running + state: RUNNING diff --git a/integration_tests/k8s_resources/opensearch_crud/dev/someteamname/opensearch_noversion.yaml b/integration_tests/k8s_resources/opensearch_crud/dev/someteamname/opensearch_noversion.yaml index b1ceebab6..127beb4ad 100644 --- a/integration_tests/k8s_resources/opensearch_crud/dev/someteamname/opensearch_noversion.yaml +++ b/integration_tests/k8s_resources/opensearch_crud/dev/someteamname/opensearch_noversion.yaml @@ -30,6 +30,7 @@ spec: tenant: nav terminationProtection: true status: + version: "3.3.0" conditions: - lastTransitionTime: "2023-11-08T10:36:06Z" message: Instance was created or update on Aiven side diff --git a/integration_tests/k8s_resources/opensearch_version/dev/myteam/opensearch-myteam-crpin.yaml b/integration_tests/k8s_resources/opensearch_version/dev/myteam/opensearch-myteam-crpin.yaml new file mode 100644 index 000000000..b2d12b10a --- /dev/null +++ b/integration_tests/k8s_resources/opensearch_version/dev/myteam/opensearch-myteam-crpin.yaml @@ -0,0 +1,29 @@ +apiVersion: aiven.io/v1alpha1 +kind: OpenSearch +metadata: + annotations: + controllers.aiven.io/generation-was-processed: "1" + controllers.aiven.io/instance-is-running: "true" + labels: + nais.io/managed-by: console + name: opensearch-myteam-crpin +spec: + cloudName: google-europe-north1 + connInfoSecretTarget: + name: "" + disk_space: 16G + plan: hobbyist + project: nav-dev + projectVpcId: fff21e17-95d5-408b-8df5-15aacf38f5de + terminationProtection: true + userConfig: + opensearch_version: "1" +status: + version: "2.19.3" + conditions: + - lastTransitionTime: "2024-01-10T09:40:58Z" + message: Instance is running on Aiven side + reason: CheckRunning + status: "True" + type: Running + state: RUNNING diff --git a/integration_tests/k8s_resources/opensearch_version/dev/myteam/opensearch-myteam-nocrd.yaml b/integration_tests/k8s_resources/opensearch_version/dev/myteam/opensearch-myteam-nocrd.yaml new file mode 100644 index 000000000..9c89e5986 --- /dev/null +++ b/integration_tests/k8s_resources/opensearch_version/dev/myteam/opensearch-myteam-nocrd.yaml @@ -0,0 +1,26 @@ +apiVersion: aiven.io/v1alpha1 +kind: OpenSearch +metadata: + annotations: + controllers.aiven.io/generation-was-processed: "1" + controllers.aiven.io/instance-is-running: "true" + labels: + nais.io/managed-by: console + name: opensearch-myteam-nocrd +spec: + cloudName: google-europe-north1 + connInfoSecretTarget: + name: "" + disk_space: 16G + plan: hobbyist + project: nav-dev + projectVpcId: fff21e17-95d5-408b-8df5-15aacf38f5de + terminationProtection: true +status: + conditions: + - lastTransitionTime: "2024-01-10T09:40:58Z" + message: Instance is running on Aiven side + reason: CheckRunning + status: "True" + type: Running + state: RUNNING diff --git a/integration_tests/k8s_resources/opensearch_version/dev/myteam/opensearch-myteam-nometa.yaml b/integration_tests/k8s_resources/opensearch_version/dev/myteam/opensearch-myteam-nometa.yaml new file mode 100644 index 000000000..b2e197917 --- /dev/null +++ b/integration_tests/k8s_resources/opensearch_version/dev/myteam/opensearch-myteam-nometa.yaml @@ -0,0 +1,28 @@ +apiVersion: aiven.io/v1alpha1 +kind: OpenSearch +metadata: + annotations: + controllers.aiven.io/generation-was-processed: "1" + controllers.aiven.io/instance-is-running: "true" + labels: + nais.io/managed-by: console + name: opensearch-myteam-nometa +spec: + cloudName: google-europe-north1 + connInfoSecretTarget: + name: "" + disk_space: 16G + plan: hobbyist + project: nav-dev + projectVpcId: fff21e17-95d5-408b-8df5-15aacf38f5de + terminationProtection: true + userConfig: + opensearch_version: "2" +status: + conditions: + - lastTransitionTime: "2024-01-10T09:40:58Z" + message: Instance is running on Aiven side + reason: CheckRunning + status: "True" + type: Running + state: RUNNING diff --git a/integration_tests/k8s_resources/simple/dev/slug-1/opensearch.yaml b/integration_tests/k8s_resources/simple/dev/slug-1/opensearch.yaml index dee7ba7a8..707271328 100644 --- a/integration_tests/k8s_resources/simple/dev/slug-1/opensearch.yaml +++ b/integration_tests/k8s_resources/simple/dev/slug-1/opensearch.yaml @@ -24,6 +24,7 @@ spec: userConfig: opensearch_version: "2" status: + version: "2.19.3" conditions: - lastTransitionTime: "2023-11-08T10:36:06Z" message: Instance was created or update on Aiven side diff --git a/integration_tests/k8s_resources/simple/dev/slug-1/valkey.yaml b/integration_tests/k8s_resources/simple/dev/slug-1/valkey.yaml index dbc784b1e..b451afbc6 100644 --- a/integration_tests/k8s_resources/simple/dev/slug-1/valkey.yaml +++ b/integration_tests/k8s_resources/simple/dev/slug-1/valkey.yaml @@ -21,6 +21,7 @@ spec: tenant: nav terminationProtection: true status: + version: "9.1.0" conditions: - lastTransitionTime: "2023-11-20T19:07:04Z" message: Instance was created or update on Aiven side diff --git a/integration_tests/k8s_resources/state/dev/myteam/opensearch-myteam-poweroff.yaml b/integration_tests/k8s_resources/state/dev/myteam/opensearch-myteam-poweroff.yaml index bc82ba64f..ea99a4d6d 100644 --- a/integration_tests/k8s_resources/state/dev/myteam/opensearch-myteam-poweroff.yaml +++ b/integration_tests/k8s_resources/state/dev/myteam/opensearch-myteam-poweroff.yaml @@ -36,3 +36,5 @@ spec: terminationProtection: true userConfig: opensearch_version: "2" +status: + state: POWEROFF diff --git a/integration_tests/k8s_resources/state/dev/myteam/opensearch-myteam-rebalancing.yaml b/integration_tests/k8s_resources/state/dev/myteam/opensearch-myteam-rebalancing.yaml index 6a109cc3f..ac5694548 100644 --- a/integration_tests/k8s_resources/state/dev/myteam/opensearch-myteam-rebalancing.yaml +++ b/integration_tests/k8s_resources/state/dev/myteam/opensearch-myteam-rebalancing.yaml @@ -36,3 +36,5 @@ spec: terminationProtection: true userConfig: opensearch_version: "2" +status: + state: REBALANCING diff --git a/integration_tests/k8s_resources/state/dev/myteam/opensearch-myteam-running.yaml b/integration_tests/k8s_resources/state/dev/myteam/opensearch-myteam-running.yaml index 3e9fd2663..3c722c425 100644 --- a/integration_tests/k8s_resources/state/dev/myteam/opensearch-myteam-running.yaml +++ b/integration_tests/k8s_resources/state/dev/myteam/opensearch-myteam-running.yaml @@ -36,3 +36,5 @@ spec: terminationProtection: true userConfig: opensearch_version: "2" +status: + state: RUNNING diff --git a/integration_tests/k8s_resources/state/dev/myteam/valkey-myteam-poweroff.yaml b/integration_tests/k8s_resources/state/dev/myteam/valkey-myteam-poweroff.yaml index 068472684..6f6191bc1 100644 --- a/integration_tests/k8s_resources/state/dev/myteam/valkey-myteam-poweroff.yaml +++ b/integration_tests/k8s_resources/state/dev/myteam/valkey-myteam-poweroff.yaml @@ -18,3 +18,5 @@ spec: team: slug-1 tenant: nav terminationProtection: true +status: + state: POWEROFF diff --git a/integration_tests/k8s_resources/state/dev/myteam/valkey-myteam-rebalancing.yaml b/integration_tests/k8s_resources/state/dev/myteam/valkey-myteam-rebalancing.yaml index 572426398..ff7794324 100644 --- a/integration_tests/k8s_resources/state/dev/myteam/valkey-myteam-rebalancing.yaml +++ b/integration_tests/k8s_resources/state/dev/myteam/valkey-myteam-rebalancing.yaml @@ -18,3 +18,5 @@ spec: team: slug-1 tenant: nav terminationProtection: true +status: + state: REBALANCING diff --git a/integration_tests/k8s_resources/state/dev/myteam/valkey-myteam-running.yaml b/integration_tests/k8s_resources/state/dev/myteam/valkey-myteam-running.yaml index d0f4482fc..0e67b1c58 100644 --- a/integration_tests/k8s_resources/state/dev/myteam/valkey-myteam-running.yaml +++ b/integration_tests/k8s_resources/state/dev/myteam/valkey-myteam-running.yaml @@ -18,3 +18,5 @@ spec: team: slug-1 tenant: nav terminationProtection: true +status: + state: RUNNING diff --git a/integration_tests/k8s_resources/valkey_version/dev/myteam/valkey-myteam-nocrd.yaml b/integration_tests/k8s_resources/valkey_version/dev/myteam/valkey-myteam-nocrd.yaml new file mode 100644 index 000000000..c0f29161d --- /dev/null +++ b/integration_tests/k8s_resources/valkey_version/dev/myteam/valkey-myteam-nocrd.yaml @@ -0,0 +1,25 @@ +apiVersion: aiven.io/v1alpha1 +kind: Valkey +metadata: + annotations: + controllers.aiven.io/generation-was-processed: "1" + controllers.aiven.io/instance-is-running: "true" + labels: + nais.io/managed-by: console + name: valkey-myteam-nocrd +spec: + cloudName: google-europe-north1 + connInfoSecretTarget: + name: "" + plan: hobbyist + project: nav-dev + projectVpcId: fff21e17-95d5-408b-8df5-15aacf38f5de + terminationProtection: true +status: + conditions: + - lastTransitionTime: "2024-01-10T09:40:58Z" + message: Instance is running on Aiven side + reason: CheckRunning + status: "True" + type: Running + state: RUNNING diff --git a/integration_tests/k8s_resources/valkey_version/dev/myteam/valkey-myteam-nometa.yaml b/integration_tests/k8s_resources/valkey_version/dev/myteam/valkey-myteam-nometa.yaml new file mode 100644 index 000000000..36b69f7f0 --- /dev/null +++ b/integration_tests/k8s_resources/valkey_version/dev/myteam/valkey-myteam-nometa.yaml @@ -0,0 +1,27 @@ +apiVersion: aiven.io/v1alpha1 +kind: Valkey +metadata: + annotations: + controllers.aiven.io/generation-was-processed: "1" + controllers.aiven.io/instance-is-running: "true" + labels: + nais.io/managed-by: console + name: valkey-myteam-nometa +spec: + cloudName: google-europe-north1 + connInfoSecretTarget: + name: "" + plan: hobbyist + project: nav-dev + projectVpcId: fff21e17-95d5-408b-8df5-15aacf38f5de + terminationProtection: true + userConfig: + valkey_version: "8.1" +status: + conditions: + - lastTransitionTime: "2024-01-10T09:40:58Z" + message: Instance is running on Aiven side + reason: CheckRunning + status: "True" + type: Running + state: RUNNING diff --git a/integration_tests/k8s_resources/valkey_version/dev/myteam/valkey-myteam-upgradable.yaml b/integration_tests/k8s_resources/valkey_version/dev/myteam/valkey-myteam-upgradable.yaml new file mode 100644 index 000000000..84212b9d7 --- /dev/null +++ b/integration_tests/k8s_resources/valkey_version/dev/myteam/valkey-myteam-upgradable.yaml @@ -0,0 +1,28 @@ +apiVersion: aiven.io/v1alpha1 +kind: Valkey +metadata: + annotations: + controllers.aiven.io/generation-was-processed: "1" + controllers.aiven.io/instance-is-running: "true" + labels: + nais.io/managed-by: console + name: valkey-myteam-upgradable +spec: + cloudName: google-europe-north1 + connInfoSecretTarget: + name: "" + plan: hobbyist + project: nav-dev + projectVpcId: fff21e17-95d5-408b-8df5-15aacf38f5de + terminationProtection: true + userConfig: + valkey_version: "8.1" +status: + version: "8.1.4" + conditions: + - lastTransitionTime: "2024-01-10T09:40:58Z" + message: Instance is running on Aiven side + reason: CheckRunning + status: "True" + type: Running + state: RUNNING diff --git a/integration_tests/k8s_resources/valkey_version/dev/myteam/valkey-myteam-versioned.yaml b/integration_tests/k8s_resources/valkey_version/dev/myteam/valkey-myteam-versioned.yaml new file mode 100644 index 000000000..0ceed9a97 --- /dev/null +++ b/integration_tests/k8s_resources/valkey_version/dev/myteam/valkey-myteam-versioned.yaml @@ -0,0 +1,28 @@ +apiVersion: aiven.io/v1alpha1 +kind: Valkey +metadata: + annotations: + controllers.aiven.io/generation-was-processed: "1" + controllers.aiven.io/instance-is-running: "true" + labels: + nais.io/managed-by: console + name: valkey-myteam-versioned +spec: + cloudName: google-europe-north1 + connInfoSecretTarget: + name: "" + plan: hobbyist + project: nav-dev + projectVpcId: fff21e17-95d5-408b-8df5-15aacf38f5de + terminationProtection: true + userConfig: + valkey_version: "8.1" +status: + version: "9.1.0" + conditions: + - lastTransitionTime: "2024-01-10T09:40:58Z" + message: Instance is running on Aiven side + reason: CheckRunning + status: "True" + type: Running + state: RUNNING diff --git a/integration_tests/opensearch_crud.lua b/integration_tests/opensearch_crud.lua index e54f2bb34..e82112138 100644 --- a/integration_tests/opensearch_crud.lua +++ b/integration_tests/opensearch_crud.lua @@ -5,6 +5,10 @@ local mainTeam = Team.new("someteamname", "purpose", "#slack_channel") mainTeam:addMember(user) Helper.readK8sResources("k8s_resources/opensearch_crud") +-- Created instances inherit their pin as status.version, the way the operator records it +-- after its first reconcile, so an update has a current version to validate against. The +-- two fixtures above 2.19 exist so a downgrade can be attempted at all: nothing below +-- 2.19 is selectable any more. Test.gql("Create opensearch in non-existing team", function(t) t.addHeader("x-user-email", user:email()) @@ -17,7 +21,7 @@ Test.gql("Create opensearch in non-existing team", function(t) teamSlug: "devteam" tier: SINGLE_NODE memory: GB_16 - version: V2 + version: V2_19 storageGB: 350 } ) { @@ -53,7 +57,7 @@ Test.gql("Create opensearch as non-team member", function(t) teamSlug: "someteamname" tier: SINGLE_NODE memory: GB_16 - version: V2 + version: V2_19 storageGB: 350 } ) { @@ -89,7 +93,7 @@ Test.gql("Create opensearch as team member", function(t) teamSlug: "someteamname" tier: SINGLE_NODE memory: GB_16 - version: V2 + version: V2_19 storageGB: 350 } ) { @@ -122,7 +126,7 @@ Test.gql("Create opensearch as team member with existing name", function(t) teamSlug: "someteamname" tier: SINGLE_NODE memory: GB_16 - version: V2 + version: V2_19 storageGB: 350 } ) { @@ -158,7 +162,7 @@ Test.gql("Create opensearch with invalid tier and memory combination", function( teamSlug: "someteamname" tier: HIGH_AVAILABILITY memory: GB_2 - version: V2 + version: V2_19 storageGB: 16 } ) { @@ -196,7 +200,7 @@ Test.gql("Create opensearch with invalid storage capacity", function(t) teamSlug: "someteamname" tier: HIGH_AVAILABILITY memory: GB_4 - version: V2 + version: V2_19 storageGB: 16 } ) { @@ -234,7 +238,7 @@ Test.gql("Create opensearch with invalid storage capacity increment", function(t teamSlug: "someteamname" tier: SINGLE_NODE memory: GB_8 - version: V2 + version: V2_19 storageGB: 180 } ) { @@ -267,6 +271,9 @@ Test.k8s("Validate OpenSearch resource", function(t) t.check("aiven.io/v1alpha1", "opensearches", "dev", mainTeam:slug(), resourceName, { apiVersion = "aiven.io/v1alpha1", kind = "OpenSearch", + status = { + version = "2.19", + }, metadata = { name = resourceName, namespace = mainTeam:slug(), @@ -292,7 +299,7 @@ Test.k8s("Validate OpenSearch resource", function(t) tenant = "some-tenant", }, userConfig = { - opensearch_version = "2", + opensearch_version = "2.19", }, }, }) @@ -344,7 +351,7 @@ Test.gql("Create opensearch with tier and memory equivalent to hobbyist plan", f teamSlug: "someteamname" tier: SINGLE_NODE memory: GB_2 - version: V2 + version: V2_19 storageGB: 16 } ) { @@ -372,6 +379,9 @@ Test.k8s("Validate hobbyist OpenSearch resource", function(t) t.check("aiven.io/v1alpha1", "opensearches", "dev", mainTeam:slug(), resourceName, { apiVersion = "aiven.io/v1alpha1", kind = "OpenSearch", + status = { + version = "2.19", + }, metadata = { name = resourceName, namespace = mainTeam:slug(), @@ -397,7 +407,7 @@ Test.k8s("Validate hobbyist OpenSearch resource", function(t) tenant = "some-tenant", }, userConfig = { - opensearch_version = "2", + opensearch_version = "2.19", }, }, }) @@ -449,7 +459,7 @@ Test.gql("Update OpenSearch in non-existing team", function(t) teamSlug: "devteam" tier: SINGLE_NODE memory: GB_16 - version: V2 + version: V2_19 storageGB: 350 } ) { @@ -485,7 +495,7 @@ Test.gql("Update OpenSearch as non-team-member", function(t) teamSlug: "devteam" tier: SINGLE_NODE memory: GB_16 - version: V2 + version: V2_19 storageGB: 350 } ) { @@ -521,7 +531,7 @@ Test.gql("Update OpenSearch as team-member", function(t) teamSlug: "someteamname" tier: HIGH_AVAILABILITY memory: GB_4 - version: V2 + version: V2_19 storageGB: 1020 } ) { @@ -549,6 +559,9 @@ Test.k8s("Validate OpenSearch resource after update", function(t) t.check("aiven.io/v1alpha1", "opensearches", "dev", mainTeam:slug(), resourceName, { apiVersion = "aiven.io/v1alpha1", kind = "OpenSearch", + status = { + version = "2.19", + }, metadata = { name = resourceName, namespace = mainTeam:slug(), @@ -574,7 +587,7 @@ Test.k8s("Validate OpenSearch resource after update", function(t) tenant = "some-tenant", }, userConfig = { - opensearch_version = "2", + opensearch_version = "2.19", }, }, }) @@ -602,6 +615,11 @@ Test.gql("List opensearches for team", function(t) team = { openSearches = { nodes = { + { + name = "downgrade", + tier = "SINGLE_NODE", + memory = "GB_2", + }, { name = "foobar", tier = "HIGH_AVAILABILITY", @@ -640,12 +658,12 @@ Test.gql("Downgrade OpenSearch as team-member", function(t) mutation UpdateOpenSearch { updateOpenSearch( input: { - name: "foobar" + name: "downgrade" environmentName: "dev" teamSlug: "someteamname" tier: HIGH_AVAILABILITY memory: GB_4 - version: V1 + version: V2_19 storageGB: 240 } ) { @@ -660,7 +678,7 @@ Test.gql("Downgrade OpenSearch as team-member", function(t) errors = { { locations = NotNull(), - message = "Cannot change OpenSearch version from V2 to V1. New version must be one of [V2_19]", + message = "Cannot change OpenSearch version from V3_6 to V2_19. No further upgrades available.", path = { "updateOpenSearch", }, @@ -670,7 +688,7 @@ Test.gql("Downgrade OpenSearch as team-member", function(t) } end) -Test.gql("Downgrade OpenSearch without explicit version set", function(t) +Test.gql("Downgrade OpenSearch without a version pinned in the CR", function(t) t.addHeader("x-user-email", user:email()) t.query [[ mutation UpdateOpenSearch { @@ -681,7 +699,7 @@ Test.gql("Downgrade OpenSearch without explicit version set", function(t) teamSlug: "someteamname" tier: HIGH_AVAILABILITY memory: GB_4 - version: V1 + version: V2_19 storageGB: 240 } ) { @@ -696,7 +714,83 @@ Test.gql("Downgrade OpenSearch without explicit version set", function(t) errors = { { locations = NotNull(), - message = "Cannot change OpenSearch version from V2 to V1. New version must be one of [V2_19]", + message = "Cannot change OpenSearch version from V3_3 to V2_19. New version must be one of [V3_6]", + path = { + "updateOpenSearch", + }, + }, + }, + data = Null, + } +end) + +Test.gql("Reject a deprecated OpenSearch version on create", function(t) + t.addHeader("x-user-email", user:email()) + t.query [[ + mutation CreateOpenSearch { + createOpenSearch( + input: { + name: "deprecated-create" + environmentName: "dev" + teamSlug: "someteamname" + tier: SINGLE_NODE + memory: GB_16 + version: V2 + storageGB: 350 + } + ) { + openSearch { + name + } + } + } + ]] + + t.check { + errors = { + { + extensions = { + field = "version", + }, + message = "OpenSearch version V2 is deprecated: use V2_19 instead.", + path = { + "createOpenSearch", + }, + }, + }, + data = Null, + } +end) + +Test.gql("Reject a deprecated OpenSearch version on update", function(t) + t.addHeader("x-user-email", user:email()) + t.query [[ + mutation UpdateOpenSearch { + updateOpenSearch( + input: { + name: "foobar" + environmentName: "dev" + teamSlug: "someteamname" + tier: HIGH_AVAILABILITY + memory: GB_4 + version: V3_3 + storageGB: 240 + } + ) { + openSearch { + name + } + } + } + ]] + + t.check { + errors = { + { + extensions = { + field = "version", + }, + message = "OpenSearch version V3_3 is deprecated: vendor support disappears 2027-02-01.", path = { "updateOpenSearch", }, @@ -717,7 +811,7 @@ Test.gql("Update non-console managed OpenSearch as team-member", function(t) teamSlug: "someteamname" tier: HIGH_AVAILABILITY memory: GB_4 - version: V2 + version: V2_19 storageGB: 240 } ) { @@ -753,7 +847,7 @@ Test.gql("Update OpenSearch with tier and memory equivalent to hobbyist plan", f teamSlug: "someteamname" tier: SINGLE_NODE memory: GB_2 - version: V2 + version: V2_19 storageGB: 16 } ) { @@ -781,6 +875,9 @@ Test.k8s("Validate hobbyist OpenSearch resource after update", function(t) t.check("aiven.io/v1alpha1", "opensearches", "dev", mainTeam:slug(), resourceName, { apiVersion = "aiven.io/v1alpha1", kind = "OpenSearch", + status = { + version = "2.19", + }, metadata = { name = resourceName, namespace = mainTeam:slug(), @@ -806,7 +903,7 @@ Test.k8s("Validate hobbyist OpenSearch resource after update", function(t) tenant = "some-tenant", }, userConfig = { - opensearch_version = "2", + opensearch_version = "2.19", }, }, }) @@ -857,7 +954,7 @@ Test.gql("Create opensearch in other team", function(t) teamSlug: "%s" tier: SINGLE_NODE memory: GB_16 - version: V2 + version: V2_19 storageGB: 350 } ) { diff --git a/integration_tests/opensearch_version_defects.lua b/integration_tests/opensearch_version_defects.lua new file mode 100644 index 000000000..6605af85a --- /dev/null +++ b/integration_tests/opensearch_version_defects.lua @@ -0,0 +1,270 @@ +local user = User.new("user", "user@usersen.com") +local team = Team.new("myteam", "purpose", "#slack_channel") +team:addMember(user) + +Helper.readK8sResources("k8s_resources/opensearch_version") + +-- status.version is the sole authority for the version. The CR pin never feeds what the +-- API reports, and a version the operator has not recorded is not reported at all. +local unreportedVersion = +"The server errored out while processing your request, and we didn't write a suitable error message. You might consider that a bug on our side. Please try again, and if the error persists, contact the Nais team." + +local function unknownVersion(name) + return string.format('The running version of "%s" is not known yet. Try again shortly.', name) +end + +local function versionQuery(name) + return string.format([[ + query { + team(slug: "myteam") { + environment(name: "dev") { + openSearch(name: "%s") { + name + version { + actual + desiredMajor + } + } + } + } + } + ]], name) +end + +Test.gql("Aiven's version wins over a CR that pins a different one", function(t) + t.addHeader("x-user-email", user:email()) + t.query(versionQuery("crpin")) + + t.check { + data = { + team = { + environment = { + openSearch = { + name = "crpin", + version = { + actual = "2.19.3", + desiredMajor = "V2_19", + }, + }, + }, + }, + }, + } +end) + +Test.gql("Aiven reports no version, CR pins one", function(t) + t.addHeader("x-user-email", user:email()) + t.query(versionQuery("nometa")) + + t.check { + data = { + team = { + environment = { + openSearch = { + name = "nometa", + version = { + actual = Null, + desiredMajor = "V2_19", + }, + }, + }, + }, + }, + } +end) + +Test.gql("Neither Aiven nor the CR reports a version", function(t) + t.addHeader("x-user-email", user:email()) + t.query(versionQuery("nocrd")) + + t.check { + errors = { + { + locations = NotNull(), + message = unreportedVersion, + path = { + "team", + "environment", + "openSearch", + "version", + }, + }, + }, + data = Null, + } +end) + +Test.gql("Update OpenSearch when no version is known anywhere", function(t) + t.addHeader("x-user-email", user:email()) + t.query [[ + mutation UpdateOpenSearch { + updateOpenSearch( + input: { + name: "nocrd" + environmentName: "dev" + teamSlug: "myteam" + tier: SINGLE_NODE + memory: GB_2 + version: V2_19 + storageGB: 16 + } + ) { + openSearch { + name + } + } + } + ]] + + t.check { + errors = { + { + locations = NotNull(), + message = unknownVersion("nocrd"), + path = { + "updateOpenSearch", + }, + }, + }, + data = Null, + } +end) + +Test.gql("Aiven silent, the CR's pinned version does not rescue the upgrade", function(t) + t.addHeader("x-user-email", user:email()) + t.query [[ + mutation UpdateOpenSearch { + updateOpenSearch( + input: { + name: "nometa" + environmentName: "dev" + teamSlug: "myteam" + tier: SINGLE_NODE + memory: GB_2 + version: V2_19 + storageGB: 16 + } + ) { + openSearch { + name + } + } + } + ]] + + t.check { + errors = { + { + locations = NotNull(), + message = unknownVersion("nometa"), + path = { + "updateOpenSearch", + }, + }, + }, + data = Null, + } +end) + +Test.gql("Upgrade OpenSearch to a version on the upgrade path", function(t) + t.addHeader("x-user-email", user:email()) + t.query [[ + mutation UpdateOpenSearch { + updateOpenSearch( + input: { + name: "crpin" + environmentName: "dev" + teamSlug: "myteam" + tier: SINGLE_NODE + memory: GB_2 + version: V3_6 + storageGB: 16 + } + ) { + openSearch { + name + } + } + } + ]] + + t.check { + data = { + updateOpenSearch = { + openSearch = { + name = "crpin", + }, + }, + }, + } +end) + +Test.gql("Creating without a version uses the newest Aiven supports", function(t) + t.addHeader("x-user-email", user:email()) + t.query [[ + mutation CreateOpenSearch { + createOpenSearch( + input: { + name: "newest" + environmentName: "dev" + teamSlug: "myteam" + tier: SINGLE_NODE + memory: GB_2 + storageGB: 16 + } + ) { + openSearch { + name + } + } + } + ]] + + t.check { + data = { + createOpenSearch = { + openSearch = { + name = "newest", + }, + }, + }, + } +end) + +Test.k8s("The created instance pins the newest version", function(t) + t.check("aiven.io/v1alpha1", "opensearches", "dev", "myteam", "opensearch-myteam-newest", { + apiVersion = "aiven.io/v1alpha1", + kind = "OpenSearch", + status = { + version = "3.6", + }, + metadata = { + name = "opensearch-myteam-newest", + namespace = "myteam", + annotations = { + ["console.nais.io/last-modified-at"] = NotNull(), + ["console.nais.io/last-modified-by"] = user:email(), + }, + labels = { + ["app.kubernetes.io/managed-by"] = "console", + ["nais.io/managed-by"] = "console", + }, + }, + spec = { + project = "aiven-dev", + projectVpcId = "aiven-vpc", + plan = "hobbyist", + cloudName = "google-europe-north1", + disk_space = "16G", + terminationProtection = true, + tags = { + environment = "dev", + team = "myteam", + tenant = "some-tenant", + }, + userConfig = { + opensearch_version = "3.6", + }, + }, + }) +end) diff --git a/integration_tests/opensearchversion.lua b/integration_tests/opensearchversion.lua index 727f0b877..22fc63cdc 100644 --- a/integration_tests/opensearchversion.lua +++ b/integration_tests/opensearchversion.lua @@ -28,8 +28,8 @@ Test.gql("Show version of OpenSearch instance", function(t) { name = "opensearch-slug-1-opensearch", version = { - actual = "2.17.2", - desiredMajor = "V2", + actual = "2.19.3", + desiredMajor = "V2_19", }, }, }, diff --git a/integration_tests/valkey_crud.lua b/integration_tests/valkey_crud.lua index 312706067..e5df0bfb3 100644 --- a/integration_tests/valkey_crud.lua +++ b/integration_tests/valkey_crud.lua @@ -6,6 +6,8 @@ mainTeam:addMember(user) local otherTeam = Team.new("someothername", "purpose", "#slack_channel") Helper.readK8sResources("k8s_resources/valkey_crud") +-- Created instances inherit their pin as status.version, the way the operator records it +-- after its first reconcile, so an update has a current version to validate against. Test.gql("Create valkey in non-existing team", function(t) t.addHeader("x-user-email", user:email()) @@ -223,6 +225,9 @@ Test.k8s("Validate Valkey resource", function(t) t.check("aiven.io/v1alpha1", "valkeys", "dev", mainTeam:slug(), resourceName, { apiVersion = "aiven.io/v1alpha1", kind = "Valkey", + status = { + version = "9.1", + }, metadata = { name = resourceName, namespace = mainTeam:slug(), @@ -243,6 +248,7 @@ Test.k8s("Validate Valkey resource", function(t) terminationProtection = true, userConfig = { valkey_number_of_databases = 32, + valkey_version = "9.1", }, tags = { environment = "dev", @@ -299,6 +305,7 @@ Test.gql("Update Valkey in non-existing team", function(t) teamSlug: "devteam" tier: SINGLE_NODE memory: GB_14 + version: V9_1 } ) { valkey { @@ -333,6 +340,7 @@ Test.gql("Update Valkey as non-team-member", function(t) teamSlug: "devteam" tier: SINGLE_NODE memory: GB_14 + version: V9_1 } ) { valkey { @@ -367,6 +375,7 @@ Test.gql("Update Valkey as team-member", function(t) teamSlug: "someteamname" tier: HIGH_AVAILABILITY memory: GB_4 + version: V9_1 maxMemoryPolicy: ALLKEYS_RANDOM notifyKeyspaceEvents: "Exd" databases: 64 @@ -398,6 +407,9 @@ Test.k8s("Validate Valkey resource after update", function(t) t.check("aiven.io/v1alpha1", "valkeys", "dev", mainTeam:slug(), resourceName, { apiVersion = "aiven.io/v1alpha1", kind = "Valkey", + status = { + version = "9.1", + }, metadata = { name = resourceName, namespace = mainTeam:slug(), @@ -420,6 +432,7 @@ Test.k8s("Validate Valkey resource after update", function(t) valkey_maxmemory_policy = "allkeys-random", valkey_notify_keyspace_events = "Exd", valkey_number_of_databases = 64, + valkey_version = "9.1", }, tags = { environment = "dev", @@ -467,6 +480,9 @@ Test.k8s("Validate hobbyist Valkey resource", function(t) t.check("aiven.io/v1alpha1", "valkeys", "dev", mainTeam:slug(), resourceName, { apiVersion = "aiven.io/v1alpha1", kind = "Valkey", + status = { + version = "9.1", + }, metadata = { name = resourceName, namespace = mainTeam:slug(), @@ -485,6 +501,9 @@ Test.k8s("Validate hobbyist Valkey resource", function(t) plan = "hobbyist", cloudName = "google-europe-north1", terminationProtection = true, + userConfig = { + valkey_version = "9.1", + }, tags = { environment = "dev", team = mainTeam:slug(), @@ -604,6 +623,7 @@ Test.gql("Update non-console managed Valkey as team-member", function(t) teamSlug: "someteamname" tier: HIGH_AVAILABILITY memory: GB_4 + version: V9_1 maxMemoryPolicy: ALLKEYS_RANDOM notifyKeyspaceEvents: "Exd" } @@ -640,6 +660,7 @@ Test.gql("Update Valkey with tier and memory equivalent to hobbyist plan", funct teamSlug: "someteamname" tier: SINGLE_NODE memory: GB_1 + version: V9_1 maxMemoryPolicy: ALLKEYS_RANDOM notifyKeyspaceEvents: "Exd" } @@ -668,6 +689,9 @@ Test.k8s("Validate hobbyist Valkey resource after update", function(t) t.check("aiven.io/v1alpha1", "valkeys", "dev", mainTeam:slug(), resourceName, { apiVersion = "aiven.io/v1alpha1", kind = "Valkey", + status = { + version = "9.1", + }, metadata = { name = resourceName, namespace = mainTeam:slug(), @@ -690,6 +714,7 @@ Test.k8s("Validate hobbyist Valkey resource after update", function(t) valkey_maxmemory_policy = "allkeys-random", valkey_notify_keyspace_events = "Exd", valkey_number_of_databases = 64, + valkey_version = "9.1", }, tags = { environment = "dev", @@ -1119,6 +1144,7 @@ Test.gql("Update Valkey labels successfully", function(t) teamSlug: "someteamname" tier: SINGLE_NODE memory: GB_14 + version: V9_1 labels: [ { key: "my-custom-key", value: "testing" } ] @@ -1160,6 +1186,7 @@ Test.gql("Update Valkey labels with reserved key -> should fail validation", fun teamSlug: "someteamname" tier: SINGLE_NODE memory: GB_14 + version: V9_1 labels: [ { key: "app", value: "invalid" } ] @@ -1199,6 +1226,7 @@ Test.gql("Update Valkey labels to specify app.kubernetes.io/managed-by: Helm", f teamSlug: "someteamname" tier: SINGLE_NODE memory: GB_14 + version: V9_1 labels: [ { key: "app.kubernetes.io/managed-by", value: "Helm" } ] @@ -1242,6 +1270,7 @@ Test.gql( teamSlug: "someteamname" tier: SINGLE_NODE memory: GB_14 + version: V9_1 labels: [ { key: "my-custom-key", value: "second-test" } ] @@ -1278,6 +1307,9 @@ Test.k8s("Validate Valkey labels after update", function(t) t.check("aiven.io/v1alpha1", "valkeys", "dev", mainTeam:slug(), resourceName, { apiVersion = "aiven.io/v1alpha1", kind = "Valkey", + status = { + version = "9.1", + }, metadata = { name = resourceName, namespace = mainTeam:slug(), diff --git a/integration_tests/valkey_version.lua b/integration_tests/valkey_version.lua new file mode 100644 index 000000000..e98432854 --- /dev/null +++ b/integration_tests/valkey_version.lua @@ -0,0 +1,179 @@ +local user = User.new("user", "user@usersen.com") +local team = Team.new("myteam", "purpose", "#slack_channel") +team:addMember(user) + +Helper.readK8sResources("k8s_resources/valkey_version") +-- "versioned" sits on the newest rung, "upgradable" on the oldest, so both ends of the +-- upgrade path have a subject. Both pin 8.1 in their CR, which Aiven overrides. + +-- The operator is the authority for the version, mirroring OpenSearch. A version it has +-- not recorded yet is not reported at all. +local function unknownVersion(name) + return string.format('The running version of "%s" is not known yet. Try again shortly.', name) +end + +local function versionQuery(name) + return string.format([[ + query { + team(slug: "myteam") { + environment(name: "dev") { + valkey(name: "%s") { + name + version { + actual + desiredMajor + } + } + } + } + } + ]], name) +end + + +Test.gql("Aiven's version wins over a CR that pins a different one", function(t) + t.addHeader("x-user-email", user:email()) + t.query(versionQuery("versioned")) + + t.check { + data = { + team = { + environment = { + valkey = { + name = "versioned", + version = { + actual = "9.1.0", + desiredMajor = "V9_1", + }, + }, + }, + }, + }, + } +end) + +Test.gql("Aiven reports no Valkey version", function(t) + t.addHeader("x-user-email", user:email()) + t.query(versionQuery("nometa")) + + t.check { + data = { + team = { + environment = { + valkey = { + name = "nometa", + version = { + actual = Null, + desiredMajor = "V8_1", + }, + }, + }, + }, + }, + } +end) + +local function updateVersion(name, version) + return string.format([[ + mutation UpdateValkey { + updateValkey( + input: { + name: "%s" + environmentName: "dev" + teamSlug: "myteam" + tier: SINGLE_NODE + memory: GB_1 + version: %s + } + ) { + valkey { + name + } + } + } + ]], name, version) +end + +Test.gql("Requesting the version Aiven already reports changes nothing", function(t) + t.addHeader("x-user-email", user:email()) + t.query(updateVersion("versioned", "V9_1")) + + t.check { + data = { + updateValkey = { + valkey = { + name = "versioned", + }, + }, + }, + } +end) + +Test.gql("Reject a Valkey downgrade below the version Aiven reports", function(t) + t.addHeader("x-user-email", user:email()) + t.query(updateVersion("versioned", "V8_1")) + + t.check { + errors = { + { + locations = NotNull(), + message = "Cannot change Valkey version from V9_1 to V8_1. No further upgrades available.", + path = { + "updateValkey", + }, + }, + }, + data = Null, + } +end) + +Test.gql("Aiven silent, the CR's pinned version does not rescue the upgrade", function(t) + t.addHeader("x-user-email", user:email()) + t.query(updateVersion("nometa", "V9_1")) + + t.check { + errors = { + { + locations = NotNull(), + message = unknownVersion("nometa"), + path = { + "updateValkey", + }, + }, + }, + data = Null, + } +end) + +Test.gql("Update Valkey when no version is known anywhere", function(t) + t.addHeader("x-user-email", user:email()) + t.query(updateVersion("nocrd", "V9_1")) + + t.check { + errors = { + { + locations = NotNull(), + message = unknownVersion("nocrd"), + path = { + "updateValkey", + }, + }, + }, + data = Null, + } +end) + +Test.gql("Upgrade Valkey to a version on the upgrade path", function(t) + t.addHeader("x-user-email", user:email()) + t.query(updateVersion("upgradable", "V9_1")) + + t.check { + data = { + updateValkey = { + valkey = { + name = "upgradable", + }, + }, + }, + } +end) diff --git a/integration_tests/valkeyversion.lua b/integration_tests/valkeyversion.lua new file mode 100644 index 000000000..7d896d841 --- /dev/null +++ b/integration_tests/valkeyversion.lua @@ -0,0 +1,40 @@ +Helper.readK8sResources("./k8s_resources/simple") +local user = User.new() +local team = Team.new("slug-1", "purpose", "#channel") + +Test.gql("Show version of Valkey instance", function(t) + t.addHeader("x-user-email", user:email()) + + t.query(string.format([[ +{ + team(slug: "%s") { + valkeys { + nodes { + name + version { + actual + desiredMajor + } + } + } + } +}]], team:slug())) + + t.check { + data = { + team = { + valkeys = { + nodes = { + { + name = "valkey-slug-1-contests", + version = { + actual = "9.1.0", + desiredMajor = "V9_1", + }, + }, + }, + }, + }, + }, + } +end) diff --git a/internal/cmd/api/http.go b/internal/cmd/api/http.go index 9864a1ed5..c38c56a52 100644 --- a/internal/cmd/api/http.go +++ b/internal/cmd/api/http.go @@ -349,8 +349,8 @@ func ConfigureGraph( ctx = config.NewLoaderContext(ctx, watchers.ConfigWatcher, log) ctx = instancegroup.NewLoaderContext(ctx, watchers.ReplicaSetWatcher, watchers.PodWatcher, watchers.AppWatcher, dynamicClients, log) ctx = aiven.NewLoaderContext(ctx, aivenProjects) - ctx = opensearch.NewLoaderContext(ctx, tenantName, watchers.OpenSearchWatcher, aivenClient, log) - ctx = valkey.NewLoaderContext(ctx, tenantName, watchers.ValkeyWatcher, aivenClient, log) + ctx = opensearch.NewLoaderContext(ctx, tenantName, watchers.OpenSearchWatcher, log) + ctx = valkey.NewLoaderContext(ctx, tenantName, watchers.ValkeyWatcher, log) ctx = price.NewLoaderContext(ctx, priceRetriever, log) ctx = utilization.NewLoaderContext(ctx, prometheusClient, log) ctx = alerts.NewLoaderContext(ctx, prometheusClient, log) diff --git a/internal/graph/gengql/opensearch.generated.go b/internal/graph/gengql/opensearch.generated.go index 6900d0aa4..5d350d336 100644 --- a/internal/graph/gengql/opensearch.generated.go +++ b/internal/graph/gengql/opensearch.generated.go @@ -2622,7 +2622,7 @@ func (ec *executionContext) unmarshalInputCreateOpenSearchInput(ctx context.Cont it.Memory = data case "version": ctx := graphql.WithPathContext(ctx, graphql.NewPathWithField("version")) - data, err := ec.unmarshalNOpenSearchMajorVersion2githubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋopensearchᚐOpenSearchMajorVersion(ctx, v) + data, err := ec.unmarshalOOpenSearchMajorVersion2ᚖgithubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋopensearchᚐOpenSearchMajorVersion(ctx, v) if err != nil { return it, err } @@ -2863,7 +2863,7 @@ func (ec *executionContext) unmarshalInputUpdateOpenSearchInput(ctx context.Cont it.Memory = data case "version": ctx := graphql.WithPathContext(ctx, graphql.NewPathWithField("version")) - data, err := ec.unmarshalNOpenSearchMajorVersion2githubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋopensearchᚐOpenSearchMajorVersion(ctx, v) + data, err := ec.unmarshalNOpenSearchMajorVersion2ᚖgithubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋopensearchᚐOpenSearchMajorVersion(ctx, v) if err != nil { return it, err } @@ -4735,6 +4735,22 @@ func (ec *executionContext) marshalNOpenSearchMajorVersion2githubᚗcomᚋnais return v } +func (ec *executionContext) unmarshalNOpenSearchMajorVersion2ᚖgithubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋopensearchᚐOpenSearchMajorVersion(ctx context.Context, v any) (*opensearch.OpenSearchMajorVersion, error) { + var res = new(opensearch.OpenSearchMajorVersion) + err := res.UnmarshalGQL(v) + return res, graphql.ErrorOnPath(ctx, err) +} + +func (ec *executionContext) marshalNOpenSearchMajorVersion2ᚖgithubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋopensearchᚐOpenSearchMajorVersion(ctx context.Context, sel ast.SelectionSet, v *opensearch.OpenSearchMajorVersion) graphql.Marshaler { + if v == nil { + if !graphql.HasFieldError(ctx, graphql.GetFieldContext(ctx)) { + graphql.AddErrorf(ctx, "the requested element is null which the schema does not allow") + } + return graphql.Null + } + return v +} + func (ec *executionContext) unmarshalNOpenSearchMemory2githubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋopensearchᚐOpenSearchMemory(ctx context.Context, v any) (opensearch.OpenSearchMemory, error) { var res opensearch.OpenSearchMemory err := res.UnmarshalGQL(v) @@ -4918,6 +4934,22 @@ func (ec *executionContext) unmarshalOOpenSearchFilter2ᚖgithubᚗcomᚋnaisᚋ return &res, graphql.ErrorOnPath(ctx, err) } +func (ec *executionContext) unmarshalOOpenSearchMajorVersion2ᚖgithubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋopensearchᚐOpenSearchMajorVersion(ctx context.Context, v any) (*opensearch.OpenSearchMajorVersion, error) { + if v == nil { + return nil, nil + } + var res = new(opensearch.OpenSearchMajorVersion) + err := res.UnmarshalGQL(v) + return res, graphql.ErrorOnPath(ctx, err) +} + +func (ec *executionContext) marshalOOpenSearchMajorVersion2ᚖgithubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋopensearchᚐOpenSearchMajorVersion(ctx context.Context, sel ast.SelectionSet, v *opensearch.OpenSearchMajorVersion) graphql.Marshaler { + if v == nil { + return graphql.Null + } + return v +} + func (ec *executionContext) unmarshalOOpenSearchOrder2ᚖgithubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋopensearchᚐOpenSearchOrder(ctx context.Context, v any) (*opensearch.OpenSearchOrder, error) { if v == nil { return nil, nil diff --git a/internal/graph/gengql/root_.generated.go b/internal/graph/gengql/root_.generated.go index bb3a21a01..a9535468e 100644 --- a/internal/graph/gengql/root_.generated.go +++ b/internal/graph/gengql/root_.generated.go @@ -3472,6 +3472,7 @@ type ComplexityRoot struct { TeamEnvironment func(childComplexity int) int TerminationProtection func(childComplexity int) int Tier func(childComplexity int) int + Version func(childComplexity int) int Workload func(childComplexity int) int } @@ -3619,6 +3620,11 @@ type ComplexityRoot struct { OldValue func(childComplexity int) int } + ValkeyVersion struct { + Actual func(childComplexity int) int + DesiredMajor func(childComplexity int) int + } + ViewSecretValuesPayload struct { Values func(childComplexity int) int } @@ -18413,6 +18419,13 @@ func (e *executableSchema) Complexity(ctx context.Context, typeName, field strin return e.ComplexityRoot.Valkey.Tier(childComplexity), true + case "Valkey.version": + if e.ComplexityRoot.Valkey.Version == nil { + break + } + + return e.ComplexityRoot.Valkey.Version(childComplexity), true + case "Valkey.workload": if e.ComplexityRoot.Valkey.Workload == nil { break @@ -18992,6 +19005,20 @@ func (e *executableSchema) Complexity(ctx context.Context, typeName, field strin return e.ComplexityRoot.ValkeyUpdatedActivityLogEntryDataUpdatedField.OldValue(childComplexity), true + case "ValkeyVersion.actual": + if e.ComplexityRoot.ValkeyVersion.Actual == nil { + break + } + + return e.ComplexityRoot.ValkeyVersion.Actual(childComplexity), true + + case "ValkeyVersion.desiredMajor": + if e.ComplexityRoot.ValkeyVersion.DesiredMajor == nil { + break + } + + return e.ComplexityRoot.ValkeyVersion.DesiredMajor(childComplexity), true + case "ViewSecretValuesPayload.values": if e.ComplexityRoot.ViewSecretValuesPayload.Values == nil { break @@ -25642,16 +25669,14 @@ enum OpenSearchMemory { } enum OpenSearchMajorVersion { - "OpenSearch Version 3.6.x" + "OpenSearch Version 3.6 LTS" V3_6 "OpenSearch Version 3.3.x" - V3_3 - "OpenSearch Version 2.19.x" + V3_3 @deprecated(reason: "Vendor support disappears 2027-02-01") + "OpenSearch Version 2.19 LTS" V2_19 - "OpenSearch Version 2.17.x" - V2 - "OpenSearch Version 1.3.x" - V1 + "OpenSearch Version 2.19 LTS - backwards compatible" + V2 @deprecated(reason: "Use ` + "`" + `V2_19` + "`" + ` instead") } input CreateOpenSearchInput { @@ -25666,7 +25691,7 @@ input CreateOpenSearchInput { "Available memory for the OpenSearch instance." memory: OpenSearchMemory! "Major version of the OpenSearch instance." - version: OpenSearchMajorVersion! + version: OpenSearchMajorVersion "Available storage in GB." storageGB: Int! } @@ -31553,6 +31578,24 @@ type TeamInventoryCountValkeys { total: Int! } +"Version information for a Valkey instance." +type ValkeyVersion { + "The full version string of the Valkey instance. This will be available after the instance is created." + actual: String + "The desired major version of the Valkey instance." + desiredMajor: ValkeyMajorVersion! +} + +"Major version of a Valkey instance." +enum ValkeyMajorVersion { + "Valkey Version 9.1.x" + V9_1 + "Valkey Version 9.0.x" + V9_0 @deprecated(reason: "No longer supported by back-end services") + "Valkey Version 8.1.x" + V8_1 +} + type Valkey implements Persistence & Node { id: ID! name: String! @@ -31581,6 +31624,8 @@ type Valkey implements Persistence & Node { notifyKeyspaceEvents: String "Number of databases the Valkey instance is configured with. Default is 16. Minimum 1, maximum 128. Changing this will cause a restart of the Valkey service." databases: Int! + "Fetch version for the Valkey instance." + version: ValkeyVersion! "Issues that affects the instance." issues( "Get the first n items in the connection. This can be used in combination with the after parameter." @@ -31778,6 +31823,8 @@ input CreateValkeyInput { tier: ValkeyTier! "Available memory for the Valkey instance." memory: ValkeyMemory! + "Major version of the Valkey instance." + version: ValkeyMajorVersion "Maximum memory policy for the Valkey instance." maxMemoryPolicy: ValkeyMaxMemoryPolicy "Configure keyspace notifications for the Valkey instance. See https://valkey.io/topics/notifications/ for details." @@ -31802,6 +31849,8 @@ input UpdateValkeyInput { tier: ValkeyTier! "Available memory for the Valkey instance." memory: ValkeyMemory! + "Major version of the Valkey instance." + version: ValkeyMajorVersion! "Maximum memory policy for the Valkey instance." maxMemoryPolicy: ValkeyMaxMemoryPolicy "Configure keyspace notifications for the Valkey instance. See https://valkey.io/topics/notifications/ for details." @@ -38019,6 +38068,8 @@ func (ec *executionContext) childFields_Valkey(ctx context.Context, field graphq return ec.fieldContext_Valkey_notifyKeyspaceEvents(ctx, field) case "databases": return ec.fieldContext_Valkey_databases(ctx, field) + case "version": + return ec.fieldContext_Valkey_version(ctx, field) case "issues": return ec.fieldContext_Valkey_issues(ctx, field) case "activityLog": @@ -38209,6 +38260,16 @@ func (ec *executionContext) childFields_ValkeyUpdatedActivityLogEntryDataUpdated return nil, fmt.Errorf("no field named %q was found under type ValkeyUpdatedActivityLogEntryDataUpdatedField", field.Name) } +func (ec *executionContext) childFields_ValkeyVersion(ctx context.Context, field graphql.CollectedField) (*graphql.FieldContext, error) { + switch field.Name { + case "actual": + return ec.fieldContext_ValkeyVersion_actual(ctx, field) + case "desiredMajor": + return ec.fieldContext_ValkeyVersion_desiredMajor(ctx, field) + } + return nil, fmt.Errorf("no field named %q was found under type ValkeyVersion", field.Name) +} + func (ec *executionContext) childFields_ViewSecretValuesPayload(ctx context.Context, field graphql.CollectedField) (*graphql.FieldContext, error) { switch field.Name { case "values": diff --git a/internal/graph/gengql/valkey.generated.go b/internal/graph/gengql/valkey.generated.go index 792bdd005..053228007 100644 --- a/internal/graph/gengql/valkey.generated.go +++ b/internal/graph/gengql/valkey.generated.go @@ -36,6 +36,7 @@ type ValkeyResolver interface { Workload(ctx context.Context, obj *valkey.Valkey) (workload.Workload, error) State(ctx context.Context, obj *valkey.Valkey) (valkey.ValkeyState, error) + Version(ctx context.Context, obj *valkey.Valkey) (*valkey.ValkeyVersion, error) Issues(ctx context.Context, obj *valkey.Valkey, first *int, after *pagination.Cursor, last *int, before *pagination.Cursor, orderBy *issue.IssueOrder, filter *issue.ResourceIssueFilter) (*issue.IssueConnection, error) ActivityLog(ctx context.Context, obj *valkey.Valkey, first *int, after *pagination.Cursor, last *int, before *pagination.Cursor, filter *activitylog.ActivityLogFilter) (*activitylog.ActivityLogEntryConnection, error) Cost(ctx context.Context, obj *valkey.Valkey) (*cost.ValkeyCost, error) @@ -727,6 +728,38 @@ func (ec *executionContext) fieldContext_Valkey_databases(_ context.Context, fie return graphql.NewScalarFieldContext("Valkey", field, false, false, errors.New("field of type Int does not have child fields")) } +func (ec *executionContext) _Valkey_version(ctx context.Context, field graphql.CollectedField, obj *valkey.Valkey) (ret graphql.Marshaler) { + return graphql.ResolveField( + ctx, + ec.OperationContext, + field, + func(ctx context.Context, field graphql.CollectedField) (*graphql.FieldContext, error) { + return ec.fieldContext_Valkey_version(ctx, field) + }, + func(ctx context.Context) (any, error) { + return ec.Resolvers.Valkey().Version(ctx, obj) + }, + nil, + func(ctx context.Context, selections ast.SelectionSet, v *valkey.ValkeyVersion) graphql.Marshaler { + return ec.marshalNValkeyVersion2ᚖgithubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋvalkeyᚐValkeyVersion(ctx, selections, v) + }, + true, + true, + ) +} +func (ec *executionContext) fieldContext_Valkey_version(_ context.Context, field graphql.CollectedField) (fc *graphql.FieldContext, err error) { + fc = &graphql.FieldContext{ + Object: "Valkey", + Field: field, + IsMethod: true, + IsResolver: true, + Child: func(ctx context.Context, field graphql.CollectedField) (*graphql.FieldContext, error) { + return ec.childFields_ValkeyVersion(ctx, field) + }, + } + return fc, nil +} + func (ec *executionContext) _Valkey_issues(ctx context.Context, field graphql.CollectedField, obj *valkey.Valkey) (ret graphql.Marshaler) { return graphql.ResolveField( ctx, @@ -2472,6 +2505,52 @@ func (ec *executionContext) fieldContext_ValkeyUpdatedActivityLogEntryDataUpdate return graphql.NewScalarFieldContext("ValkeyUpdatedActivityLogEntryDataUpdatedField", field, false, false, errors.New("field of type String does not have child fields")) } +func (ec *executionContext) _ValkeyVersion_actual(ctx context.Context, field graphql.CollectedField, obj *valkey.ValkeyVersion) (ret graphql.Marshaler) { + return graphql.ResolveField( + ctx, + ec.OperationContext, + field, + func(ctx context.Context, field graphql.CollectedField) (*graphql.FieldContext, error) { + return ec.fieldContext_ValkeyVersion_actual(ctx, field) + }, + func(ctx context.Context) (any, error) { + return obj.Actual, nil + }, + nil, + func(ctx context.Context, selections ast.SelectionSet, v *string) graphql.Marshaler { + return ec.marshalOString2ᚖstring(ctx, selections, v) + }, + true, + false, + ) +} +func (ec *executionContext) fieldContext_ValkeyVersion_actual(_ context.Context, field graphql.CollectedField) (fc *graphql.FieldContext, err error) { + return graphql.NewScalarFieldContext("ValkeyVersion", field, false, false, errors.New("field of type String does not have child fields")) +} + +func (ec *executionContext) _ValkeyVersion_desiredMajor(ctx context.Context, field graphql.CollectedField, obj *valkey.ValkeyVersion) (ret graphql.Marshaler) { + return graphql.ResolveField( + ctx, + ec.OperationContext, + field, + func(ctx context.Context, field graphql.CollectedField) (*graphql.FieldContext, error) { + return ec.fieldContext_ValkeyVersion_desiredMajor(ctx, field) + }, + func(ctx context.Context) (any, error) { + return obj.DesiredMajor, nil + }, + nil, + func(ctx context.Context, selections ast.SelectionSet, v valkey.ValkeyMajorVersion) graphql.Marshaler { + return ec.marshalNValkeyMajorVersion2githubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋvalkeyᚐValkeyMajorVersion(ctx, selections, v) + }, + true, + true, + ) +} +func (ec *executionContext) fieldContext_ValkeyVersion_desiredMajor(_ context.Context, field graphql.CollectedField) (fc *graphql.FieldContext, err error) { + return graphql.NewScalarFieldContext("ValkeyVersion", field, false, false, errors.New("field of type ValkeyMajorVersion does not have child fields")) +} + // endregion **************************** field.gotpl ***************************** // region **************************** input.gotpl ***************************** @@ -2545,7 +2624,7 @@ func (ec *executionContext) unmarshalInputCreateValkeyInput(ctx context.Context, asMap[k] = v } - fieldsInOrder := [...]string{"name", "environmentName", "teamSlug", "tier", "memory", "maxMemoryPolicy", "notifyKeyspaceEvents", "databases"} + fieldsInOrder := [...]string{"name", "environmentName", "teamSlug", "tier", "memory", "version", "maxMemoryPolicy", "notifyKeyspaceEvents", "databases"} for _, k := range fieldsInOrder { v, ok := asMap[k] if !ok { @@ -2587,6 +2666,13 @@ func (ec *executionContext) unmarshalInputCreateValkeyInput(ctx context.Context, return it, err } it.Memory = data + case "version": + ctx := graphql.WithPathContext(ctx, graphql.NewPathWithField("version")) + data, err := ec.unmarshalOValkeyMajorVersion2ᚖgithubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋvalkeyᚐValkeyMajorVersion(ctx, v) + if err != nil { + return it, err + } + it.Version = data case "maxMemoryPolicy": ctx := graphql.WithPathContext(ctx, graphql.NewPathWithField("maxMemoryPolicy")) data, err := ec.unmarshalOValkeyMaxMemoryPolicy2ᚖgithubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋvalkeyᚐValkeyMaxMemoryPolicy(ctx, v) @@ -2668,7 +2754,7 @@ func (ec *executionContext) unmarshalInputUpdateValkeyInput(ctx context.Context, asMap[k] = v } - fieldsInOrder := [...]string{"name", "environmentName", "teamSlug", "tier", "memory", "maxMemoryPolicy", "notifyKeyspaceEvents", "databases", "labels"} + fieldsInOrder := [...]string{"name", "environmentName", "teamSlug", "tier", "memory", "version", "maxMemoryPolicy", "notifyKeyspaceEvents", "databases", "labels"} for _, k := range fieldsInOrder { v, ok := asMap[k] if !ok { @@ -2710,6 +2796,13 @@ func (ec *executionContext) unmarshalInputUpdateValkeyInput(ctx context.Context, return it, err } it.Memory = data + case "version": + ctx := graphql.WithPathContext(ctx, graphql.NewPathWithField("version")) + data, err := ec.unmarshalNValkeyMajorVersion2githubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋvalkeyᚐValkeyMajorVersion(ctx, v) + if err != nil { + return it, err + } + it.Version = data case "maxMemoryPolicy": ctx := graphql.WithPathContext(ctx, graphql.NewPathWithField("maxMemoryPolicy")) data, err := ec.unmarshalOValkeyMaxMemoryPolicy2ᚖgithubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋvalkeyᚐValkeyMaxMemoryPolicy(ctx, v) @@ -3295,6 +3388,42 @@ func (ec *executionContext) _Valkey(ctx context.Context, sel ast.SelectionSet, o if out.Values[i] == graphql.Null { atomic.AddUint32(&out.Invalids, 1) } + case "version": + field := field + + innerFunc := func(ctx context.Context, fs *graphql.FieldSet) (res graphql.Marshaler) { + defer func() { + if r := recover(); r != nil { + ec.Error(ctx, ec.Recover(ctx, r)) + } + }() + res = ec._Valkey_version(ctx, field, obj) + if res == graphql.Null { + atomic.AddUint32(&fs.Invalids, 1) + } + return res + } + + if field.Deferrable != nil { + dfs, ok := deferred[field.Deferrable.Label] + di := 0 + if ok { + dfs.AddField(field) + di = len(dfs.Values) - 1 + } else { + dfs = graphql.NewFieldSet([]graphql.CollectedField{field}) + deferred[field.Deferrable.Label] = dfs + } + dfs.Concurrently(di, func(ctx context.Context) graphql.Marshaler { + return innerFunc(ctx, dfs) + }) + + // don't run the out.Concurrently() call below + out.Values[i] = graphql.Null + continue + } + + out.Concurrently(i, func(ctx context.Context) graphql.Marshaler { return innerFunc(ctx, out) }) case "issues": field := field @@ -4418,6 +4547,47 @@ func (ec *executionContext) _ValkeyUpdatedActivityLogEntryDataUpdatedField(ctx c return out } +var valkeyVersionImplementors = []string{"ValkeyVersion"} + +func (ec *executionContext) _ValkeyVersion(ctx context.Context, sel ast.SelectionSet, obj *valkey.ValkeyVersion) graphql.Marshaler { + fields := graphql.CollectFields(ec.OperationContext, sel, valkeyVersionImplementors) + + out := graphql.NewFieldSet(fields) + deferred := make(map[string]*graphql.FieldSet) + for i, field := range fields { + switch field.Name { + case "__typename": + out.Values[i] = graphql.MarshalString("ValkeyVersion") + case "actual": + out.Values[i] = ec._ValkeyVersion_actual(ctx, field, obj) + case "desiredMajor": + out.Values[i] = ec._ValkeyVersion_desiredMajor(ctx, field, obj) + if out.Values[i] == graphql.Null { + out.Invalids++ + } + default: + panic("unknown field " + strconv.Quote(field.Name)) + } + } + out.Dispatch(ctx) + if out.Invalids > 0 { + return graphql.Null + } + + atomic.AddInt32(&ec.Deferred, int32(min(len(deferred), math.MaxInt32))) + + for label, dfs := range deferred { + ec.ProcessDeferredGroup(graphql.DeferredGroup{ + Label: label, + Path: graphql.GetPath(ctx), + FieldSet: dfs, + Context: ctx, + }) + } + + return out +} + // endregion **************************** object.gotpl **************************** // region ***************************** type.gotpl ***************************** @@ -4666,6 +4836,16 @@ func (ec *executionContext) marshalNValkeyEdge2ᚕgithubᚗcomᚋnaisᚋapiᚋin return ret } +func (ec *executionContext) unmarshalNValkeyMajorVersion2githubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋvalkeyᚐValkeyMajorVersion(ctx context.Context, v any) (valkey.ValkeyMajorVersion, error) { + var res valkey.ValkeyMajorVersion + err := res.UnmarshalGQL(v) + return res, graphql.ErrorOnPath(ctx, err) +} + +func (ec *executionContext) marshalNValkeyMajorVersion2githubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋvalkeyᚐValkeyMajorVersion(ctx context.Context, sel ast.SelectionSet, v valkey.ValkeyMajorVersion) graphql.Marshaler { + return v +} + func (ec *executionContext) unmarshalNValkeyMemory2githubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋvalkeyᚐValkeyMemory(ctx context.Context, v any) (valkey.ValkeyMemory, error) { var res valkey.ValkeyMemory err := res.UnmarshalGQL(v) @@ -4772,6 +4952,20 @@ func (ec *executionContext) marshalNValkeyUpdatedActivityLogEntryDataUpdatedFiel return ec._ValkeyUpdatedActivityLogEntryDataUpdatedField(ctx, sel, v) } +func (ec *executionContext) marshalNValkeyVersion2githubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋvalkeyᚐValkeyVersion(ctx context.Context, sel ast.SelectionSet, v valkey.ValkeyVersion) graphql.Marshaler { + return ec._ValkeyVersion(ctx, sel, &v) +} + +func (ec *executionContext) marshalNValkeyVersion2ᚖgithubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋvalkeyᚐValkeyVersion(ctx context.Context, sel ast.SelectionSet, v *valkey.ValkeyVersion) graphql.Marshaler { + if v == nil { + if !graphql.HasFieldError(ctx, graphql.GetFieldContext(ctx)) { + graphql.AddErrorf(ctx, "the requested element is null which the schema does not allow") + } + return graphql.Null + } + return ec._ValkeyVersion(ctx, sel, v) +} + func (ec *executionContext) unmarshalOValkeyAccessOrder2ᚖgithubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋvalkeyᚐValkeyAccessOrder(ctx context.Context, v any) (*valkey.ValkeyAccessOrder, error) { if v == nil { return nil, nil @@ -4795,6 +4989,22 @@ func (ec *executionContext) unmarshalOValkeyFilter2ᚖgithubᚗcomᚋnaisᚋapi return &res, graphql.ErrorOnPath(ctx, err) } +func (ec *executionContext) unmarshalOValkeyMajorVersion2ᚖgithubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋvalkeyᚐValkeyMajorVersion(ctx context.Context, v any) (*valkey.ValkeyMajorVersion, error) { + if v == nil { + return nil, nil + } + var res = new(valkey.ValkeyMajorVersion) + err := res.UnmarshalGQL(v) + return res, graphql.ErrorOnPath(ctx, err) +} + +func (ec *executionContext) marshalOValkeyMajorVersion2ᚖgithubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋvalkeyᚐValkeyMajorVersion(ctx context.Context, sel ast.SelectionSet, v *valkey.ValkeyMajorVersion) graphql.Marshaler { + if v == nil { + return graphql.Null + } + return v +} + func (ec *executionContext) unmarshalOValkeyMaxMemoryPolicy2githubᚗcomᚋnaisᚋapiᚋinternalᚋpersistenceᚋvalkeyᚐValkeyMaxMemoryPolicy(ctx context.Context, v any) (valkey.ValkeyMaxMemoryPolicy, error) { var res valkey.ValkeyMaxMemoryPolicy err := res.UnmarshalGQL(v) diff --git a/internal/graph/opensearch.resolvers.go b/internal/graph/opensearch.resolvers.go index 06e8dd9ac..3890c7913 100644 --- a/internal/graph/opensearch.resolvers.go +++ b/internal/graph/opensearch.resolvers.go @@ -59,7 +59,7 @@ func (r *openSearchResolver) TeamEnvironment(ctx context.Context, obj *opensearc } func (r *openSearchResolver) State(ctx context.Context, obj *opensearch.OpenSearch) (opensearch.OpenSearchState, error) { - return opensearch.State(ctx, obj) + return opensearch.State(obj), nil } func (r *openSearchResolver) Workload(ctx context.Context, obj *opensearch.OpenSearch) (workload.Workload, error) { @@ -76,7 +76,7 @@ func (r *openSearchResolver) Access(ctx context.Context, obj *opensearch.OpenSea } func (r *openSearchResolver) Version(ctx context.Context, obj *opensearch.OpenSearch) (*opensearch.OpenSearchVersion, error) { - return opensearch.GetOpenSearchVersion(ctx, obj) + return opensearch.GetOpenSearchVersion(obj) } func (r *openSearchResolver) Issues(ctx context.Context, obj *opensearch.OpenSearch, first *int, after *pagination.Cursor, last *int, before *pagination.Cursor, orderBy *issue.IssueOrder, filter *issue.ResourceIssueFilter) (*issue.IssueConnection, error) { diff --git a/internal/graph/schema/opensearch.graphqls b/internal/graph/schema/opensearch.graphqls index 0c6887512..03da51aa4 100644 --- a/internal/graph/schema/opensearch.graphqls +++ b/internal/graph/schema/opensearch.graphqls @@ -242,16 +242,14 @@ enum OpenSearchMemory { } enum OpenSearchMajorVersion { - "OpenSearch Version 3.6.x" + "OpenSearch Version 3.6 LTS" V3_6 "OpenSearch Version 3.3.x" - V3_3 - "OpenSearch Version 2.19.x" + V3_3 @deprecated(reason: "Vendor support disappears 2027-02-01") + "OpenSearch Version 2.19 LTS" V2_19 - "OpenSearch Version 2.17.x" - V2 - "OpenSearch Version 1.3.x" - V1 + "OpenSearch Version 2.19 LTS - backwards compatible" + V2 @deprecated(reason: "Use `V2_19` instead") } input CreateOpenSearchInput { @@ -266,7 +264,7 @@ input CreateOpenSearchInput { "Available memory for the OpenSearch instance." memory: OpenSearchMemory! "Major version of the OpenSearch instance." - version: OpenSearchMajorVersion! + version: OpenSearchMajorVersion "Available storage in GB." storageGB: Int! } diff --git a/internal/graph/schema/valkey.graphqls b/internal/graph/schema/valkey.graphqls index 371979d4a..b4a2c0147 100644 --- a/internal/graph/schema/valkey.graphqls +++ b/internal/graph/schema/valkey.graphqls @@ -59,6 +59,24 @@ type TeamInventoryCountValkeys { total: Int! } +"Version information for a Valkey instance." +type ValkeyVersion { + "The full version string of the Valkey instance. This will be available after the instance is created." + actual: String + "The desired major version of the Valkey instance." + desiredMajor: ValkeyMajorVersion! +} + +"Major version of a Valkey instance." +enum ValkeyMajorVersion { + "Valkey Version 9.1.x" + V9_1 + "Valkey Version 9.0.x" + V9_0 @deprecated(reason: "No longer supported by back-end services") + "Valkey Version 8.1.x" + V8_1 +} + type Valkey implements Persistence & Node { id: ID! name: String! @@ -87,6 +105,8 @@ type Valkey implements Persistence & Node { notifyKeyspaceEvents: String "Number of databases the Valkey instance is configured with. Default is 16. Minimum 1, maximum 128. Changing this will cause a restart of the Valkey service." databases: Int! + "Fetch version for the Valkey instance." + version: ValkeyVersion! "Issues that affects the instance." issues( "Get the first n items in the connection. This can be used in combination with the after parameter." @@ -284,6 +304,8 @@ input CreateValkeyInput { tier: ValkeyTier! "Available memory for the Valkey instance." memory: ValkeyMemory! + "Major version of the Valkey instance." + version: ValkeyMajorVersion "Maximum memory policy for the Valkey instance." maxMemoryPolicy: ValkeyMaxMemoryPolicy "Configure keyspace notifications for the Valkey instance. See https://valkey.io/topics/notifications/ for details." @@ -308,6 +330,8 @@ input UpdateValkeyInput { tier: ValkeyTier! "Available memory for the Valkey instance." memory: ValkeyMemory! + "Major version of the Valkey instance." + version: ValkeyMajorVersion! "Maximum memory policy for the Valkey instance." maxMemoryPolicy: ValkeyMaxMemoryPolicy "Configure keyspace notifications for the Valkey instance. See https://valkey.io/topics/notifications/ for details." diff --git a/internal/graph/valkey.resolvers.go b/internal/graph/valkey.resolvers.go index a276b6973..e6b6f198a 100644 --- a/internal/graph/valkey.resolvers.go +++ b/internal/graph/valkey.resolvers.go @@ -91,7 +91,11 @@ func (r *valkeyResolver) Workload(ctx context.Context, obj *valkey.Valkey) (work } func (r *valkeyResolver) State(ctx context.Context, obj *valkey.Valkey) (valkey.ValkeyState, error) { - return valkey.State(ctx, obj) + return valkey.State(obj), nil +} + +func (r *valkeyResolver) Version(ctx context.Context, obj *valkey.Valkey) (*valkey.ValkeyVersion, error) { + return valkey.GetValkeyVersion(obj) } func (r *valkeyResolver) Issues(ctx context.Context, obj *valkey.Valkey, first *int, after *pagination.Cursor, last *int, before *pagination.Cursor, orderBy *issue.IssueOrder, filter *issue.ResourceIssueFilter) (*issue.IssueConnection, error) { diff --git a/internal/kubernetes/fake/fake.go b/internal/kubernetes/fake/fake.go index 234fdfc1f..094023849 100644 --- a/internal/kubernetes/fake/fake.go +++ b/internal/kubernetes/fake/fake.go @@ -219,6 +219,17 @@ func NewDynamicClient(scheme *runtime.Scheme) *dynfake.FakeDynamicClient { nais_io_v1alpha1.GroupVersion.WithResource("tunnels"): "TunnelList", }) + client.PrependReactor("create", "*", func(action k8stesting.Action) (handled bool, ret runtime.Object, err error) { + createAction, ok := action.(k8stesting.CreateAction) + if !ok { + return false, nil, nil + } + if o, ok := createAction.GetObject().(*unstructured.Unstructured); ok { + observeAivenVersion(o) + } + return false, nil, nil + }) + client.PrependReactor("patch", "*", func(action k8stesting.Action) (handled bool, ret runtime.Object, err error) { patchAction, ok := action.(k8stesting.PatchAction) if !ok { @@ -379,3 +390,22 @@ func newDynamicClient(scheme *runtime.Scheme, objs ...runtime.Object) dynamic.In AddObjectToDynamicClient(scheme, fc, objs...) return fc } + +// observeAivenVersion stands in for the Aiven operator's first reconcile, which records +// the version the service reports in status. Only the create path needs it: fixtures +// declare status.version themselves, and the ones that deliberately omit it are the +// cases where the operator has not observed the service yet. +func observeAivenVersion(o *unstructured.Unstructured) { + kind := strings.ToLower(o.GetKind()) + if kind != "valkey" && kind != "opensearch" { + return + } + if v, _, _ := unstructured.NestedString(o.Object, "status", "version"); v != "" { + return + } + pinned, ok, _ := unstructured.NestedString(o.Object, "spec", "userConfig", kind+"_version") + if !ok || pinned == "" { + return + } + _ = unstructured.SetNestedField(o.Object, pinned, "status", "version") +} diff --git a/internal/persistence/aivenversion/match.go b/internal/persistence/aivenversion/match.go new file mode 100644 index 000000000..83346d5c4 --- /dev/null +++ b/internal/persistence/aivenversion/match.go @@ -0,0 +1,206 @@ +// Package aivenversion resolves the version string Aiven reports for a service to the +// GraphQL enum value that names it. OpenSearch and Valkey declare separate enums but +// answer the same question, so the rule lives here rather than in both. +package aivenversion + +import ( + "fmt" + "slices" + "strconv" + "strings" +) + +// StatusVersion is where the Aiven operator records the version a service is actually +// running, as of aiven/aiven-operator#1280. Both kinds read the same path, so it is +// declared once here rather than beside each kind's spec paths. +var StatusVersion = []string{"status", "version"} + +// Declared returns every version the upgrade-path map knows about, sorted. Deriving the +// list rather than maintaining a second one beside it is what keeps validation, matching +// and the newest-version lookup from disagreeing. The sort matters: callers scan the +// result in order, so a raw map range would make matching depend on iteration order. +func Declared[T ~string, U any](paths map[T]U) []T { + versions := make([]T, 0, len(paths)) + for v := range paths { + versions = append(versions, v) + } + slices.Sort(versions) + return versions +} + +// Match picks the enum value that names the version Aiven reports. +// +// Enum names carry the version they stand for: "V3_6" is 3.6, and a bare "V2" is major +// 2 with no minor of its own. Matching never crosses a major, because a major bump is +// never a silent equivalence. Within one major: +// - an exact minor match wins; +// - where two or more minors are declared, a reported minor between two of them +// resolves to the lower, and one outside their range to the closest; +// - a bare VX catches reported minors below the VX_Y it upgrades to. +// +// Anything left over is an error, so the API never claims a version it cannot name. +func Match[T ~string, U ~[]T](reported string, declared []T, upgradesTo map[T]U) (T, error) { + var zero T + + major, minor, err := parseReported(reported) + if err != nil { + return zero, err + } + + var withMinor []version[T] + var bare *version[T] + for _, value := range declared { + v, err := parseDeclared(value) + if err != nil { + return zero, err + } + if v.major != major { + continue + } + if v.hasMinor { + withMinor = append(withMinor, v) + } else { + bare = &v + } + } + + slices.SortFunc(withMinor, func(a, b version[T]) int { return a.minor - b.minor }) + + for _, v := range withMinor { + if v.minor == minor { + return v.value, nil + } + } + + // A lone declared minor is a point, not a range, so there is nothing to interpolate + // between or fall off the end of. + if len(withMinor) >= 2 { + lowest, highest := withMinor[0], withMinor[len(withMinor)-1] + switch { + case minor < lowest.minor: + return lowest.value, nil + case minor > highest.minor: + return highest.value, nil + } + for i := len(withMinor) - 1; i >= 0; i-- { + if withMinor[i].minor < minor { + return withMinor[i].value, nil + } + } + } + + if bare != nil { + if boundary, ok := bareBoundary(*bare, upgradesTo); ok && minor < boundary { + return bare.value, nil + } + } + + return zero, fmt.Errorf("unsupported Aiven version: %q", reported) +} + +// MatchPin resolves the version pinned in a Kubernetes resource. A pin is written by this +// API rather than reported by Aiven, so it is matched against the declared versions +// directly instead of through Match. Older releases pinned a bare major such as "2", so +// both that shape and the current "2.19" occur in the wild. +func MatchPin[T ~string](pin string, declared []T, aivenString func(T) (string, error)) (T, error) { + var zero T + + for _, value := range declared { + s, err := aivenString(value) + if err == nil && s == pin { + return qualified(value, declared, aivenString), nil + } + } + + if bare := T("V" + pin); slices.Contains(declared, bare) { + return qualified(bare, declared, aivenString), nil + } + + return zero, fmt.Errorf("unsupported pinned version: %q", pin) +} + +// qualified returns the VX_Y spelling of a bare VX where one is declared for the same +// Aiven version. The two are one version under two names, and the qualified one is what +// callers should see. Anything already carrying a minor is returned untouched. +func qualified[T ~string](value T, declared []T, aivenString func(T) (string, error)) T { + if v, err := parseDeclared(value); err != nil || v.hasMinor { + return value + } + + want, err := aivenString(value) + if err != nil { + return value + } + + for _, other := range declared { + if other == value { + continue + } + if v, err := parseDeclared(other); err != nil || !v.hasMinor { + continue + } + if s, err := aivenString(other); err == nil && s == want { + return other + } + } + return value +} + +type version[T ~string] struct { + value T + major int + minor int + hasMinor bool +} + +// bareBoundary returns the minor below which a bare VX claims a reported version. VX +// has no minor itself, so the VX_Y it upgrades to supplies the boundary. +func bareBoundary[T ~string, U ~[]T](bare version[T], upgradesTo map[T]U) (int, bool) { + boundary := -1 + for _, target := range upgradesTo[bare.value] { + v, err := parseDeclared(target) + if err != nil || !v.hasMinor || v.major != bare.major { + continue + } + if boundary < 0 || v.minor < boundary { + boundary = v.minor + } + } + return boundary, boundary >= 0 +} + +func parseDeclared[T ~string](value T) (version[T], error) { + majorPart, minorPart, hasMinor := strings.Cut(strings.TrimPrefix(string(value), "V"), "_") + + major, err := strconv.Atoi(majorPart) + if err != nil { + return version[T]{}, fmt.Errorf("parsing major version from %q: %w", value, err) + } + + v := version[T]{value: value, major: major, hasMinor: hasMinor} + if hasMinor { + if v.minor, err = strconv.Atoi(minorPart); err != nil { + return version[T]{}, fmt.Errorf("parsing minor version from %q: %w", value, err) + } + } + return v, nil +} + +func parseReported(reported string) (int, int, error) { + parts := strings.SplitN(reported, ".", 3) + if len(parts) < 2 { + return 0, 0, fmt.Errorf("unsupported Aiven version %q: no minor version", reported) + } + + major, err := strconv.Atoi(parts[0]) + if err != nil { + return 0, 0, fmt.Errorf("parsing major version from %q: %w", reported, err) + } + + minor, err := strconv.Atoi(parts[1]) + if err != nil { + return 0, 0, fmt.Errorf("parsing minor version from %q: %w", reported, err) + } + + return major, minor, nil +} diff --git a/internal/persistence/opensearch/dataloader.go b/internal/persistence/opensearch/dataloader.go index 72a6ab567..b9d058293 100644 --- a/internal/persistence/opensearch/dataloader.go +++ b/internal/persistence/opensearch/dataloader.go @@ -3,12 +3,8 @@ package opensearch import ( "context" - "github.com/nais/api/internal/graph/loader" "github.com/nais/api/internal/kubernetes/watcher" - "github.com/nais/api/internal/thirdparty/aiven" "github.com/sirupsen/logrus" - "github.com/sourcegraph/conc/pool" - "github.com/vikstrous/dataloadgen" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/runtime/schema" ) @@ -17,13 +13,8 @@ type ctxKey int const loadersKey ctxKey = iota -type AivenDataLoaderKey struct { - Project string - ServiceName string -} - -func NewLoaderContext(ctx context.Context, tenantName string, watcher *watcher.Watcher[*OpenSearch], aivenClient aiven.AivenClient, logger logrus.FieldLogger) context.Context { - return context.WithValue(ctx, loadersKey, newLoaders(tenantName, watcher, aivenClient, logger)) +func NewLoaderContext(ctx context.Context, tenantName string, watcher *watcher.Watcher[*OpenSearch], logger logrus.FieldLogger) context.Context { + return context.WithValue(ctx, loadersKey, newLoaders(tenantName, watcher, logger)) } func NewWatcher(ctx context.Context, mgr *watcher.Manager) *watcher.Watcher[*OpenSearch] { @@ -47,60 +38,21 @@ func fromContext(ctx context.Context) *loaders { } type loaders struct { - client *client - watcher *watcher.Watcher[*OpenSearch] - versionLoader *dataloadgen.Loader[*AivenDataLoaderKey, string] - tenantName string - aivenClient aiven.AivenClient - log logrus.FieldLogger + client *client + watcher *watcher.Watcher[*OpenSearch] + tenantName string + log logrus.FieldLogger } -func newLoaders(tenantName string, watcher *watcher.Watcher[*OpenSearch], aivenClient aiven.AivenClient, logger logrus.FieldLogger) *loaders { +func newLoaders(tenantName string, watcher *watcher.Watcher[*OpenSearch], logger logrus.FieldLogger) *loaders { client := &client{ watcher: watcher, } - versionLoader := &dataloader{aivenClient: aivenClient, log: logger} - return &loaders{ - client: client, - watcher: watcher, - tenantName: tenantName, - versionLoader: dataloadgen.NewLoader(versionLoader.getVersions, loader.DefaultDataLoaderOptions...), - aivenClient: aivenClient, - log: logger, - } -} - -type dataloader struct { - aivenClient aiven.AivenClient - log logrus.FieldLogger -} - -func (l dataloader) getVersions(ctx context.Context, aivenDataLoaderKeys []*AivenDataLoaderKey) ([]string, []error) { - wg := pool.New().WithContext(ctx) - rets := make([]string, len(aivenDataLoaderKeys)) - errs := make([]error, len(aivenDataLoaderKeys)) - - for i, pair := range aivenDataLoaderKeys { - wg.Go(func(ctx context.Context) error { - res, err := l.aivenClient.ServiceGet(ctx, pair.Project, pair.ServiceName) - if err != nil { - errs[i] = err - } else { - if res.Metadata != nil { - if version, ok := res.Metadata["opensearch_version"]; ok { - rets[i] = version.(string) - } - } - } - return nil - }) + client: client, + watcher: watcher, + tenantName: tenantName, + log: logger, } - - if err := wg.Wait(); err != nil { - l.log.WithError(err).Error("error waiting for dataloader") - } - - return rets, errs } diff --git a/internal/persistence/opensearch/models.go b/internal/persistence/opensearch/models.go index 6ca38f8a3..3912081c1 100644 --- a/internal/persistence/opensearch/models.go +++ b/internal/persistence/opensearch/models.go @@ -10,17 +10,18 @@ import ( "sync" "github.com/99designs/gqlgen/graphql" + "github.com/Masterminds/semver/v3" "github.com/nais/api/internal/graph/apierror" "github.com/nais/api/internal/graph/ident" "github.com/nais/api/internal/graph/model" "github.com/nais/api/internal/graph/pagination" "github.com/nais/api/internal/kubernetes" "github.com/nais/api/internal/persistence/aivencredentials" + "github.com/nais/api/internal/persistence/aivenversion" "github.com/nais/api/internal/slug" "github.com/nais/api/internal/validate" "github.com/nais/api/internal/workload" aiven_io_v1alpha1 "github.com/nais/liberator/pkg/apis/aiven.io/v1alpha1" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/runtime/schema" @@ -69,6 +70,7 @@ type OpenSearch struct { WorkloadReference *workload.Reference `json:"-"` AivenProject string `json:"-"` MajorVersion OpenSearchMajorVersion `json:"-"` + ReportedVersion string `json:"-"` } func (OpenSearch) IsPersistence() {} @@ -109,8 +111,7 @@ type OpenSearchAccess struct { } type OpenSearchStatus struct { - State string `json:"state"` - Conditions []metav1.Condition `json:"conditions"` + State string `json:"state"` } type OpenSearchOrder struct { @@ -197,12 +198,12 @@ func toOpenSearch(u *unstructured.Unstructured, envName string) (*OpenSearch, er name = strings.TrimPrefix(obj.GetName(), NamePrefix(slug.Slug(obj.GetNamespace()))) } - majorVersion := OpenSearchMajorVersion("") - if v, found, _ := unstructured.NestedString(u.Object, specOpenSearchVersion...); found { - version, err := OpenSearchMajorVersionFromAivenString(v) - if err == nil { - majorVersion = version - } + reportedVersion, _, _ := unstructured.NestedString(u.Object, aivenversion.StatusVersion...) + + pinnedVersion, _, _ := unstructured.NestedString(u.Object, specOpenSearchVersion...) + majorVersion, err := aivenversion.MatchPin(pinnedVersion, allMajorVersions, OpenSearchMajorVersion.ToAivenString) + if err != nil { + majorVersion = "" } // default to minimum storage capacity for the selected plan, in case the field is not set explicitly @@ -219,16 +220,16 @@ func toOpenSearch(u *unstructured.Unstructured, envName string) (*OpenSearch, er EnvironmentName: envName, TerminationProtection: terminationProtection, Status: &OpenSearchStatus{ - Conditions: obj.Status.Conditions, - State: obj.Status.State, + State: obj.Status.State, }, TeamSlug: slug.Slug(obj.GetNamespace()), WorkloadReference: workload.ReferenceFromOwnerReferences(obj.GetOwnerReferences()), AivenProject: obj.Spec.Project, Tier: machine.Tier, Memory: machine.Memory, - MajorVersion: majorVersion, StorageGB: storageGB, + MajorVersion: majorVersion, + ReportedVersion: reportedVersion, Labels: model.UserLabels(obj.GetLabels()), }, nil } @@ -270,10 +271,10 @@ func (o *OpenSearchMetadataInput) ValidationErrors(ctx context.Context) *validat type OpenSearchInput struct { OpenSearchMetadataInput - Tier OpenSearchTier `json:"tier"` - Memory OpenSearchMemory `json:"memory"` - Version OpenSearchMajorVersion `json:"version"` - StorageGB StorageGB `json:"storageGB"` + Tier OpenSearchTier `json:"tier"` + Memory OpenSearchMemory `json:"memory"` + Version *OpenSearchMajorVersion `json:"version,omitempty"` + StorageGB StorageGB `json:"storageGB"` } func (o *OpenSearchInput) Validate(ctx context.Context) error { @@ -289,8 +290,12 @@ func (o *OpenSearchInput) ValidationErrors(ctx context.Context) *validate.Valida if !o.Memory.IsValid() { verr.Add("memory", "Invalid OpenSearch memory: %s.", o.Memory) } - if !o.Version.IsValid() { - verr.Add("version", "Invalid OpenSearch version: %s.", o.Version.String()) + if o.Version != nil { + if !o.Version.IsValid() { + verr.Add("version", "Invalid OpenSearch version: %s.", o.Version.String()) + } else if reason, deprecated := o.Version.DeprecationReason(); deprecated { + verr.Add("version", "OpenSearch version %s is deprecated: %s.", o.Version.String(), reason) + } } machine, err := machineTypeFromTierAndMemory(o.Tier, o.Memory) @@ -371,7 +376,6 @@ type CreateOpenSearchPayload struct { type OpenSearchMajorVersion string const ( - OpenSearchMajorVersionV1 OpenSearchMajorVersion = "V1" OpenSearchMajorVersionV2 OpenSearchMajorVersion = "V2" OpenSearchMajorVersionV2_19 OpenSearchMajorVersion = "V2_19" OpenSearchMajorVersionV3_3 OpenSearchMajorVersion = "V3_3" @@ -389,7 +393,6 @@ func (u upgradePath) String() string { } var upgradePaths = map[OpenSearchMajorVersion]upgradePath{ - OpenSearchMajorVersionV1: {OpenSearchMajorVersionV2, OpenSearchMajorVersionV2_19}, OpenSearchMajorVersionV2: {OpenSearchMajorVersionV2_19}, OpenSearchMajorVersionV2_19: {OpenSearchMajorVersionV3_3, OpenSearchMajorVersionV3_6}, OpenSearchMajorVersionV3_3: {OpenSearchMajorVersionV3_6}, @@ -413,17 +416,52 @@ func (e OpenSearchMajorVersion) ValidateUpgradePath(other OpenSearchMajorVersion return apierror.Errorf("Cannot change OpenSearch version from %v to %v. New version must be one of [%s]", other, e, path) } -func (e OpenSearchMajorVersion) IsValid() bool { - switch e { - case - OpenSearchMajorVersionV1, - OpenSearchMajorVersionV2, - OpenSearchMajorVersionV2_19, - OpenSearchMajorVersionV3_3, - OpenSearchMajorVersionV3_6: - return true +// newestMajorVersion returns the highest version the codebase knows about, derived +// from the versions themselves so that adding one to upgradePaths is enough. +func newestMajorVersion() (OpenSearchMajorVersion, error) { + var newest OpenSearchMajorVersion + var highest *semver.Version + + for v := range upgradePaths { + s, err := v.ToAivenString() + if err != nil { + return "", err + } + parsed, err := semver.NewVersion(s) + if err != nil { + return "", fmt.Errorf("parsing OpenSearch version %q: %w", s, err) + } + if highest == nil || parsed.GreaterThan(highest) { + highest, newest = parsed, v + } } - return false + + if newest == "" { + return "", fmt.Errorf("no OpenSearch versions defined") + } + return newest, nil +} + +// allMajorVersions is derived from upgradePaths rather than maintained beside it, so the +// two cannot disagree. +var allMajorVersions = aivenversion.Declared(upgradePaths) + +func (e OpenSearchMajorVersion) IsValid() bool { + return slices.Contains(allMajorVersions, e) +} + +// deprecatedMajorVersions name versions Aiven still runs, so they stay valid for reads +// and remain in the GraphQL enum. No client may choose one. Each carries its own reason +// because versions are deprecated for different causes, and refusing without saying why +// leaves the caller guessing. +var deprecatedMajorVersions = map[OpenSearchMajorVersion]string{ + OpenSearchMajorVersionV2: "use V2_19 instead", + OpenSearchMajorVersionV3_3: "vendor support disappears 2027-02-01", +} + +func (e OpenSearchMajorVersion) DeprecationReason() (string, bool) { + reason, ok := deprecatedMajorVersions[e] + return reason, ok } func (e OpenSearchMajorVersion) String() string { @@ -447,14 +485,10 @@ func (e OpenSearchMajorVersion) MarshalGQL(w io.Writer) { fmt.Fprint(w, strconv.Quote(e.String())) } -// ToAivenString returns the version string without the "V" prefix, e.g. "2" or "1". +// ToAivenString returns the version as Aiven writes it in userConfig, e.g. "2.19". func (e OpenSearchMajorVersion) ToAivenString() (string, error) { switch e { - case OpenSearchMajorVersionV1: - return "1", nil - case OpenSearchMajorVersionV2: - return "2", nil - case OpenSearchMajorVersionV2_19: + case OpenSearchMajorVersionV2, OpenSearchMajorVersionV2_19: return "2.19", nil case OpenSearchMajorVersionV3_3: return "3.3", nil @@ -466,20 +500,7 @@ func (e OpenSearchMajorVersion) ToAivenString() (string, error) { } func OpenSearchMajorVersionFromAivenString(s string) (OpenSearchMajorVersion, error) { - switch { - case strings.HasPrefix(s, "1"): - return OpenSearchMajorVersionV1, nil - case strings.HasPrefix(s, "2.19"): - return OpenSearchMajorVersionV2_19, nil - case strings.HasPrefix(s, "2"): - return OpenSearchMajorVersionV2, nil - case strings.HasPrefix(s, "3.3"): - return OpenSearchMajorVersionV3_3, nil - case strings.HasPrefix(s, "3.6"): - return OpenSearchMajorVersionV3_6, nil - default: - return "", fmt.Errorf("unsupported Aiven OpenSearch version: %q", s) - } + return aivenversion.Match(s, allMajorVersions, upgradePaths) } type OpenSearchMemory string diff --git a/internal/persistence/opensearch/queries.go b/internal/persistence/opensearch/queries.go index f29857ece..3ceeae070 100644 --- a/internal/persistence/opensearch/queries.go +++ b/internal/persistence/opensearch/queries.go @@ -48,28 +48,20 @@ func Get(ctx context.Context, teamSlug slug.Slug, environment, name string) (*Op return fromContext(ctx).client.watcher.Get(environment, teamSlug.String(), name) } -func State(ctx context.Context, os *OpenSearch) (OpenSearchState, error) { - s, err := fromContext(ctx).aivenClient.ServiceGet(ctx, os.AivenProject, os.FullyQualifiedName()) - if err != nil { - // The OpenSearch instance may not have been created in Aiven yet, or it has been deleted. - // In both cases, we return "unknown" state rather than an error. - if aiven.IsNotFound(err) { - return OpenSearchStateUnknown, nil - } - return OpenSearchStateUnknown, err - } - - switch s.State { +// State reports the state the operator last observed. An instance Aiven has not created +// yet, or has already deleted, has no state recorded and reads as unknown. +func State(os *OpenSearch) OpenSearchState { + switch os.Status.State { case "RUNNING": - return OpenSearchStateRunning, nil + return OpenSearchStateRunning case "REBALANCING": - return OpenSearchStateRebalancing, nil + return OpenSearchStateRebalancing case "REBUILDING": - return OpenSearchStateRebuilding, nil + return OpenSearchStateRebuilding case "POWEROFF": - return OpenSearchStatePoweroff, nil + return OpenSearchStatePoweroff default: - return OpenSearchStateUnknown, nil + return OpenSearchStateUnknown } } @@ -120,33 +112,28 @@ func ListAccess(ctx context.Context, openSearch *OpenSearch, page *pagination.Pa return pagination.NewConnection(ret, page, len(all)), nil } -func GetOpenSearchVersion(ctx context.Context, os *OpenSearch) (*OpenSearchVersion, error) { - key := AivenDataLoaderKey{ - Project: os.AivenProject, - ServiceName: os.FullyQualifiedName(), - } - - major := os.MajorVersion - var versionString *string - v, err := fromContext(ctx).versionLoader.Load(ctx, &key) - if err == nil { - versionString = new(v) - if major == "" { - mv, err := OpenSearchMajorVersionFromAivenString(v) - if err != nil { - return nil, err - } - major = mv +func GetOpenSearchVersion(os *OpenSearch) (*OpenSearchVersion, error) { + if os.ReportedVersion == "" { + // Kubernetes is only eventually consistent with Aiven: the CR exists here before + // the service exists there, and lingers after it is deleted. Aiven has no version + // to report in that window, so fall back to the version pinned in the CR rather + // than failing the whole query. Actual stays nil, which the schema documents as + // "available after the instance is created". + if os.MajorVersion == "" { + return nil, fmt.Errorf("no OpenSearch version known for service %q", os.FullyQualifiedName()) } + return &OpenSearchVersion{DesiredMajor: os.MajorVersion}, nil } - if major == "" { - major = OpenSearchMajorVersionV2 + actual := os.ReportedVersion + major, err := OpenSearchMajorVersionFromAivenString(actual) + if err != nil { + return nil, err } return &OpenSearchVersion{ DesiredMajor: major, - Actual: versionString, + Actual: &actual, }, nil } @@ -186,7 +173,15 @@ func Create(ctx context.Context, input CreateOpenSearchInput) (*CreateOpenSearch if err != nil { return nil, err } - version, err := input.Version.ToAivenString() + desired := input.Version + if desired == nil { + newest, err := newestMajorVersion() + if err != nil { + return nil, err + } + desired = &newest + } + version, err := desired.ToAivenString() if err != nil { return nil, err } @@ -269,7 +264,7 @@ func Update(ctx context.Context, input UpdateOpenSearchInput) (*UpdateOpenSearch } changes = append(changes, res...) - res, err = updateVersion(ctx, openSearch, input) + res, err = updateVersion(openSearch, input) if err != nil { return nil, err } @@ -442,32 +437,32 @@ func updatePlan(openSearch *unstructured.Unstructured, input UpdateOpenSearchInp return changes, nil } -func updateVersion(ctx context.Context, openSearch *unstructured.Unstructured, input UpdateOpenSearchInput) ([]*OpenSearchUpdatedActivityLogEntryDataUpdatedField, error) { +func updateVersion(openSearch *unstructured.Unstructured, input UpdateOpenSearchInput) ([]*OpenSearchUpdatedActivityLogEntryDataUpdatedField, error) { changes := make([]*OpenSearchUpdatedActivityLogEntryDataUpdatedField, 0) - oldVersion, found, err := unstructured.NestedString(openSearch.Object, specOpenSearchVersion...) + os, err := toOpenSearch(openSearch, input.EnvironmentName) if err != nil { return nil, err } - if !found { - os, err := toOpenSearch(openSearch, input.EnvironmentName) - if err != nil { - return nil, err - } - version, err := GetOpenSearchVersion(ctx, os) - if err != nil { - return nil, err - } - oldVersion = *version.Actual + // Only the running version can judge whether the requested change is legal, and the + // operator records it once it has observed the service. The CR's own pin says what + // was asked for, so it cannot stand in. + if os.ReportedVersion == "" { + return nil, apierror.Errorf("The running version of %q is not known yet. Try again shortly.", os.Name) } + oldVersion := os.ReportedVersion oldMajorVersion, err := OpenSearchMajorVersionFromAivenString(oldVersion) if err != nil { return nil, err } - if oldMajorVersion == input.Version { + if input.Version == nil { + return nil, fmt.Errorf("no OpenSearch version supplied") + } + + if oldMajorVersion == *input.Version { return changes, nil } @@ -476,13 +471,8 @@ func updateVersion(ctx context.Context, openSearch *unstructured.Unstructured, i } changes = append(changes, &OpenSearchUpdatedActivityLogEntryDataUpdatedField{ - Field: "version", - OldValue: func() *string { - if found { - return new(oldVersion) - } - return nil - }(), + Field: "version", + OldValue: &oldVersion, NewValue: new(input.Version.String()), }) diff --git a/internal/persistence/opensearch/sortfilter.go b/internal/persistence/opensearch/sortfilter.go index c9a059808..63d86ff49 100644 --- a/internal/persistence/opensearch/sortfilter.go +++ b/internal/persistence/opensearch/sortfilter.go @@ -22,12 +22,7 @@ func init() { return strings.Compare(a.EnvironmentName, b.EnvironmentName) }, "NAME") SortFilterOpenSearch.RegisterConcurrentSort("STATE", func(ctx context.Context, a *OpenSearch) int { - s, err := State(ctx, a) - if err != nil { - return int(OpenSearchStateUnknown) - } - - return int(s) + return int(State(a)) }, "NAME") SortFilterOpenSearch.RegisterFilter(func(ctx context.Context, v *OpenSearch, filter *OpenSearchFilter) bool { diff --git a/internal/persistence/valkey/dataloader.go b/internal/persistence/valkey/dataloader.go index 31f30a200..6e1a79f1e 100644 --- a/internal/persistence/valkey/dataloader.go +++ b/internal/persistence/valkey/dataloader.go @@ -4,7 +4,6 @@ import ( "context" "github.com/nais/api/internal/kubernetes/watcher" - "github.com/nais/api/internal/thirdparty/aiven" "github.com/sirupsen/logrus" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/runtime/schema" @@ -14,8 +13,8 @@ type ctxKey int const loadersKey ctxKey = iota -func NewLoaderContext(ctx context.Context, tenantName string, valkeyWatcher *watcher.Watcher[*Valkey], aivenClient aiven.AivenClient, logger logrus.FieldLogger) context.Context { - return context.WithValue(ctx, loadersKey, newLoaders(tenantName, valkeyWatcher, aivenClient, logger)) +func NewLoaderContext(ctx context.Context, tenantName string, valkeyWatcher *watcher.Watcher[*Valkey], logger logrus.FieldLogger) context.Context { + return context.WithValue(ctx, loadersKey, newLoaders(tenantName, valkeyWatcher, logger)) } func NewWatcher(ctx context.Context, mgr *watcher.Manager) *watcher.Watcher[*Valkey] { @@ -39,23 +38,21 @@ func fromContext(ctx context.Context) *loaders { } type loaders struct { - client *client - tenantName string - watcher *watcher.Watcher[*Valkey] - aivenClient aiven.AivenClient - log logrus.FieldLogger + client *client + tenantName string + watcher *watcher.Watcher[*Valkey] + log logrus.FieldLogger } -func newLoaders(tenantName string, watcher *watcher.Watcher[*Valkey], aivenClient aiven.AivenClient, logger logrus.FieldLogger) *loaders { +func newLoaders(tenantName string, watcher *watcher.Watcher[*Valkey], logger logrus.FieldLogger) *loaders { client := &client{ watcher: watcher, } return &loaders{ - client: client, - tenantName: tenantName, - watcher: watcher, - aivenClient: aivenClient, - log: logger, + client: client, + tenantName: tenantName, + watcher: watcher, + log: logger, } } diff --git a/internal/persistence/valkey/models.go b/internal/persistence/valkey/models.go index 2b18eb851..5006cbffe 100644 --- a/internal/persistence/valkey/models.go +++ b/internal/persistence/valkey/models.go @@ -4,20 +4,23 @@ import ( "context" "fmt" "io" + "slices" "strconv" "strings" "sync" + "github.com/Masterminds/semver/v3" + "github.com/nais/api/internal/graph/apierror" "github.com/nais/api/internal/graph/ident" "github.com/nais/api/internal/graph/model" "github.com/nais/api/internal/graph/pagination" "github.com/nais/api/internal/kubernetes" "github.com/nais/api/internal/persistence/aivencredentials" + "github.com/nais/api/internal/persistence/aivenversion" "github.com/nais/api/internal/slug" "github.com/nais/api/internal/validate" "github.com/nais/api/internal/workload" aiven_io_v1alpha1 "github.com/nais/liberator/pkg/apis/aiven.io/v1alpha1" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/runtime/schema" @@ -67,6 +70,8 @@ type Valkey struct { EnvironmentName string `json:"-"` WorkloadReference *workload.Reference `json:"-"` AivenProject string `json:"-"` + MajorVersion ValkeyMajorVersion `json:"-"` + ReportedVersion string `json:"-"` } func (Valkey) IsPersistence() {} @@ -107,8 +112,7 @@ type ValkeyAccess struct { } type ValkeyStatus struct { - State string `json:"state"` - Conditions []metav1.Condition `json:"conditions"` + State string `json:"state"` } type ValkeyOrder struct { @@ -191,6 +195,14 @@ func toValkey(u *unstructured.Unstructured, envName string) (*Valkey, error) { maxMemoryPolicy = "" } + reportedVersion, _, _ := unstructured.NestedString(u.Object, aivenversion.StatusVersion...) + + pinnedVersion, _, _ := unstructured.NestedString(u.Object, specValkeyVersion...) + majorVersion, err := aivenversion.MatchPin(pinnedVersion, allMajorVersions, ValkeyMajorVersion.ToAivenString) + if err != nil { + majorVersion = "" + } + notifyKeyspaceEvents, _, _ := unstructured.NestedString(u.Object, specNotifyKeyspaceEvents...) numberOfDatabases, found, _ := unstructured.NestedNumberAsFloat64(u.Object, specNumberOfDatabases...) @@ -213,8 +225,7 @@ func toValkey(u *unstructured.Unstructured, envName string) (*Valkey, error) { EnvironmentName: envName, TerminationProtection: terminationProtection, Status: &ValkeyStatus{ - Conditions: obj.Status.Conditions, - State: obj.Status.State, + State: obj.Status.State, }, TeamSlug: slug.Slug(obj.GetNamespace()), WorkloadReference: workload.ReferenceFromOwnerReferences(obj.GetOwnerReferences()), @@ -224,6 +235,8 @@ func toValkey(u *unstructured.Unstructured, envName string) (*Valkey, error) { MaxMemoryPolicy: maxMemoryPolicy, NotifyKeyspaceEvents: notifyKeyspaceEvents, Databases: int(numberOfDatabases), + MajorVersion: majorVersion, + ReportedVersion: reportedVersion, Labels: model.UserLabels(obj.GetLabels()), }, nil } @@ -272,6 +285,133 @@ type ValkeyInput struct { Databases *int `json:"databases,omitempty"` } +type ValkeyMajorVersion string + +const ( + ValkeyMajorVersionV8_1 ValkeyMajorVersion = "V8_1" + ValkeyMajorVersionV9_0 ValkeyMajorVersion = "V9_0" + ValkeyMajorVersionV9_1 ValkeyMajorVersion = "V9_1" +) + +type upgradePath []ValkeyMajorVersion + +func (u upgradePath) String() string { + versions := make([]string, len(u)) + for i, v := range u { + versions[i] = v.String() + } + return strings.Join(versions, ",") +} + +var upgradePaths = map[ValkeyMajorVersion]upgradePath{ + ValkeyMajorVersionV8_1: {ValkeyMajorVersionV9_1}, + ValkeyMajorVersionV9_0: {ValkeyMajorVersionV9_1}, + ValkeyMajorVersionV9_1: {}, +} + +func (e ValkeyMajorVersion) ValidateUpgradePath(other ValkeyMajorVersion) error { + path, ok := upgradePaths[other] + if !ok { + return fmt.Errorf("unknown Valkey major version: %q", other) + } + + if len(path) == 0 { + return apierror.Errorf("Cannot change Valkey version from %v to %v. No further upgrades available.", other, e) + } + + if slices.Contains(path, e) { + return nil + } + + return apierror.Errorf("Cannot change Valkey version from %v to %v. New version must be one of [%s]", other, e, path) +} + +// newestMajorVersion returns the highest version the codebase knows about, derived +// from the versions themselves so that adding one to upgradePaths is enough. +func newestMajorVersion() (ValkeyMajorVersion, error) { + var newest ValkeyMajorVersion + var highest *semver.Version + + for v := range upgradePaths { + s, err := v.ToAivenString() + if err != nil { + return "", err + } + parsed, err := semver.NewVersion(s) + if err != nil { + return "", fmt.Errorf("parsing Valkey version %q: %w", s, err) + } + if highest == nil || parsed.GreaterThan(highest) { + highest, newest = parsed, v + } + } + + if newest == "" { + return "", fmt.Errorf("no Valkey versions defined") + } + return newest, nil +} + +// allMajorVersions is derived from upgradePaths rather than maintained beside it, so the +// two cannot disagree. +var allMajorVersions = aivenversion.Declared(upgradePaths) + +func (e ValkeyMajorVersion) IsValid() bool { + return slices.Contains(allMajorVersions, e) +} + +// deprecatedMajorVersions name versions Aiven still runs, so they stay valid for reads +// and remain in the GraphQL enum. No client may choose one. Each carries its own reason +// because versions are deprecated for different causes, and refusing without saying why +// leaves the caller guessing. +var deprecatedMajorVersions = map[ValkeyMajorVersion]string{ + ValkeyMajorVersionV9_0: "no longer supported by back-end services", +} + +func (e ValkeyMajorVersion) DeprecationReason() (string, bool) { + reason, ok := deprecatedMajorVersions[e] + return reason, ok +} + +func (e ValkeyMajorVersion) String() string { + return string(e) +} + +func (e *ValkeyMajorVersion) UnmarshalGQL(v any) error { + str, ok := v.(string) + if !ok { + return fmt.Errorf("enums must be strings") + } + + *e = ValkeyMajorVersion(str) + if !e.IsValid() { + return fmt.Errorf("%s is not a valid ValkeyMajorVersion", str) + } + return nil +} + +func (e ValkeyMajorVersion) MarshalGQL(w io.Writer) { + fmt.Fprint(w, strconv.Quote(e.String())) +} + +// ToAivenString returns the version as Aiven writes it in userConfig, e.g. "8.1". +func (e ValkeyMajorVersion) ToAivenString() (string, error) { + switch e { + case ValkeyMajorVersionV8_1: + return "8.1", nil + case ValkeyMajorVersionV9_0: + return "9.0", nil + case ValkeyMajorVersionV9_1: + return "9.1", nil + default: + return "", fmt.Errorf("unexpected Valkey major version: %q", e) + } +} + +func ValkeyMajorVersionFromAivenString(s string) (ValkeyMajorVersion, error) { + return aivenversion.Match(s, allMajorVersions, upgradePaths) +} + func (v *ValkeyInput) Validate(ctx context.Context) error { return v.ValidationErrors(ctx).NilIfEmpty() } @@ -286,8 +426,9 @@ func (v *ValkeyInput) ValidationErrors(ctx context.Context) *validate.Validation if !v.Memory.IsValid() { verr.Add("memory", "Invalid Valkey memory: %s.", v.Memory) } + if v.MaxMemoryPolicy != nil && !v.MaxMemoryPolicy.IsValid() { - verr.Add("version", "Invalid Valkey max memory policy: %s.", v.MaxMemoryPolicy.String()) + verr.Add("maxMemoryPolicy", "Invalid Valkey max memory policy: %s.", v.MaxMemoryPolicy.String()) } if v.Databases != nil { @@ -301,6 +442,19 @@ func (v *ValkeyInput) ValidationErrors(ctx context.Context) *validate.Validation type CreateValkeyInput struct { ValkeyInput + Version *ValkeyMajorVersion `json:"version,omitempty"` +} + +func (i *CreateValkeyInput) Validate(ctx context.Context) error { + verr := i.ValkeyInput.ValidationErrors(ctx) + if i.Version != nil { + if !i.Version.IsValid() { + verr.Add("version", "Invalid Valkey version: %s.", i.Version) + } else if reason, deprecated := i.Version.DeprecationReason(); deprecated { + verr.Add("version", "Valkey version %s is deprecated: %s.", i.Version, reason) + } + } + return verr.NilIfEmpty() } type CreateValkeyPayload struct { @@ -479,11 +633,17 @@ func (e ValkeyTier) MarshalGQL(w io.Writer) { type UpdateValkeyInput struct { ValkeyInput - Labels []*model.ResourceLabel `json:"labels,omitempty"` + Version ValkeyMajorVersion `json:"version"` + Labels []*model.ResourceLabel `json:"labels,omitempty"` } func (i *UpdateValkeyInput) Validate(ctx context.Context) error { verr := i.ValkeyInput.ValidationErrors(ctx) + if !i.Version.IsValid() { + verr.Add("version", "Invalid Valkey version: %s.", i.Version) + } else if reason, deprecated := i.Version.DeprecationReason(); deprecated { + verr.Add("version", "Valkey version %s is deprecated: %s.", i.Version, reason) + } validateUserLabels(verr, i.Labels) return verr.NilIfEmpty() } @@ -506,6 +666,11 @@ type DeleteValkeyPayload struct { ValkeyDeleted *bool `json:"valkeyDeleted,omitempty"` } +type ValkeyVersion struct { + Actual *string `json:"actual,omitempty"` + DesiredMajor ValkeyMajorVersion `json:"desiredMajor"` +} + type ValkeyState int const ( diff --git a/internal/persistence/valkey/queries.go b/internal/persistence/valkey/queries.go index 9d5a4ac0f..4a01cadb5 100644 --- a/internal/persistence/valkey/queries.go +++ b/internal/persistence/valkey/queries.go @@ -31,6 +31,7 @@ var ( specMaxMemoryPolicy = []string{"spec", "userConfig", "valkey_maxmemory_policy"} specNotifyKeyspaceEvents = []string{"spec", "userConfig", "valkey_notify_keyspace_events"} specNumberOfDatabases = []string{"spec", "userConfig", "valkey_number_of_databases"} + specValkeyVersion = []string{"spec", "userConfig", "valkey_version"} ) func GetByIdent(ctx context.Context, id ident.Ident) (*Valkey, error) { @@ -97,6 +98,31 @@ func ListAccess(ctx context.Context, valkey *Valkey, page *pagination.Pagination return pagination.NewConnection(ret, page, len(all)), nil } +func GetValkeyVersion(v *Valkey) (*ValkeyVersion, error) { + if v.ReportedVersion == "" { + // Kubernetes is only eventually consistent with Aiven: the CR exists here before + // the service exists there, and lingers after it is deleted. The operator has no + // version to record in that window, so fall back to the version pinned in the CR + // rather than failing the whole query. Actual stays nil, which the schema + // documents as "available after the instance is created". + if v.MajorVersion == "" { + return nil, fmt.Errorf("no Valkey version known for service %q", v.FullyQualifiedName()) + } + return &ValkeyVersion{DesiredMajor: v.MajorVersion}, nil + } + + actual := v.ReportedVersion + major, err := ValkeyMajorVersionFromAivenString(actual) + if err != nil { + return nil, err + } + + return &ValkeyVersion{ + Actual: &actual, + DesiredMajor: major, + }, nil +} + func ListForWorkload(ctx context.Context, teamSlug slug.Slug, environmentName string, references []nais_io_v1.Valkey, orderBy *ValkeyOrder) (*ValkeyConnection, error) { all := fromContext(ctx).client.watcher.GetByNamespace(teamSlug.String(), watcher.InCluster(environmentName)) ret := make([]*Valkey, 0) @@ -154,6 +180,19 @@ func Create(ctx context.Context, input CreateValkeyInput) (*CreateValkeyPayload, return nil, err } + desired := input.Version + if desired == nil { + newest, err := newestMajorVersion() + if err != nil { + return nil, err + } + desired = &newest + } + version, err := desired.ToAivenString() + if err != nil { + return nil, err + } + res.Object["spec"] = map[string]any{ "cloudName": "google-europe-north1", "plan": machine.AivenPlan, @@ -167,6 +206,10 @@ func Create(ctx context.Context, input CreateValkeyInput) (*CreateValkeyPayload, }, } + if err := unstructured.SetNestedField(res.Object, version, specValkeyVersion...); err != nil { + return nil, err + } + if input.MaxMemoryPolicy != nil { maxMemoryPolicy := input.MaxMemoryPolicy.ToAivenString() err := unstructured.SetNestedField(res.Object, maxMemoryPolicy, specMaxMemoryPolicy...) @@ -251,6 +294,12 @@ func Update(ctx context.Context, input UpdateValkeyInput) (*UpdateValkeyPayload, } changes = append(changes, res...) + res, err = updateValkeyVersion(valkey, input) + if err != nil { + return nil, err + } + changes = append(changes, res...) + res, err = updateMaxMemoryPolicy(valkey, input) if err != nil { return nil, err @@ -483,6 +532,53 @@ func updateMaxMemoryPolicy(valkey *unstructured.Unstructured, input UpdateValkey return changes, nil } +func updateValkeyVersion(valkey *unstructured.Unstructured, input UpdateValkeyInput) ([]*ValkeyUpdatedActivityLogEntryDataUpdatedField, error) { + changes := make([]*ValkeyUpdatedActivityLogEntryDataUpdatedField, 0) + + newValue, err := input.Version.ToAivenString() + if err != nil { + return nil, err + } + + v, err := toValkey(valkey, input.EnvironmentName) + if err != nil { + return nil, err + } + + // Only the running version can judge whether the requested change is legal, and the + // operator records it once it has observed the service. The CR's own pin says what + // was asked for, so it cannot stand in. + if v.ReportedVersion == "" { + return nil, apierror.Errorf("The running version of %q is not known yet. Try again shortly.", v.Name) + } + + oldValue := v.ReportedVersion + oldMajor, err := ValkeyMajorVersionFromAivenString(oldValue) + if err != nil { + return nil, err + } + + if oldMajor == input.Version { + return changes, nil + } + + if err := input.Version.ValidateUpgradePath(oldMajor); err != nil { + return nil, err + } + + changes = append(changes, &ValkeyUpdatedActivityLogEntryDataUpdatedField{ + Field: "version", + OldValue: &oldValue, + NewValue: new(newValue), + }) + + if err := unstructured.SetNestedField(valkey.Object, newValue, specValkeyVersion...); err != nil { + return nil, err + } + + return changes, nil +} + func updateNotifyKeyspaceEvents(valkey *unstructured.Unstructured, input UpdateValkeyInput) ([]*ValkeyUpdatedActivityLogEntryDataUpdatedField, error) { changes := make([]*ValkeyUpdatedActivityLogEntryDataUpdatedField, 0) @@ -638,27 +734,19 @@ func Delete(ctx context.Context, input DeleteValkeyInput) (*DeleteValkeyPayload, }, nil } -func State(ctx context.Context, v *Valkey) (ValkeyState, error) { - s, err := fromContext(ctx).aivenClient.ServiceGet(ctx, v.AivenProject, v.FullyQualifiedName()) - if err != nil { - // The Valkey instance may not have been created in Aiven yet, or it has been deleted. - // In both cases, we return "unknown" state rather than an error. - if aiven.IsNotFound(err) { - return ValkeyStateUnknown, nil - } - return ValkeyStateUnknown, err - } - - switch s.State { +// State reports the state the operator last observed. An instance Aiven has not created +// yet, or has already deleted, has no state recorded and reads as unknown. +func State(v *Valkey) ValkeyState { + switch v.Status.State { case "RUNNING": - return ValkeyStateRunning, nil + return ValkeyStateRunning case "REBALANCING": - return ValkeyStateRebalancing, nil + return ValkeyStateRebalancing case "REBUILDING": - return ValkeyStateRebuilding, nil + return ValkeyStateRebuilding case "POWEROFF": - return ValkeyStatePoweroff, nil + return ValkeyStatePoweroff default: - return ValkeyStateUnknown, nil + return ValkeyStateUnknown } } diff --git a/internal/persistence/valkey/sortfilter.go b/internal/persistence/valkey/sortfilter.go index 393925d44..8aa771134 100644 --- a/internal/persistence/valkey/sortfilter.go +++ b/internal/persistence/valkey/sortfilter.go @@ -23,12 +23,7 @@ func init() { }, "NAME") SortFilterValkey.RegisterConcurrentSort("STATE", func(ctx context.Context, a *Valkey) int { - s, err := State(ctx, a) - if err != nil { - return int(ValkeyStateUnknown) - } - - return int(s) + return int(State(a)) }, "NAME") SortFilterValkey.RegisterFilter(func(ctx context.Context, v *Valkey, filter *ValkeyFilter) bool { diff --git a/internal/thirdparty/aiven/dataloader.go b/internal/thirdparty/aiven/dataloader.go index 2d8b037e4..f7fc5037e 100644 --- a/internal/thirdparty/aiven/dataloader.go +++ b/internal/thirdparty/aiven/dataloader.go @@ -9,7 +9,7 @@ type ctxKey int const loadersKey ctxKey = iota func NewLoaderContext(ctx context.Context, projects Projects) context.Context { - return context.WithValue(ctx, loadersKey, newLoaders(projects)) + return context.WithValue(ctx, loadersKey, &loaders{projects: projects}) } func fromContext(ctx context.Context) *loaders { @@ -19,9 +19,3 @@ func fromContext(ctx context.Context) *loaders { type loaders struct { projects Projects } - -func newLoaders(projects Projects) *loaders { - return &loaders{ - projects: projects, - } -} diff --git a/internal/thirdparty/aiven/error.go b/internal/thirdparty/aiven/error.go deleted file mode 100644 index a146c3376..000000000 --- a/internal/thirdparty/aiven/error.go +++ /dev/null @@ -1,10 +0,0 @@ -package aiven - -import ( - aiven "github.com/aiven/go-client-codegen" -) - -// IsNotFound re-exports [aiven.IsNotFound] -func IsNotFound(err error) bool { - return aiven.IsNotFound(err) -} diff --git a/internal/thirdparty/aiven/fake.go b/internal/thirdparty/aiven/fake.go index 95b364981..e0bc54d3f 100644 --- a/internal/thirdparty/aiven/fake.go +++ b/internal/thirdparty/aiven/fake.go @@ -71,7 +71,7 @@ func (f *FakeAivenClient) ProjectAlertsList(ctx context.Context, p string) ([]pr } // ServiceGet returns hardcoded example dataset -func (f *FakeAivenClient) ServiceGet(_ context.Context, _ string, serviceName string, _ ...[2]string) (*aiven.ServiceGetOut, error) { +func (f *FakeAivenClient) ServiceGet(_ context.Context, _ string, _ string, _ ...[2]string) (*aiven.ServiceGetOut, error) { description := "This is a description (Nais API call it title)" link := "https://nais.io" impact := "This is the impact (Nais API call it description)" @@ -79,15 +79,7 @@ func (f *FakeAivenClient) ServiceGet(_ context.Context, _ string, serviceName st deadline := startAt.Add(24 * time.Hour).Format(time.RFC3339) startAfter := startAt.Add(1 * time.Hour).Format(time.RFC3339) - state := aiven.ServiceStateTypeRunning - if strings.HasSuffix(serviceName, "poweroff") { - state = aiven.ServiceStateTypePoweroff - } else if strings.HasSuffix(serviceName, "rebalancing") { - state = aiven.ServiceStateTypeRebalancing - } - return &aiven.ServiceGetOut{ - State: state, Maintenance: &aiven.MaintenanceOut{ Updates: []aiven.UpdateOut{ { @@ -106,8 +98,5 @@ func (f *FakeAivenClient) ServiceGet(_ context.Context, _ string, serviceName st Dow: "sunday", Time: "12:34:56", }, - Metadata: map[string]any{ - "opensearch_version": "2.17.2", - }, }, nil }